Skip to content

Commit 504c1c1

Browse files
Merge pull request #84 from SourcePointUSA/refactor_repository_settings
improve concurrency safety of `Repository` and `Settings`
2 parents 6cdd1bf + c62017f commit 504c1c1

5 files changed

Lines changed: 452 additions & 42 deletions

File tree

‎core/src/commonMain/kotlin/com/sourcepoint/mobile_core/storage/Repository.kt‎

Lines changed: 14 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -16,51 +16,31 @@ class Repository(private val storage: Settings) {
1616
}
1717

1818
var tcData: IABData
19-
get() = storage.withLock {
20-
keys
21-
.filter { it.startsWith(TCF_PREFIX) }
22-
.associateWith { this[it]!! }
23-
}
24-
set(value) {
25-
storage.withLock {
26-
removeKeysStartingWith(prefix = TCF_PREFIX)
27-
value.entries.forEach { this[it.key] = it.value }
28-
}
29-
}
19+
get() = storage.getKeysWithPrefix(TCF_PREFIX)
20+
set(value) = storage.replaceKeysWithPrefix(TCF_PREFIX, value)
3021

3122
var gppData: IABData
32-
get() = storage.withLock {
33-
keys
34-
.filter { it.startsWith(GPP_PREFIX) }
35-
.associateWith { this[it]!! }
36-
}
37-
set(value) {
38-
storage.withLock {
39-
removeKeysStartingWith(prefix = GPP_PREFIX)
40-
value.entries.forEach { this[it.key] = it.value }
41-
}
42-
}
23+
get() = storage.getKeysWithPrefix(GPP_PREFIX)
24+
set(value) = storage.replaceKeysWithPrefix(GPP_PREFIX, value)
4325

4426
var uspString: String?
45-
get() = storage.withLock { this[USPSTRING_KEY] }
46-
set(value) { storage.withLock { this[USPSTRING_KEY] = value } }
27+
get() = storage.readStringOrNull(USPSTRING_KEY)
28+
set(value) = storage.writeString(USPSTRING_KEY, value)
4729

4830
var state: State?
4931
get() = runCatching {
50-
storage.withLock {
51-
Json.decodeFromString<State>(getString(SP_STATE_KEY, defaultValue = ""))
52-
}
32+
val json = storage.readString(SP_STATE_KEY, defaultValue = "")
33+
if (json.isBlank()) null else Json.decodeFromString<State>(json)
5334
}.getOrNull()
5435
set(value) {
55-
storage.withLock { this[SP_STATE_KEY] = Json.encodeToString(value) }
36+
val json = value?.let { Json.encodeToString(it) }
37+
storage.writeString(SP_STATE_KEY, json)
5638
}
5739

5840
fun clear() {
59-
storage.withLock {
60-
removeKeysStartingWith(prefix = TCF_PREFIX)
61-
removeKeysStartingWith(prefix = GPP_PREFIX)
62-
remove(USPSTRING_KEY)
63-
remove(SP_STATE_KEY)
64-
}
41+
storage.removeKeysStartingWith(TCF_PREFIX)
42+
storage.removeKeysStartingWith(GPP_PREFIX)
43+
storage.delete(USPSTRING_KEY)
44+
storage.delete(SP_STATE_KEY)
6545
}
6646
}

‎core/src/commonMain/kotlin/com/sourcepoint/mobile_core/storage/SettingsExt.kt‎

Lines changed: 28 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,7 @@
11
package com.sourcepoint.mobile_core.storage
22

3+
import com.russhwolf.settings.MapSettings
34
import com.russhwolf.settings.Settings
4-
import kotlinx.coroutines.runBlocking
5-
import kotlinx.coroutines.sync.withLock
65
import kotlinx.coroutines.InternalCoroutinesApi
76
import kotlinx.coroutines.internal.SynchronizedObject
87
import kotlinx.coroutines.internal.synchronized
@@ -37,6 +36,33 @@ internal fun Settings.removeKeysStartingWith(prefix: String) = withLock {
3736
toRemove.forEach { remove(it) }
3837
}
3938

39+
internal fun Settings.getKeysWithPrefix(prefix: String): Map<String, JsonPrimitive> = withLock {
40+
keys.filter { it.startsWith(prefix) }
41+
.associateWith { getJsonPrimitive(it) }
42+
}
43+
44+
internal fun Settings.replaceKeysWithPrefix(prefix: String, data: Map<String, JsonPrimitive>) = withLock {
45+
val toRemove = keys.filter { it.startsWith(prefix) }
46+
toRemove.forEach { remove(it) }
47+
data.forEach { (key, value) -> putJsonPrimitive(key, value) }
48+
}
49+
50+
internal fun Settings.readString(key: String, defaultValue: String = ""): String = withLock {
51+
getString(key, defaultValue)
52+
}
53+
54+
internal fun Settings.readStringOrNull(key: String): String? = withLock {
55+
getStringOrNull(key)
56+
}
57+
58+
internal fun Settings.writeString(key: String, value: String?) = withLock {
59+
if (value == null) remove(key) else putString(key, value)
60+
}
61+
62+
internal fun Settings.delete(key: String) = withLock {
63+
remove(key)
64+
}
65+
4066
internal operator fun Settings.set(key: String, value: JsonPrimitive) = putJsonPrimitive(key, value)
4167

4268
internal fun Settings.putJsonPrimitive(key: String, value: JsonPrimitive) {
Lines changed: 177 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,177 @@
1+
package com.sourcepoint.mobile_core
2+
3+
import com.russhwolf.settings.MapSettings
4+
import com.sourcepoint.mobile_core.models.consents.State
5+
import com.sourcepoint.mobile_core.storage.Repository
6+
import kotlinx.coroutines.Dispatchers
7+
import kotlinx.coroutines.async
8+
import kotlinx.coroutines.awaitAll
9+
import kotlinx.coroutines.delay
10+
import kotlinx.coroutines.test.runTest
11+
import kotlinx.serialization.json.JsonPrimitive
12+
import com.sourcepoint.mobile_core.asserters.assertDoesNotContain
13+
import com.sourcepoint.mobile_core.asserters.assertIsEmpty
14+
import com.sourcepoint.mobile_core.asserters.assertNotEmpty
15+
import kotlin.test.Test
16+
import kotlin.test.assertEquals
17+
import kotlin.test.assertNotNull
18+
import kotlin.time.Duration.Companion.milliseconds
19+
20+
class RepositoryConcurrencyTest {
21+
22+
@Test
23+
fun concurrentWritesToTcDataProduceSingleEntry() = runTest {
24+
val repository = Repository(MapSettings())
25+
repository.tcData = mapOf("IABTCF_initial" to JsonPrimitive("should_be_removed"))
26+
27+
val jobs = (0 until 10).map { threadId ->
28+
async(Dispatchers.Default) {
29+
repeat(100) { iteration ->
30+
repository.tcData = mapOf(
31+
"IABTCF_key1" to JsonPrimitive("thread${threadId}_iter${iteration}"),
32+
"IABTCF_key2" to JsonPrimitive("thread${threadId}_iter${iteration}")
33+
)
34+
}
35+
}
36+
}
37+
jobs.awaitAll()
38+
39+
val tcData = repository.tcData
40+
assertEquals(2, tcData.size, "Expected exactly 2 keys after concurrent writes")
41+
assertDoesNotContain(tcData.keys, "IABTCF_initial")
42+
assertEquals(tcData["IABTCF_key1"], tcData["IABTCF_key2"])
43+
}
44+
45+
@Test
46+
fun concurrentWritesToGppDataProduceSingleEntry() = runTest {
47+
val repository = Repository(MapSettings())
48+
repository.gppData = mapOf("IABGPP_initial" to JsonPrimitive("should_be_removed"))
49+
50+
val jobs = (0 until 10).map { threadId ->
51+
async(Dispatchers.Default) {
52+
repeat(100) { iteration ->
53+
repository.gppData = mapOf(
54+
"IABGPP_key1" to JsonPrimitive("thread${threadId}_iter${iteration}"),
55+
"IABGPP_key2" to JsonPrimitive("thread${threadId}_iter${iteration}")
56+
)
57+
}
58+
}
59+
}
60+
jobs.awaitAll()
61+
62+
val gppData = repository.gppData
63+
assertEquals(2, gppData.size, "Expected exactly 2 keys after concurrent writes")
64+
assertDoesNotContain(gppData.keys, "IABGPP_initial")
65+
assertEquals(gppData["IABGPP_key1"], gppData["IABGPP_key2"])
66+
}
67+
68+
@Test
69+
fun concurrentReadsNeverSeeInconsistentTcData() = runTest {
70+
val repository = Repository(MapSettings())
71+
repository.tcData = mapOf("IABTCF_initial" to JsonPrimitive("initial_value"))
72+
73+
val writeJobs = (0 until 50).map { writeId ->
74+
async(Dispatchers.Default) {
75+
repository.tcData = mapOf("IABTCF_key" to JsonPrimitive("value_$writeId"))
76+
}
77+
}
78+
79+
val readJobs = (0 until 50).map {
80+
async(Dispatchers.Default) {
81+
assertNotEmpty(repository.tcData)
82+
}
83+
}
84+
85+
(writeJobs + readJobs).awaitAll()
86+
assertEquals(1, repository.tcData.size)
87+
}
88+
89+
@Test
90+
fun concurrentReadsNeverSeeInconsistentGppData() = runTest {
91+
val repository = Repository(MapSettings())
92+
repository.gppData = mapOf("IABGPP_initial" to JsonPrimitive("initial_value"))
93+
94+
val writeJobs = (0 until 50).map { writeId ->
95+
async(Dispatchers.Default) {
96+
repository.gppData = mapOf("IABGPP_key" to JsonPrimitive("value_$writeId"))
97+
}
98+
}
99+
100+
val readJobs = (0 until 50).map {
101+
async(Dispatchers.Default) {
102+
assertNotEmpty(repository.gppData)
103+
}
104+
}
105+
106+
(writeJobs + readJobs).awaitAll()
107+
assertEquals(1, repository.gppData.size)
108+
}
109+
110+
@Test
111+
fun concurrentClearOperationsLeaveStorageEmpty() = runTest {
112+
val storage = MapSettings()
113+
val repository = Repository(storage)
114+
115+
val jobs = (0 until 100).map { iteration ->
116+
async(Dispatchers.Default) {
117+
repository.tcData = mapOf("IABTCF_key" to JsonPrimitive("value"))
118+
repository.gppData = mapOf("IABGPP_key" to JsonPrimitive("value"))
119+
repository.uspString = "uspString_$iteration"
120+
delay(1.milliseconds)
121+
repository.clear()
122+
}
123+
}
124+
jobs.awaitAll()
125+
126+
assertIsEmpty(storage.keys)
127+
}
128+
129+
@Test
130+
fun multiKeyWritesAreAtomic() = runTest {
131+
val repository = Repository(MapSettings())
132+
133+
val jobs = (0 until 100).map { id ->
134+
async(Dispatchers.Default) {
135+
repository.tcData = mapOf(
136+
"IABTCF_key1" to JsonPrimitive("value1_$id"),
137+
"IABTCF_key2" to JsonPrimitive("value2_$id"),
138+
"IABTCF_key3" to JsonPrimitive("value3_$id")
139+
)
140+
}
141+
}
142+
jobs.awaitAll()
143+
144+
val tcData = repository.tcData
145+
assertEquals(3, tcData.size)
146+
147+
val ids = tcData.values.map { it.content.substringAfterLast("_") }.toSet()
148+
assertEquals(1, ids.size, "All values must be from same write operation")
149+
}
150+
151+
@Test
152+
fun stateSerializationIsThreadSafe() = runTest {
153+
val repository = Repository(MapSettings())
154+
155+
val writeJobs = (0 until 100).map { id ->
156+
async(Dispatchers.Default) {
157+
repository.state = State(accountId = id, propertyId = id * 2)
158+
}
159+
}
160+
161+
val readJobs = (0 until 100).map {
162+
async(Dispatchers.Default) {
163+
repository.state?.let { state ->
164+
assertNotNull(state.accountId)
165+
assertNotNull(state.propertyId)
166+
assertEquals(state.accountId * 2, state.propertyId)
167+
}
168+
}
169+
}
170+
171+
(writeJobs + readJobs).awaitAll()
172+
173+
val finalState = repository.state
174+
assertNotNull(finalState)
175+
assertEquals(finalState.accountId * 2, finalState.propertyId)
176+
}
177+
}

0 commit comments

Comments
 (0)