Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 32 additions & 10 deletions app/src/main/java/to/bitkit/repositories/ActivityRepo.kt
Original file line number Diff line number Diff line change
Expand Up @@ -986,19 +986,41 @@ class ActivityRepo @Inject constructor(
}
}

/**
* Applies each slice of the backup envelope on its own so one rejected record cannot discard the others.
* Core fails a bulk write as a whole, so applying all three slices together would let a single unusable tag cost
* the activities and the closed channels too. The overall result still fails when any slice failed, keeping
* [BackupRepo] from treating a partial restore as authoritative and rewriting a good backup with it.
* Observers are notified whenever at least one slice was applied, and the tag signal only fires when the
* tags slice itself was applied, so a rejected tags slice cannot mark the metadata backup as changed.
*/
suspend fun restoreFromBackup(payload: ActivityBackupV1): Result<Unit> = withContext(bgDispatcher) {
runCatching {
coreService.activity.upsertList(payload.activities)
coreService.activity.upsertTags(payload.activityTags)
val activities = runSuspendCatching { coreService.activity.upsertList(payload.activities) }
val activityTags = runSuspendCatching { coreService.activity.upsertTags(payload.activityTags) }
val closedChannels = runSuspendCatching {
coreService.activity.upsertClosedChannelList(payload.closedChannels)
}.onSuccess {
Logger.debug(
"Restored ${payload.activities.size} activities, ${payload.activityTags.size} activity tags, " +
"${payload.closedChannels.size} closed channels",
context = TAG,
)
notifyActivitiesChanged(tagsChanged = true)
}
val results = listOf(
"activities" to activities,
"activityTags" to activityTags,
"closedChannels" to closedChannels,
)
val failures = results.mapNotNull { (slice, result) ->
result.exceptionOrNull()?.also {
Logger.error("Failed to restore '$slice' activity backup slice", it, context = TAG)
}
}

if (failures.size < results.size) notifyActivitiesChanged(tagsChanged = activityTags.isSuccess)

failures.firstOrNull()?.let { return@withContext Result.failure(it) }

Logger.debug(
"Restored ${payload.activities.size} activities, ${payload.activityTags.size} activity tags, " +
"${payload.closedChannels.size} closed channels",
context = TAG,
)
return@withContext Result.success(Unit)
}

suspend fun markAllUnseenActivitiesAsSeen(): Result<Unit> = withContext(bgDispatcher) {
Expand Down
30 changes: 28 additions & 2 deletions app/src/main/java/to/bitkit/repositories/BackupRepo.kt
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ import javax.inject.Inject
import javax.inject.Provider
import javax.inject.Singleton
import kotlin.time.Clock
import kotlin.time.Duration.Companion.minutes
import kotlin.time.Duration.Companion.seconds
import kotlin.time.ExperimentalTime

Expand Down Expand Up @@ -124,15 +125,34 @@ class BackupRepo @Inject constructor(
private val _isWiping = MutableStateFlow(false)
val isWiping: StateFlow<Boolean> = _isWiping.asStateFlow()

private val restorePendingUntil = MutableStateFlow(0L)

fun reset() {
stopObservingBackups()
vssBackupClient.reset()
vssBackupClientLdk.reset()
}

fun setWiping(isWiping: Boolean) = _isWiping.update { isWiping }

/**
* Holds ordinary uploads from the moment a restore is announced, closing the window between the
* restore flow starting and [performFullRestoreFromLatestBackup] raising [_isRestoring]: the node
* starts and syncs while the backup is still being read, and the resulting activity upload would
* replace the stored envelope with the fresh wallet's state before it is read.
*
* The gate cannot suppress uploads indefinitely: it expires after [RESTORE_PENDING_TIMEOUT_MS], it is
* held in memory only so a process death clears it, and it never blocks an explicit
* [triggerBackup], including the migration rewrite a restore ends with.
*/
fun setRestorePending(isPending: Boolean) {
restorePendingUntil.update { if (isPending) currentTimeMillis() + RESTORE_PENDING_TIMEOUT_MS else 0L }
Logger.debug("Set restore pending to '$isPending'", context = TAG)
}

private fun currentTimeMillis(): Long = nowMillis(clock)
private fun shouldSkipBackup(): Boolean = _isRestoring.value || _isWiping.value
private fun isRestorePending(): Boolean = currentTimeMillis() < restorePendingUntil.value
private fun shouldSkipBackup(): Boolean = _isRestoring.value || _isWiping.value || isRestorePending()
private fun BackupItemStatus.shouldBackup(category: BackupCategory) =
this.isRequired &&
!this.running &&
Expand Down Expand Up @@ -708,7 +728,7 @@ class BackupRepo @Inject constructor(
)
val parsed = json.decodeFromString<ActivityBackupV1>(migration.json)
val persisted = activityRepo.restoreFromBackup(parsed)
.onFailure { Logger.warn("Failed to restore activity backup", it, context = TAG) }
.onFailure { Logger.warn("Skipped activity backup rewrite after a failed restore", context = TAG) }
.isSuccess

return RestoredCoreBackup(createdAt = parsed.createdAt, needsRewrite = migration.changed && persisted)
Expand Down Expand Up @@ -864,6 +884,12 @@ class BackupRepo @Inject constructor(
private const val FAILED_BACKUP_NOTIFICATION_INTERVAL = 10 * 60 * 1000L // 10 minutes
private const val SYNC_STATUS_DEBOUNCE = 500L // 500ms debounce for sync status updates
private val VSS_TIMESTAMP_TIMEOUT = 60.seconds

/**
* How long a pending restore gates ordinary uploads for. Longer than any restore in practice,
* short enough that a restore that never returns cannot hold the gate for the session.
*/
private val RESTORE_PENDING_TIMEOUT_MS = 10.minutes.inWholeMilliseconds
}
}

Expand Down
6 changes: 6 additions & 0 deletions app/src/main/java/to/bitkit/viewmodels/WalletViewModel.kt
Original file line number Diff line number Diff line change
Expand Up @@ -219,6 +219,7 @@ class WalletViewModel @Inject constructor(
Logger.error("Restore from backup failed", it, context = TAG)
}
_restoreState.update { RestoreState.Completed }
backupRepo.setRestorePending(false)
}

private suspend fun restoreFromMostRecentBackup() {
Expand Down Expand Up @@ -560,11 +561,16 @@ class WalletViewModel @Inject constructor(
suspend fun restoreWallet(mnemonic: String, bip39Passphrase: String?) {
setInitNodeLifecycleState()
_restoreState.update { RestoreState.InProgress.Wallet }
// The node starts and syncs long before the backup is read, so ordinary uploads are held from
// here rather than from the restore itself, which would upload over the backup it has not read.
backupRepo.setRestorePending(true)

walletRepo.restoreWallet(
mnemonic = mnemonic,
bip39Passphrase = bip39Passphrase,
).onFailure {
// Nothing reaches restoreFromBackup when the wallet was never created, so release here.
backupRepo.setRestorePending(false)
ToastEventBus.send(it)
}
}
Expand Down
110 changes: 110 additions & 0 deletions app/src/test/java/to/bitkit/repositories/ActivityRepoTest.kt
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package to.bitkit.repositories
import com.synonym.bitkitcore.Activity
import com.synonym.bitkitcore.ActivityFilter
import com.synonym.bitkitcore.ActivityTags
import com.synonym.bitkitcore.ClosedChannelDetails
import com.synonym.bitkitcore.IcJitEntry
import com.synonym.bitkitcore.LightningActivity
import com.synonym.bitkitcore.OnchainActivity
Expand Down Expand Up @@ -31,6 +32,7 @@ import to.bitkit.data.dto.PendingBoostActivity
import to.bitkit.ext.create
import to.bitkit.ext.createChannelDetails
import to.bitkit.ext.mock
import to.bitkit.models.ActivityBackupV1
import to.bitkit.models.WalletScope
import to.bitkit.services.CoreService
import to.bitkit.services.HwSnapshotResult
Expand Down Expand Up @@ -70,6 +72,10 @@ class ActivityRepoTest : BaseUnitTest() {
on { v1 } doReturn testActivityV1
}

private val backupTags by lazy { ActivityTags(WalletScope.default, "activity1", listOf("daily")) }

private val backupClosedChannel = mock<ClosedChannelDetails>()

private val baseOnchainActivity = OnchainActivity.create(
walletId = "wallet0",
id = "base_activity_id",
Expand Down Expand Up @@ -1024,6 +1030,110 @@ class ActivityRepoTest : BaseUnitTest() {
assertEquals(listOf("hw-txid"), result.map { it.paymentId })
}

@Test
fun `restoreFromBackup applies every slice and signals tag changes`() = test {
val tagsBefore = sut.activityTagsChanged.value

val result = sut.restoreFromBackup(backupPayload())

assertTrue(result.isSuccess)
verify(coreService.activity).upsertList(listOf(testActivity))
verify(coreService.activity).upsertTags(listOf(backupTags))
verify(coreService.activity).upsertClosedChannelList(listOf(backupClosedChannel))
assertTrue(sut.activityTagsChanged.value > tagsBefore)
}

@Test
fun `restoreFromBackup applies remaining slices when the tags slice fails`() = test {
whenever(coreService.activity.upsertTags(any()))
.thenThrow(RuntimeException("Failed to insert tag: FOREIGN KEY constraint failed"))
val activitiesBefore = sut.activitiesChanged.value
val tagsBefore = sut.activityTagsChanged.value

val result = sut.restoreFromBackup(backupPayload())

// One unusable tag must not cost the activities or the closed channels.
verify(coreService.activity).upsertList(listOf(testActivity))
verify(coreService.activity).upsertClosedChannelList(listOf(backupClosedChannel))
// Still a failure, so BackupRepo never rewrites a good backup with partial state.
assertTrue(result.isFailure)
assertTrue(sut.activitiesChanged.value > activitiesBefore)
// No tag was stored, so the metadata backup must not be marked as changed.
assertEquals(tagsBefore, sut.activityTagsChanged.value)
}

@Test
fun `restoreFromBackup applies remaining slices when the activities slice fails`() = test {
whenever(coreService.activity.upsertList(any())).thenThrow(RuntimeException("upsert failed"))

val result = sut.restoreFromBackup(backupPayload())

verify(coreService.activity).upsertTags(listOf(backupTags))
verify(coreService.activity).upsertClosedChannelList(listOf(backupClosedChannel))
assertTrue(result.isFailure)
}

@Test
fun `restoreFromBackup applies remaining slices when the closed channels slice fails`() = test {
val failure = RuntimeException("closed channels upsert failed")
whenever(coreService.activity.upsertClosedChannelList(any())).thenThrow(failure)
val activitiesBefore = sut.activitiesChanged.value

val result = sut.restoreFromBackup(backupPayload())

verify(coreService.activity).upsertList(listOf(testActivity))
verify(coreService.activity).upsertTags(listOf(backupTags))
assertEquals(failure, result.exceptionOrNull())
assertTrue(sut.activitiesChanged.value > activitiesBefore)
}

@Test
fun `restoreFromBackup returns the first failure when several slices fail`() = test {
val activitiesFailure = RuntimeException("activities upsert failed")
whenever(coreService.activity.upsertList(any())).thenThrow(activitiesFailure)
whenever(coreService.activity.upsertTags(any())).thenThrow(RuntimeException("tags upsert failed"))

val result = sut.restoreFromBackup(backupPayload())

verify(coreService.activity).upsertClosedChannelList(listOf(backupClosedChannel))
assertEquals(activitiesFailure, result.exceptionOrNull())
}

@Test
fun `restoreFromBackup does not notify observers when every slice fails`() = test {
whenever(coreService.activity.upsertList(any())).thenThrow(RuntimeException("activities upsert failed"))
whenever(coreService.activity.upsertTags(any())).thenThrow(RuntimeException("tags upsert failed"))
whenever(coreService.activity.upsertClosedChannelList(any()))
.thenThrow(RuntimeException("closed channels upsert failed"))
val activitiesBefore = sut.activitiesChanged.value
val tagsBefore = sut.activityTagsChanged.value

val result = sut.restoreFromBackup(backupPayload())

assertTrue(result.isFailure)
assertEquals(activitiesBefore, sut.activitiesChanged.value)
assertEquals(tagsBefore, sut.activityTagsChanged.value)
}

@Test
fun `restoreFromBackup rethrows cancellation`() = test {
val cancellation = CancellationException("cancelled")
whenever(coreService.activity.upsertTags(any())).thenThrow(cancellation)

val thrown = assertFailsWith<CancellationException> {
sut.restoreFromBackup(backupPayload())
}

assertEquals(cancellation.message, thrown.message)
}

private fun backupPayload() = ActivityBackupV1(
createdAt = 1234567890L,
activities = listOf(testActivity),
activityTags = listOf(backupTags),
closedChannels = listOf(backupClosedChannel),
)

private suspend fun stubHardwareTagLookup(activity: Activity.Onchain) {
whenever { coreService.activity.getAllActivitiesTags() }
.thenReturn(listOf(ActivityTags(HARDWARE_WALLET_ID, "hw-activity", listOf("cold"))))
Expand Down
Loading
Loading