-
Notifications
You must be signed in to change notification settings - Fork 23
Expand file tree
/
Copy pathBucketStorageImpl.kt
More file actions
147 lines (119 loc) · 5.01 KB
/
Copy pathBucketStorageImpl.kt
File metadata and controls
147 lines (119 loc) · 5.01 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
package com.powersync.bucket
import co.touchlab.kermit.Logger
import co.touchlab.stately.concurrency.AtomicBoolean
import com.powersync.db.SqlCursor
import com.powersync.db.crud.CrudEntry
import com.powersync.db.crud.CrudRow
import com.powersync.db.internal.InternalDatabase
import com.powersync.db.internal.InternalTable
import com.powersync.db.internal.PowerSyncTransaction
import com.powersync.sync.Instruction
import com.powersync.utils.JsonUtil
internal class BucketStorageImpl(
private val db: InternalDatabase,
private val logger: Logger,
) : BucketStorage {
private var hasCompletedSync = AtomicBoolean(false)
companion object {
const val MAX_OP_ID = "9223372036854775807"
}
override fun getMaxOpId(): String = MAX_OP_ID
override suspend fun getClientId(): String {
val id =
db.getOptional("SELECT powersync_client_id() as client_id") {
it.getString(0)!!
}
return id ?: throw IllegalStateException("Client ID not found")
}
override suspend fun nextCrudItem(): CrudEntry? = db.getOptional(sql = nextCrudQuery, mapper = ::mapCrudEntry)
override suspend fun nextCrudItem(transaction: PowerSyncTransaction): CrudEntry? =
transaction.getOptionalAsync(sql = nextCrudQuery, mapper = ::mapCrudEntry)
private val nextCrudQuery = "SELECT id, tx_id, data FROM ${InternalTable.CRUD} ORDER BY id ASC LIMIT 1"
override fun mapCrudEntry(row: SqlCursor): CrudEntry =
CrudEntry.fromRow(
CrudRow(
id = row.getString(0)!!,
txId = row.getString(1)?.toInt(),
data = row.getString(2)!!,
),
)
override suspend fun hasCrud(): Boolean {
val res = db.getOptional(sql = hasCrudQuery, mapper = hasCrudMapper)
return res == 1L
}
override suspend fun hasCrud(transaction: PowerSyncTransaction): Boolean {
val res = transaction.getOptionalAsync(sql = hasCrudQuery, mapper = hasCrudMapper)
return res == 1L
}
private val hasCrudQuery = "SELECT 1 FROM ps_crud LIMIT 1"
private val hasCrudMapper: (SqlCursor) -> Long = {
it.getLong(0)!!
}
override suspend fun updateLocalTarget(checkpointCallback: suspend () -> String): Boolean {
db.getOptional(
"SELECT target_op FROM ${InternalTable.BUCKETS} WHERE name = '\$local' AND target_op = ?",
parameters = listOf(MAX_OP_ID),
mapper = { cursor -> cursor.getLong(0)!! },
)
?: // Nothing to update
return false
val seqBefore =
db.getOptional("SELECT seq FROM main.sqlite_sequence WHERE name = '${InternalTable.CRUD}'") {
it.getLong(0)!!
} ?: // Nothing to update
return false
val opId = checkpointCallback()
logger.i { "[updateLocalTarget] Updating target to checkpoint $opId" }
return db.writeTransactionAsync { tx ->
if (hasCrud(tx)) {
logger.w { "[updateLocalTarget] ps crud is not empty" }
return@writeTransactionAsync false
}
val seqAfter =
tx.getOptionalAsync("SELECT seq FROM main.sqlite_sequence WHERE name = '${InternalTable.CRUD}'") {
it.getLong(0)!!
}
?: // assert isNotEmpty
throw AssertionError("Sqlite Sequence should not be empty")
if (seqAfter != seqBefore) {
logger.d("seqAfter != seqBefore seqAfter: $seqAfter seqBefore: $seqBefore")
// New crud data may have been uploaded since we got the checkpoint. Abort.
return@writeTransactionAsync false
}
tx.executeAsync(
"UPDATE ${InternalTable.BUCKETS} SET target_op = CAST(? as INTEGER) WHERE name='\$local'",
listOf(opId),
)
return@writeTransactionAsync true
}
}
override suspend fun hasCompletedSync(): Boolean {
if (hasCompletedSync.value) {
return true
}
val completedSync =
db.getOptional(
"SELECT powersync_last_synced_at()",
mapper = { cursor ->
cursor.getString(0)!!
},
)
return if (completedSync != null) {
hasCompletedSync.value = true
true
} else {
false
}
}
private fun handleControlResult(cursor: SqlCursor): List<Instruction> {
val result = cursor.getString(0)!!
logger.v { "control result: $result" }
return JsonUtil.json.decodeFromString<List<Instruction>>(result)
}
override suspend fun control(args: PowerSyncControlArguments): List<Instruction> =
db.writeTransactionAsync { tx ->
logger.v { "powersync_control: $args" }
val (op: String, data: Any?) = args.sqlArguments
tx.getAsync("SELECT powersync_control(?, ?) AS r", listOf(op, data), ::handleControlResult)
}
}