// SPDX-License-Identifier: MIT // Copyright (c) 2026 Thibault Ducray // // This file is part of MyPwdTool's open-source sync/encryption core — see // LICENSE-SYNC-CRYPTO.md at the repo root and https://mypwdtool.com/open-source/. // The rest of this application is proprietary and NOT covered by this license. package fr.tducray.mypwdtool.sync import fr.tducray.mypwdtool.crypto.CryptoManager import fr.tducray.mypwdtool.relay.Base64Url import fr.tducray.mypwdtool.relay.OutboundRelayMessage import fr.tducray.mypwdtool.relay.RelayAPIClient import fr.tducray.mypwdtool.relay.RelayError import kotlin.random.Random import fr.tducray.mypwdtool.vault.Entry import fr.tducray.mypwdtool.vault.EntryType import fr.tducray.mypwdtool.vault.PasswordHistoryEntry import fr.tducray.mypwdtool.vault.VaultInstantFormat import fr.tducray.mypwdtool.vault.VaultStorage import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.currentCoroutineContext import kotlinx.coroutines.delay import kotlinx.coroutines.isActive import kotlinx.coroutines.launch import kotlinx.serialization.json.Json import java.time.Instant import java.util.UUID /** * Port of the join-only subset of `Sync/SyncManager.swift` (see the "not ported" list in * `project_windows_client_progress.md`/README — sync group *creation*, device revocation * *initiation*, and vault-transfer are Apple-only per the agreed v1 scope; this client only * ever joins an existing group, but must correctly receive+apply device_hello/device_leave/ * device_revoke/key_rotate sent by any device in the group). **Exception, added 2026-09-08**: * [createGroup] exists for a tester-only Android build flavor (`BuildConfig. * ENABLE_TESTER_GROUP_CREATION`) that needs to exercise the create/invite/revoke side of sync * during Play Console closed testing, ahead of Play Billing existing — see that function's own * doc comment. * * Requires the vault to be unlocked — [dek] is the vault's data-encryption key, needed both to * decrypt entry fields before packing them into outbound sync payloads and to encrypt fields * from inbound payloads before writing them to [db]. Caller is responsible for calling * [stopEngine] on lock and constructing a fresh [SyncManager] (or providing an updated DEK) on * unlock — this class doesn't watch vault-lock state itself the way the Swift `@MainActor` * singleton observes `MasterPasswordManager.$isUnlocked`. */ class SyncManager( private val db: VaultStorage, private val store: SyncStore, private val dek: ByteArray, private val platform: String, private val clientVersion: String, private val outbox: SyncOutbox = SyncOutbox(), private val scope: CoroutineScope = CoroutineScope(SupervisorJob()), /** Invoked after any inbound change actually lands in [db] (upsert/delete apply, or * self-revocation wiping the vault) — the JVM-only equivalent of Swift's `.databaseDidOpen` * notification. Called from whatever coroutine dispatcher the poll loop runs on, not * necessarily a UI thread; Compose's snapshot state system tolerates writes from background * threads, but other UI toolkits calling into this may need to hop threads themselves. */ private val onVaultChanged: () -> Unit = {}, /** Invoked when this device is remotely revoked from the group, after the vault has been * wiped and sync credentials cleared — matches Swift's `performSelfRevocation()` calling * `MasterPasswordManager.shared.lock()` as its step 3. Without this, the wiped (now-empty) * vault stayed unlocked and browsable until the user manually locked it — reproduced live * (2026-09-05): revoking an Android device from macOS correctly deleted its entries but * left it sitting unlocked. UI wires this to the same lock action as the manual "Lock * vault" button (clears the in-memory DEK). */ private val onSelfRevoked: () -> Unit = {}, /** Optional — records join/leave/revoke/key-rotation/undecryptable-message events to the * local activity log, mirroring `SecurityEventLogger` calls scattered through * `Sync/SyncManager.swift`. Nullable so existing tests don't need to construct one. */ private val securityLogger: fr.tducray.mypwdtool.security.SecurityEventLogger? = null, ) { companion object { // Matches Sync/SyncStore.swift's hardcoded serverURL — not persisted, same reasoning: // this is a fixed, developer-hosted relay, not a user setting. const val SERVER_URL = "https://relay.mypwdtool.com/relay/v3" // Long-poll hold duration passed as `waitMs` to the relay — matches Swift's // `longPollWaitMs`. Switched from short-polling (wait_ms=0 + a local delay(), previously // split 60s desktop / 120s Android) 2026-10-01, once the relay (a Go service, not a // thread-per-connection server) added real long-poll support (MAX_WAIT_MS=60000) — // holding an idle connection open costs it about the same as a closed one, so the old // per-platform battery-vs-latency split no longer applies either; same value everywhere. private const val LONG_POLL_WAIT_MS = 55_000 // Short-hold long poll used during known-active bursts — matches Swift's // `fastPollWaitMs`. Still long-polling, just with a short hold, so a burst feels // responsive without an artificial local delay. private const val FAST_POLL_WAIT_MS = 2_000 private const val FAST_POLL_DURATION_MS = 30_000L private const val HEARTBEAT_INTERVAL_MS = 30 * 60_000L private const val SEND_DELAY_MS = 60L private const val POLL_BATCH_LIMIT = 100 } private val json = Json { ignoreUnknownKeys = true } private var client: RelayAPIClient = RelayAPIClient(SERVER_URL, platform, clientVersion) private var pollJob: Job? = null private var drainJob: Job? = null private var kickDrainJob: Job? = null private var fastPollUntil: Instant? = null private var pollFailureCount = 0 private var lastHeartbeatSentAt: Instant? = null // MARK: - Joining a group suspend fun joinGroup(pairingCodeString: String, deviceName: String) { val decoded = Base64Url.decode(pairingCodeString.trim()) ?: throw SyncError.InvalidPairingCode val pairing = try { json.decodeFromString(PairingCode.serializer(), decoded.decodeToString()) } catch (e: Exception) { throw SyncError.InvalidPairingCode } val sgkBytes = Base64Url.decode(pairing.sgk) ?: throw SyncError.InvalidPairingCode require(sgkBytes.size == 32) { "SGK must be 32 bytes" } val joinClient = RelayAPIClient(pairing.url, platform, clientVersion) val response = joinClient.registerDevice(deviceName) println("[Sync] registered device_id=${response.deviceId} inbox_id=${response.inboxId} on ${pairing.url}") client = joinClient store.streamId = pairing.streamId store.deviceId = response.deviceId store.inboxId = response.inboxId store.deviceName = deviceName store.sendToken = response.sendToken store.recvToken = response.recvToken store.sgk = sgkBytes store.isEnabled = true store.upsertPeer(SyncPeer(inboxId = pairing.inboxId, deviceName = pairing.deviceName, addedAt = Instant.now().toString())) securityLogger?.log(fr.tducray.mypwdtool.security.SecurityEventType.SYNC_GROUP_JOINED, detail = pairing.deviceName) startEngine() enableFastPoll() scope.launch { println("[Sync] sending hello + bulk sync to initiator inbox=${pairing.inboxId}") sendHello(listOf(pairing.inboxId)) bulkSyncAllEntries(listOf(pairing.inboxId)) println("[Sync] hello + bulk sync sent") } } /** "Disconnect + revoke" — unlike Swift's `disableSync()` (which only sends `device_leave` * and deregisters, leaving the SGK usable by anyone who kept a copy), this also rotates the * SGK for the remaining peers first, mirroring what `revokeDevice(peer:)` does when revoking * *another* device. Requested explicitly (2026-07-24): a voluntary leave should burn the key * the leaving device held, not just poof out of the peer list. Worth porting the same * key-rotation-on-leave behavior back into Swift's `disableSync()` for parity — currently * only this Kotlin client does it. */ /** Unlike [revokeDevice] (which must apply immediately regardless of network — you may not * trust or control the target device, so the local suppression can't wait on a round trip), * leaving the group yourself has no such urgency, and an optimistic "leave locally no matter * what" design has a real downside: if this device can never reach the relay again, the * group's SGK never gets rotated away from a key this device still holds. So self-leave * requires actually reaching the relay first: if the `device_leave` broadcast fails right now * (most commonly, offline), this device stays fully in the group — sync keeps running, local * state is untouched — and this throws [SyncError.LeaveRequiresConnection] instead of * pretending to have left. A retry is queued for the rest of *this* app session only (see * [SyncControlOutbox], in-memory only): if connectivity returns before the process dies, the * leave completes automatically in the background; if not, the next launch finds the device * still in the group, exactly as if `leaveGroup()` had never been called — the caller just * calls it again once actually connected. Changed 2026-09-19 from the previous * always-succeeds-locally design. Key rotation stays best-effort/non-blocking, same as * [revokeDevice]'s own — a lost rotation only affects whether the *remaining* peers' key * changes, not whether this device successfully left. */ suspend fun leaveGroup() { val sendToken = store.sendToken val sgk = store.sgk val streamId = store.streamId val deviceId = store.deviceId val recipients = store.peers.map { it.inboxId } if (sendToken == null || sgk == null || streamId == null || deviceId == null || recipients.isEmpty()) { // Either not configured enough to notify anyone, or no one to notify — safe to leave // immediately either way. finishLeavingGroup(sendToken) return } // Frozen (built once, not re-derived on retry) — finishLeavingGroup()'s clearAll() wipes // this device's own sync state right after attempting this, so by retry time there is no // "current state" left to re-derive from; whatever will ever be sent has to be built now. val leavePayload = SyncPayload( eventId = UUID.randomUUID().toString(), senderDeviceId = deviceId, senderCounter = store.incrementSenderCounter(), op = SyncOp.DEVICE_LEAVE.rawValue, inboxId = store.inboxId, deviceName = store.deviceName, ) val leaveMsg = buildEncryptedMessage(leavePayload, recipients, sgk, streamId, deviceId) // No queued retry here (unlike device_revoke/key_rotate) — a background retry firing once // connectivity returns can't be verified end-to-end (the relay accepting the send doesn't // guarantee the other devices actually applied it), and a false "succeeded" would // silently desync this device from the group's own view of membership with no way for the // user to notice. Reproduced live 2026-09-19: two mobile devices went offline, disabled // sync, got queued for retry, and once back online considered themselves out of the group // — but macOS still saw them as active members, and the relay still had their device // registrations. Requiring a fresh, user-visible manual retry keeps this observable and // simple instead of silently wrong. if (!attemptResend(leaveMsg, sendToken)) { println("[Sync] device_leave send failed — staying in the group, no automatic retry") throw SyncError.LeaveRequiresConnection } // Burns the key this leaving device held, same as revokeDevice does for another device — // a voluntary leave shouldn't just poof out of the peer list while still holding a usable // copy of the group's key. Best-effort: doesn't block finishing the leave below, and // store.sgk is never updated from the returned new key here (unlike revokeDevice's own // use of key rotation) since clearAll() wipes it a moment later regardless. val (rotateMsg, _) = buildKeyRotationMessage(sgk, recipients, streamId, deviceId) if (!attemptResend(rotateMsg, sendToken)) { println("[Sync] key rotation on leave failed — queued for retry") SyncControlOutbox.enqueue(SyncControlOutbox.Item("key_rotate(on leave)") { attemptResend(rotateMsg, sendToken) }) } finishLeavingGroup(sendToken) } /** The actual local-leave step, shared between [leaveGroup]'s immediate-success path and its * queued retry's eventual success path. */ private suspend fun finishLeavingGroup(sendToken: String?) { if (sendToken != null) { try { client.deregisterDevice(sendToken) } catch (e: Exception) { /* best-effort */ } } stopEngine() store.clearAll() } private fun buildKeyRotationMessage( oldSgk: ByteArray, recipients: List, streamId: String, deviceId: String, ): Pair { val newSgk = SyncCrypto.generateSGK() val payload = SyncPayload( eventId = UUID.randomUUID().toString(), senderDeviceId = deviceId, senderCounter = store.incrementSenderCounter(), op = SyncOp.KEY_ROTATE.rawValue, newSgk = Base64Url.encode(newSgk), ) // Sealed under the OLD sgk deliberately — that's the key every remaining peer still has. val msg = buildEncryptedMessage(payload, recipients, oldSgk, streamId, deviceId) return msg to newSgk } private suspend fun sendKeyRotationMessage(oldSgk: ByteArray, recipients: List, sendToken: String) { val streamId = store.streamId ?: return val deviceId = store.deviceId ?: return val (msg, newSgk) = buildKeyRotationMessage(oldSgk, recipients, streamId, deviceId) client.sendMessage(msg, sendToken) println("[Sync] sent op=${SyncOp.KEY_ROTATE} to ${recipients.size} recipient(s): $recipients") // Only switch our own key once the announcement is durably queued on the relay — matches // Swift's sendKeyRotation. leaveGroup builds its own key rotation message separately // (never touching store.sgk) since it clears the whole store right after anyway; // revokeDevice needs the update here so this device keeps syncing with the remaining // peers under the new key. store.sgk = newSgk } /** The current device's own pairing code — lets an already-joined device invite another one, * without needing a fresh relay round-trip (unlike joining, generating an invite is purely * local: it's just this device's own streamId/sgk/inboxId/deviceName, re-encoded). Port of * Swift's `SyncManager.currentPairingCode()`. * * Deliberately no entitlement check in here (mirrors [createGroup], which never had one * either) — as of 2026-09-14 every platform has a real purchase flow, so the old "one free * invite per vault without paying" compromise (`sync_one_shot_invite_used`, designed * 2026-07-30 specifically for Kotlin having no purchase flow at all — see git history) no * longer has a reason to exist. Policy now lives entirely at the UI layer, same as group * creation: the caller checks its own local entitlement (`canCreateGroup()`/ * `LicenseManager.isPurchased`/`StoreManager.isSyncPurchased`) before ever calling this, and * shows the paywall instead of calling it at all when not entitled — consistent with how * [createGroup] has always worked, and with option (c) chosen 2026-09-14: "Add a device" * stays visible and clickable regardless of entitlement, leading to the same paywall * "Create a group" already uses rather than a disabled control or a numeric "X left" label. */ fun currentPairingCode(): String { val sgk = store.sgk ?: throw SyncError.InvalidPairingCode val streamId = store.streamId ?: throw SyncError.InvalidPairingCode val inboxId = store.inboxId ?: throw SyncError.InvalidPairingCode val code = encodePairingCode(streamId, sgk, inboxId, store.deviceName ?: "") // The inviting device otherwise only discovers the new joiner once its current long-poll // hold (up to 55s) expires — the joiner itself gets a fast 2s-hold poll immediately via // joinGroup()'s own enableFastPoll(), but until now the inviter side had no equivalent, // so it looked like the joiner "wasn't seen" for up to that long after actually joining. // Starts the moment the code is shown, not whenever the dialog happens to be dismissed — // the join can complete well before that. Added 2026-09-10. // // startEngine() first is required, not just enableFastPoll() alone: the poll loop only // re-checks its fastPollUntil deadline once its *current* long-poll hold finishes (see // pollLoop), so calling enableFastPoll() while a 55s hold is already in flight would // otherwise not take effect until that hold's own up-to-55s delay elapsed anyway — same // restart-to-take-effect-now pattern joinGroup() already relies on. startEngine() enableFastPoll() return code } /** Shared by [currentPairingCode] (one-shot-gated, for inviting an *additional* device into * an already-configured group) and [createGroup] (ungated — a brand-new group's very first * code must not burn the one-shot-invite flag, which is scoped to invites sent *after* * creating/joining, not the creation act itself). */ private fun encodePairingCode(streamId: String, sgk: ByteArray, inboxId: String, deviceName: String): String { val pairing = PairingCode( url = SERVER_URL, streamId = streamId, sgk = Base64Url.encode(sgk), inboxId = inboxId, deviceName = deviceName, ) return Base64Url.encode(json.encodeToString(PairingCode.serializer(), pairing).encodeToByteArray()) } // MARK: - Creating a group (tester-flavor-only on Android; see BuildConfig.ENABLE_TESTER_GROUP_CREATION) /** Starts a brand-new sync group with this device as its first member, and returns a pairing * code another device can join with. Port of Swift's `SyncManager.createGroup(serverURL: * deviceName:)`, minus the `serverURL` parameter — unlike Swift, this client always talks to * the one fixed [SERVER_URL], so there's no arbitrary relay to plumb through. There is no * receipt/entitlement check anywhere in this function or in the relay's `registerDevice` * call — the same is true of Swift's own implementation; the real gate is purely the caller * deciding whether to expose this function's UI entry point at all (see * `BuildConfig.ENABLE_TESTER_GROUP_CREATION` and the hardcoded-expiry self-revoke in * `VaultScreen`, both Android-only). */ suspend fun createGroup(deviceName: String): String { val registration = client.registerDevice(deviceName) val sgk = SyncCrypto.generateSGK() val streamId = UUID.randomUUID().toString() store.isEnabled = true store.streamId = streamId store.deviceId = registration.deviceId store.inboxId = registration.inboxId store.deviceName = deviceName store.sendToken = registration.sendToken store.recvToken = registration.recvToken store.sgk = sgk store.createdByThisDevice = true startEngine() // Same reasoning as currentPairingCode(): the pairing code returned below is about to be // displayed to invite the group's first joiner, so this device should fast-poll for it // immediately rather than waiting up to a full normal-interval poll away. enableFastPoll() securityLogger?.log(fr.tducray.mypwdtool.security.SecurityEventType.SYNC_GROUP_CREATED) return encodePairingCode(streamId, sgk, registration.inboxId, deviceName) } /** Removes `peer` from the local group immediately and unconditionally — the local removal * is what actually enforces the revoke from this device's point of view (every outbound path * only ever sends to `store.peers`, so once `peer` is gone from that list this device never * queues anything addressed to it again, regardless of whether the relay is reachable right * now). The broadcast + key rotation below are still attempted, best-effort, since they're * what lets the target itself (self-wipe) and other peers learn about the revoke — but their * failure doesn't block the local suppression the user actually needs. Port of Swift's * `SyncManager.revokeDevice(peer:)`. */ suspend fun revokeDevice(peer: SyncPeer) { store.peers = store.peers.filterNot { it.inboxId == peer.inboxId } securityLogger?.log(fr.tducray.mypwdtool.security.SecurityEventType.DEVICE_REVOKED, detail = peer.deviceName) // Both attempts below are non-fatal to the local suppression above (already applied, // unconditionally) — they're what lets OTHER devices and the target itself learn about // the revoke. If either fails right now (most commonly: offline at this exact moment), // it's queued on SyncControlOutbox and retried automatically on its own internal timer // (see its doc comment) — no more "revoke locally regardless, with nothing else ever // finding out." Found live 2026-09-18. if (!attemptDeviceRevokeSend(peer.inboxId, peer.deviceName)) { SyncControlOutbox.enqueue(SyncControlOutbox.Item("device_revoke(${peer.deviceName})") { attemptDeviceRevokeSend(peer.inboxId, peer.deviceName) }) } if (!attemptKeyRotationSend()) { SyncControlOutbox.enqueue(SyncControlOutbox.Item("key_rotate(after revoking ${peer.deviceName})") { attemptKeyRotationSend() }) } } /** Re-derives recipients/SGK fresh from current [SyncStore] state on every call (first * attempt and every retry) — see [SyncControlOutbox]'s doc comment for why that matters * here. Always includes `targetInboxId` even though it's no longer in `store.peers`, so the * revoked device itself still receives the instruction to self-wipe once it reconnects. * Returns `true` for "nothing to send" (not configured) — not a failure worth retrying. */ private suspend fun attemptDeviceRevokeSend(targetInboxId: String, targetDeviceName: String?): Boolean { val sendToken = store.sendToken ?: return true val sgk = store.sgk ?: return true val streamId = store.streamId ?: return true val deviceId = store.deviceId ?: return true val recipients = (store.peers.map { it.inboxId } + targetInboxId).distinct() if (recipients.isEmpty()) return true val payload = SyncPayload( eventId = UUID.randomUUID().toString(), senderDeviceId = deviceId, senderCounter = store.incrementSenderCounter(), op = SyncOp.DEVICE_REVOKE.rawValue, inboxId = targetInboxId, // target's inbox — receivers remove this device deviceName = targetDeviceName, ) return try { sendPayload(payload, recipients, sgk, streamId, deviceId, sendToken) true } catch (e: CancellationException) { throw e } catch (e: Exception) { println("[Sync] device_revoke send attempt failed: $e") false } } /** Re-derives its old SGK/recipients fresh from current [SyncStore] state on every call. * Returns `true` for "nothing to rotate for" (no peers left, or not configured) — not a * failure worth retrying. */ private suspend fun attemptKeyRotationSend(): Boolean { val sendToken = store.sendToken ?: return true val oldSgk = store.sgk ?: return true val recipients = store.peers.map { it.inboxId } if (recipients.isEmpty()) return true return try { sendKeyRotationMessage(oldSgk, recipients, sendToken) true } catch (e: CancellationException) { throw e } catch (e: Exception) { println("[Sync] key rotation send attempt failed: $e") false } } // MARK: - Lifecycle fun startEngine() { if (!store.isConfigured) return stopEngine() // Pick up any entries VaultTransferService wrote into this vault while it was inactive — // see SyncStore.addPendingTransferEntryID's doc comment. val pendingIds = store.drainPendingTransferEntryIDs() if (pendingIds.isNotEmpty()) { val entries = db.fetchAllEntries().associateBy { it.id } pendingIds.forEach { id -> entries[id]?.let { enqueueUpsert(it) } } } pollJob = scope.launch { pollLoop() } drainJob = scope.launch { sendOutboxItems() } } fun stopEngine() { pollJob?.cancel(); pollJob = null drainJob?.cancel(); drainJob = null kickDrainJob?.cancel(); kickDrainJob = null } fun forceSyncNow() { if (store.isConfigured) startEngine() } /** Renames this device: updates the relay's own `device_name` record (§4.4, diagnostics * only) via PATCH /devices, and separately pushes the new name to paired peers on the next * heartbeat via the encrypted `sender_device_name` payload field — the two are unrelated on * the wire (one is a plaintext relay-visible field, the other is inside the E2EE payload * peers actually display), so both need updating for the rename to be visible everywhere. * Resetting `lastHeartbeatSentAt` and kicking the engine makes the peer-facing heartbeat go * out now instead of up to 30 minutes away (matches the same fix in Swift's * SyncManager.swift). */ fun renameDevice(newName: String) { val trimmed = newName.trim() if (trimmed.isEmpty()) return store.deviceName = trimmed lastHeartbeatSentAt = null store.sendToken?.let { sendToken -> scope.launch { try { client.renameDevice(sendToken, trimmed) } catch (e: Exception) { // Best-effort: the relay's copy is diagnostics-only, so a failed PATCH here // (offline, transient error) shouldn't block the peer-facing rename below. println("[Sync] relay device rename failed (non-fatal): $e") } } } forceSyncNow() } fun resumeIfNeeded() { if (pollJob == null) startEngine() } private fun enableFastPoll() { fastPollUntil = Instant.now().plusMillis(FAST_POLL_DURATION_MS) } // MARK: - Outbox fun enqueueUpsert(entry: Entry) { println("[Sync] enqueueUpsert entry=${entry.id} name=${entry.name} updatedAt=${entry.updatedAt}") outbox.enqueue(SyncOutbox.Item(SyncOp.UPSERT, entry.id, VaultInstantFormat.format(entry.updatedAt), entry)) kickDrain() } fun enqueueDelete(entryId: UUID, updatedAt: Instant = Instant.now()) { println("[Sync] enqueueDelete entry=$entryId") outbox.enqueue(SyncOutbox.Item(SyncOp.DELETE, entryId, VaultInstantFormat.format(updatedAt), null)) kickDrain() } // enqueueUpsert (and therefore kickDrain) can now be called from the Chrome-extension // socket server's background threads, not just the Compose UI thread/sync coroutines — this // lock keeps the cancel-then-reassign of kickDrainJob atomic against that, so two // near-simultaneous callers can't race and silently drop the scheduled drain (same class of // bug as VaultDatabase's — see its own lock's doc comment for the live-observed symptom this // general pattern causes). private val kickDrainLock = Any() private fun kickDrain() { synchronized(kickDrainLock) { kickDrainJob?.cancel() kickDrainJob = scope.launch { delay(100) sendOutboxItems() } } } private suspend fun sendOutboxItems() { // Control messages (device_revoke, device_leave, hello, peer introductions) are no // longer drained here — SyncControlOutbox is a self-draining app-lifetime singleton (see // its own doc comment for why: this SyncManager instance, and everything a queued retry's // closure captures, must survive being GC'd once the vault locks and this instance is // discarded). val sendToken = store.sendToken ?: run { println("[Sync] sendOutboxItems: no sendToken, skipping drain (items stay queued)") return } if (store.sgk == null) { println("[Sync] sendOutboxItems: no sgk, skipping drain (items stay queued)") return } val recipients = store.peers.filterNot { it.isSuspended }.map { it.inboxId } if (recipients.isEmpty()) { println("[Sync] sendOutboxItems: no non-suspended peers, skipping drain (items stay queued)") return } sendItems(outbox.drain(), recipients, sendToken) } private suspend fun bulkSyncAllEntries(recipientInboxIds: List) { val sendToken = store.sendToken ?: return if (recipientInboxIds.isEmpty()) return val items = db.fetchAllEntries().map { SyncOutbox.Item(SyncOp.UPSERT, it.id, VaultInstantFormat.format(it.updatedAt), it) } sendItems(items, recipientInboxIds, sendToken) } private suspend fun sendItems(items: List, recipients: List, sendToken: String) { val sgk = store.sgk ?: return val streamId = store.streamId ?: return val deviceId = store.deviceId ?: return for (item in items) { try { val payload = buildOutboundPayload(item, deviceId) sendPayload(payload, recipients, sgk, streamId, deviceId, sendToken) delay(SEND_DELAY_MS) } catch (e: CancellationException) { throw e } catch (e: Exception) { // Previously silent — a permanently-failing item (e.g. a decrypt failure on a // corrupted field) would re-enqueue here forever with zero diagnostic output, // which is exactly what made a real live bug (2026-07-29: an entry updated via // the Chrome extension stopped syncing, with "no sync log in IntelliJ, no // message on the relay") impossible to diagnose after the fact. println("[Sync] send failed for entry=${item.entryId} op=${item.op}, re-enqueuing: $e") outbox.enqueue(item) } } } private fun buildOutboundPayload(item: SyncOutbox.Item, deviceId: String): SyncPayload { val vaultId = db.vaultId?.toString() ?: "unknown-vault" val entry = item.entry fun dec(fieldName: String, blob: String?): String? = blob?.let { try { CryptoManager.decrypt(it, dek, CryptoManager.aad(vaultId, item.entryId.toString(), fieldName)) } catch (e: Exception) { println("[Sync] decrypt failed for entry=${item.entryId} field=$fieldName: $e") throw e } } return SyncPayload( eventId = UUID.randomUUID().toString(), senderDeviceId = deviceId, senderDeviceName = store.deviceName, // Authoritative identity for last-seen matching (added 2026-07-27) — see // touchPeerLastSeen: unlike senderDeviceId (relay-assigned, observed live to drift // out of sync with a peer's stored copy), inboxId is the same identifier that // already, correctly, routes every message, so it can't drift. inboxId = store.inboxId, senderCounter = store.incrementSenderCounter(), op = item.op.rawValue, entryId = item.entryId.toString(), updatedAt = item.updatedAt, createdAt = entry?.let { VaultInstantFormat.format(it.createdAt) }, name = entry?.name, username = entry?.username, websites = entry?.websites, password = entry?.let { dec("password", it.encryptedPassword) }, previousPassword = entry?.let { dec("previousPassword", it.encryptedPreviousPassword) }, // Send the complete history, not just this change — see SyncPasswordHistoryItem. // Well within the relay's 262144-byte message cap even for a long history. passwordHistory = entry?.let { e -> db.fetchPasswordHistory(e.id).mapNotNull { h -> val pwd = runCatching { CryptoManager.decrypt(h.encryptedPassword, dek, CryptoManager.aad(vaultId, e.id.toString(), "password")) }.getOrNull() ?: return@mapNotNull null SyncPasswordHistoryItem( id = h.id.toString(), changedAt = VaultInstantFormat.format(h.changedAt), password = pwd, ) } }, note = entry?.let { dec("note", it.encryptedNote) }, cardNumber = entry?.let { dec("cardNumber", it.encryptedCardNumber) }, cardExpiry = entry?.let { dec("cardExpiry", it.encryptedCardExpiry) }, cardCvv = entry?.let { dec("cardCvv", it.encryptedCardCVV) }, cardPin = entry?.let { dec("cardPin", it.encryptedCardPin) }, totpSecret = entry?.let { dec("totpSecret", it.encryptedTOTPSecret) }, entryType = entry?.entryType?.rawValue, isFavorite = entry?.isFavorite, deletedAt = entry?.deletedAt?.let { VaultInstantFormat.format(it) }, passkeyRpId = entry?.passkeyRelyingPartyId, passkeyCredentialId = entry?.let { dec("passkeyCredentialId", it.encryptedPasskeyCredentialId) }, passkeyUserHandle = entry?.let { dec("passkeyUserHandle", it.encryptedPasskeyUserHandle) }, ) } private fun buildEncryptedMessage( payload: SyncPayload, recipients: List, sgk: ByteArray, streamId: String, deviceId: String, ): OutboundRelayMessage { val plaintext = json.encodeToString(SyncPayload.serializer(), payload).toByteArray() val messageId = UUID.randomUUID().toString() val ciphertext = SyncCrypto.encrypt(plaintext, sgk, messageId, streamId, deviceId, recipients) return OutboundRelayMessage( messageId = messageId, streamId = streamId, senderDeviceId = deviceId, recipientInboxIds = recipients, ciphertext = ciphertext, ) } private suspend fun sendPayload( payload: SyncPayload, recipients: List, sgk: ByteArray, streamId: String, deviceId: String, sendToken: String, ) { val msg = buildEncryptedMessage(payload, recipients, sgk, streamId, deviceId) client.sendMessage(msg, sendToken) println("[Sync] sent op=${payload.op} to ${recipients.size} recipient(s): $recipients") } /** Resends an already-built message as-is — used for control messages whose content doesn't * depend on group state that could keep changing after the first attempt (device_leave and * its paired key rotation), unlike device_revoke/key_rotate-after-revoke/hello/peer * introductions, which re-derive fresh on every attempt instead (see * [attemptDeviceRevokeSend]/[attemptKeyRotationSend]/[attemptProtocolMessage]). */ private suspend fun attemptResend(msg: OutboundRelayMessage, sendToken: String): Boolean = try { client.sendMessage(msg, sendToken) true } catch (e: CancellationException) { throw e } catch (e: Exception) { println("[Sync] resend attempt failed: $e") false } /** Re-derives sgk/streamId/deviceId/sendToken/counter fresh from current [SyncStore] state on * every call (first attempt and every retry) — `recipients`/`op`/`configure` are frozen by * the caller, since for hello/peer-introductions they're already a fixed, specific target, * not "whatever the group currently looks like" (unlike revoke/key-rotation's own recipient * lists). Returns `true` for "not configured"/"no recipients" too — that's not a failure * worth retrying. */ private suspend fun attemptProtocolMessage( op: SyncOp, recipients: List, configure: SyncPayload.() -> SyncPayload, ): Boolean { val sgk = store.sgk ?: return true val streamId = store.streamId ?: return true val deviceId = store.deviceId ?: return true val sendToken = store.sendToken ?: return true if (recipients.isEmpty()) return true val base = SyncPayload( eventId = UUID.randomUUID().toString(), senderDeviceId = deviceId, senderDeviceName = store.deviceName, senderCounter = store.incrementSenderCounter(), op = op.rawValue, ).configure() return try { sendPayload(base, recipients, sgk, streamId, deviceId, sendToken) true } catch (e: CancellationException) { throw e } catch (e: Exception) { println("[Sync] $op send attempt failed: $e") false } } /** [retryLabel] gates whether a failed attempt gets queued on [SyncControlOutbox] — `null` (only * used by heartbeat) means "don't retry," per explicit direction: a heartbeat is a pure * liveness ping, superseded by the next one 30 minutes later regardless. */ private suspend fun sendProtocolMessage( op: SyncOp, recipients: List, retryLabel: String? = null, configure: SyncPayload.() -> SyncPayload = { this }, ) { val delivered = attemptProtocolMessage(op, recipients, configure) if (!delivered && retryLabel != null) { SyncControlOutbox.enqueue(SyncControlOutbox.Item(retryLabel) { attemptProtocolMessage(op, recipients, configure) }) } } private suspend fun sendHello(recipients: List) { sendProtocolMessage(SyncOp.DEVICE_HELLO, recipients, retryLabel = "device_hello") { copy(inboxId = store.inboxId, deviceName = store.deviceName) } } private suspend fun sendHeartbeat(recipients: List) { sendProtocolMessage(SyncOp.HEARTBEAT, recipients) { copy(inboxId = store.inboxId) } } /** Announces `announcedInboxId`/`announcedDeviceName` (a peer) to `recipients`. */ private suspend fun sendPeerIntroduction(recipients: List, announcedInboxId: String, announcedDeviceName: String) { sendProtocolMessage(SyncOp.DEVICE_HELLO, recipients, retryLabel = "peer_introduction($announcedDeviceName)") { copy(inboxId = announcedInboxId, deviceName = announcedDeviceName) } } // MARK: - Poll loop private suspend fun pollLoop() { println("[Sync] poll loop started (long-poll ${LONG_POLL_WAIT_MS}ms normal / ${FAST_POLL_WAIT_MS}ms fast)") while (currentCoroutineContext().isActive) { val fastPoll = fastPollUntil?.isAfter(Instant.now()) == true try { pollOnce(waitMs = if (fastPoll) FAST_POLL_WAIT_MS else LONG_POLL_WAIT_MS) } catch (e: CancellationException) { throw e } catch (e: RelayError.TokenExpired) { println("[Sync] poll stopped: token expired") break } catch (e: RelayError.Unauthorized) { // Needs re-registration (re-join) — stop polling until that happens. println("[Sync] poll stopped: unauthorized (needs re-join)") break } catch (e: Exception) { val retryMs = pollRetryDelayMs(e) println("[Sync] poll failed, retrying in ${retryMs}ms: $e") delay(retryMs) continue } pollFailureCount = 0 val now = Instant.now() if (lastHeartbeatSentAt == null || now.isAfter(lastHeartbeatSentAt!!.plusMillis(HEARTBEAT_INTERVAL_MS))) { val recipients = store.peers.filterNot { it.isSuspended }.map { it.inboxId } if (recipients.isNotEmpty()) sendHeartbeat(recipients) lastHeartbeatSentAt = now } // Safety-net retry — see SyncOutbox's doc comment for why this exists: without it, an // outbox item whose very first drain attempt failed its guards (e.g. no peers yet // right after a fresh join) would sit stuck forever, since nothing else periodically // retries it. if (!outbox.isEmpty) sendOutboxItems() // No delay() here: the long poll just made in pollOnce() already held the connection // open for up to `waitMs`, which is what used to pace this loop — the next // iteration's hold does that job now. } } /** Honours the relay's own Retry-After on 429; otherwise backs off exponentially from 30s to a * 5-min cap, with ±20% jitter so devices that all lost the relay together don't retry in * lockstep. Reset by the next successful poll. */ private fun pollRetryDelayMs(e: Exception): Long { if (e is RelayError.RateLimited) return e.retryAfterSeconds * 1000L pollFailureCount++ val base = minOf(30_000L * (1L shl minOf(pollFailureCount - 1, 4)), 300_000L) return (base * (0.8 + Random.nextDouble() * 0.4)).toLong() } private suspend fun pollOnce(waitMs: Int) { val recvToken = store.recvToken ?: return val inboxId = store.inboxId ?: return val response = client.pollMessages(recvToken, inboxId, store.cursor, POLL_BATCH_LIMIT, waitMs) if (response.messages.isEmpty()) return println("[Sync] poll: received ${response.messages.size} message(s)") var allApplied = true val toAck = mutableListOf() for (msg in response.messages) { if (store.hasApplied(msg.messageId)) { toAck.add(msg.messageId) touchPeerLastSeen(msg.senderDeviceId, senderInboxId = null, senderDeviceName = null) continue } val plaintext = try { SyncCrypto.decrypt(msg.ciphertext, store.sgk ?: byteArrayOf(), msg.messageId, msg.streamId, msg.senderDeviceId, msg.recipientInboxIds) } catch (e: Exception) { // Permanent (AEAD is deterministic) — most likely SGK rotated since this was // sent. ACK it to clear it from the relay rather than retrying forever. securityLogger?.log(fr.tducray.mypwdtool.security.SecurityEventType.SYNC_MESSAGE_UNDECRYPTABLE) toAck.add(msg.messageId) continue } val payload = try { json.decodeFromString(SyncPayload.serializer(), plaintext.decodeToString()) } catch (e: Exception) { toAck.add(msg.messageId) continue } // Allow-list, not a deny-list: inboxId only reliably means "the sender's own inbox" // for heartbeat/upsert/delete (the three ops it's populated for in // buildOutboundPayload/sendHeartbeat). device_revoke repurposes it as the TARGET // being revoked; device_hello can be a *relayed* introduction naming a third peer, // not whoever actually relayed the message — trusting it there would credit the // wrong peer's last-seen. val senderOp = SyncOp.fromRawValue(payload.op) val senderInboxId = if (senderOp == SyncOp.HEARTBEAT || senderOp == SyncOp.UPSERT || senderOp == SyncOp.DELETE) payload.inboxId else null if (store.hasApplied(payload.eventId)) { toAck.add(msg.messageId) touchPeerLastSeen(msg.senderDeviceId, senderInboxId, payload.senderDeviceName) continue } val applied = applyPayload(payload) if (applied) { store.markApplied(payload.eventId) store.markApplied(msg.messageId) touchPeerLastSeen(msg.senderDeviceId, senderInboxId, payload.senderDeviceName) toAck.add(msg.messageId) } else { allApplied = false } } if (toAck.isNotEmpty()) client.ackMessages(recvToken, inboxId, toAck) if (allApplied && response.nextCursor.isNotEmpty()) store.cursor = response.nextCursor } /** Called on *every* ACKed message, including the two dedupe-and-skip paths above (a * message already marked applied, e.g. a redelivery after a prior ACK attempt failed) — * not just freshly-applied ones. A dedupe-skipped message still proves the sender is alive * and reachable right now; skipping the update there let a peer show a stale "last seen" * timestamp days in the past even while its redelivered messages were being ACKed * successfully in the present (matches the same fix in Swift's SyncManager.swift). * * `senderInboxId` is the authoritative match — it's the same identifier that already, * correctly, routes every message, so unlike `relayDeviceId` it can't drift. `relayDeviceId`, * then `senderDeviceName`, are fallbacks only for messages that predate `inboxId` being * carried on heartbeat/upsert/delete — matching by name risks a false match if two peers * share a display name, which `inboxId` can't, since it's relay-assigned and unique. */ private fun touchPeerLastSeen(relayDeviceId: String, senderInboxId: String?, senderDeviceName: String?) { var peer = senderInboxId?.let { id -> store.peers.firstOrNull { it.inboxId == id } } val matchedByInboxId = peer != null if (peer == null) peer = store.peers.firstOrNull { it.relayDeviceId == relayDeviceId } if (peer == null && senderDeviceName != null) { peer = store.peers.firstOrNull { it.deviceName == senderDeviceName } } if (peer == null) { println("[Sync] last-seen update skipped: no peer matched senderDeviceId=$relayDeviceId senderInboxId=$senderInboxId senderDeviceName=$senderDeviceName — known peers: ${store.peers.map { Triple(it.deviceName, it.inboxId, it.relayDeviceId) }}") return } // A device can rename itself locally (see AppSettings/Main.kt) and has no other way to // tell its peers — it just rides along on whatever it next sends, heartbeat included, in // the already-optional `sender_device_name` field. Older peer builds that never populate // this field pass `null` here, which correctly leaves the stored name untouched (matches // the same fix in Swift's SyncManager.swift). Only applied when matched via `inboxId` — // the one unique, relay-assigned identifier that can't collide. Live-tested 2026-09-02: // pairing an old build (relayDeviceId not yet reliably populated) hit the `relayDeviceId` // fallback match here against the WRONG existing peer, overwriting that peer's name with // the new device's name — exactly the false-match risk already called out above for name // overwriting, just not previously guarded against. val newName = senderDeviceName?.takeIf { matchedByInboxId && it.isNotEmpty() && it != peer.deviceName } store.upsertPeer(peer.copy( lastSeenAt = Instant.now().toString(), relayDeviceId = relayDeviceId, deviceName = newName ?: peer.deviceName )) } // MARK: - Applying inbound payloads internal suspend fun applyPayload(payload: SyncPayload): Boolean { println("[Sync] applying op=${payload.op} from sender_device_id=${payload.senderDeviceId}") return when (SyncOp.fromRawValue(payload.op)) { SyncOp.UPSERT -> applyUpsert(payload) SyncOp.DELETE -> applyDelete(payload) SyncOp.DEVICE_HELLO -> applyDeviceHello(payload) SyncOp.DEVICE_LEAVE -> applyDeviceLeave(payload) SyncOp.DEVICE_REVOKE -> applyDeviceRevoke(payload) SyncOp.KEY_ROTATE -> applyKeyRotate(payload) SyncOp.HEARTBEAT, null -> true } } internal fun applyUpsert(payload: SyncPayload): Boolean { val entryId = payload.entryId?.let { runCatching { UUID.fromString(it) }.getOrNull() } ?: return true val updatedAt = payload.updatedAt?.let { runCatching { VaultInstantFormat.parse(it) }.getOrNull() } ?: Instant.now() // Falls back to a fixed placeholder when the vault row doesn't exist yet (matches // buildOutboundPayload's fallback) rather than dropping the update — vaultId only scopes // AAD, it isn't a correctness requirement for applying an inbound entry. val vaultId = db.vaultId?.toString() ?: "unknown-vault" fun enc(fieldName: String, plain: String?): String? = plain?.let { CryptoManager.encrypt(it, dek, CryptoManager.aad(vaultId, entryId.toString(), fieldName)) } val resolvedWebsites = payload.websites ?: payload.website?.let { listOf(it) } val existing = db.fetchEntry(entryId) if (existing != null) { if (!updatedAt.isAfter(existing.updatedAt)) return true // stale — already applied a newer version val updated = existing.copy( name = payload.name ?: existing.name, username = payload.username ?: existing.username, websites = resolvedWebsites ?: existing.websites, encryptedPassword = enc("password", payload.password) ?: existing.encryptedPassword, encryptedPreviousPassword = enc("previousPassword", payload.previousPassword) ?: existing.encryptedPreviousPassword, encryptedNote = enc("note", payload.note) ?: existing.encryptedNote, encryptedCardNumber = enc("cardNumber", payload.cardNumber) ?: existing.encryptedCardNumber, encryptedCardExpiry = enc("cardExpiry", payload.cardExpiry) ?: existing.encryptedCardExpiry, encryptedCardCVV = enc("cardCvv", payload.cardCvv) ?: existing.encryptedCardCVV, encryptedCardPin = enc("cardPin", payload.cardPin) ?: existing.encryptedCardPin, encryptedTOTPSecret = enc("totpSecret", payload.totpSecret) ?: existing.encryptedTOTPSecret, passkeyRelyingPartyId = payload.passkeyRpId ?: existing.passkeyRelyingPartyId, encryptedPasskeyCredentialId = enc("passkeyCredentialId", payload.passkeyCredentialId) ?: existing.encryptedPasskeyCredentialId, encryptedPasskeyUserHandle = enc("passkeyUserHandle", payload.passkeyUserHandle) ?: existing.encryptedPasskeyUserHandle, entryType = payload.entryType?.let { EntryType.fromRawValue(it) } ?: existing.entryType, isFavorite = payload.isFavorite ?: existing.isFavorite, // Direct replacement, not a fallback-to-existing like the fields above: a // restore needs to explicitly clear this to null, which `?: existing.deletedAt` // could never do (nil-coalescing can't distinguish "absent" from "restored"). deletedAt = payload.deletedAt?.let { runCatching { VaultInstantFormat.parse(it) }.getOrNull() }, updatedAt = updatedAt, ) db.updateEntry(updated) onVaultChanged() } else { val entry = Entry( id = entryId, name = payload.name ?: "", username = payload.username, websites = resolvedWebsites ?: emptyList(), encryptedPassword = enc("password", payload.password), encryptedPreviousPassword = enc("previousPassword", payload.previousPassword), encryptedNote = enc("note", payload.note), encryptedCardNumber = enc("cardNumber", payload.cardNumber), encryptedCardExpiry = enc("cardExpiry", payload.cardExpiry), encryptedCardCVV = enc("cardCvv", payload.cardCvv), encryptedCardPin = enc("cardPin", payload.cardPin), encryptedTOTPSecret = enc("totpSecret", payload.totpSecret), passkeyRelyingPartyId = payload.passkeyRpId, encryptedPasskeyCredentialId = enc("passkeyCredentialId", payload.passkeyCredentialId), encryptedPasskeyUserHandle = enc("passkeyUserHandle", payload.passkeyUserHandle), entryType = payload.entryType?.let { EntryType.fromRawValue(it) } ?: EntryType.LOGIN, isFavorite = payload.isFavorite ?: false, createdAt = payload.createdAt?.let { runCatching { VaultInstantFormat.parse(it) }.getOrNull() } ?: updatedAt, updatedAt = updatedAt, deletedAt = payload.deletedAt?.let { runCatching { VaultInstantFormat.parse(it) }.getOrNull() }, ) db.insertEntry(entry) onVaultChanged() } // Merge the sender's full history rather than just this one change, so a device that // joined the group late (or was offline for several password changes) still ends up with // every entry another device already has. insertPasswordHistory is INSERT OR IGNORE on // id (see VaultDatabase), so re-applying the same message is a no-op. payload.passwordHistory?.forEach { item -> val historyId = runCatching { UUID.fromString(item.id) }.getOrNull() ?: return@forEach val encPwd = enc("password", item.password) ?: return@forEach val changedAt = runCatching { VaultInstantFormat.parse(item.changedAt) }.getOrNull() ?: Instant.now() db.insertPasswordHistory(PasswordHistoryEntry(id = historyId, entryId = entryId, encryptedPassword = encPwd, changedAt = changedAt)) } return true } internal fun applyDelete(payload: SyncPayload): Boolean { val entryId = payload.entryId?.let { runCatching { UUID.fromString(it) }.getOrNull() } ?: return true db.deleteEntry(entryId) onVaultChanged() return true } private suspend fun applyDeviceHello(payload: SyncPayload): Boolean { val newInboxId = payload.inboxId ?: return true if (newInboxId in store.revokedInboxIds) return true val newDeviceName = payload.deviceName ?: "" val isNewPeer = store.peers.none { it.inboxId == newInboxId } store.upsertPeer(SyncPeer(inboxId = newInboxId, deviceName = newDeviceName, addedAt = Instant.now().toString())) enableFastPoll() if (isNewPeer && newInboxId != store.inboxId) { bulkSyncAllEntries(listOf(newInboxId)) // Full-mesh bootstrap: a joiner's pairing code only names the initiator, so relay // introductions both ways so every device eventually learns about every peer. val otherPeers = store.peers.filter { it.inboxId != newInboxId } for (peer in otherPeers) { sendPeerIntroduction(listOf(peer.inboxId), newInboxId, newDeviceName) sendPeerIntroduction(listOf(newInboxId), peer.inboxId, peer.deviceName) } } return true } private fun applyDeviceLeave(payload: SyncPayload): Boolean { val inboxId = payload.inboxId ?: return true store.removePeer(inboxId) securityLogger?.log(fr.tducray.mypwdtool.security.SecurityEventType.PEER_LEFT_SYNC_GROUP, detail = payload.deviceName) return true } private suspend fun applyDeviceRevoke(payload: SyncPayload): Boolean { // Swift's revokeDevice(peer:) reuses the device_hello fields (inbox_id/device_name) to // carry the revoked device's identity for device_revoke too — there is no dedicated // "target" field on the wire. val targetInboxId = payload.inboxId ?: return true if (targetInboxId == store.inboxId) { performSelfRevocation() } else { store.removePeer(targetInboxId) } store.markRevoked(targetInboxId) return true } private fun applyKeyRotate(payload: SyncPayload): Boolean { val newSgk = payload.newSgk?.let { Base64Url.decode(it) } ?: return true if (newSgk.size != 32) return true store.sgk = newSgk securityLogger?.log(fr.tducray.mypwdtool.security.SecurityEventType.SYNC_KEY_ROTATED) return true } private suspend fun performSelfRevocation() { db.deleteAllEntries() securityLogger?.log(fr.tducray.mypwdtool.security.SecurityEventType.SELF_REVOKED) val sendToken = store.sendToken if (sendToken != null) { try { client.deregisterDevice(sendToken) } catch (e: Exception) { /* best-effort */ } } stopEngine() store.clearAll() // Fired after clearAll(), not before — onVaultChanged is also what the UI uses to // refresh its mirrored "is sync configured" flag (see desktopApp's isSyncConfigured); // firing it while the store was still fully configured (only entries had been wiped so // far) left that flag stuck stale until some unrelated recomposition happened to catch up // — confirmed live: the sidebar cloud icon and its dialog kept showing the pre-revoke // configured state (stale device name, stale peer list) until the dialog was closed. onVaultChanged() onSelfRevoked() } }