Skip to content

Commit e958aef

Browse files
committed
Use Dapper as internal query engine
1 parent 19574b6 commit e958aef

6 files changed

Lines changed: 128 additions & 156 deletions

File tree

PowerSync/PowerSync.Common/Client/PowerSyncDatabase.cs

Lines changed: 17 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -422,7 +422,8 @@ await Database.WriteTransaction(async tx =>
422422
Closed = true;
423423
}
424424

425-
private record UploadQueueStatsResult(int size, int count);
425+
private record UploadQueueStatsSizeCountResult(long size, long count);
426+
private record UploadQueueStatsCountResult(long count);
426427
/// <summary>
427428
/// Get upload queue size estimate and count.
428429
/// </summary>
@@ -432,23 +433,22 @@ public async Task<UploadQueueStats> GetUploadQueueStats(bool includeSize = false
432433
{
433434
if (includeSize)
434435
{
435-
var result = await tx.Get<UploadQueueStatsResult>(
436+
var result = await tx.Get<UploadQueueStatsSizeCountResult>(
436437
$"SELECT SUM(cast(data as blob) + 20) as size, count(*) as count FROM {PSInternalTable.CRUD}"
437438
);
438439

439440
return new UploadQueueStats(result.count, result.size);
440441
}
441442
else
442443
{
443-
var result = await tx.Get<UploadQueueStatsResult>(
444+
var result = await tx.Get<UploadQueueStatsCountResult>(
444445
$"SELECT count(*) as count FROM {PSInternalTable.CRUD}"
445446
);
446447
return new UploadQueueStats(result.count);
447448
}
448449
});
449450
}
450451

451-
452452
/// <summary>
453453
/// Get a batch of crud data to upload.
454454
/// <para />
@@ -468,7 +468,7 @@ public async Task<UploadQueueStats> GetUploadQueueStats(bool includeSize = false
468468
/// </summary>
469469
public async Task<CrudBatch?> GetCrudBatch(int limit = 100)
470470
{
471-
var crudResult = await GetAll<CrudEntryJSON>($"SELECT id, tx_id, data FROM {PSInternalTable.CRUD} ORDER BY id ASC LIMIT ?", [limit + 1]);
471+
var crudResult = await GetAll<CrudEntryJSON>($"SELECT id, tx_id AS TransactionId, data FROM {PSInternalTable.CRUD} ORDER BY id ASC LIMIT ?", [limit + 1]);
472472

473473
var all = crudResult.Select(CrudEntry.FromRow).ToList();
474474

@@ -510,7 +510,7 @@ public async Task<UploadQueueStats> GetUploadQueueStats(bool includeSize = false
510510
return await ReadTransaction(async tx =>
511511
{
512512
var first = await tx.GetOptional<CrudEntryJSON>(
513-
$"SELECT id, tx_id, data FROM {PSInternalTable.CRUD} ORDER BY id ASC LIMIT 1");
513+
$"SELECT id, tx_id AS TransactionId, data FROM {PSInternalTable.CRUD} ORDER BY id ASC LIMIT 1");
514514

515515
if (first == null)
516516
{
@@ -528,7 +528,7 @@ public async Task<UploadQueueStats> GetUploadQueueStats(bool includeSize = false
528528
else
529529
{
530530
var result = await tx.GetAll<CrudEntryJSON>(
531-
$"SELECT id, tx_id, data FROM {PSInternalTable.CRUD} WHERE tx_id = ? ORDER BY id ASC",
531+
$"SELECT id, tx_id AS TransactionId, data FROM {PSInternalTable.CRUD} WHERE tx_id = ? ORDER BY id ASC",
532532
[txId]);
533533

534534
all = result.Select(CrudEntry.FromRow).ToList();
@@ -691,7 +691,16 @@ public Task Watch<T>(string query, object?[]? parameters, WatchHandler<T> handle
691691
return tcs.Task;
692692
}
693693

694-
private record ExplainedResult(string opcode, int p2, int p3);
694+
private class ExplainedResult
695+
{
696+
public int addr;
697+
public string opcode;
698+
public int p1;
699+
public int p2;
700+
public int p3;
701+
public string p4;
702+
public int p5;
703+
}
695704
private record TableSelectResult(string tbl_name);
696705
public async Task<string[]> ResolveTables(string sql, object?[]? parameters = null, SQLWatchOptions? options = null)
697706
{
@@ -718,7 +727,6 @@ public async Task<string[]> ResolveTables(string sql, object?[]? parameters = nu
718727
resolvedTables.Add(POWERSYNC_TABLE_MATCH.Replace(table.tbl_name, ""));
719728
}
720729
}
721-
722730
return [.. resolvedTables];
723731
}
724732

PowerSync/PowerSync.Common/DB/Crud/UploadQueueStatus.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,8 @@
11
namespace PowerSync.Common.DB.Crud;
22

3-
public class UploadQueueStats(int count, long? size = null)
3+
public class UploadQueueStats(long count, long? size = null)
44
{
5-
public int Count { get; set; } = count;
5+
public long Count { get; set; } = count;
66

77
public long? Size { get; set; } = size;
88

PowerSync/PowerSync.Common/MDSQLite/MDSQLiteConnection.cs

Lines changed: 80 additions & 125 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,11 @@
11
namespace PowerSync.Common.MDSQLite;
22

3+
using System.Data;
34
using System.Threading.Tasks;
45

5-
using Microsoft.Data.Sqlite;
6+
using Dapper;
67

7-
using Newtonsoft.Json;
8+
using Microsoft.Data.Sqlite;
89

910
using PowerSync.Common.DB;
1011
using PowerSync.Common.Utils;
@@ -72,21 +73,14 @@ public void FlushUpdates()
7273
Emit(new DBAdapterEvent { TablesUpdated = batchedUpdate });
7374
}
7475

75-
/// <summary>
76-
/// Replaces ? placeholders with named parameters and sets up the command.
77-
/// Returns the parameter names for reference.
78-
/// </summary>
79-
private static List<string> PrepareCommandParameters(SqliteCommand command, string query, int parameterCount)
76+
private static List<string> PrepareQueryString(ref string query, int parameterCount)
8077
{
81-
var parameterNames = new List<string>();
82-
78+
var parameterList = new List<string>();
8379
if (parameterCount == 0)
8480
{
85-
command.CommandText = query;
86-
return parameterNames;
81+
return parameterList;
8782
}
8883

89-
// Count placeholders
9084
int placeholderCount = query.Count(c => c == '?');
9185
if (placeholderCount != parameterCount)
9286
{
@@ -97,173 +91,134 @@ private static List<string> PrepareCommandParameters(SqliteCommand command, stri
9791
for (int i = 0; i < parameterCount; i++)
9892
{
9993
string paramName = $"@param{i}";
100-
parameterNames.Add(paramName);
94+
parameterList.Add(paramName);
10195

10296
int index = query.IndexOf('?');
10397
if (index == -1)
10498
{
10599
throw new ArgumentException("Mismatch between placeholders and parameters.");
106100
}
107101

102+
// TODO inefficient, but maybe not noticeably so
108103
query = string.Concat(query.Substring(0, index), paramName, query.Substring(index + 1));
109-
110104
}
111105

112-
command.CommandText = query;
113-
114-
// Create empty parameter objects
115-
foreach (var paramName in parameterNames)
116-
{
117-
var parameter = command.CreateParameter();
118-
parameter.ParameterName = paramName;
119-
command.Parameters.Add(parameter);
120-
}
121-
122-
return parameterNames;
106+
return parameterList;
123107
}
124108

125-
private static void PrepareCommand(SqliteCommand command, string query, object?[]? parameters)
109+
private static DynamicParameters? PrepareQuery(ref string query, object?[]? parameters)
126110
{
127-
int paramCount = parameters?.Length ?? 0;
128-
PrepareCommandParameters(command, query, paramCount);
129-
130-
// Set the values
131-
if (parameters != null)
111+
if (parameters is null)
132112
{
133-
for (int i = 0; i < parameters.Length; i++)
134-
{
135-
command.Parameters[i].Value = parameters[i] ?? DBNull.Value;
136-
}
113+
return null;
137114
}
138-
}
139115

116+
int parameterCount = parameters.Length;
117+
if (parameterCount == 0)
118+
{
119+
return null;
120+
}
140121

141-
public async Task<NonQueryResult> Execute(string query, object?[]? parameters = null)
142-
{
143-
using var command = Db.CreateCommand();
144-
PrepareCommand(command, query, parameters);
122+
var parameterNames = PrepareQueryString(ref query, parameterCount);
145123

146-
int rowsAffected = await command.ExecuteNonQueryAsync();
124+
var dynamicParams = new DynamicParameters();
147125

148-
return new NonQueryResult
126+
for (int i = 0; i < parameterCount; i++)
149127
{
150-
InsertId = raw.sqlite3_last_insert_rowid(Db.Handle),
151-
RowsAffected = rowsAffected
152-
};
128+
dynamicParams.Add(parameterNames[i], parameters[i]);
129+
}
130+
131+
return dynamicParams;
153132
}
154133

155-
public async Task<NonQueryResult> ExecuteBatch(string query, object?[][]? parameters = null)
134+
private static List<DynamicParameters>? PrepareQuery(ref string query, object?[][]? parameters)
156135
{
157-
parameters ??= [];
158-
159-
if (parameters.Length == 0)
136+
if (parameters is null || parameters.Length == 0)
160137
{
161-
return new NonQueryResult { RowsAffected = 0 };
138+
return null;
162139
}
140+
var parameterCount = parameters[0].Length;
141+
var parameterNames = PrepareQueryString(ref query, parameterCount);
163142

164-
int totalRowsAffected = 0;
165-
166-
var command = Db.CreateCommand();
143+
var dynamicParamsList = new List<DynamicParameters>();
167144

168-
// Prepare command once with parameter placeholders
169-
int paramCount = parameters[0]?.Length ?? 0;
170-
PrepareCommandParameters(command, query, paramCount);
171-
172-
// Execute for each parameter set (reuses compiled statement)
173145
foreach (var paramSet in parameters)
174146
{
175-
if (paramSet != null)
147+
if (paramSet.Length != parameterCount)
176148
{
177-
for (int i = 0; i < paramSet.Length; i++)
178-
{
179-
command.Parameters[i].Value = paramSet[i] ?? DBNull.Value;
180-
}
149+
throw new ArgumentException("Parameter sets have different number of arguments.");
181150
}
182151

183-
totalRowsAffected += await command.ExecuteNonQueryAsync();
152+
var dynamicParams = new DynamicParameters();
153+
for (int i = 0; i < parameterCount; i++)
154+
{
155+
dynamicParams.Add(parameterNames[i], paramSet[i]);
156+
}
157+
dynamicParamsList.Add(dynamicParams);
184158
}
185159

186-
return new NonQueryResult
187-
{
188-
RowsAffected = totalRowsAffected,
189-
InsertId = raw.sqlite3_last_insert_rowid(Db.Handle)
190-
};
160+
return dynamicParamsList;
191161
}
192162

193-
public async Task<QueryResult> ExecuteQuery(string query, object?[]? parameters = null)
163+
public async Task<T?> GetOptional<T>(string query, object?[]? parameters = null)
194164
{
195-
var result = new QueryResult();
196-
using var command = Db.CreateCommand();
197-
PrepareCommand(command, query, parameters);
198-
199-
var rows = new List<Dictionary<string, object>>();
200-
201-
using var reader = await command.ExecuteReaderAsync();
202-
203-
while (await reader.ReadAsync())
204-
{
205-
var row = new Dictionary<string, object>();
206-
for (int i = 0; i < reader.FieldCount; i++)
207-
{
208-
row[reader.GetName(i)] = reader.IsDBNull(i) ? null! : reader.GetValue(i);
209-
}
210-
rows.Add(row);
211-
}
212-
213-
result.Rows.Array = rows;
214-
return result;
165+
DynamicParameters? dynamicParams = PrepareQuery(ref query, parameters);
166+
return dynamicParams == null
167+
? await Db.QueryFirstOrDefaultAsync<T>(query, commandType: CommandType.Text)
168+
: await Db.QueryFirstOrDefaultAsync<T>(query, dynamicParams, commandType: CommandType.Text);
215169
}
216170

217-
public async Task<T[]> GetAll<T>(string sql, object?[]? parameters = null)
171+
public async Task<T> Get<T>(string query, object?[]? parameters = null)
218172
{
219-
var result = await ExecuteQuery(sql, parameters);
173+
DynamicParameters? dynamicParams = PrepareQuery(ref query, parameters);
174+
return dynamicParams == null
175+
? await Db.QueryFirstAsync<T>(query, commandType: CommandType.Text)
176+
: await Db.QueryFirstAsync<T>(query, dynamicParams, commandType: CommandType.Text);
177+
}
220178

221-
// If there are no rows, return an empty array.
222-
if (result.Rows.Array.Count == 0)
223-
{
224-
return [];
225-
}
179+
public async Task<T[]> GetAll<T>(string query, object?[]? parameters = null)
180+
{
181+
DynamicParameters? dynamicParams = PrepareQuery(ref query, parameters);
182+
return [..dynamicParams == null
183+
? await Db.QueryAsync<T>(query, commandType: CommandType.Text)
184+
: await Db.QueryAsync<T>(query, dynamicParams, commandType: CommandType.Text)];
185+
}
226186

227-
var items = new List<T>();
187+
public async Task<NonQueryResult> Execute(string query, object?[]? parameters = null)
188+
{
189+
DynamicParameters? dynamicParams = PrepareQuery(ref query, parameters);
190+
int rowsAffected = dynamicParams == null
191+
? await Db.ExecuteAsync(query, commandType: CommandType.Text)
192+
: await Db.ExecuteAsync(query, dynamicParams, commandType: CommandType.Text);
228193

229-
foreach (var row in result.Rows.Array)
194+
return new NonQueryResult
230195
{
231-
if (row != null)
232-
{
233-
// Serialize the row to JSON and then deserialize it into type T.
234-
string json = JsonConvert.SerializeObject(row);
235-
T item = JsonConvert.DeserializeObject<T>(json)!;
236-
items.Add(item);
237-
}
238-
}
239-
240-
return [.. items];
196+
InsertId = raw.sqlite3_last_insert_rowid(Db.Handle),
197+
RowsAffected = rowsAffected,
198+
};
241199
}
242200

243-
public async Task<T?> GetOptional<T>(string sql, object?[]? parameters = null)
201+
public async Task<NonQueryResult> ExecuteBatch(string query, object?[][]? parameters = null)
244202
{
245-
var result = await ExecuteQuery(sql, parameters);
246-
247-
// If there are no rows, return null
248-
if (result.Rows.Array.Count == 0)
203+
if (parameters is null || parameters.Length == 0)
249204
{
250-
return default;
205+
return new NonQueryResult { RowsAffected = 0 };
251206
}
252207

253-
var firstRow = result.Rows.Array[0];
254-
255-
if (firstRow == null)
208+
List<DynamicParameters>? dynamicParamsList = PrepareQuery(ref query, parameters);
209+
if (dynamicParamsList is null)
256210
{
257-
return default;
211+
// Should be unreachable but you never know
212+
return new NonQueryResult { RowsAffected = 0 };
258213
}
259214

260-
string json = JsonConvert.SerializeObject(firstRow);
261-
return JsonConvert.DeserializeObject<T>(json);
262-
}
215+
int rowsAffected = await Db.ExecuteAsync(query, dynamicParamsList, commandType: CommandType.Text);
263216

264-
public async Task<T> Get<T>(string sql, object?[]? parameters = null)
265-
{
266-
return await GetOptional<T>(sql, parameters) ?? throw new InvalidOperationException("Result set is empty");
217+
return new NonQueryResult
218+
{
219+
InsertId = raw.sqlite3_last_insert_rowid(Db.Handle),
220+
RowsAffected = rowsAffected,
221+
};
267222
}
268223

269224
public new void Close()

0 commit comments

Comments
 (0)