Coverage Report

Created: 2026-08-15 06:25

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/tinysparql/src/libtinysparql/direct/tracker-direct.c
Line
Count
Source
1
/*
2
 * Copyright (C) 2010, Nokia <ivan.frade@nokia.com>
3
 * Copyright (C) 2017, Red Hat, Inc.
4
 *
5
 * This library is free software; you can redistribute it and/or
6
 * modify it under the terms of the GNU Lesser General Public
7
 * License as published by the Free Software Foundation; either
8
 * version 2.1 of the License, or (at your option) any later version.
9
 *
10
 * This library is distributed in the hope that it will be useful,
11
 * but WITHOUT ANY WARRANTY; without even the implied warranty of
12
 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the GNU
13
 * Lesser General Public License for more details.
14
 *
15
 * You should have received a copy of the GNU Lesser General Public
16
 * License along with this library; if not, write to the
17
 * Free Software Foundation, Inc., 51 Franklin Street, Fifth Floor,
18
 * Boston, MA  02110-1301, USA.
19
 */
20
21
#include "config.h"
22
23
#include "tracker-direct.h"
24
#include "tracker-direct-batch.h"
25
#include "tracker-direct-statement.h"
26
27
#include <tracker-common.h>
28
29
#include "core/tracker-data.h"
30
#include "tracker-deserializer-directory.h"
31
#include "tracker-notifier-private.h"
32
#include "tracker-private.h"
33
#include "tracker-serializer.h"
34
35
typedef struct _TrackerDirectConnectionPrivate TrackerDirectConnectionPrivate;
36
37
struct _TrackerDirectConnectionPrivate
38
{
39
  TrackerSparqlConnectionFlags flags;
40
  GFile *store;
41
  GFile *ontology;
42
  GInputStream *ontology_rdf;
43
44
  TrackerNamespaceManager *namespace_manager;
45
  TrackerDataManager *data_manager;
46
  GMutex update_mutex;
47
48
  GThreadPool *update_thread; /* Contains 1 exclusive thread */
49
  GThreadPool *select_pool;
50
51
  GList *notifiers;
52
  GMutex notifiers_mutex;
53
54
  gint64 timestamp;
55
  gint64 cleanup_timestamp;
56
57
  guint cleanup_timeout_id;
58
  TrackerRdfFormat ontology_rdf_format;
59
60
  guint initialized : 1;
61
  guint closing     : 1;
62
};
63
64
typedef enum {
65
  TASK_TYPE_QUERY,
66
  TASK_TYPE_QUERY_STATEMENT,
67
  TASK_TYPE_SERIALIZE,
68
  TASK_TYPE_SERIALIZE_STATEMENT,
69
  TASK_TYPE_UPDATE,
70
  TASK_TYPE_UPDATE_BLANK,
71
  TASK_TYPE_UPDATE_RESOURCE,
72
  TASK_TYPE_UPDATE_BATCH,
73
  TASK_TYPE_UPDATE_STATEMENT,
74
  TASK_TYPE_DESERIALIZE,
75
  TASK_TYPE_RELEASE_MEMORY,
76
} TaskType;
77
78
typedef struct {
79
  TaskType type;
80
81
  union {
82
    gchar *sparql;
83
84
    TrackerBatch *batch;
85
86
    struct {
87
      TrackerSparqlStatement *stmt;
88
      GHashTable *parameters;
89
    } statement;
90
91
    struct {
92
      gchar *graph;
93
      TrackerResource *resource;
94
    } update_resource;
95
96
    struct {
97
      gchar *sparql;
98
      TrackerRdfFormat format;
99
      TrackerSerializeFlags flags;
100
    } serialize;
101
102
    struct {
103
      TrackerSparqlStatement *stmt;
104
      GHashTable *parameters;
105
      TrackerRdfFormat format;
106
      TrackerSerializeFlags flags;
107
    } serialize_statement;
108
109
    struct {
110
      GInputStream *stream;
111
      gchar *default_graph;
112
      TrackerRdfFormat format;
113
      TrackerDeserializeFlags flags;
114
    } deserialize;
115
  } d;
116
} TaskData;
117
118
enum {
119
  PROP_0,
120
  PROP_FLAGS,
121
  PROP_STORE_LOCATION,
122
  PROP_ONTOLOGY_LOCATION,
123
  PROP_ONTOLOGY_STREAM,
124
  PROP_ONTOLOGY_STREAM_FORMAT,
125
  N_PROPS
126
};
127
128
static GParamSpec *props[N_PROPS] = { NULL };
129
130
static void tracker_direct_connection_initable_iface_init (GInitableIface *iface);
131
static void tracker_direct_connection_async_initable_iface_init (GAsyncInitableIface *iface);
132
133
G_DEFINE_QUARK (TrackerDirectNotifier, tracker_direct_notifier)
134
135
44.6k
G_DEFINE_TYPE_WITH_CODE (TrackerDirectConnection, tracker_direct_connection,
136
44.6k
                         TRACKER_TYPE_SPARQL_CONNECTION,
137
44.6k
                         G_ADD_PRIVATE (TrackerDirectConnection)
138
44.6k
                         G_IMPLEMENT_INTERFACE (G_TYPE_INITABLE,
139
44.6k
                                                tracker_direct_connection_initable_iface_init)
140
44.6k
                         G_IMPLEMENT_INTERFACE (G_TYPE_ASYNC_INITABLE,
141
44.6k
                                                tracker_direct_connection_async_initable_iface_init))
142
44.6k
143
44.6k
static TaskData *
144
44.6k
task_data_new (TaskType type)
145
44.6k
{
146
0
  TaskData *task;
147
148
0
  task = g_new0 (TaskData, 1);
149
0
  task->type = type;
150
151
0
  return task;
152
0
}
153
154
static void
155
task_data_free (TaskData *task)
156
0
{
157
0
  switch (task->type) {
158
0
  case TASK_TYPE_QUERY:
159
0
  case TASK_TYPE_UPDATE:
160
0
  case TASK_TYPE_UPDATE_BLANK:
161
0
    g_free (task->d.sparql);
162
0
    break;
163
0
  case TASK_TYPE_SERIALIZE:
164
0
    g_free (task->d.serialize.sparql);
165
0
    break;
166
0
  case TASK_TYPE_UPDATE_RESOURCE:
167
0
    g_free (task->d.update_resource.graph);
168
0
    g_object_unref (task->d.update_resource.resource);
169
0
    break;
170
0
  case TASK_TYPE_UPDATE_BATCH:
171
0
    g_clear_object (&task->d.batch);
172
0
    break;
173
0
  case TASK_TYPE_UPDATE_STATEMENT:
174
0
  case TASK_TYPE_QUERY_STATEMENT:
175
0
    g_clear_object (&task->d.statement.stmt);
176
0
    g_clear_pointer (&task->d.statement.parameters, g_hash_table_unref);
177
0
    break;
178
0
  case TASK_TYPE_SERIALIZE_STATEMENT:
179
0
    g_clear_object (&task->d.serialize_statement.stmt);
180
0
    g_clear_pointer (&task->d.serialize_statement.parameters,
181
0
                     g_hash_table_unref);
182
0
    break;
183
0
  case TASK_TYPE_RELEASE_MEMORY:
184
0
    break;
185
0
  case TASK_TYPE_DESERIALIZE:
186
0
    g_clear_object (&task->d.deserialize.stream);
187
0
    g_free (task->d.deserialize.default_graph);
188
0
    break;
189
0
  }
190
0
  g_free (task);
191
0
}
192
193
static gboolean
194
cleanup_timeout_cb (gpointer user_data)
195
0
{
196
0
  TrackerDirectConnection *conn = user_data;
197
0
  TrackerDirectConnectionPrivate *priv;
198
0
  gint64 timestamp;
199
0
  GTask *task;
200
201
0
  priv = tracker_direct_connection_get_instance_private (conn);
202
0
  timestamp = g_get_monotonic_time ();
203
204
  /* If we already cleaned up */
205
0
  if (priv->timestamp < priv->cleanup_timestamp)
206
0
    return G_SOURCE_CONTINUE;
207
  /* If the connection was used less than 10s ago */
208
0
  if (timestamp - priv->timestamp < 10 * G_USEC_PER_SEC)
209
0
    return G_SOURCE_CONTINUE;
210
211
0
  priv->cleanup_timestamp = timestamp;
212
213
0
  task = g_task_new (conn, NULL, NULL, NULL);
214
0
  g_task_set_task_data (task,
215
0
                        task_data_new (TASK_TYPE_RELEASE_MEMORY),
216
0
                        (GDestroyNotify) task_data_free);
217
218
0
  g_thread_pool_push (priv->update_thread, task, NULL);
219
220
0
  return G_SOURCE_CONTINUE;
221
0
}
222
223
gboolean
224
update_resource (TrackerData      *data,
225
                 const gchar      *graph,
226
                 TrackerResource  *resource,
227
                 GError          **error)
228
0
{
229
0
  GError *inner_error = NULL;
230
0
  GHashTable *visited;
231
232
0
  if (!tracker_data_begin_transaction (data, &inner_error))
233
0
    goto error;
234
235
0
  visited = g_hash_table_new_full (NULL, NULL, NULL,
236
0
                                   (GDestroyNotify) tracker_rowid_free);
237
238
0
  tracker_data_update_resource (data,
239
0
                                graph,
240
0
                                resource,
241
0
                                NULL,
242
0
                                visited,
243
0
                                &inner_error);
244
245
0
  g_hash_table_unref (visited);
246
247
0
  if (inner_error) {
248
0
    tracker_data_rollback_transaction (data);
249
0
    goto error;
250
0
  }
251
252
0
  if (!tracker_data_commit_transaction (data, &inner_error))
253
0
    goto error;
254
255
0
  return TRUE;
256
257
0
error:
258
0
  g_propagate_error (error, inner_error);
259
0
  return FALSE;
260
0
}
261
262
static TrackerSerializerFormat
263
convert_format (TrackerRdfFormat format)
264
0
{
265
0
  switch (format) {
266
0
  case TRACKER_RDF_FORMAT_TURTLE:
267
0
    return TRACKER_SERIALIZER_FORMAT_TTL;
268
0
  case TRACKER_RDF_FORMAT_TRIG:
269
0
    return TRACKER_SERIALIZER_FORMAT_TRIG;
270
0
  case TRACKER_RDF_FORMAT_JSON_LD:
271
0
    return TRACKER_SERIALIZER_FORMAT_JSON_LD;
272
0
  default:
273
0
    g_assert_not_reached ();
274
0
  }
275
0
}
276
277
static void
278
update_thread_func (gpointer data,
279
                    gpointer user_data)
280
0
{
281
0
  TrackerDirectConnectionPrivate *priv;
282
0
  TrackerDirectConnection *conn;
283
0
  GTask *task = data;
284
0
  TaskData *task_data = g_task_get_task_data (task);
285
0
  TrackerData *tracker_data;
286
0
  GError *error = NULL;
287
0
  gpointer retval = NULL;
288
0
  GDestroyNotify destroy_notify = NULL;
289
0
  gboolean update_timestamp = TRUE;
290
291
0
  conn = user_data;
292
0
  priv = tracker_direct_connection_get_instance_private (conn);
293
294
0
  g_mutex_lock (&priv->update_mutex);
295
0
  tracker_data = tracker_data_manager_get_data (priv->data_manager);
296
297
0
  switch (task_data->type) {
298
0
  case TASK_TYPE_QUERY:
299
0
  case TASK_TYPE_QUERY_STATEMENT:
300
0
  case TASK_TYPE_SERIALIZE:
301
0
  case TASK_TYPE_SERIALIZE_STATEMENT:
302
0
    g_warning ("Queries don't go through this thread");
303
0
    break;
304
0
  case TASK_TYPE_UPDATE:
305
0
    tracker_data_update_sparql (tracker_data, task_data->d.sparql, &error);
306
0
    break;
307
0
  case TASK_TYPE_UPDATE_BLANK:
308
0
    retval = tracker_data_update_sparql_blank (tracker_data, task_data->d.sparql, &error);
309
0
    destroy_notify = (GDestroyNotify) g_variant_unref;
310
0
    break;
311
0
  case TASK_TYPE_UPDATE_RESOURCE:
312
0
    update_resource (tracker_data,
313
0
                     task_data->d.update_resource.graph,
314
0
                     task_data->d.update_resource.resource,
315
0
                     &error);
316
0
    break;
317
0
  case TASK_TYPE_DESERIALIZE: {
318
0
    TrackerSparqlCursor *deserializer;
319
320
0
    if (!tracker_data_begin_transaction (tracker_data, &error))
321
0
      break;
322
323
0
    deserializer = tracker_deserializer_new (task_data->d.deserialize.stream,
324
0
               NULL,
325
0
                                             convert_format (task_data->d.deserialize.format));
326
327
0
    if (tracker_data_load_from_deserializer (tracker_data,
328
0
                                             TRACKER_DESERIALIZER (deserializer),
329
0
                                             task_data->d.deserialize.default_graph,
330
0
                                             NULL,
331
0
                                             &error)) {
332
0
      tracker_data_commit_transaction (tracker_data, &error);
333
0
    } else {
334
0
      tracker_data_rollback_transaction (tracker_data);
335
0
    }
336
0
    g_object_unref (deserializer);
337
0
    break;
338
0
  }
339
0
  case TASK_TYPE_UPDATE_BATCH:
340
0
    tracker_direct_batch_update (TRACKER_DIRECT_BATCH (task_data->d.batch),
341
0
                                 priv->data_manager, &error);
342
0
    break;
343
0
  case TASK_TYPE_UPDATE_STATEMENT:
344
0
    if (!tracker_data_begin_transaction (tracker_data, &error))
345
0
      break;
346
347
0
    if (tracker_direct_statement_execute_update (task_data->d.statement.stmt,
348
0
                                                 task_data->d.statement.parameters,
349
0
                                                 NULL,
350
0
                                                 &error)) {
351
0
      tracker_data_commit_transaction (tracker_data, &error);
352
0
    } else {
353
0
      tracker_data_rollback_transaction (tracker_data);
354
0
    }
355
0
    break;
356
0
  case TASK_TYPE_RELEASE_MEMORY:
357
0
    tracker_data_manager_release_memory (priv->data_manager);
358
0
    update_timestamp = FALSE;
359
0
    break;
360
0
  }
361
362
0
  if (error)
363
0
    g_task_return_error (task, error);
364
0
  else if (retval)
365
0
    g_task_return_pointer (task, retval, destroy_notify);
366
0
  else
367
0
    g_task_return_boolean (task, TRUE);
368
369
0
  g_object_unref (task);
370
371
0
  if (update_timestamp)
372
0
    tracker_direct_connection_update_timestamp (conn);
373
374
0
  g_mutex_unlock (&priv->update_mutex);
375
0
}
376
377
static void
378
execute_query_in_thread (GTask    *task,
379
                         TaskData *task_data)
380
0
{
381
0
  TrackerSparqlConnection *conn;
382
0
  TrackerSparqlCursor *cursor;
383
0
  GError *error = NULL;
384
385
0
  if (g_task_return_error_if_cancelled (task))
386
0
    return;
387
388
0
  conn = g_task_get_source_object (task);
389
390
0
  if (task_data->type == TASK_TYPE_QUERY) {
391
0
    cursor = tracker_sparql_connection_query (conn,
392
0
                                              task_data->d.sparql,
393
0
                                              g_task_get_cancellable (task),
394
0
                                              &error);
395
0
  } else if (task_data->type == TASK_TYPE_QUERY_STATEMENT) {
396
0
    TrackerSparql *sparql;
397
398
0
    sparql = tracker_direct_statement_get_sparql (task_data->d.statement.stmt);
399
0
    cursor = tracker_sparql_execute_cursor (sparql,
400
0
                                            task_data->d.statement.parameters,
401
0
                                            &error);
402
0
  } else {
403
0
    g_assert_not_reached ();
404
0
  }
405
406
0
  if (cursor) {
407
0
    tracker_direct_connection_update_timestamp (TRACKER_DIRECT_CONNECTION (conn));
408
0
    g_task_return_pointer (task, cursor, g_object_unref);
409
0
  } else {
410
0
    g_task_return_error (task, error);
411
0
  }
412
0
}
413
414
static void
415
serialize_in_thread (GTask    *task,
416
                     TaskData *task_data)
417
0
{
418
0
  TrackerDirectConnectionPrivate *priv;
419
0
  TrackerDirectConnection *conn;
420
0
  TrackerSparql *query = NULL;
421
0
  TrackerSparqlCursor *cursor = NULL;
422
0
  TrackerNamespaceManager *namespaces;
423
0
  TrackerRdfFormat format;
424
0
  GInputStream *istream = NULL;
425
0
  GHashTable *parameters = NULL;
426
0
  GError *error = NULL;
427
428
0
  conn = g_task_get_source_object (task);
429
0
  priv = tracker_direct_connection_get_instance_private (conn);
430
431
0
  if (task_data->type == TASK_TYPE_SERIALIZE) {
432
0
    format = task_data->d.serialize.format;
433
0
    query = tracker_sparql_new (priv->data_manager,
434
0
                                task_data->d.serialize.sparql,
435
0
                                &error);
436
0
    if (!query)
437
0
      goto out;
438
0
  } else if (task_data->type == TASK_TYPE_SERIALIZE_STATEMENT) {
439
0
    TrackerSparqlStatement *stmt;
440
441
0
    format = task_data->d.serialize_statement.format;
442
0
    stmt = task_data->d.serialize_statement.stmt;
443
0
    query = g_object_ref (tracker_direct_statement_get_sparql (stmt));
444
0
    parameters = task_data->d.serialize_statement.parameters;
445
0
  } else {
446
0
    g_assert_not_reached ();
447
0
  }
448
449
0
  if (!tracker_sparql_is_serializable (query)) {
450
0
    g_set_error (&error,
451
0
                 TRACKER_SPARQL_ERROR,
452
0
                 TRACKER_SPARQL_ERROR_PARSE,
453
0
                 "Query is not DESCRIBE or CONSTRUCT");
454
0
    goto out;
455
0
  }
456
457
0
  cursor = tracker_sparql_execute_cursor (query, parameters, &error);
458
0
  if (!cursor)
459
0
    goto out;
460
461
0
  tracker_direct_connection_update_timestamp (conn);
462
0
  tracker_sparql_cursor_set_connection (cursor, TRACKER_SPARQL_CONNECTION (conn));
463
0
  namespaces = tracker_sparql_connection_get_namespace_manager (TRACKER_SPARQL_CONNECTION (conn));
464
0
  istream = tracker_serializer_new (cursor, namespaces, convert_format (format));
465
466
0
 out:
467
0
  g_clear_object (&query);
468
0
  g_clear_object (&cursor);
469
470
0
  if (istream)
471
0
    g_task_return_pointer (task, istream, g_object_unref);
472
0
  else
473
0
    g_task_return_error (task, error);
474
0
}
475
476
static void
477
query_thread_pool_func (gpointer data,
478
                        gpointer user_data)
479
0
{
480
0
  TrackerDirectConnection *conn = user_data;
481
0
  TrackerDirectConnectionPrivate *priv;
482
0
  GTask *task = data;
483
0
  TaskData *task_data = g_task_get_task_data (task);
484
485
0
  priv = tracker_direct_connection_get_instance_private (conn);
486
487
0
  if (priv->closing) {
488
0
    g_task_return_new_error (task,
489
0
                             G_IO_ERROR,
490
0
                             G_IO_ERROR_CONNECTION_CLOSED,
491
0
                             "Connection is closed");
492
0
    g_object_unref (task);
493
0
    return;
494
0
  }
495
496
0
  switch (task_data->type) {
497
0
  case TASK_TYPE_QUERY:
498
0
  case TASK_TYPE_QUERY_STATEMENT:
499
0
    execute_query_in_thread (task, task_data);
500
0
    break;
501
0
  case TASK_TYPE_SERIALIZE:
502
0
  case TASK_TYPE_SERIALIZE_STATEMENT:
503
0
    serialize_in_thread (task, task_data);
504
0
    break;
505
0
  default:
506
0
    g_assert_not_reached ();
507
0
  }
508
509
0
  g_object_unref (task);
510
0
}
511
512
static gboolean
513
set_up_thread_pools (TrackerDirectConnection  *conn,
514
         GError                  **error)
515
14.8k
{
516
14.8k
  TrackerDirectConnectionPrivate *priv;
517
518
14.8k
  priv = tracker_direct_connection_get_instance_private (conn);
519
520
14.8k
  priv->select_pool = g_thread_pool_new (query_thread_pool_func,
521
14.8k
                                         conn, 16, FALSE, error);
522
14.8k
  if (!priv->select_pool)
523
0
    return FALSE;
524
525
14.8k
  priv->update_thread = g_thread_pool_new (update_thread_func,
526
14.8k
                                           conn, 1, TRUE, error);
527
14.8k
  if (!priv->update_thread)
528
0
    return FALSE;
529
530
14.8k
  return TRUE;
531
14.8k
}
532
533
static TrackerDBManagerFlags
534
translate_flags (TrackerSparqlConnectionFlags flags)
535
14.8k
{
536
14.8k
  TrackerDBManagerFlags db_flags = TRACKER_DB_MANAGER_FLAGS_NONE;
537
538
14.8k
  if ((flags & TRACKER_SPARQL_CONNECTION_FLAGS_READONLY) != 0)
539
0
    db_flags |= TRACKER_DB_MANAGER_READONLY;
540
14.8k
  if ((flags & TRACKER_SPARQL_CONNECTION_FLAGS_FTS_ENABLE_STEMMER) != 0)
541
0
    db_flags |= TRACKER_DB_MANAGER_FTS_ENABLE_STEMMER;
542
14.8k
  if ((flags & TRACKER_SPARQL_CONNECTION_FLAGS_FTS_ENABLE_UNACCENT) != 0)
543
0
    db_flags |= TRACKER_DB_MANAGER_FTS_ENABLE_UNACCENT;
544
14.8k
  if ((flags & TRACKER_SPARQL_CONNECTION_FLAGS_FTS_ENABLE_STOP_WORDS) != 0)
545
0
    db_flags |= TRACKER_DB_MANAGER_FTS_ENABLE_STOP_WORDS;
546
14.8k
  if ((flags & TRACKER_SPARQL_CONNECTION_FLAGS_FTS_IGNORE_NUMBERS) != 0)
547
0
    db_flags |= TRACKER_DB_MANAGER_FTS_IGNORE_NUMBERS;
548
14.8k
  if ((flags & TRACKER_SPARQL_CONNECTION_FLAGS_ANONYMOUS_BNODES) != 0)
549
5.84k
    db_flags |= TRACKER_DB_MANAGER_ANONYMOUS_BNODES;
550
551
  /* This flag is inverted */
552
14.8k
  if ((flags & TRACKER_SPARQL_CONNECTION_FLAGS_DISABLE_SYNTAX_EXTENSIONS) == 0)
553
9.05k
    db_flags |= TRACKER_DB_MANAGER_ENABLE_SYNTAX_EXTENSIONS;
554
555
14.8k
  return db_flags;
556
14.8k
}
557
558
static gboolean
559
tracker_direct_connection_initable_init (GInitable     *initable,
560
                                         GCancellable  *cancellable,
561
                                         GError       **error)
562
14.8k
{
563
14.8k
  TrackerDirectConnectionPrivate *priv;
564
14.8k
  TrackerDirectConnection *conn;
565
14.8k
  TrackerDBManagerFlags db_flags;
566
14.8k
  TrackerSparqlCursor *ontology_data = NULL;
567
14.8k
  GHashTable *namespaces;
568
14.8k
  GHashTableIter iter;
569
14.8k
  gchar *prefix, *ns;
570
14.8k
  GError *inner_error = NULL;
571
572
14.8k
  conn = TRACKER_DIRECT_CONNECTION (initable);
573
14.8k
  priv = tracker_direct_connection_get_instance_private (conn);
574
575
14.8k
  if (!set_up_thread_pools (conn, error))
576
0
    return FALSE;
577
578
14.8k
  if (priv->ontology &&
579
5.84k
      g_file_query_file_type (priv->ontology, G_FILE_QUERY_INFO_NONE, NULL) != G_FILE_TYPE_DIRECTORY) {
580
0
    gchar *uri;
581
582
0
    uri = g_file_get_uri (priv->ontology);
583
0
    g_set_error (error, TRACKER_DATA_ONTOLOGY_ERROR,
584
0
                 TRACKER_DATA_ONTOLOGY_NOT_FOUND,
585
0
                 "'%s' is not a ontology location", uri);
586
0
    g_free (uri);
587
0
    return FALSE;
588
0
  }
589
590
14.8k
  db_flags = translate_flags (priv->flags);
591
592
14.8k
  if (!priv->store) {
593
14.8k
    db_flags |= TRACKER_DB_MANAGER_IN_MEMORY;
594
14.8k
  }
595
596
14.8k
  if (priv->ontology) {
597
5.84k
    ontology_data = tracker_deserializer_directory_new (priv->ontology, NULL);
598
9.05k
  } else if (priv->ontology_rdf) {
599
9.05k
    TrackerSerializerFormat format;
600
601
9.05k
    switch (priv->ontology_rdf_format) {
602
9.05k
    case TRACKER_RDF_FORMAT_TURTLE:
603
9.05k
      format = TRACKER_SERIALIZER_FORMAT_TTL;
604
9.05k
      break;
605
0
    case TRACKER_RDF_FORMAT_TRIG:
606
0
      format = TRACKER_SERIALIZER_FORMAT_TRIG;
607
0
      break;
608
0
    case TRACKER_RDF_FORMAT_JSON_LD:
609
0
      format = TRACKER_SERIALIZER_FORMAT_JSON_LD;
610
0
      break;
611
0
    default:
612
0
      g_assert_not_reached ();
613
0
      break;
614
9.05k
    }
615
616
9.05k
    ontology_data = tracker_deserializer_new (priv->ontology_rdf, NULL, format);
617
9.05k
  }
618
619
14.8k
  priv->data_manager = tracker_data_manager_new (db_flags, priv->store,
620
14.8k
                                                 TRACKER_DESERIALIZER (ontology_data),
621
14.8k
                                                 100);
622
14.8k
  g_clear_object (&ontology_data);
623
624
14.8k
  if (!g_initable_init (G_INITABLE (priv->data_manager), cancellable, &inner_error)) {
625
5.84k
    g_propagate_error (error, _translate_internal_error (inner_error));
626
5.84k
    g_clear_object (&priv->data_manager);
627
5.84k
    return FALSE;
628
5.84k
  }
629
630
  /* Initialize namespace manager */
631
9.05k
  priv->namespace_manager = tracker_namespace_manager_new ();
632
9.05k
  namespaces = tracker_data_manager_get_namespaces (priv->data_manager);
633
9.05k
  g_hash_table_iter_init (&iter, namespaces);
634
635
75.3k
  while (g_hash_table_iter_next (&iter, (gpointer*) &prefix, (gpointer*) &ns)) {
636
66.3k
    tracker_namespace_manager_add_prefix (priv->namespace_manager,
637
66.3k
                                          prefix, ns);
638
66.3k
  }
639
640
9.05k
  g_hash_table_unref (namespaces);
641
642
9.05k
  priv->cleanup_timeout_id =
643
9.05k
    g_timeout_add_seconds (30, cleanup_timeout_cb, conn);
644
645
9.05k
  return TRUE;
646
14.8k
}
647
648
static void
649
tracker_direct_connection_initable_iface_init (GInitableIface *iface)
650
5
{
651
5
  iface->init = tracker_direct_connection_initable_init;
652
5
}
653
654
static void
655
async_initable_thread_func (GTask        *task,
656
                            gpointer      source_object,
657
                            gpointer      task_data,
658
                            GCancellable *cancellable)
659
0
{
660
0
  GError *error = NULL;
661
662
0
  if (!g_initable_init (G_INITABLE (source_object), cancellable, &error))
663
0
    g_task_return_error (task, error);
664
0
  else
665
0
    g_task_return_boolean (task, TRUE);
666
667
0
  g_object_unref (task);
668
0
}
669
670
static void
671
tracker_direct_connection_async_initable_init_async (GAsyncInitable      *async_initable,
672
                                                     gint                 priority,
673
                                                     GCancellable        *cancellable,
674
                                                     GAsyncReadyCallback  callback,
675
                                                     gpointer             user_data)
676
0
{
677
0
  GTask *task;
678
679
0
  task = g_task_new (async_initable, cancellable, callback, user_data);
680
0
  g_task_set_priority (task, priority);
681
0
  g_task_run_in_thread (task, async_initable_thread_func);
682
0
}
683
684
static gboolean
685
tracker_direct_connection_async_initable_init_finish (GAsyncInitable  *async_initable,
686
                                                      GAsyncResult    *res,
687
                                                      GError         **error)
688
0
{
689
0
  return g_task_propagate_boolean (G_TASK (res), error);
690
0
}
691
692
static void
693
tracker_direct_connection_async_initable_iface_init (GAsyncInitableIface *iface)
694
5
{
695
5
  iface->init_async = tracker_direct_connection_async_initable_init_async;
696
5
  iface->init_finish = tracker_direct_connection_async_initable_init_finish;
697
5
}
698
699
static void
700
tracker_direct_connection_init (TrackerDirectConnection *conn)
701
14.8k
{
702
14.8k
  TrackerDirectConnectionPrivate *priv =
703
14.8k
    tracker_direct_connection_get_instance_private (conn);
704
705
14.8k
  g_mutex_init (&priv->update_mutex);
706
14.8k
  g_mutex_init (&priv->notifiers_mutex);
707
14.8k
}
708
709
static GHashTable *
710
get_event_cache_ht (TrackerNotifier *notifier)
711
0
{
712
0
  GHashTable *events;
713
714
0
  events = g_object_get_qdata (G_OBJECT (notifier), tracker_direct_notifier_quark ());
715
0
  if (!events) {
716
0
    events = g_hash_table_new_full (g_str_hash, g_str_equal, NULL,
717
0
                                    (GDestroyNotify) _tracker_notifier_event_cache_free);
718
0
    g_object_set_qdata_full (G_OBJECT (notifier), tracker_direct_notifier_quark (),
719
0
                             events, (GDestroyNotify) g_hash_table_unref);
720
0
  }
721
722
0
  return events;
723
0
}
724
725
static TrackerNotifierEventCache *
726
lookup_event_cache (TrackerNotifier *notifier,
727
                    const gchar     *graph)
728
0
{
729
0
  TrackerNotifierEventCache *cache;
730
0
  GHashTable *events;
731
732
0
  if (!graph)
733
0
    graph = "";
734
735
0
  events = get_event_cache_ht (notifier);
736
0
  cache = g_hash_table_lookup (events, graph);
737
738
0
  if (!cache) {
739
0
    cache = _tracker_notifier_event_cache_new (notifier, graph);
740
0
    g_hash_table_insert (events,
741
0
                         (gpointer) tracker_notifier_event_cache_get_graph (cache),
742
0
                         cache);
743
0
  }
744
745
0
  return cache;
746
0
}
747
748
/* These callbacks will be called from a different thread
749
 * (always the same one though), handle with care.
750
 */
751
static void
752
statement_cb (TrackerDataUpdateType  type,
753
              const gchar           *graph,
754
              TrackerRowid           subject_id,
755
              TrackerRowid           predicate_id,
756
              TrackerRowid           object_id,
757
              GPtrArray             *rdf_types,
758
              gpointer               user_data)
759
0
{
760
0
  TrackerNotifier *notifier = user_data;
761
0
  TrackerSparqlConnection *conn = _tracker_notifier_get_connection (notifier);
762
0
  TrackerDirectConnection *direct = TRACKER_DIRECT_CONNECTION (conn);
763
0
  TrackerDirectConnectionPrivate *priv = tracker_direct_connection_get_instance_private (direct);
764
0
  TrackerOntologies *ontologies = tracker_data_manager_get_ontologies (priv->data_manager);
765
0
  TrackerProperty *rdf_type = tracker_ontologies_get_rdf_type (ontologies);
766
0
  TrackerNotifierEventCache *cache;
767
0
  TrackerClass *rdf_type_class = NULL;
768
0
  guint i;
769
770
0
  cache = lookup_event_cache (notifier, graph);
771
772
0
  if (predicate_id == tracker_property_get_id (rdf_type)) {
773
0
    const gchar *uri;
774
775
0
    uri = tracker_ontologies_get_uri_by_id (ontologies, object_id);
776
0
    rdf_type_class = tracker_ontologies_get_class_by_uri (ontologies, uri);
777
0
  }
778
779
0
  for (i = 0; i < rdf_types->len; i++) {
780
0
    TrackerClass *class = g_ptr_array_index (rdf_types, i);
781
0
    TrackerNotifierEventType event_type;
782
783
0
    if (!tracker_class_get_notify (class))
784
0
      continue;
785
786
0
    if (rdf_type_class && class == rdf_type_class) {
787
0
      if (type == TRACKER_DATA_INSERT)
788
0
        event_type = TRACKER_NOTIFIER_EVENT_CREATE;
789
0
      else
790
0
        event_type = TRACKER_NOTIFIER_EVENT_DELETE;
791
0
    } else {
792
0
      event_type = TRACKER_NOTIFIER_EVENT_UPDATE;
793
0
    }
794
795
0
    _tracker_notifier_event_cache_push_event (cache, subject_id, event_type);
796
0
  }
797
0
}
798
799
static void
800
transaction_cb (TrackerDataTransactionType type,
801
                gpointer                   user_data)
802
0
{
803
0
  TrackerNotifier *notifier = user_data;
804
0
  GHashTable *events;
805
806
0
  events = get_event_cache_ht (notifier);
807
808
0
  if (type == TRACKER_DATA_COMMIT) {
809
0
    TrackerNotifierEventCache *cache;
810
0
    GHashTableIter iter;
811
812
0
    g_hash_table_iter_init (&iter, events);
813
814
0
    while (g_hash_table_iter_next (&iter, NULL, (gpointer *) &cache)) {
815
0
      g_hash_table_iter_steal (&iter);
816
0
      _tracker_notifier_event_cache_flush_events (notifier, cache);
817
0
    }
818
0
  } else {
819
0
    g_hash_table_remove_all (events);
820
0
  }
821
0
}
822
823
static void
824
detach_notifier (TrackerDirectConnection *conn,
825
                 TrackerNotifier         *notifier)
826
0
{
827
0
  TrackerDirectConnectionPrivate *priv;
828
0
  TrackerData *tracker_data;
829
830
0
  priv = tracker_direct_connection_get_instance_private (conn);
831
832
0
  priv->notifiers = g_list_remove (priv->notifiers, notifier);
833
834
0
  tracker_notifier_stop (notifier);
835
0
  tracker_data = tracker_data_manager_get_data (priv->data_manager);
836
837
0
  tracker_data_remove_callbacks (tracker_data,
838
0
                                 statement_cb,
839
0
                                 transaction_cb,
840
0
                                 notifier);
841
0
}
842
843
static void
844
weak_ref_notify (gpointer  data,
845
                 GObject  *prev_location)
846
0
{
847
0
  TrackerDirectConnection *conn = data;
848
0
  TrackerDirectConnectionPrivate *priv =
849
0
    tracker_direct_connection_get_instance_private (conn);
850
851
0
  g_mutex_lock (&priv->notifiers_mutex);
852
0
  detach_notifier (conn, (TrackerNotifier *) prev_location);
853
0
  g_mutex_unlock (&priv->notifiers_mutex);
854
0
}
855
856
static void
857
tracker_direct_connection_finalize (GObject *object)
858
14.8k
{
859
14.8k
  TrackerDirectConnectionPrivate *priv;
860
14.8k
  TrackerDirectConnection *conn;
861
862
14.8k
  conn = TRACKER_DIRECT_CONNECTION (object);
863
14.8k
  priv = tracker_direct_connection_get_instance_private (conn);
864
865
14.8k
  if (!priv->closing)
866
0
    tracker_sparql_connection_close (TRACKER_SPARQL_CONNECTION (object));
867
868
14.8k
  g_clear_object (&priv->store);
869
14.8k
  g_clear_object (&priv->ontology);
870
14.8k
  g_clear_object (&priv->namespace_manager);
871
14.8k
  g_clear_object (&priv->ontology_rdf);
872
14.8k
  g_mutex_clear (&priv->update_mutex);
873
14.8k
  g_mutex_clear (&priv->notifiers_mutex);
874
875
14.8k
  G_OBJECT_CLASS (tracker_direct_connection_parent_class)->finalize (object);
876
14.8k
}
877
878
static void
879
tracker_direct_connection_set_property (GObject      *object,
880
                                        guint         prop_id,
881
                                        const GValue *value,
882
                                        GParamSpec   *pspec)
883
74.4k
{
884
74.4k
  TrackerDirectConnectionPrivate *priv;
885
74.4k
  TrackerDirectConnection *conn;
886
887
74.4k
  conn = TRACKER_DIRECT_CONNECTION (object);
888
74.4k
  priv = tracker_direct_connection_get_instance_private (conn);
889
890
74.4k
  switch (prop_id) {
891
14.8k
  case PROP_FLAGS:
892
14.8k
    priv->flags = g_value_get_flags (value);
893
14.8k
    break;
894
14.8k
  case PROP_STORE_LOCATION:
895
14.8k
    priv->store = g_value_dup_object (value);
896
14.8k
    break;
897
14.8k
  case PROP_ONTOLOGY_LOCATION:
898
14.8k
    priv->ontology = g_value_dup_object (value);
899
14.8k
    break;
900
14.8k
  case PROP_ONTOLOGY_STREAM:
901
14.8k
    priv->ontology_rdf = g_value_dup_object (value);
902
14.8k
    break;
903
14.8k
  case PROP_ONTOLOGY_STREAM_FORMAT:
904
14.8k
    priv->ontology_rdf_format = g_value_get_enum (value);
905
14.8k
    break;
906
0
  default:
907
0
    G_OBJECT_WARN_INVALID_PROPERTY_ID (object, prop_id, pspec);
908
0
    break;
909
74.4k
  }
910
74.4k
}
911
912
static void
913
tracker_direct_connection_get_property (GObject    *object,
914
                                        guint       prop_id,
915
                                        GValue     *value,
916
                                        GParamSpec *pspec)
917
0
{
918
0
  TrackerDirectConnectionPrivate *priv;
919
0
  TrackerDirectConnection *conn;
920
921
0
  conn = TRACKER_DIRECT_CONNECTION (object);
922
0
  priv = tracker_direct_connection_get_instance_private (conn);
923
924
0
  switch (prop_id) {
925
0
  case PROP_FLAGS:
926
0
    g_value_set_flags (value, priv->flags);
927
0
    break;
928
0
  case PROP_STORE_LOCATION:
929
0
    g_value_set_object (value, priv->store);
930
0
    break;
931
0
  case PROP_ONTOLOGY_LOCATION:
932
0
    g_value_set_object (value, priv->ontology);
933
0
    break;
934
0
  default:
935
0
    G_OBJECT_WARN_INVALID_PROPERTY_ID (object, prop_id, pspec);
936
0
    break;
937
0
  }
938
0
}
939
940
static TrackerSparqlCursor *
941
tracker_direct_connection_query (TrackerSparqlConnection  *self,
942
                                 const gchar              *sparql,
943
                                 GCancellable             *cancellable,
944
                                 GError                  **error)
945
0
{
946
0
  TrackerDirectConnectionPrivate *priv;
947
0
  TrackerDirectConnection *conn;
948
0
  TrackerSparql *query;
949
0
  TrackerSparqlCursor *cursor = NULL;
950
0
  GError *inner_error = NULL;
951
952
0
  conn = TRACKER_DIRECT_CONNECTION (self);
953
0
  priv = tracker_direct_connection_get_instance_private (conn);
954
955
0
  query = tracker_sparql_new (priv->data_manager, sparql, &inner_error);
956
0
  if (query) {
957
0
    cursor = tracker_sparql_execute_cursor (query, NULL, &inner_error);
958
0
    tracker_direct_connection_update_timestamp (conn);
959
0
    g_object_unref (query);
960
0
  }
961
962
0
  if (inner_error)
963
0
    g_propagate_error (error, _translate_internal_error (inner_error));
964
965
0
  return cursor;
966
0
}
967
968
static void
969
tracker_direct_connection_query_async (TrackerSparqlConnection *self,
970
                                       const gchar             *sparql,
971
                                       GCancellable            *cancellable,
972
                                       GAsyncReadyCallback      callback,
973
                                       gpointer                 user_data)
974
0
{
975
0
  TrackerDirectConnectionPrivate *priv;
976
0
  TrackerDirectConnection *conn;
977
0
  TaskData *task_data;
978
0
  GError *error = NULL;
979
0
  GTask *task;
980
981
0
  conn = TRACKER_DIRECT_CONNECTION (self);
982
0
  priv = tracker_direct_connection_get_instance_private (conn);
983
984
0
  task_data = task_data_new (TASK_TYPE_QUERY);
985
0
  task_data->d.sparql = g_strdup (sparql);
986
987
0
  task = g_task_new (self, cancellable, callback, user_data);
988
0
  g_task_set_task_data (task, task_data,
989
0
                        (GDestroyNotify) task_data_free);
990
991
0
  if (!g_thread_pool_push (priv->select_pool, task, &error)) {
992
0
    g_task_return_error (task, _translate_internal_error (error));
993
0
    g_object_unref (task);
994
0
  }
995
0
}
996
997
static TrackerSparqlCursor *
998
tracker_direct_connection_query_finish (TrackerSparqlConnection  *self,
999
                                        GAsyncResult             *res,
1000
                                        GError                  **error)
1001
0
{
1002
0
  return g_task_propagate_pointer (G_TASK (res), error);
1003
0
}
1004
1005
static TrackerSparqlStatement *
1006
tracker_direct_connection_query_statement (TrackerSparqlConnection  *self,
1007
                                            const gchar              *query,
1008
                                            GCancellable             *cancellable,
1009
                                            GError                  **error)
1010
5.59k
{
1011
5.59k
  return TRACKER_SPARQL_STATEMENT (tracker_direct_statement_new (self, query, error));
1012
5.59k
}
1013
1014
static TrackerSparqlStatement *
1015
tracker_direct_connection_update_statement (TrackerSparqlConnection  *self,
1016
                                            const gchar              *query,
1017
                                            GCancellable             *cancellable,
1018
                                            GError                  **error)
1019
5.59k
{
1020
5.59k
  return TRACKER_SPARQL_STATEMENT (tracker_direct_statement_new_update (self, query, error));
1021
5.59k
}
1022
1023
static void
1024
tracker_direct_connection_update (TrackerSparqlConnection  *self,
1025
                                  const gchar              *sparql,
1026
                                  GCancellable             *cancellable,
1027
                                  GError                  **error)
1028
0
{
1029
0
  TrackerDirectConnectionPrivate *priv;
1030
0
  TrackerDirectConnection *conn;
1031
0
  TrackerData *data;
1032
0
  GError *inner_error = NULL;
1033
1034
0
  conn = TRACKER_DIRECT_CONNECTION (self);
1035
0
  priv = tracker_direct_connection_get_instance_private (conn);
1036
1037
0
  g_mutex_lock (&priv->update_mutex);
1038
0
  data = tracker_data_manager_get_data (priv->data_manager);
1039
0
  tracker_data_update_sparql (data, sparql, &inner_error);
1040
0
  tracker_direct_connection_update_timestamp (conn);
1041
0
  g_mutex_unlock (&priv->update_mutex);
1042
1043
0
  if (inner_error)
1044
0
    g_propagate_error (error, inner_error);
1045
0
}
1046
1047
static void
1048
tracker_direct_connection_update_async (TrackerSparqlConnection *self,
1049
                                        const gchar             *sparql,
1050
                                        GCancellable            *cancellable,
1051
                                        GAsyncReadyCallback      callback,
1052
                                        gpointer                 user_data)
1053
0
{
1054
0
  TrackerDirectConnectionPrivate *priv;
1055
0
  TrackerDirectConnection *conn;
1056
0
  TaskData *task_data;
1057
0
  GTask *task;
1058
1059
0
  conn = TRACKER_DIRECT_CONNECTION (self);
1060
0
  priv = tracker_direct_connection_get_instance_private (conn);
1061
1062
0
  task_data = task_data_new (TASK_TYPE_UPDATE);
1063
0
  task_data->d.sparql = g_strdup (sparql);
1064
1065
0
  task = g_task_new (self, cancellable, callback, user_data);
1066
0
  g_task_set_task_data (task, task_data,
1067
0
                        (GDestroyNotify) task_data_free);
1068
1069
0
  g_thread_pool_push (priv->update_thread, task, NULL);
1070
0
}
1071
1072
static void
1073
tracker_direct_connection_update_finish (TrackerSparqlConnection  *self,
1074
                                         GAsyncResult             *res,
1075
                                         GError                  **error)
1076
0
{
1077
0
  GError *inner_error = NULL;
1078
1079
0
  g_task_propagate_boolean (G_TASK (res), &inner_error);
1080
0
  if (inner_error)
1081
0
    g_propagate_error (error, _translate_internal_error (inner_error));
1082
0
}
1083
1084
static void
1085
on_batch_finished (GObject      *source,
1086
                   GAsyncResult *result,
1087
                   gpointer      user_data)
1088
0
{
1089
0
  TrackerBatch *batch = TRACKER_BATCH (source);
1090
0
  GTask *task = user_data;
1091
0
  GError *error = NULL;
1092
0
  gboolean retval;
1093
1094
0
  retval = tracker_batch_execute_finish (batch, result, &error);
1095
1096
0
  if (retval)
1097
0
    g_task_return_boolean (task, TRUE);
1098
0
  else
1099
0
    g_task_return_error (task, error);
1100
1101
0
  g_object_unref (task);
1102
0
}
1103
1104
static void
1105
tracker_direct_connection_update_array_async (TrackerSparqlConnection  *self,
1106
                                              gchar                   **updates,
1107
                                              gint                      n_updates,
1108
                                              GCancellable             *cancellable,
1109
                                              GAsyncReadyCallback       callback,
1110
                                              gpointer                  user_data)
1111
0
{
1112
0
  TrackerBatch *batch;
1113
0
  GTask *task;
1114
0
  gint i;
1115
1116
0
  batch = tracker_sparql_connection_create_batch (self);
1117
1118
0
  for (i = 0; i < n_updates; i++)
1119
0
    tracker_batch_add_sparql (batch, updates[i]);
1120
1121
0
  task = g_task_new (self, cancellable, callback, user_data);
1122
0
  tracker_batch_execute_async (batch, cancellable, on_batch_finished, task);
1123
0
  g_object_unref (batch);
1124
0
}
1125
1126
static gboolean
1127
tracker_direct_connection_update_array_finish (TrackerSparqlConnection  *self,
1128
                                               GAsyncResult             *res,
1129
                                               GError                  **error)
1130
0
{
1131
0
  GError *inner_error = NULL;
1132
0
  gboolean result;
1133
1134
0
  result = g_task_propagate_boolean (G_TASK (res), &inner_error);
1135
0
  if (inner_error)
1136
0
    g_propagate_error (error, _translate_internal_error (inner_error));
1137
1138
0
  return result;
1139
0
}
1140
1141
static GVariant *
1142
tracker_direct_connection_update_blank (TrackerSparqlConnection  *self,
1143
                                        const gchar              *sparql,
1144
                                        GCancellable             *cancellable,
1145
                                        GError                  **error)
1146
0
{
1147
0
  TrackerDirectConnectionPrivate *priv;
1148
0
  TrackerDirectConnection *conn;
1149
0
  TrackerData *data;
1150
0
  GVariant *blank_nodes;
1151
0
  GError *inner_error = NULL;
1152
1153
0
  conn = TRACKER_DIRECT_CONNECTION (self);
1154
0
  priv = tracker_direct_connection_get_instance_private (conn);
1155
1156
0
  g_mutex_lock (&priv->update_mutex);
1157
0
  data = tracker_data_manager_get_data (priv->data_manager);
1158
0
  blank_nodes = tracker_data_update_sparql_blank (data, sparql, &inner_error);
1159
0
  tracker_direct_connection_update_timestamp (conn);
1160
0
  g_mutex_unlock (&priv->update_mutex);
1161
1162
0
  if (inner_error)
1163
0
    g_propagate_error (error, _translate_internal_error (inner_error));
1164
0
  return blank_nodes;
1165
0
}
1166
1167
static void
1168
tracker_direct_connection_update_blank_async (TrackerSparqlConnection *self,
1169
                                              const gchar             *sparql,
1170
                                              GCancellable            *cancellable,
1171
                                              GAsyncReadyCallback      callback,
1172
                                              gpointer                 user_data)
1173
0
{
1174
0
  TrackerDirectConnectionPrivate *priv;
1175
0
  TrackerDirectConnection *conn;
1176
0
  TaskData *task_data;
1177
0
  GTask *task;
1178
1179
0
  conn = TRACKER_DIRECT_CONNECTION (self);
1180
0
  priv = tracker_direct_connection_get_instance_private (conn);
1181
1182
0
  task_data = task_data_new (TASK_TYPE_UPDATE_BLANK);
1183
0
  task_data->d.sparql = g_strdup (sparql);
1184
1185
0
  task = g_task_new (self, cancellable, callback, user_data);
1186
0
  g_task_set_task_data (task, task_data,
1187
0
                        (GDestroyNotify) task_data_free);
1188
1189
0
  g_thread_pool_push (priv->update_thread, task, NULL);
1190
0
}
1191
1192
static GVariant *
1193
tracker_direct_connection_update_blank_finish (TrackerSparqlConnection  *self,
1194
                                               GAsyncResult             *res,
1195
                                               GError                  **error)
1196
0
{
1197
0
  GError *inner_error = NULL;
1198
0
  GVariant *result;
1199
1200
0
  result = g_task_propagate_pointer (G_TASK (res), &inner_error);
1201
0
  if (inner_error)
1202
0
    g_propagate_error (error, _translate_internal_error (inner_error));
1203
1204
0
  return result;
1205
0
}
1206
1207
static TrackerNamespaceManager *
1208
tracker_direct_connection_get_namespace_manager (TrackerSparqlConnection *self)
1209
0
{
1210
0
  TrackerDirectConnectionPrivate *priv;
1211
1212
0
  priv = tracker_direct_connection_get_instance_private (TRACKER_DIRECT_CONNECTION (self));
1213
1214
0
  return priv->namespace_manager;
1215
0
}
1216
1217
static TrackerNotifier *
1218
tracker_direct_connection_create_notifier (TrackerSparqlConnection *self)
1219
0
{
1220
0
  TrackerDirectConnectionPrivate *priv;
1221
0
  TrackerNotifier *notifier;
1222
0
  TrackerData *tracker_data;
1223
1224
0
  priv = tracker_direct_connection_get_instance_private (TRACKER_DIRECT_CONNECTION (self));
1225
1226
0
  notifier = g_object_new (TRACKER_TYPE_NOTIFIER,
1227
0
                           "connection", self,
1228
0
         NULL);
1229
1230
0
  g_mutex_lock (&priv->notifiers_mutex);
1231
1232
0
  if (!priv->closing) {
1233
0
    tracker_data = tracker_data_manager_get_data (priv->data_manager);
1234
0
    tracker_data_add_callbacks (tracker_data,
1235
0
                                statement_cb,
1236
0
                                transaction_cb,
1237
0
                                notifier);
1238
1239
0
    g_object_weak_ref (G_OBJECT (notifier), weak_ref_notify, self);
1240
0
    priv->notifiers = g_list_prepend (priv->notifiers, notifier);
1241
0
  }
1242
1243
0
  g_mutex_unlock (&priv->notifiers_mutex);
1244
1245
0
  return notifier;
1246
0
}
1247
1248
static void
1249
tracker_direct_connection_close (TrackerSparqlConnection *self)
1250
14.8k
{
1251
14.8k
  TrackerDirectConnectionPrivate *priv;
1252
14.8k
  TrackerDirectConnection *conn;
1253
1254
14.8k
  conn = TRACKER_DIRECT_CONNECTION (self);
1255
14.8k
  priv = tracker_direct_connection_get_instance_private (conn);
1256
14.8k
  priv->closing = TRUE;
1257
1258
14.8k
  if (priv->cleanup_timeout_id) {
1259
9.05k
    g_source_remove (priv->cleanup_timeout_id);
1260
9.05k
    priv->cleanup_timeout_id = 0;
1261
9.05k
  }
1262
1263
14.8k
  if (priv->update_thread) {
1264
14.8k
    g_thread_pool_free (priv->update_thread, TRUE, TRUE);
1265
14.8k
    priv->update_thread = NULL;
1266
14.8k
  }
1267
1268
14.8k
  if (priv->select_pool) {
1269
14.8k
    g_thread_pool_free (priv->select_pool, TRUE, TRUE);
1270
14.8k
    priv->select_pool = NULL;
1271
14.8k
  }
1272
1273
14.8k
  g_mutex_lock (&priv->notifiers_mutex);
1274
1275
14.8k
  while (priv->notifiers) {
1276
0
    TrackerNotifier *notifier = priv->notifiers->data;
1277
1278
0
    g_object_weak_unref (G_OBJECT (notifier),
1279
0
                         weak_ref_notify,
1280
0
                         conn);
1281
0
    detach_notifier (conn, notifier);
1282
0
  }
1283
1284
14.8k
  g_mutex_unlock (&priv->notifiers_mutex);
1285
1286
14.8k
  if (priv->data_manager) {
1287
9.05k
    tracker_data_manager_shutdown (priv->data_manager);
1288
9.05k
    g_clear_object (&priv->data_manager);
1289
9.05k
  }
1290
14.8k
}
1291
1292
static void
1293
async_close_thread_func (GTask        *task,
1294
                         gpointer      source_object,
1295
                         gpointer      task_data,
1296
                         GCancellable *cancellable)
1297
0
{
1298
0
  if (g_task_return_error_if_cancelled (task))
1299
0
    return;
1300
1301
0
  tracker_sparql_connection_close (source_object);
1302
0
  g_task_return_boolean (task, TRUE);
1303
0
}
1304
1305
static void
1306
tracker_direct_connection_close_async (TrackerSparqlConnection *connection,
1307
                                       GCancellable            *cancellable,
1308
                                       GAsyncReadyCallback      callback,
1309
                                       gpointer                 user_data)
1310
0
{
1311
0
  GTask *task;
1312
1313
0
  task = g_task_new (connection, cancellable, callback, user_data);
1314
0
  g_task_run_in_thread (task, async_close_thread_func);
1315
0
  g_object_unref (task);
1316
0
}
1317
1318
static gboolean
1319
tracker_direct_connection_close_finish (TrackerSparqlConnection  *connection,
1320
                                        GAsyncResult             *res,
1321
                                        GError                  **error)
1322
0
{
1323
0
  return g_task_propagate_boolean (G_TASK (res), error);
1324
0
}
1325
1326
static gboolean
1327
tracker_direct_connection_update_resource (TrackerSparqlConnection  *self,
1328
                                           const gchar              *graph,
1329
                                           TrackerResource          *resource,
1330
                                           GCancellable             *cancellable,
1331
                                           GError                  **error)
1332
0
{
1333
0
  TrackerDirectConnectionPrivate *priv;
1334
0
  TrackerDirectConnection *conn;
1335
0
  TrackerData *data;
1336
0
  GError *inner_error = NULL;
1337
1338
0
  conn = TRACKER_DIRECT_CONNECTION (self);
1339
0
  priv = tracker_direct_connection_get_instance_private (conn);
1340
1341
0
  g_mutex_lock (&priv->update_mutex);
1342
0
  data = tracker_data_manager_get_data (priv->data_manager);
1343
0
  update_resource (data, graph, resource, &inner_error);
1344
0
  tracker_direct_connection_update_timestamp (conn);
1345
0
  g_mutex_unlock (&priv->update_mutex);
1346
1347
0
  if (inner_error) {
1348
0
    g_propagate_error (error, inner_error);
1349
0
    return FALSE;
1350
0
  }
1351
1352
0
  return TRUE;
1353
0
}
1354
1355
static void
1356
tracker_direct_connection_update_resource_async (TrackerSparqlConnection *self,
1357
                                                 const gchar             *graph,
1358
                                                 TrackerResource         *resource,
1359
                                                 GCancellable            *cancellable,
1360
                                                 GAsyncReadyCallback      callback,
1361
                                                 gpointer                 user_data)
1362
0
{
1363
0
  TrackerDirectConnectionPrivate *priv;
1364
0
  TrackerDirectConnection *conn;
1365
0
  TaskData *task_data;
1366
0
  GTask *task;
1367
1368
0
  conn = TRACKER_DIRECT_CONNECTION (self);
1369
0
  priv = tracker_direct_connection_get_instance_private (conn);
1370
1371
0
  task_data = task_data_new (TASK_TYPE_UPDATE_RESOURCE);
1372
0
  task_data->d.update_resource.graph = g_strdup (graph);
1373
0
  task_data->d.update_resource.resource = g_object_ref (resource);
1374
1375
0
  task = g_task_new (self, cancellable, callback, user_data);
1376
0
  g_task_set_task_data (task, task_data,
1377
0
                        (GDestroyNotify) task_data_free);
1378
1379
0
  g_thread_pool_push (priv->update_thread, task, NULL);
1380
0
}
1381
1382
static gboolean
1383
tracker_direct_connection_update_resource_finish (TrackerSparqlConnection  *connection,
1384
                                                  GAsyncResult             *res,
1385
                                                  GError                  **error)
1386
0
{
1387
0
  return g_task_propagate_boolean (G_TASK (res), error);
1388
0
}
1389
1390
static TrackerBatch *
1391
tracker_direct_connection_create_batch (TrackerSparqlConnection *connection)
1392
24.9k
{
1393
24.9k
  TrackerDirectConnectionPrivate *priv;
1394
24.9k
  TrackerDirectConnection *conn;
1395
1396
24.9k
  conn = TRACKER_DIRECT_CONNECTION (connection);
1397
24.9k
  priv = tracker_direct_connection_get_instance_private (conn);
1398
1399
24.9k
  if (priv->flags & TRACKER_SPARQL_CONNECTION_FLAGS_READONLY)
1400
0
    return NULL;
1401
1402
24.9k
  return tracker_direct_batch_new (connection);
1403
24.9k
}
1404
1405
static gboolean
1406
tracker_direct_connection_lookup_dbus_service (TrackerSparqlConnection  *connection,
1407
                                               const gchar              *dbus_name,
1408
                                               const gchar              *dbus_path,
1409
                                               gchar                   **name,
1410
                                               gchar                   **path)
1411
0
{
1412
0
  TrackerDirectConnectionPrivate *priv;
1413
0
  TrackerDirectConnection *conn;
1414
0
  TrackerSparqlConnection *remote;
1415
0
  GError *error = NULL;
1416
0
  gchar *uri;
1417
1418
0
  conn = TRACKER_DIRECT_CONNECTION (connection);
1419
0
  priv = tracker_direct_connection_get_instance_private (conn);
1420
1421
0
  uri = tracker_util_build_dbus_uri (G_BUS_TYPE_SESSION,
1422
0
                                     dbus_name, dbus_path);
1423
0
  remote = tracker_data_manager_get_remote_connection (priv->data_manager,
1424
0
                                                       uri, &error);
1425
0
  if (error) {
1426
0
    g_warning ("Error getting remote connection '%s': %s", uri, error->message);
1427
0
    g_error_free (error);
1428
0
  }
1429
1430
0
  g_free (uri);
1431
1432
0
  if (!remote)
1433
0
    return FALSE;
1434
0
  if (!g_object_class_find_property (G_OBJECT_GET_CLASS (remote), "bus-name"))
1435
0
    return FALSE;
1436
1437
0
  g_object_get (remote,
1438
0
                "bus-name", name,
1439
0
                "bus-object-path", path,
1440
0
                NULL);
1441
1442
0
  return TRUE;
1443
0
}
1444
1445
static void
1446
tracker_direct_connection_serialize_async (TrackerSparqlConnection  *self,
1447
                                           TrackerSerializeFlags     flags,
1448
                                           TrackerRdfFormat          format,
1449
                                           const gchar              *query,
1450
                                           GCancellable             *cancellable,
1451
                                           GAsyncReadyCallback      callback,
1452
                                           gpointer                 user_data)
1453
0
{
1454
0
  TrackerDirectConnectionPrivate *priv;
1455
0
  TrackerDirectConnection *conn;
1456
0
  GError *error = NULL;
1457
0
  TaskData *task_data;
1458
0
  GTask *task;
1459
1460
0
  conn = TRACKER_DIRECT_CONNECTION (self);
1461
0
  priv = tracker_direct_connection_get_instance_private (conn);
1462
1463
0
  task_data = task_data_new (TASK_TYPE_SERIALIZE);
1464
0
  task_data->d.serialize.sparql = g_strdup (query);
1465
0
  task_data->d.serialize.format = format;
1466
0
  task_data->d.serialize.flags = flags;
1467
1468
0
  task = g_task_new (self, cancellable, callback, user_data);
1469
0
  g_task_set_task_data (task, task_data,
1470
0
                        (GDestroyNotify) task_data_free);
1471
1472
0
  if (!g_thread_pool_push (priv->select_pool, task, &error)) {
1473
0
    g_task_return_error (task, _translate_internal_error (error));
1474
0
    g_object_unref (task);
1475
0
  }
1476
0
}
1477
1478
static GInputStream *
1479
tracker_direct_connection_serialize_finish (TrackerSparqlConnection  *connection,
1480
                                            GAsyncResult             *res,
1481
                                            GError                  **error)
1482
0
{
1483
0
  return g_task_propagate_pointer (G_TASK (res), error);
1484
0
}
1485
1486
static void
1487
tracker_direct_connection_deserialize_async (TrackerSparqlConnection *self,
1488
                                             TrackerDeserializeFlags  flags,
1489
                                             TrackerRdfFormat         format,
1490
                                             const gchar             *default_graph,
1491
                                             GInputStream            *stream,
1492
                                             GCancellable            *cancellable,
1493
                                             GAsyncReadyCallback      callback,
1494
                                             gpointer                 user_data)
1495
0
{
1496
0
  TrackerDirectConnectionPrivate *priv;
1497
0
  TrackerDirectConnection *conn;
1498
0
  TaskData *task_data;
1499
0
  GTask *task;
1500
1501
0
  conn = TRACKER_DIRECT_CONNECTION (self);
1502
0
  priv = tracker_direct_connection_get_instance_private (conn);
1503
1504
0
  task_data = task_data_new (TASK_TYPE_DESERIALIZE);
1505
0
  task_data->d.deserialize.stream = g_object_ref (stream);
1506
0
  task_data->d.deserialize.default_graph = g_strdup (default_graph);
1507
0
  task_data->d.deserialize.format = format;
1508
0
  task_data->d.deserialize.flags = flags;
1509
1510
0
  task = g_task_new (self, cancellable, callback, user_data);
1511
0
  g_task_set_task_data (task, task_data,
1512
0
                        (GDestroyNotify) task_data_free);
1513
1514
0
  g_thread_pool_push (priv->update_thread, task, NULL);
1515
0
}
1516
1517
static gboolean
1518
tracker_direct_connection_deserialize_finish (TrackerSparqlConnection  *connection,
1519
                                              GAsyncResult             *res,
1520
                                              GError                  **error)
1521
0
{
1522
0
  return g_task_propagate_boolean (G_TASK (res), error);
1523
0
}
1524
1525
static void
1526
tracker_direct_connection_map_connection (TrackerSparqlConnection *connection,
1527
            const gchar             *handle_name,
1528
            TrackerSparqlConnection *service_connection)
1529
0
{
1530
0
  TrackerDirectConnectionPrivate *priv;
1531
0
  TrackerDirectConnection *conn;
1532
1533
0
  conn = TRACKER_DIRECT_CONNECTION (connection);
1534
0
  priv = tracker_direct_connection_get_instance_private (conn);
1535
1536
0
  tracker_data_manager_map_connection (priv->data_manager,
1537
0
                                       handle_name,
1538
0
                                       service_connection);
1539
0
}
1540
1541
static void
1542
tracker_direct_connection_class_init (TrackerDirectConnectionClass *klass)
1543
5
{
1544
5
  TrackerSparqlConnectionClass *sparql_connection_class;
1545
5
  GObjectClass *object_class;
1546
1547
5
  object_class = G_OBJECT_CLASS (klass);
1548
5
  sparql_connection_class = TRACKER_SPARQL_CONNECTION_CLASS (klass);
1549
1550
5
  object_class->finalize = tracker_direct_connection_finalize;
1551
5
  object_class->set_property = tracker_direct_connection_set_property;
1552
5
  object_class->get_property = tracker_direct_connection_get_property;
1553
1554
5
  sparql_connection_class->query = tracker_direct_connection_query;
1555
5
  sparql_connection_class->query_async = tracker_direct_connection_query_async;
1556
5
  sparql_connection_class->query_finish = tracker_direct_connection_query_finish;
1557
5
  sparql_connection_class->query_statement = tracker_direct_connection_query_statement;
1558
5
  sparql_connection_class->update_statement = tracker_direct_connection_update_statement;
1559
5
  sparql_connection_class->update = tracker_direct_connection_update;
1560
5
  sparql_connection_class->update_async = tracker_direct_connection_update_async;
1561
5
  sparql_connection_class->update_finish = tracker_direct_connection_update_finish;
1562
5
  sparql_connection_class->update_array_async = tracker_direct_connection_update_array_async;
1563
5
  sparql_connection_class->update_array_finish = tracker_direct_connection_update_array_finish;
1564
5
  sparql_connection_class->update_blank = tracker_direct_connection_update_blank;
1565
5
  sparql_connection_class->update_blank_async = tracker_direct_connection_update_blank_async;
1566
5
  sparql_connection_class->update_blank_finish = tracker_direct_connection_update_blank_finish;
1567
5
  sparql_connection_class->get_namespace_manager = tracker_direct_connection_get_namespace_manager;
1568
5
  sparql_connection_class->create_notifier = tracker_direct_connection_create_notifier;
1569
5
  sparql_connection_class->close = tracker_direct_connection_close;
1570
5
  sparql_connection_class->close_async = tracker_direct_connection_close_async;
1571
5
  sparql_connection_class->close_finish = tracker_direct_connection_close_finish;
1572
5
  sparql_connection_class->update_resource = tracker_direct_connection_update_resource;
1573
5
  sparql_connection_class->update_resource_async = tracker_direct_connection_update_resource_async;
1574
5
  sparql_connection_class->update_resource_finish = tracker_direct_connection_update_resource_finish;
1575
5
  sparql_connection_class->create_batch = tracker_direct_connection_create_batch;
1576
5
  sparql_connection_class->lookup_dbus_service = tracker_direct_connection_lookup_dbus_service;
1577
5
  sparql_connection_class->serialize_async = tracker_direct_connection_serialize_async;
1578
5
  sparql_connection_class->serialize_finish = tracker_direct_connection_serialize_finish;
1579
5
  sparql_connection_class->deserialize_async = tracker_direct_connection_deserialize_async;
1580
5
  sparql_connection_class->deserialize_finish = tracker_direct_connection_deserialize_finish;
1581
5
  sparql_connection_class->map_connection = tracker_direct_connection_map_connection;
1582
1583
5
  props[PROP_FLAGS] =
1584
5
    g_param_spec_flags ("flags",
1585
5
                        "Flags",
1586
5
                        "Flags",
1587
5
                        TRACKER_TYPE_SPARQL_CONNECTION_FLAGS,
1588
5
                        TRACKER_SPARQL_CONNECTION_FLAGS_NONE,
1589
5
                        G_PARAM_READWRITE |
1590
5
                        G_PARAM_CONSTRUCT_ONLY);
1591
5
  props[PROP_STORE_LOCATION] =
1592
5
    g_param_spec_object ("store-location",
1593
5
                         "Store location",
1594
5
                         "Store location",
1595
5
                         G_TYPE_FILE,
1596
5
                         G_PARAM_READWRITE |
1597
5
                         G_PARAM_CONSTRUCT_ONLY);
1598
5
  props[PROP_ONTOLOGY_LOCATION] =
1599
5
    g_param_spec_object ("ontology-location",
1600
5
                         "Ontology location",
1601
5
                         "Ontology location",
1602
5
                         G_TYPE_FILE,
1603
5
                         G_PARAM_READWRITE |
1604
5
                         G_PARAM_CONSTRUCT_ONLY);
1605
5
  props[PROP_ONTOLOGY_STREAM] =
1606
5
    g_param_spec_object ("ontology-stream", NULL, NULL,
1607
5
                         G_TYPE_INPUT_STREAM,
1608
5
                         G_PARAM_WRITABLE |
1609
5
                         G_PARAM_CONSTRUCT_ONLY);
1610
5
  props[PROP_ONTOLOGY_STREAM_FORMAT] =
1611
5
    g_param_spec_enum ("ontology-stream-format", NULL, NULL,
1612
5
                       TRACKER_TYPE_RDF_FORMAT,
1613
5
                       TRACKER_RDF_FORMAT_TURTLE,
1614
5
                       G_PARAM_WRITABLE |
1615
5
                       G_PARAM_CONSTRUCT_ONLY);
1616
1617
5
  g_object_class_install_properties (object_class, N_PROPS, props);
1618
5
}
1619
1620
TrackerSparqlConnection *
1621
tracker_direct_connection_new (TrackerSparqlConnectionFlags   flags,
1622
                               GFile                         *store,
1623
                               GFile                         *ontology,
1624
                               GError                       **error)
1625
5.84k
{
1626
5.84k
  return g_initable_new (TRACKER_TYPE_DIRECT_CONNECTION,
1627
5.84k
                         NULL, error,
1628
5.84k
                         "flags", flags,
1629
5.84k
                         "store-location", store,
1630
5.84k
                         "ontology-location", ontology,
1631
5.84k
                         NULL);
1632
5.84k
}
1633
1634
void
1635
tracker_direct_connection_new_async (TrackerSparqlConnectionFlags  flags,
1636
                                     GFile                        *store,
1637
                                     GFile                        *ontology,
1638
                                     GCancellable                 *cancellable,
1639
                                     GAsyncReadyCallback           cb,
1640
                                     gpointer                      user_data)
1641
0
{
1642
0
  g_async_initable_new_async (TRACKER_TYPE_DIRECT_CONNECTION,
1643
0
                              G_PRIORITY_DEFAULT,
1644
0
                              cancellable,
1645
0
                              cb,
1646
0
                              user_data,
1647
0
                              "flags", flags,
1648
0
                              "store-location", store,
1649
0
                              "ontology-location", ontology,
1650
0
                              NULL);
1651
0
}
1652
1653
TrackerSparqlConnection *
1654
tracker_direct_connection_new_finish (GAsyncResult  *res,
1655
                                      GError       **error)
1656
0
{
1657
0
  GAsyncInitable *initable;
1658
1659
0
  initable = g_task_get_source_object (G_TASK (res));
1660
1661
0
  return TRACKER_SPARQL_CONNECTION (g_async_initable_new_finish (initable,
1662
0
                                                                 res,
1663
0
                                                                 error));
1664
0
}
1665
1666
TrackerSparqlConnection *
1667
tracker_direct_connection_new_from_rdf (TrackerSparqlConnectionFlags   flags,
1668
                                        GFile                         *store,
1669
                                        TrackerDeserializeFlags        deserialize_flags,
1670
                                        TrackerRdfFormat               rdf_format,
1671
                                        GInputStream                  *rdf_stream,
1672
                                        GCancellable                  *cancellable,
1673
                                        GError                       **error)
1674
9.05k
{
1675
9.05k
  return g_initable_new (TRACKER_TYPE_DIRECT_CONNECTION,
1676
9.05k
                         cancellable, error,
1677
9.05k
                         "flags", flags,
1678
9.05k
                         "store-location", store,
1679
9.05k
                         "ontology-stream-format", rdf_format,
1680
9.05k
                         "ontology-stream", rdf_stream,
1681
9.05k
                         NULL);
1682
9.05k
}
1683
1684
void
1685
tracker_direct_connection_new_from_rdf_async (TrackerSparqlConnectionFlags  flags,
1686
                                              GFile                        *store,
1687
                                              TrackerDeserializeFlags       deserialize_flags,
1688
                                              TrackerRdfFormat              rdf_format,
1689
                                              GInputStream                 *rdf_stream,
1690
                                              GCancellable                 *cancellable,
1691
                                              GAsyncReadyCallback           cb,
1692
                                              gpointer                      user_data)
1693
0
{
1694
0
  g_async_initable_new_async (TRACKER_TYPE_DIRECT_CONNECTION,
1695
0
                              G_PRIORITY_DEFAULT,
1696
0
                              cancellable,
1697
0
                              cb,
1698
0
                              user_data,
1699
0
                              "flags", flags,
1700
0
                              "store-location", store,
1701
0
                              "ontology-stream-format", rdf_format,
1702
0
                              "ontology-stream", rdf_stream,
1703
0
                              NULL);
1704
0
}
1705
1706
TrackerSparqlConnection *
1707
tracker_direct_connection_new_from_rdf_finish (GAsyncResult  *res,
1708
                                               GError       **error)
1709
0
{
1710
0
  GAsyncInitable *initable;
1711
1712
0
  initable = g_task_get_source_object (G_TASK (res));
1713
1714
0
  return TRACKER_SPARQL_CONNECTION (g_async_initable_new_finish (initable,
1715
0
                                                                 res,
1716
0
                                                                 error));
1717
0
}
1718
1719
TrackerDataManager *
1720
tracker_direct_connection_get_data_manager (TrackerDirectConnection *conn)
1721
11.1k
{
1722
11.1k
  TrackerDirectConnectionPrivate *priv;
1723
1724
11.1k
  priv = tracker_direct_connection_get_instance_private (conn);
1725
11.1k
  return priv->data_manager;
1726
11.1k
}
1727
1728
void
1729
tracker_direct_connection_update_timestamp (TrackerDirectConnection *conn)
1730
24.9k
{
1731
24.9k
  TrackerDirectConnectionPrivate *priv;
1732
1733
24.9k
  priv = tracker_direct_connection_get_instance_private (conn);
1734
24.9k
  priv->timestamp = g_get_monotonic_time ();
1735
24.9k
}
1736
1737
gboolean
1738
tracker_direct_connection_update_batch (TrackerDirectConnection  *conn,
1739
                                        TrackerBatch             *batch,
1740
                                        GError                  **error)
1741
24.9k
{
1742
24.9k
  TrackerDirectConnectionPrivate *priv;
1743
24.9k
  GError *inner_error = NULL;
1744
1745
24.9k
  priv = tracker_direct_connection_get_instance_private (conn);
1746
1747
24.9k
  g_mutex_lock (&priv->update_mutex);
1748
24.9k
  tracker_direct_batch_update (TRACKER_DIRECT_BATCH (batch),
1749
24.9k
                               priv->data_manager, &inner_error);
1750
24.9k
  tracker_direct_connection_update_timestamp (conn);
1751
24.9k
  g_mutex_unlock (&priv->update_mutex);
1752
1753
24.9k
  if (inner_error) {
1754
19.1k
    g_propagate_error (error, inner_error);
1755
19.1k
    return FALSE;
1756
19.1k
  }
1757
1758
5.83k
  return TRUE;
1759
24.9k
}
1760
1761
void
1762
tracker_direct_connection_update_batch_async (TrackerDirectConnection  *conn,
1763
                                              TrackerBatch             *batch,
1764
                                              GCancellable             *cancellable,
1765
                                              GAsyncReadyCallback       callback,
1766
                                              gpointer                  user_data)
1767
0
{
1768
0
  TrackerDirectConnectionPrivate *priv;
1769
0
  TaskData *task_data;
1770
0
  GTask *task;
1771
1772
0
  priv = tracker_direct_connection_get_instance_private (conn);
1773
1774
0
  task_data = task_data_new (TASK_TYPE_UPDATE_BATCH);
1775
0
  task_data->d.batch = g_object_ref (batch);
1776
1777
0
  task = g_task_new (batch, cancellable, callback, user_data);
1778
0
  g_task_set_task_data (task, task_data,
1779
0
                        (GDestroyNotify) task_data_free);
1780
1781
0
  g_thread_pool_push (priv->update_thread, task, NULL);
1782
0
}
1783
1784
gboolean
1785
tracker_direct_connection_update_batch_finish (TrackerDirectConnection  *conn,
1786
                                               GAsyncResult             *res,
1787
                                               GError                  **error)
1788
0
{
1789
0
  GError *inner_error = NULL;
1790
1791
0
  g_task_propagate_boolean (G_TASK (res), &inner_error);
1792
0
  if (inner_error) {
1793
0
    g_propagate_error (error, _translate_internal_error (inner_error));
1794
0
    return FALSE;
1795
0
  }
1796
1797
0
  return TRUE;
1798
0
}
1799
1800
gboolean
1801
tracker_direct_connection_execute_update_statement (TrackerDirectConnection  *conn,
1802
                                                    TrackerSparqlStatement   *stmt,
1803
                                                    GHashTable               *parameters,
1804
                                                    GError                  **error)
1805
0
{
1806
0
  TrackerDirectConnectionPrivate *priv;
1807
0
  TrackerData *tracker_data;
1808
0
  GError *inner_error = NULL;
1809
1810
0
  priv = tracker_direct_connection_get_instance_private (conn);
1811
1812
0
  g_mutex_lock (&priv->update_mutex);
1813
1814
0
  tracker_data = tracker_data_manager_get_data (priv->data_manager);
1815
0
  if (!tracker_data_begin_transaction (tracker_data, &inner_error))
1816
0
    goto out;
1817
1818
0
  if (tracker_direct_statement_execute_update (stmt, parameters, NULL, &inner_error)) {
1819
0
    if (tracker_data_commit_transaction (tracker_data, &inner_error))
1820
0
      tracker_direct_connection_update_timestamp (conn);
1821
0
  } else {
1822
0
    tracker_data_rollback_transaction (tracker_data);
1823
0
  }
1824
1825
0
 out:
1826
0
  g_mutex_unlock (&priv->update_mutex);
1827
1828
0
  if (inner_error) {
1829
0
    g_propagate_error (error, inner_error);
1830
0
    return FALSE;
1831
0
  }
1832
1833
0
  return TRUE;
1834
0
}
1835
1836
void
1837
tracker_direct_connection_execute_query_statement_async (TrackerDirectConnection *conn,
1838
                                                         TrackerSparqlStatement  *stmt,
1839
                                                         GHashTable              *parameters,
1840
                                                         GCancellable            *cancellable,
1841
                                                         GAsyncReadyCallback      callback,
1842
                                                         gpointer                 user_data)
1843
0
{
1844
0
  TrackerDirectConnectionPrivate *priv =
1845
0
    tracker_direct_connection_get_instance_private (conn);
1846
0
  GError *error = NULL;
1847
0
  TaskData *task_data;
1848
0
  GTask *task;
1849
1850
0
  task_data = task_data_new (TASK_TYPE_QUERY_STATEMENT);
1851
0
  task_data->d.statement.stmt = g_object_ref (stmt);
1852
0
  task_data->d.statement.parameters =
1853
0
    parameters ? g_hash_table_ref (parameters) : NULL;
1854
1855
0
  task = g_task_new (conn, cancellable, callback, user_data);
1856
0
  g_task_set_task_data (task, task_data,
1857
0
                        (GDestroyNotify) task_data_free);
1858
1859
0
  if (!g_thread_pool_push (priv->select_pool, task, &error)) {
1860
0
    g_task_return_error (task, _translate_internal_error (error));
1861
0
    g_object_unref (task);
1862
0
  }
1863
0
}
1864
1865
TrackerSparqlCursor *
1866
tracker_direct_connection_execute_query_statement_finish (TrackerDirectConnection  *conn,
1867
                                                          GAsyncResult             *res,
1868
                                                          GError                  **error)
1869
0
{
1870
0
  return g_task_propagate_pointer (G_TASK (res), error);
1871
0
}
1872
1873
void
1874
tracker_direct_connection_execute_serialize_statement_async (TrackerDirectConnection *conn,
1875
                                                             TrackerSparqlStatement  *stmt,
1876
                                                             GHashTable              *parameters,
1877
                                                             TrackerSerializeFlags    flags,
1878
                                                             TrackerRdfFormat         format,
1879
                                                             GCancellable            *cancellable,
1880
                                                             GAsyncReadyCallback      callback,
1881
                                                             gpointer                 user_data)
1882
0
{
1883
0
  TrackerDirectConnectionPrivate *priv =
1884
0
    tracker_direct_connection_get_instance_private (conn);
1885
0
  GError *error = NULL;
1886
0
  TaskData *task_data;
1887
0
  GTask *task;
1888
1889
0
  task_data = task_data_new (TASK_TYPE_SERIALIZE_STATEMENT);
1890
0
  task_data->d.serialize_statement.stmt = g_object_ref (stmt);
1891
0
  task_data->d.serialize_statement.parameters =
1892
0
    parameters ? g_hash_table_ref (parameters) : NULL;
1893
0
  task_data->d.serialize_statement.flags = flags;
1894
0
  task_data->d.serialize_statement.format = format;
1895
1896
0
  task = g_task_new (conn, cancellable, callback, user_data);
1897
0
  g_task_set_task_data (task, task_data,
1898
0
                        (GDestroyNotify) task_data_free);
1899
1900
0
  if (!g_thread_pool_push (priv->select_pool, task, &error)) {
1901
0
    g_task_return_error (task, _translate_internal_error (error));
1902
0
    g_object_unref (task);
1903
0
  }
1904
0
}
1905
1906
GInputStream *
1907
tracker_direct_connection_execute_serialize_statement_finish (TrackerDirectConnection  *conn,
1908
                                                              GAsyncResult             *res,
1909
                                                              GError                  **error)
1910
0
{
1911
0
  return g_task_propagate_pointer (G_TASK (res), error);
1912
0
}
1913
1914
1915
void
1916
tracker_direct_connection_execute_update_statement_async (TrackerDirectConnection  *conn,
1917
                                                          TrackerSparqlStatement   *stmt,
1918
                                                          GHashTable               *parameters,
1919
                                                          GCancellable             *cancellable,
1920
                                                          GAsyncReadyCallback       callback,
1921
                                                          gpointer                  user_data)
1922
0
{
1923
0
  TrackerDirectConnectionPrivate *priv;
1924
0
  TaskData *task_data;
1925
0
  GTask *task;
1926
1927
0
  priv = tracker_direct_connection_get_instance_private (conn);
1928
1929
0
  task_data = task_data_new (TASK_TYPE_UPDATE_STATEMENT);
1930
0
  task_data->d.statement.stmt = g_object_ref (stmt);
1931
0
  task_data->d.statement.parameters =
1932
0
    parameters ? g_hash_table_ref (parameters) : NULL;
1933
1934
0
  task = g_task_new (stmt, cancellable, callback, user_data);
1935
0
  g_task_set_task_data (task, task_data,
1936
0
                        (GDestroyNotify) task_data_free);
1937
1938
0
  g_thread_pool_push (priv->update_thread, task, NULL);
1939
0
}
1940
1941
gboolean
1942
tracker_direct_connection_execute_update_statement_finish (TrackerDirectConnection  *conn,
1943
                                                           GAsyncResult             *res,
1944
                                                           GError                  **error)
1945
0
{
1946
0
  GError *inner_error = NULL;
1947
1948
0
  g_task_propagate_boolean (G_TASK (res), &inner_error);
1949
0
  if (inner_error) {
1950
0
    g_propagate_error (error, _translate_internal_error (inner_error));
1951
0
    return FALSE;
1952
0
  }
1953
1954
0
  return TRUE;
1955
0
}