feat(sync): a fetch queue per identity, and the inbox reopened on activation
Phase 6 of docs/multiple-profiles.md. The one data problem, closed from both ends. The two queues that fetch -- SynchronizeNostrEventRequest and NegentropySynchronizeRequest -- gain a nullable ownerPublicKey, the identity that asked, because what a fetch brings back is opened with the fetcher's key. The pumps ask for their own: the head of the queue among this identity's requests and the ones nobody owns, so that a request another identity queued waits for that identity. Rows from before the column read back null, meaning "the device's", which is what they were, and any identity may drain them. The broadcast queue deliberately gets no owner: a signed event is anyone's to carry, and holding A's outgoing message until A is opened again would be a delivery failure the user would never be told about. Room version 20, an AutoMigration for two nullable columns, with the schema export committed beside its predecessors. The stamp is the repository's, not the call site's. Eighteen view models and two DAO paths queue requests, and every one of them does so as the active identity -- a screen cannot queue anything as anyone else. DatabaseNostrRepository takes the identity flow at construction, which meant moving the view model above the repositories in the nav host, and stamps every request it queues where the caller stamped nothing; a negentropy request carries its owner into the REQ it becomes. The DAO's own requests -- the placeholder-profile syncs it plants while indexing, and the participant syncs of a Marmot join -- stay unowned, on purpose and against the plan's sketch: they fetch public kinds that need no key to open, and any profile that is open may as well fetch them. The other end: wraps that were fetched under the wrong key -- before this, or by a read-only identity whose nsec was pasted later, the case the npub plan's DAO guard left with the words "the key that opens this one may be signed in later" and no code behind them. storeNostrEvent never indexes an event it already holds, so a wrap that arrived under the wrong key stayed closed for good: the live subscription re-received it, the DAO saw a known id, and returned. NostrDao.reopenInbox finds every wrap addressed to the active key that no seal names and runs each through indexNostrEvent again under the right key, one transaction per wrap so that one that cannot be opened rolls back its own changes and the next is still tried; one pass, since wraps do not depend on one another. It runs from the sync pumps' collector once per activation, for an identity that can sign, before the live inbox and the broadcast pump are launched -- so the sweep and the live subscription are not opening the same wrap at once; both are idempotent by event id. Found on the way: GiftWrapSeal.giftWrapMessageId, a foreign key to the wrap a seal came out of, existed and was never written, so nothing in the database could say which wraps had been opened. decryptGiftWrapSeal now sets it. A seal from before reads back null, is swept once, and comes back with the link -- the reindex is idempotent, and persistInboundChatMessage already files a message it has once. Tests: ReadOnlyGiftWrapDaoJvmTest -- a wrap for us stored under another profile's key is stored and not opened, a re-delivery under our key changes nothing, the sweep opens it once and finds nothing the second time; the other profile's sweep and a read-only pair open nothing; a wrap stored read-only is opened once the key is here. OwnedRequestQueueJvmTest -- A's request is offered to A and not to B, nobody's to both, oldest first; the repository stamps what it queues with the identity that is open, nothing when none is, and keeps an owner a caller named. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Pulled-From: curated/curated@759199f2af
This commit is contained in:
File diff suppressed because it is too large
Load Diff
@@ -177,7 +177,7 @@ val GENESIS_AT = Instant.fromEpochMilliseconds(1231006505000L)
|
||||
UnsignedNostrEvent::class,
|
||||
Zap::class
|
||||
],
|
||||
version = 19,
|
||||
version = 20,
|
||||
autoMigrations = [
|
||||
// v2 only adds the DkgSession/DkgParticipantMessage tables, so Room can
|
||||
// generate the migration itself — nothing existing changes shape.
|
||||
@@ -312,6 +312,15 @@ val GENESIS_AT = Instant.fromEpochMilliseconds(1231006505000L)
|
||||
// room, which is correct for them: a ceremony held in its own NIP-17 room
|
||||
// ran over exactly that room's members and was named after it.
|
||||
AutoMigration(from = 18, to = 19),
|
||||
// v20 adds the nullable ownerPublicKey to SynchronizeNostrEventRequest and
|
||||
// NegentropySynchronizeRequest: the identity that asked, so that only its
|
||||
// pump answers. The two queues used to be the device's, drained by whoever
|
||||
// was open with their key pair, and a gift wrap fetched for one profile
|
||||
// while another was open was stored under the wrong key and never opened
|
||||
// again. Rows written before this read back null, meaning "the device's",
|
||||
// which is what they were, and any identity may drain them. The broadcast
|
||||
// queue deliberately gets no owner: a signed event is anyone's to carry.
|
||||
AutoMigration(from = 19, to = 20),
|
||||
]
|
||||
)
|
||||
@ColumnTypeConverters(MantraConverters::class)
|
||||
|
||||
@@ -9,6 +9,18 @@ interface GiftWrapMessageDao {
|
||||
@Query("SELECT * FROM GiftWrapMessage")
|
||||
suspend fun getAllGiftWrapMessages(): List<press.mantra.compose.database.model.GiftWrapMessage>
|
||||
|
||||
/**
|
||||
* Every wrap addressed to [receiverPublicKey] that the device stored and never
|
||||
* opened -- no seal names it. Fetched under another identity's key, or under a
|
||||
* read-only identity's none; the inbox sweep at activation opens them.
|
||||
*/
|
||||
@Query(
|
||||
"SELECT * FROM GiftWrapMessage WHERE receiverPublicKey = :receiverPublicKey " +
|
||||
"AND id NOT IN (SELECT giftWrapMessageId FROM GiftWrapSeal WHERE giftWrapMessageId IS NOT NULL) " +
|
||||
"ORDER BY createdAt ASC"
|
||||
)
|
||||
suspend fun getUnopenedAddressedTo(receiverPublicKey: String): List<press.mantra.compose.database.model.GiftWrapMessage>
|
||||
|
||||
@Upsert
|
||||
suspend fun upsert(giftWrapMessage: press.mantra.compose.database.model.GiftWrapMessage)
|
||||
|
||||
|
||||
@@ -9,6 +9,9 @@ interface GiftWrapSealDao {
|
||||
@Query("SELECT * FROM GiftWrapSeal")
|
||||
suspend fun getAllGiftWrapSeals(): List<press.mantra.compose.database.model.GiftWrapSeal>
|
||||
|
||||
@Query("SELECT COUNT(*) FROM GiftWrapSeal WHERE giftWrapMessageId = :giftWrapMessageId")
|
||||
suspend fun countForMessage(giftWrapMessageId: String): Int
|
||||
|
||||
@Upsert
|
||||
suspend fun upsert(giftWrapSeal: press.mantra.compose.database.model.GiftWrapSeal)
|
||||
|
||||
|
||||
@@ -14,6 +14,14 @@ interface NegentropySynchronizeRequestDao {
|
||||
@Query("SELECT * FROM NegentropySynchronizeRequest WHERE status = :status ORDER BY createdAt ASC, id ASC LIMIT 1")
|
||||
fun observeNegentropySynchronizeRequestsByStatus(status: String): Flow<NegentropySynchronizeRequest?>
|
||||
|
||||
/** The same head-of-queue, for one identity; see `SynchronizeNostrEventRequestDao.observeSynchronizeNostrEventRequestsByStatusFor`. */
|
||||
@Query(
|
||||
"SELECT * FROM NegentropySynchronizeRequest " +
|
||||
"WHERE status = :status AND (ownerPublicKey = :ownerPublicKey OR ownerPublicKey IS NULL) " +
|
||||
"ORDER BY createdAt ASC, id ASC LIMIT 1"
|
||||
)
|
||||
fun observeNegentropySynchronizeRequestsByStatusFor(status: String, ownerPublicKey: String): Flow<NegentropySynchronizeRequest?>
|
||||
|
||||
@Query("SELECT COUNT(*) FROM NegentropySynchronizeRequest WHERE purpose = :purpose AND status IN (:status)")
|
||||
fun observeNegentropySynchronizeRequestByPurposeAndStatusCount(purpose: String, status: List<String>): Flow<Int>
|
||||
|
||||
@@ -29,4 +37,7 @@ interface NegentropySynchronizeRequestDao {
|
||||
*/
|
||||
@Upsert
|
||||
suspend fun insert(negentropySynchronizeRequests: List<NegentropySynchronizeRequest>)
|
||||
|
||||
@Query("SELECT * FROM NegentropySynchronizeRequest")
|
||||
suspend fun getAll(): List<NegentropySynchronizeRequest>
|
||||
}
|
||||
@@ -31,6 +31,7 @@ import press.mantra.compose.exceptions.MarmotNotMemberOfChatGroupException
|
||||
import press.mantra.compose.exceptions.MarmotWelcomeEventMissingKeyPackageEventIdException
|
||||
import press.mantra.compose.extensions.toHex
|
||||
import press.mantra.compose.managers.ChillDkgRitualManager
|
||||
import press.mantra.compose.nostr.Relays
|
||||
import press.mantra.compose.nostr.dkg.DkgRitualEvents
|
||||
import press.mantra.compose.nostr.frost.FrostSigningEvents
|
||||
import press.mantra.compose.managers.FrostSigningManager
|
||||
@@ -146,6 +147,58 @@ abstract class NostrDao(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Opens every gift wrap addressed to [activeKeyPair] that the device stored and could
|
||||
* not open at the time -- because another identity fetched it, or because this one
|
||||
* held no key yet. Phase 6 of docs/multiple-profiles.md.
|
||||
*
|
||||
* [storeNostrEvent] never indexes an event it already holds, so a wrap that arrived
|
||||
* under the wrong key would otherwise stay closed for good: the live subscription
|
||||
* re-receives it, the DAO sees a known id, and returns. This runs the wrap through
|
||||
* [indexNostrEvent] again, under the right key, one transaction per wrap so that a
|
||||
* wrap that cannot be opened rolls back its own changes and the next is still tried.
|
||||
* One pass: wraps do not depend on one another, so a wrap that cannot be opened now
|
||||
* will not be helped by opening its neighbours.
|
||||
*
|
||||
* @return how many were opened.
|
||||
*/
|
||||
open suspend fun reopenInbox(activeKeyPair: KeyPair): Int {
|
||||
if (activeKeyPair.privKey == null) return 0
|
||||
val me = activeKeyPair.pubKey.toHex()
|
||||
val unopened = database.giftWrapMessageDao().getUnopenedAddressedTo(me)
|
||||
if (unopened.isEmpty()) return 0
|
||||
logger.i("Reopening ${unopened.size} gift wrap(s) addressed to $me")
|
||||
|
||||
var opened = 0
|
||||
for (wrap in unopened) {
|
||||
val nostrEvent = database.nostrEventDao().getNostrEventById(wrap.nostrEventId) ?: continue
|
||||
val relayURL = database.nostrEventRelayDao().getRelayUrls(nostrEvent.id).firstOrNull()
|
||||
?: Relays.DefaultDMRelayList.first().url
|
||||
try {
|
||||
reopenGiftWrap(nostrEvent, relayURL, activeKeyPair)
|
||||
if (database.giftWrapSealDao().countForMessage(wrap.id) > 0) opened++
|
||||
} catch (e: CancellationException) {
|
||||
throw e
|
||||
} catch (e: Throwable) {
|
||||
logger.w("Gift wrap ${wrap.id} could not be opened on the second try", e)
|
||||
}
|
||||
}
|
||||
logger.i("Reopened $opened of ${unopened.size} gift wrap(s) for $me")
|
||||
return opened
|
||||
}
|
||||
|
||||
/** One wrap, indexed again, in a transaction of its own; see [reopenInbox]. */
|
||||
@Transaction
|
||||
open suspend fun reopenGiftWrap(nostrEvent: NostrEvent, relayURL: String, activeKeyPair: KeyPair) {
|
||||
indexNostrEvent(
|
||||
nostrEvent = nostrEvent,
|
||||
relayURL = relayURL,
|
||||
synchronizationRelayURLs = listOf(relayURL),
|
||||
level = 0,
|
||||
activeKeyPair = activeKeyPair,
|
||||
)
|
||||
}
|
||||
|
||||
@Transaction
|
||||
open suspend fun storeNostrEvent(
|
||||
nostrEvent: NostrEvent,
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package press.mantra.compose.database.dao
|
||||
|
||||
import androidx.room3.Dao
|
||||
import androidx.room3.Query
|
||||
import androidx.room3.Upsert
|
||||
|
||||
@Dao
|
||||
@@ -9,4 +10,8 @@ interface NostrEventRelayDao {
|
||||
@Upsert
|
||||
suspend fun upsert(nostrEventRelay: press.mantra.compose.database.model.NostrEventRelay)
|
||||
|
||||
/** The relays an event was seen on, oldest first. */
|
||||
@Query("SELECT relayURL FROM NostrEventRelay WHERE nostrEventId = :nostrEventId ORDER BY createdAt ASC")
|
||||
suspend fun getRelayUrls(nostrEventId: String): List<String>
|
||||
|
||||
}
|
||||
@@ -14,6 +14,18 @@ interface SynchronizeNostrEventRequestDao {
|
||||
@Query("SELECT * FROM SynchronizeNostrEventRequest WHERE status = :status ORDER BY createdAt ASC, id ASC LIMIT 1")
|
||||
fun observeSynchronizeNostrEventRequestsByStatus(status: String): Flow<press.mantra.compose.database.model.SynchronizeNostrEventRequest?>
|
||||
|
||||
/**
|
||||
* The same head-of-queue, for one identity: its own requests and the ones nobody
|
||||
* owns. A request another identity queued waits for that identity, because what it
|
||||
* brings back is opened with that identity's key.
|
||||
*/
|
||||
@Query(
|
||||
"SELECT * FROM SynchronizeNostrEventRequest " +
|
||||
"WHERE status = :status AND (ownerPublicKey = :ownerPublicKey OR ownerPublicKey IS NULL) " +
|
||||
"ORDER BY createdAt ASC, id ASC LIMIT 1"
|
||||
)
|
||||
fun observeSynchronizeNostrEventRequestsByStatusFor(status: String, ownerPublicKey: String): Flow<press.mantra.compose.database.model.SynchronizeNostrEventRequest?>
|
||||
|
||||
@Query("SELECT COUNT(*) FROM SynchronizeNostrEventRequest WHERE purpose = :purpose AND status IN (:status)")
|
||||
fun observeSynchronizeNostrEventRequestByPurposeAndStatusCount(purpose: String, status: List<String>): Flow<Int>
|
||||
|
||||
@@ -31,4 +43,7 @@ interface SynchronizeNostrEventRequestDao {
|
||||
|
||||
@Insert
|
||||
suspend fun insert(synchronizeNostrEventRequests: List<press.mantra.compose.database.model.SynchronizeNostrEventRequest>)
|
||||
|
||||
@Query("SELECT * FROM SynchronizeNostrEventRequest")
|
||||
suspend fun getAll(): List<press.mantra.compose.database.model.SynchronizeNostrEventRequest>
|
||||
}
|
||||
@@ -137,7 +137,12 @@ data class GiftWrapMessage(
|
||||
content = giftWrapSeal.content,
|
||||
tags = giftWrapSeal.tags,
|
||||
createdAt = Instant.fromEpochSeconds(giftWrapSeal.createdAt),
|
||||
signature = giftWrapSeal.sig
|
||||
signature = giftWrapSeal.sig,
|
||||
// The wrap this came out of. The column existed and was never
|
||||
// written, so nothing could say which wraps had been opened;
|
||||
// the inbox sweep asks exactly that. A seal from before this
|
||||
// reads back null, is swept once, and comes back with the link.
|
||||
giftWrapMessageId = id,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -42,6 +42,9 @@ data class NegentropySynchronizeRequest(
|
||||
|
||||
val synchronizationFilter: press.mantra.compose.database.model.types.SynchronizationFilter,
|
||||
|
||||
/** The identity that asked; null for rows from before, which any identity may drain. See [SynchronizeNostrEventRequest.ownerPublicKey]. */
|
||||
val ownerPublicKey: HexKey? = null,
|
||||
|
||||
override val nostrEventId: HexKey? = null,
|
||||
override val createdAt: Instant = Clock.System.now(),
|
||||
override val updatedAt: Instant = createdAt,
|
||||
@@ -64,6 +67,7 @@ data class NegentropySynchronizeRequest(
|
||||
),
|
||||
relayURL = relayURL,
|
||||
level = level,
|
||||
ownerPublicKey = ownerPublicKey,
|
||||
)
|
||||
}
|
||||
companion object {
|
||||
|
||||
@@ -43,6 +43,13 @@ data class SynchronizeNostrEventRequest(
|
||||
|
||||
val synchronizationFilters: Array<press.mantra.compose.database.model.types.SynchronizationFilter>,
|
||||
|
||||
/**
|
||||
* The identity that asked, so that only its pump answers: what a fetch brings back
|
||||
* is opened with the fetcher's key. Null for rows written before there was one to
|
||||
* record, which any identity may drain. See docs/multiple-profiles.md, Phase 6.
|
||||
*/
|
||||
val ownerPublicKey: HexKey? = null,
|
||||
|
||||
override val nostrEventId: HexKey? = null,
|
||||
override val unsignedNostrEventId: Long? = null,
|
||||
override val createdAt: Instant = Clock.System.now(),
|
||||
@@ -61,6 +68,7 @@ data class SynchronizeNostrEventRequest(
|
||||
if (status != other.status) return false
|
||||
if (relayURL != other.relayURL) return false
|
||||
if (!synchronizationFilters.contentEquals(other.synchronizationFilters)) return false
|
||||
if (ownerPublicKey != other.ownerPublicKey) return false
|
||||
if (nostrEventId != other.nostrEventId) return false
|
||||
if (createdAt != other.createdAt) return false
|
||||
if (updatedAt != other.updatedAt) return false
|
||||
@@ -74,6 +82,7 @@ data class SynchronizeNostrEventRequest(
|
||||
result = 31 * result + status.hashCode()
|
||||
result = 31 * result + relayURL.hashCode()
|
||||
result = 31 * result + synchronizationFilters.contentHashCode()
|
||||
result = 31 * result + (ownerPublicKey?.hashCode() ?: 0)
|
||||
result = 31 * result + (nostrEventId?.hashCode() ?: 0)
|
||||
result = 31 * result + createdAt.hashCode()
|
||||
result = 31 * result + updatedAt.hashCode()
|
||||
|
||||
@@ -50,14 +50,25 @@ import com.vitorpamplona.quartz.nip57Zaps.LnZapEvent
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import press.mantra.compose.identity.Identity
|
||||
import kotlin.time.Clock
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* @param activeIdentity whose requests the two fetch queues are stamped with. Every
|
||||
* request queued through here is queued *as* the active identity -- a screen cannot
|
||||
* queue anything as anyone else -- so the stamp is the repository's, not the call
|
||||
* site's, and no call site can forget it. Defaults to nobody, which stamps nothing
|
||||
* and is what a test that has no identity gets.
|
||||
*/
|
||||
class DatabaseNostrRepository(
|
||||
private val database: press.mantra.compose.database.MantraDatabase,
|
||||
private val scope: CoroutineScope
|
||||
private val scope: CoroutineScope,
|
||||
private val activeIdentity: StateFlow<Identity?> = MutableStateFlow(null),
|
||||
): press.mantra.compose.repository.NostrRepository, press.mantra.compose.repository.RelayRepository {
|
||||
companion object {
|
||||
const val TAG = "DatabaseNostrRepository"
|
||||
@@ -108,12 +119,19 @@ class DatabaseNostrRepository(
|
||||
return database.broadcastNostrEventRequestDao().observeBroadcastNostrEventRequestsByStatus("pending")
|
||||
}
|
||||
|
||||
override suspend fun observePendingSynchronizeNostrEventRequests(): Flow<press.mantra.compose.database.model.SynchronizeNostrEventRequest?> {
|
||||
return database.synchronizeNostrEventRequestDao().observeSynchronizeNostrEventRequestsByStatus("pending")
|
||||
override suspend fun observePendingSynchronizeNostrEventRequests(ownerPublicKey: HexKey): Flow<press.mantra.compose.database.model.SynchronizeNostrEventRequest?> {
|
||||
return database.synchronizeNostrEventRequestDao().observeSynchronizeNostrEventRequestsByStatusFor("pending", ownerPublicKey)
|
||||
}
|
||||
|
||||
override suspend fun observePendingNegentropySynchronizeRequests(): Flow<press.mantra.compose.database.model.NegentropySynchronizeRequest?> {
|
||||
return database.negentropySynchronizeRequestDao().observeNegentropySynchronizeRequestsByStatus("pending")
|
||||
override suspend fun observePendingNegentropySynchronizeRequests(ownerPublicKey: HexKey): Flow<press.mantra.compose.database.model.NegentropySynchronizeRequest?> {
|
||||
return database.negentropySynchronizeRequestDao().observeNegentropySynchronizeRequestsByStatusFor("pending", ownerPublicKey)
|
||||
}
|
||||
|
||||
override suspend fun reopenInbox(activeKeyPair: KeyPair): Int {
|
||||
// Under the same lock as every other write through indexNostrEvent.
|
||||
return storeNostrEventMutex.withLock {
|
||||
database.nostrDao().reopenInbox(activeKeyPair)
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun createNewProfile(
|
||||
@@ -661,17 +679,23 @@ class DatabaseNostrRepository(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Stamped with the active identity where the caller stamped nothing. A request
|
||||
* queued while nothing is active stays nobody's, which any identity may drain.
|
||||
*/
|
||||
override suspend fun queueSynchronizeNostrEvent(
|
||||
synchronizeNostrEventRequests: List<press.mantra.compose.database.model.SynchronizeNostrEventRequest>,
|
||||
) {
|
||||
val owner = activeIdentity.value?.nostrPublicKey
|
||||
database.synchronizeNostrEventRequestDao().insert(
|
||||
synchronizeNostrEventRequests
|
||||
synchronizeNostrEventRequests.map { it.copy(ownerPublicKey = it.ownerPublicKey ?: owner) }
|
||||
)
|
||||
}
|
||||
|
||||
override suspend fun queueNegentropySynchronizeRequest(negentropySynchronizeRequests: List<press.mantra.compose.database.model.NegentropySynchronizeRequest>) {
|
||||
val owner = activeIdentity.value?.nostrPublicKey
|
||||
database.negentropySynchronizeRequestDao().insert(
|
||||
negentropySynchronizeRequests
|
||||
negentropySynchronizeRequests.map { it.copy(ownerPublicKey = it.ownerPublicKey ?: owner) }
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -45,9 +45,17 @@ interface NostrRepository {
|
||||
|
||||
suspend fun observePendingBroadcastNostrEventRequests(): Flow<LocalBroadcastNostrEventRequest?>
|
||||
|
||||
suspend fun observePendingSynchronizeNostrEventRequests(): Flow<SynchronizeNostrEventRequest?>
|
||||
/** The head of the fetch queue for one identity: its own requests and the ones nobody owns. */
|
||||
suspend fun observePendingSynchronizeNostrEventRequests(ownerPublicKey: HexKey): Flow<SynchronizeNostrEventRequest?>
|
||||
|
||||
suspend fun observePendingNegentropySynchronizeRequests(): Flow<NegentropySynchronizeRequest?>
|
||||
suspend fun observePendingNegentropySynchronizeRequests(ownerPublicKey: HexKey): Flow<NegentropySynchronizeRequest?>
|
||||
|
||||
/**
|
||||
* Opens every gift wrap addressed to [activeKeyPair] that the device stored and could
|
||||
* not open at the time. Run once per activation, before the live inbox is opened.
|
||||
* Returns how many were opened.
|
||||
*/
|
||||
suspend fun reopenInbox(activeKeyPair: KeyPair): Int
|
||||
|
||||
suspend fun createNewProfile(
|
||||
publicKey: HexKey,
|
||||
@@ -237,11 +245,13 @@ interface NostrRepository {
|
||||
TODO("Not yet implemented")
|
||||
}
|
||||
|
||||
override suspend fun observePendingSynchronizeNostrEventRequests(): Flow<SynchronizeNostrEventRequest?> {
|
||||
override suspend fun reopenInbox(activeKeyPair: KeyPair): Int = 0
|
||||
|
||||
override suspend fun observePendingSynchronizeNostrEventRequests(ownerPublicKey: HexKey): Flow<SynchronizeNostrEventRequest?> {
|
||||
TODO("Not yet implemented")
|
||||
}
|
||||
|
||||
override suspend fun observePendingNegentropySynchronizeRequests(): Flow<NegentropySynchronizeRequest?> {
|
||||
override suspend fun observePendingNegentropySynchronizeRequests(ownerPublicKey: HexKey): Flow<NegentropySynchronizeRequest?> {
|
||||
TODO("Not yet implemented")
|
||||
}
|
||||
|
||||
|
||||
@@ -177,10 +177,19 @@ fun MantraNavHost(
|
||||
AuxDatabaseManager(mantraGlobal)
|
||||
}
|
||||
|
||||
// Before the repositories: the nostr repository stamps every fetch request it queues
|
||||
// with the active identity, and takes the flow at construction.
|
||||
val sovereignWalletViewModel: SovereignWalletViewModel = viewModel(
|
||||
factory = SovereignWalletViewModel.factory(
|
||||
phoenixGlobal = phoenixGlobal
|
||||
)
|
||||
)
|
||||
|
||||
val databaseNostrRepository = remember {
|
||||
DatabaseNostrRepository(
|
||||
database = auxDatabaseManager.auxDatabase,
|
||||
applicationIOScope
|
||||
applicationIOScope,
|
||||
activeIdentity = sovereignWalletViewModel.activeIdentity,
|
||||
)
|
||||
}
|
||||
|
||||
@@ -224,12 +233,6 @@ fun MantraNavHost(
|
||||
)
|
||||
}
|
||||
|
||||
val sovereignWalletViewModel: SovereignWalletViewModel = viewModel(
|
||||
factory = SovereignWalletViewModel.factory(
|
||||
phoenixGlobal = phoenixGlobal
|
||||
)
|
||||
)
|
||||
|
||||
val navigationViewModel: NavigationViewModel = viewModel (
|
||||
factory = NavigationViewModel.factory(
|
||||
activeIdentity = sovereignWalletViewModel.activeIdentity,
|
||||
|
||||
@@ -22,6 +22,7 @@ import press.mantra.compose.repository.RelayRepository
|
||||
import co.touchlab.kermit.Logger
|
||||
import com.vitorpamplona.negentropy.Negentropy
|
||||
import com.vitorpamplona.negentropy.storage.StorageVector
|
||||
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.crypto.KeyPair
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.CloseCmd
|
||||
import com.vitorpamplona.quartz.nip01Core.relay.commands.toRelay.ReqCmd
|
||||
@@ -227,6 +228,14 @@ class SynchronizationViewModel(
|
||||
// things this identity cannot open. Not asking is not an
|
||||
// optimisation; it is the identity not asking for what it cannot use.
|
||||
if (identity.canSign) {
|
||||
// The wraps that arrived while another profile was open, or
|
||||
// while this one held no key: opened now, before the live inbox
|
||||
// starts, so the two are not opening the same wrap at once. Both
|
||||
// are idempotent by event id; the order is about not doing the
|
||||
// work twice. See docs/multiple-profiles.md, Phase 6.
|
||||
runCatching { nostrRepository.reopenInbox(keyPair) }
|
||||
.onSuccess { if (it > 0) logger.i { "Reopened $it gift wrap(s) that arrived while this profile was not open" } }
|
||||
.onFailure { logger.e("Could not reopen the inbox", it) }
|
||||
launch(Dispatchers.IO) { observePendingBroadcastNostrEventRequests(keyPair) }
|
||||
// Runs until cancelled rather than draining a queue, but it
|
||||
// belongs here for the same reason the pumps do: it needs the
|
||||
@@ -245,7 +254,8 @@ class SynchronizationViewModel(
|
||||
): Unit = coroutineScope {
|
||||
logger.i { "observePendingSyncNostrEventRequests" }
|
||||
|
||||
nostrRepository.observePendingSynchronizeNostrEventRequests().distinctUntilChanged().collect { synchronizeNostrEventRequestOrNull ->
|
||||
// This identity's requests and the ones nobody owns; another identity's wait for it.
|
||||
nostrRepository.observePendingSynchronizeNostrEventRequests(keyPair.pubKey.toHexKey()).distinctUntilChanged().collect { synchronizeNostrEventRequestOrNull ->
|
||||
synchronizeNostrEventRequestOrNull?.let { synchronizeNostrEventRequest ->
|
||||
guardPump("sync request ${synchronizeNostrEventRequest.id}") {
|
||||
logger.i("synchronizeNostrEventRequest: $synchronizeNostrEventRequest")
|
||||
@@ -406,7 +416,7 @@ class SynchronizationViewModel(
|
||||
): Unit = coroutineScope {
|
||||
logger.i { "observePendingNegentropySynchronizeRequests" }
|
||||
|
||||
nostrRepository.observePendingNegentropySynchronizeRequests().distinctUntilChanged().collect { negentropySynchronizeRequestOrNull ->
|
||||
nostrRepository.observePendingNegentropySynchronizeRequests(keyPair.pubKey.toHexKey()).distinctUntilChanged().collect { negentropySynchronizeRequestOrNull ->
|
||||
negentropySynchronizeRequestOrNull?.let { negentropySynchronizeRequest ->
|
||||
guardPump("negentropy request ${negentropySynchronizeRequest.id}") {
|
||||
mutex.withLock {
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
package press.mantra.compose.database.dao
|
||||
|
||||
import androidx.datastore.preferences.core.PreferenceDataStoreFactory
|
||||
import androidx.room3.Room
|
||||
import fr.acinq.bitcoin.ByteVector32
|
||||
import fr.acinq.bitcoin.PrivateKey
|
||||
import fr.acinq.phoenix.managers.nostrPublicKeyHex
|
||||
import fr.acinq.phoenix.utils.preferences.InternalPrefs
|
||||
import fr.acinq.phoenix.utils.preferences.UserPrefs
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.first
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import okio.Path.Companion.toPath
|
||||
import press.mantra.compose.database.MantraDatabase
|
||||
import press.mantra.compose.database.builder.getRoomDatabase
|
||||
import press.mantra.compose.database.model.NegentropySynchronizeRequest
|
||||
import press.mantra.compose.database.model.SynchronizeNostrEventRequest
|
||||
import press.mantra.compose.database.model.types.SynchronizationFilter
|
||||
import press.mantra.compose.database.repository.DatabaseNostrRepository
|
||||
import press.mantra.compose.identity.Identity
|
||||
import press.mantra.compose.identity.IdentityKind
|
||||
import press.mantra.compose.identity.toWalletId
|
||||
import java.nio.file.Files
|
||||
import kotlin.test.AfterTest
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* The two fetch queues, owned -- Phase 6 of docs/multiple-profiles.md.
|
||||
*
|
||||
* A request another identity queued waits for that identity, because what it brings
|
||||
* back is opened with that identity's key; a request nobody owns -- one from before the
|
||||
* column, or queued while nothing was open -- is anyone's. And the owner is the
|
||||
* repository's to stamp, from the identity that is open when the request is queued,
|
||||
* so that no call site can forget it.
|
||||
*/
|
||||
class OwnedRequestQueueJvmTest {
|
||||
|
||||
private val db: MantraDatabase = getRoomDatabase(Room.inMemoryDatabaseBuilder<MantraDatabase>())
|
||||
private val scope = CoroutineScope(Job() + Dispatchers.IO)
|
||||
|
||||
private val keyA = PrivateKey(ByteVector32("0a".repeat(32)))
|
||||
private val keyB = PrivateKey(ByteVector32("0b".repeat(32)))
|
||||
private val a = keyA.nostrPublicKeyHex()
|
||||
private val b = keyB.nostrPublicKeyHex()
|
||||
private val at = Instant.fromEpochSeconds(1_700_000_000)
|
||||
|
||||
@AfterTest
|
||||
fun tearDown() {
|
||||
scope.cancel()
|
||||
db.close()
|
||||
}
|
||||
|
||||
private fun sync(id: String, owner: String?, createdAt: Instant = at) = SynchronizeNostrEventRequest(
|
||||
id = id,
|
||||
purpose = "test",
|
||||
relayURL = "wss://relay.example",
|
||||
level = 0,
|
||||
synchronizationFilters = arrayOf(SynchronizationFilter(kinds = arrayOf(1))),
|
||||
ownerPublicKey = owner,
|
||||
createdAt = createdAt,
|
||||
updatedAt = createdAt,
|
||||
)
|
||||
|
||||
private fun negentropy(id: String, owner: String?, createdAt: Instant = at) = NegentropySynchronizeRequest(
|
||||
id = id,
|
||||
purpose = "test",
|
||||
relayURL = "wss://relay.example",
|
||||
level = 0,
|
||||
synchronizationFilter = SynchronizationFilter(kinds = arrayOf(1)),
|
||||
ownerPublicKey = owner,
|
||||
createdAt = createdAt,
|
||||
updatedAt = createdAt,
|
||||
)
|
||||
|
||||
@Test
|
||||
fun `a sync request is offered to its owner and to nobody else, and an unowned one to everyone`() = runBlocking {
|
||||
val dao = db.synchronizeNostrEventRequestDao()
|
||||
dao.insert(listOf(sync("a-1", owner = a)))
|
||||
|
||||
assertEquals("a-1", dao.observeSynchronizeNostrEventRequestsByStatusFor("pending", a).first()?.id)
|
||||
assertNull(dao.observeSynchronizeNostrEventRequestsByStatusFor("pending", b).first(), "A's request is not B's to answer")
|
||||
|
||||
dao.insert(listOf(sync("anyone", owner = null, createdAt = at.plus(kotlin.time.Duration.parse("1s")))))
|
||||
assertEquals("a-1", dao.observeSynchronizeNostrEventRequestsByStatusFor("pending", a).first()?.id, "oldest first, still")
|
||||
assertEquals("anyone", dao.observeSynchronizeNostrEventRequestsByStatusFor("pending", b).first()?.id, "nobody's is anyone's")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `a negentropy request is offered the same way`() = runBlocking {
|
||||
val dao = db.negentropySynchronizeRequestDao()
|
||||
dao.insert(listOf(negentropy("a-1", owner = a), negentropy("anyone", owner = null, createdAt = at.plus(kotlin.time.Duration.parse("1s")))))
|
||||
|
||||
assertEquals("a-1", dao.observeNegentropySynchronizeRequestsByStatusFor("pending", a).first()?.id)
|
||||
assertEquals("anyone", dao.observeNegentropySynchronizeRequestsByStatusFor("pending", b).first()?.id)
|
||||
}
|
||||
|
||||
private fun prefsFile(name: String) =
|
||||
Files.createTempFile(name, ".preferences_pb").also { it.toFile().delete() }.toString().toPath()
|
||||
|
||||
private fun identity(key: PrivateKey) = Identity.signing(
|
||||
id = key.publicKey().xOnly().toWalletId(),
|
||||
kind = IdentityKind.NostrSecret,
|
||||
nostrPrivateKey = key,
|
||||
userPrefs = UserPrefs(PreferenceDataStoreFactory.createWithPath { prefsFile("user") }),
|
||||
internalPrefs = InternalPrefs(PreferenceDataStoreFactory.createWithPath { prefsFile("internal") }),
|
||||
business = null,
|
||||
)
|
||||
|
||||
@Test
|
||||
fun `the repository stamps what it queues with the identity that is open, and nothing when none is`() = runBlocking {
|
||||
val active = MutableStateFlow<Identity?>(null)
|
||||
val repository = DatabaseNostrRepository(db, scope, activeIdentity = active)
|
||||
|
||||
repository.queueSynchronizeNostrEvent(listOf(sync("nobody", owner = null)))
|
||||
repository.queueNegentropySynchronizeRequest(listOf(negentropy("nobody-n", owner = null)))
|
||||
|
||||
active.value = identity(keyA)
|
||||
repository.queueSynchronizeNostrEvent(listOf(sync("as-a", owner = null)))
|
||||
repository.queueNegentropySynchronizeRequest(listOf(negentropy("as-a-n", owner = null)))
|
||||
// A caller that named an owner keeps it.
|
||||
repository.queueSynchronizeNostrEvent(listOf(sync("named-b", owner = b)))
|
||||
|
||||
val syncs = db.synchronizeNostrEventRequestDao().getAll().associate { it.id to it.ownerPublicKey }
|
||||
assertEquals(mapOf("nobody" to null, "as-a" to a, "named-b" to b), syncs)
|
||||
val negs = db.negentropySynchronizeRequestDao().getAll().associate { it.id to it.ownerPublicKey }
|
||||
assertEquals(mapOf("nobody-n" to null, "as-a-n" to a), negs)
|
||||
|
||||
// A negentropy request carries its owner into the REQ it becomes.
|
||||
assertEquals(a, negentropy("x", owner = a).toSynchronizeNostrEventRequest().ownerPublicKey)
|
||||
}
|
||||
}
|
||||
@@ -17,7 +17,9 @@ import kotlin.test.assertNotNull
|
||||
import kotlin.time.Instant
|
||||
|
||||
/**
|
||||
* A gift wrap addressed to a read-only identity is kept and not opened.
|
||||
* A gift wrap addressed to a read-only identity is kept and not opened -- and, since
|
||||
* docs/multiple-profiles.md, opened later by the inbox sweep, along with any wrap that
|
||||
* arrived while another profile on the device was the open one.
|
||||
*
|
||||
* `storeNostrEvent` is `@Transaction` and indexes inside it, and a wrap addressed to the
|
||||
* active key that cannot be unsealed throws -- so before the guard in `indexNostrEvent`
|
||||
@@ -103,4 +105,63 @@ class ReadOnlyGiftWrapDaoJvmTest {
|
||||
|
||||
assertEquals(1, db.giftWrapSealDao().getAllGiftWrapSeals().size)
|
||||
}
|
||||
|
||||
// --- The inbox, reopened: Phase 6 of docs/multiple-profiles.md ---
|
||||
|
||||
/** Another profile on the device, open when the wrap for us arrived. */
|
||||
private val other = KeyPair()
|
||||
|
||||
/**
|
||||
* The case the sweep exists for. A wrap for us, fetched while another profile was
|
||||
* open: stored, found not to be addressed to the active key, and left -- and
|
||||
* `storeNostrEvent` never indexes an event it already holds, so the live subscription
|
||||
* re-delivering it later would change nothing. The sweep runs it through the index
|
||||
* again under our key.
|
||||
*/
|
||||
@Test
|
||||
fun `a wrap fetched under another profile's key is opened by the sweep, once`() = runBlocking {
|
||||
val wrap = wrapTo(us)
|
||||
db.nostrDao().storeNostrEvent(
|
||||
nostrEvent = wrap,
|
||||
relayURL = relay,
|
||||
synchronizationRelayURLs = listOf(relay),
|
||||
level = 0,
|
||||
activeKeyPair = other,
|
||||
)
|
||||
assertEquals(listOf(wrap.id), db.giftWrapMessageDao().getAllGiftWrapMessages().map { it.id }, "stored")
|
||||
assertEquals(emptyList(), db.giftWrapSealDao().getAllGiftWrapSeals(), "and not opened")
|
||||
|
||||
// Re-delivery changes nothing: the id is known, and the DAO returns.
|
||||
db.nostrDao().storeNostrEvent(wrap, relay, listOf(relay), 0, us)
|
||||
assertEquals(emptyList(), db.giftWrapSealDao().getAllGiftWrapSeals(), "re-delivery does not reopen a known event")
|
||||
|
||||
assertEquals(1, db.nostrDao().reopenInbox(us))
|
||||
assertEquals(1, db.giftWrapSealDao().getAllGiftWrapSeals().size, "the seal is out")
|
||||
assertEquals(1, db.giftWrapPayloadDao().getAllGiftWrapUnsigned().size, "and the payload with it")
|
||||
|
||||
assertEquals(0, db.nostrDao().reopenInbox(us), "nothing left to open")
|
||||
assertEquals(1, db.giftWrapSealDao().getAllGiftWrapSeals().size)
|
||||
}
|
||||
|
||||
/** The other profile's sweep finds nothing of ours to open, and a read-only pair opens nothing. */
|
||||
@Test
|
||||
fun `the sweep opens only what is addressed to the key it runs with`() = runBlocking {
|
||||
val wrap = wrapTo(us)
|
||||
db.nostrDao().storeNostrEvent(wrap, relay, listOf(relay), 0, other)
|
||||
|
||||
assertEquals(0, db.nostrDao().reopenInbox(other))
|
||||
assertEquals(0, db.nostrDao().reopenInbox(readOnlyUs))
|
||||
assertEquals(emptyList(), db.giftWrapSealDao().getAllGiftWrapSeals())
|
||||
}
|
||||
|
||||
/** The upgrade the npub plan left open: stored read-only, opened once the nsec is here. */
|
||||
@Test
|
||||
fun `a wrap stored while the identity was read-only is opened once it holds its key`() = runBlocking {
|
||||
val wrap = wrapTo(us)
|
||||
db.nostrDao().storeNostrEvent(wrap, relay, listOf(relay), 0, readOnlyUs)
|
||||
assertEquals(emptyList(), db.giftWrapSealDao().getAllGiftWrapSeals())
|
||||
|
||||
assertEquals(1, db.nostrDao().reopenInbox(us))
|
||||
assertEquals(1, db.giftWrapSealDao().getAllGiftWrapSeals().size)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user