diff --git a/changes/unreleased/386-support-cancellation-recovery.md b/changes/unreleased/386-support-cancellation-recovery.md new file mode 100644 index 000000000..c31977ed0 --- /dev/null +++ b/changes/unreleased/386-support-cancellation-recovery.md @@ -0,0 +1,7 @@ +category: fix +issue: 351 +pull: 386 +platforms: android, desktop +user-facing: yes + +Finish cancelling a private support report only after Obiente Support confirms that its upload key is terminally cancelled. diff --git a/ui/src/desktopTest/kotlin/dev/obiente/nextcloudnative/app/JvmSupportIntakeTest.kt b/ui/src/desktopTest/kotlin/dev/obiente/nextcloudnative/app/JvmSupportIntakeTest.kt index b9fc39f99..a6c165a77 100644 --- a/ui/src/desktopTest/kotlin/dev/obiente/nextcloudnative/app/JvmSupportIntakeTest.kt +++ b/ui/src/desktopTest/kotlin/dev/obiente/nextcloudnative/app/JvmSupportIntakeTest.kt @@ -16,6 +16,7 @@ import kotlin.test.Test import kotlin.test.assertEquals import kotlin.test.assertFalse import kotlin.test.assertIs +import kotlin.test.assertNull import kotlin.test.assertTrue import kotlinx.coroutines.CancellationException import kotlinx.coroutines.Dispatchers @@ -558,6 +559,113 @@ class JvmSupportIntakeTest { } } + @Test + fun cancellationAtTheUploadMarkerCannotBeOverwrittenByAStorageRejection() = runBlocking { + val markerEntered = CountDownLatch(1) + val allowMarkerTransition = CountDownLatch(1) + testFixture( + beforeUploadMarker = { + markerEntered.countDown() + check(allowMarkerTransition.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + assertTrue(markerEntered.await(5, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + assertIs(fixture.intake.states().value) + + allowMarkerTransition.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + assertEquals(0, fixture.server.requestCount) + assertFalse(File(fixture.temporaryRoot, "pending.json").exists()) + } + } + + @Test + fun transportGateClosingBeforePostStillAllowsLocalCancellation() = runBlocking { + var gateChecks = 0 + testFixture( + supportMutationsAllowed = { + gateChecks += 1 + gateChecks <= 2 + }, + ).use { fixture -> + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + + assertIs(fixture.intake.states().value) + assertEquals(0, fixture.server.requestCount) + assertTrue(File(fixture.temporaryRoot, "pending.json").isFile) + + assertTrue(fixture.intake.cancel()) + + assertIs(fixture.intake.states().value) + assertEquals(0, fixture.server.requestCount) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + } + } + + @Test + fun closedTransportGateCannotReinsertALocallyCancelledSubmission() = runBlocking { + val gateFailureDispositionEntered = CountDownLatch(1) + val allowGateFailureDisposition = CountDownLatch(1) + var gateChecks = 0 + testFixture( + supportMutationsAllowed = { + gateChecks += 1 + gateChecks <= 2 + }, + beforeTransportGateFailureDisposition = { + gateFailureDispositionEntered.countDown() + check(allowGateFailureDisposition.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + assertTrue(gateFailureDispositionEntered.await(5, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + assertIs(fixture.intake.states().value) + allowGateFailureDisposition.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + assertEquals(0, fixture.server.requestCount) + assertFalse(File(fixture.temporaryRoot, "pending.json").exists()) + } + } + + @Test + fun cancellationDuringArchiveValidationCompletesCleanly() = runBlocking { + val archiveValidationEntered = CountDownLatch(1) + val allowArchiveValidation = CountDownLatch(1) + testFixture( + beforeUploadArchiveValidation = { + archiveValidationEntered.countDown() + check(allowArchiveValidation.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + assertTrue(archiveValidationEntered.await(5, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + assertIs(fixture.intake.states().value) + allowArchiveValidation.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + assertEquals(0, fixture.server.requestCount) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + } + } + @Test fun cancellationWinsBeforeTheUploadCallIsRegistered() = runBlocking { val registrationEntered = CountDownLatch(1) @@ -584,6 +692,218 @@ class JvmSupportIntakeTest { } } + @Test + fun cancellationPendingStopsRetryBeforeAnotherUploadStarts() = runBlocking { + val retryTransitionEntered = CountDownLatch(1) + val allowRetryTransition = CountDownLatch(1) + testFixture( + beforeRetryUploadTransition = { + retryTransitionEntered.countDown() + check(allowRetryTransition.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fixture.server.enqueue(MockResponse.Builder().code(404).build()) + + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + + val retryable = assertIs( + fixture.intake.states().value, + ) + assertFalse(retryable.outcomeAmbiguous) + assertEquals("POST", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + assertEquals("GET", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + val retry = launch(Dispatchers.Default) { fixture.intake.retry() } + assertTrue(retryTransitionEntered.await(5, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + assertIs(fixture.intake.states().value) + allowRetryTransition.countDown() + retry.join() + + assertIs(fixture.intake.states().value) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) + assertEquals(3, fixture.server.requestCount) + assertNull(fixture.server.takeRequest(200, TimeUnit.MILLISECONDS)) + } + } + + @Test + fun tombstoneRecoveryAgesFromTheLatestUploadAttempt() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fixture.server.enqueue(MockResponse.Builder().code(404).build()) + + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + + val retryable = assertIs( + fixture.intake.states().value, + ) + assertFalse(retryable.outcomeAmbiguous) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("GET", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + val descriptor = File(fixture.temporaryRoot, "pending.json") + val oldCreatedAt = System.currentTimeMillis() - TimeUnit.DAYS.toMillis(31) + descriptor.writeText( + descriptor.readText().replace( + Regex("\\\"createdAtEpochMillis\\\":\\d+"), + "\"createdAtEpochMillis\":$oldCreatedAt", + ), + ) + fixture.intake.close() + fixture.server.enqueue(receiptResponse(fixture.statusUrl)) + + fixture.newIntake().use { restored -> + assertIs(restored.states().value) + + restored.retry() + + assertIs(restored.states().value) + val retriedUpload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("POST", retriedUpload.method) + assertEquals(upload.headers["Idempotency-Key"], retriedUpload.headers["Idempotency-Key"]) + } + } + } + + @Test + fun cancellationWinsAgainstConcurrentRecoveryExpiry() = runBlocking { + val expiryDispositionEntered = CountDownLatch(1) + val allowExpiryDisposition = CountDownLatch(1) + var nowEpochMillis = System.currentTimeMillis() + testFixture( + currentTimeMillis = { nowEpochMillis }, + beforeRecoveryExpiryDisposition = { + expiryDispositionEntered.countDown() + check(allowExpiryDisposition.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fixture.server.enqueue(MockResponse.Builder().code(404).build()) + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("GET", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + nowEpochMillis += TimeUnit.DAYS.toMillis(31) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + val retry = launch(Dispatchers.Default) { fixture.intake.retry() } + assertTrue(expiryDispositionEntered.await(5, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + assertIs(fixture.intake.states().value) + allowExpiryDisposition.countDown() + retry.join() + + assertIs(fixture.intake.states().value) + val tombstone = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", tombstone.method) + assertEquals(upload.headers["Idempotency-Key"], tombstone.headers["Idempotency-Key"]) + assertFalse(File(fixture.temporaryRoot, "pending.json").exists()) + } + } + + @Test + fun localCancellationDuringRetryPersistenceRemainsTerminal() = runBlocking { + val retryTransitionCompleted = CountDownLatch(1) + val allowRetryPersistence = CountDownLatch(1) + var allowAllMutations = false + var initialGateChecks = 0 + testFixture( + supportMutationsAllowed = { + allowAllMutations || ++initialGateChecks == 1 + }, + afterRetryUploadTransition = { + retryTransitionCompleted.countDown() + check(allowRetryPersistence.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + assertIs(fixture.intake.states().value) + assertEquals(0, fixture.server.requestCount) + allowAllMutations = true + + val retry = launch(Dispatchers.Default) { fixture.intake.retry() } + assertTrue(retryTransitionCompleted.await(5, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + assertIs(fixture.intake.states().value) + allowRetryPersistence.countDown() + retry.join() + + assertIs(fixture.intake.states().value) + assertEquals(0, fixture.server.requestCount) + assertFalse(File(fixture.temporaryRoot, "pending.json").exists()) + } + } + + @Test + fun markerWriteFailurePreservesThePriorTombstoneCapability() = runBlocking { + val rejectNextDirectorySync = AtomicBoolean(false) + var markerCount = 0 + testFixture( + directorySync = { + if (rejectNextDirectorySync.compareAndSet(true, false)) { + throw IOException("Synthetic marker persistence failure.") + } + }, + beforeUploadMarker = { + markerCount += 1 + if (markerCount == 2) rejectNextDirectorySync.set(true) + }, + ).use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fixture.server.enqueue(MockResponse.Builder().code(404).build()) + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("GET", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + + fixture.intake.retry() + + assertIs(fixture.intake.states().value) + assertTrue(File(fixture.temporaryRoot, "pending.json").isFile) + assertNull(fixture.server.takeRequest(200, TimeUnit.MILLISECONDS)) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + assertTrue(fixture.intake.cancel()) + + assertIs(fixture.intake.states().value) + val tombstone = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", tombstone.method) + assertEquals(upload.headers["Idempotency-Key"], tombstone.headers["Idempotency-Key"]) + assertFalse(File(fixture.temporaryRoot, "pending.json").exists()) + } + } + + @Test + fun cancellationBeforePendingInstallRemainsLocal() = runBlocking { + val pendingInstallEntered = CountDownLatch(1) + val allowPendingInstall = CountDownLatch(1) + testFixture( + beforePendingSubmissionInstall = { + pendingInstallEntered.countDown() + check(allowPendingInstall.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + assertTrue(pendingInstallEntered.await(5, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + assertIs(fixture.intake.states().value) + allowPendingInstall.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + assertEquals(0, fixture.server.requestCount) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + } + } + @Test fun publishesBusyStateAndCancelsBeforeSubmissionPreparationCompletes() = runBlocking { val preparationEntered = CountDownLatch(1) @@ -611,6 +931,32 @@ class JvmSupportIntakeTest { } } + @Test + fun cancellationDuringArchivePromotionRemainsTerminal() = runBlocking { + val promotionEntered = CountDownLatch(1) + val allowPromotion = CountDownLatch(1) + testFixture( + beforeArchivePromotion = { + promotionEntered.countDown() + check(allowPromotion.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + assertTrue(promotionEntered.await(5, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + assertIs(fixture.intake.states().value) + allowPromotion.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + assertEquals(0, fixture.server.requestCount) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + } + } + @Test fun preservesPreparationBlockAcrossAccountSwitchesUntilTheOperationEnds() = runBlocking { val preparationEntered = CountDownLatch(1) @@ -643,38 +989,85 @@ class JvmSupportIntakeTest { @Test fun cancellationStopsTheActiveCallWhenIntentPersistenceFails() = runBlocking { - var directorySyncs = 0 + var rejectCancellationWrites = false testFixture( directorySync = { - directorySyncs += 1 - if (directorySyncs == 4) throw IOException("Synthetic cancellation persistence failure.") + if (rejectCancellationWrites) throw IOException("Synthetic cancellation persistence failure.") }, ).use { fixture -> fixture.server.enqueue( receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build(), ) - fixture.server.enqueue(MockResponse.Builder().code(404).build()) val submission = launch(Dispatchers.Default) { fixture.intake.submit("A refresh failed.", "nightly", emptyList()) } - assertEquals("POST", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("POST", upload.method) + rejectCancellationWrites = true assertFalse(fixture.intake.cancel()) withTimeout(5_000) { submission.join() } - assertEquals("GET", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) - assertIs(fixture.intake.states().value) + val retryable = assertIs(fixture.intake.states().value) + assertTrue(retryable.message.contains("was not sent")) + assertEquals(1, fixture.server.requestCount) + assertNull(fixture.server.takeRequest(200, TimeUnit.MILLISECONDS)) assertTrue(File(fixture.temporaryRoot, "pending.json").isFile) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + } + Unit + } - fixture.intake.retry() + @Test + fun concurrentCancellationPersistenceNeverRepopulatesPrivatePayload() = runBlocking { + val rejectNextCancellationWrite = AtomicBoolean(false) + val cancellationWriteEntered = CountDownLatch(1) + val allowCancellationWriteFailure = CountDownLatch(1) + val uploadResponseCompleted = CountDownLatch(1) + testFixture( + directorySync = { + if (rejectNextCancellationWrite.compareAndSet(true, false)) { + cancellationWriteEntered.countDown() + assertTrue(allowCancellationWriteFailure.await(3, TimeUnit.SECONDS)) + throw IOException("Synthetic cancellation persistence failure.") + } + }, + afterUploadResponse = { uploadResponseCompleted.countDown() }, + ).use { fixture -> + fixture.server.enqueue( + receiptResponse(fixture.statusUrl).newBuilder().headersDelay(1, TimeUnit.SECONDS).build(), + ) + fixture.server.enqueue(MockResponse.Builder().code(503).build()) + val submission = launch(Dispatchers.Default) { + fixture.intake.submit( + "Private concurrent cancellation note.", + "nightly", + listOf(SupportDiagnosticFieldDraft("private_concurrent_field", "Private concurrent value")), + ) + } + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + rejectNextCancellationWrite.set(true) + val cancellation = launch(Dispatchers.Default) { + assertFalse(fixture.intake.cancel()) + } + assertTrue(cancellationWriteEntered.await(2, TimeUnit.SECONDS)) + assertTrue(uploadResponseCompleted.await(3, TimeUnit.SECONDS)) - assertIs(fixture.intake.states().value) - assertEquals("GET", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) - assertEquals("DELETE", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + allowCancellationWriteFailure.countDown() + cancellation.join() + submission.join() + + assertIs(fixture.intake.states().value) + val tombstone = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", tombstone.method) + assertEquals(upload.headers["Idempotency-Key"], tombstone.headers["Idempotency-Key"]) + assertFalse(fixture.temporaryRoot.listFiles().orEmpty().any { it.extension == "zip" }) + val descriptor = File(fixture.temporaryRoot, "pending.json").readText() + assertTrue(descriptor.contains(upload.headers["Idempotency-Key"].orEmpty())) + assertTrue(descriptor.contains("\"cancellationPending\":true")) + assertFalse(descriptor.contains("Private concurrent cancellation note.")) + assertFalse(descriptor.contains("private_concurrent_field")) + assertFalse(descriptor.contains("Private concurrent value")) } - Unit } @Test @@ -691,9 +1084,12 @@ class JvmSupportIntakeTest { assertTrue(fixture.intake.cancel()) assertIs(fixture.intake.states().value) - fixture.server.enqueue(MockResponse.Builder().code(404).build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) withTimeout(5_000) { submission.join() } - assertIs(fixture.intake.states().value) + assertIs(fixture.intake.states().value) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) } Unit } @@ -972,35 +1368,332 @@ class JvmSupportIntakeTest { } @Test - fun cancellingAmbiguousSubmissionRequiresDeletionReconciliation() = runBlocking { - testFixture().use { fixture -> - fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fun cancellingAmbiguousSubmissionUsesAuthoritativeServerTombstone() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fixture.server.enqueue(MockResponse.Builder().code(503).build()) + + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + + assertIs(fixture.intake.states().value) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + assertTrue(fixture.intake.cancel()) + + assertIs(fixture.intake.states().value) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + } + } + + @Test + fun cancellationDuringReceiptReconciliationSendsAuthoritativeTombstone() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fixture.server.enqueue( + MockResponse.Builder().code(404).headersDelay(10, TimeUnit.SECONDS).build(), + ) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + val reconciliation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("GET", reconciliation.method) + + assertTrue(fixture.intake.cancel()) + submission.join() + + assertIs(fixture.intake.states().value) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertEquals(3, fixture.server.requestCount) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + } + } + + @Test + fun cancellationBeforeReceiptRegistrationStartsTheTombstoneImmediately() = runBlocking { + val receiptRegistrationEntered = CountDownLatch(1) + val allowReceiptRegistration = CountDownLatch(1) + testFixture( + beforeReceiptCallRegistration = { + receiptRegistrationEntered.countDown() + check(allowReceiptRegistration.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertTrue(receiptRegistrationEntered.await(5, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + assertIs(fixture.intake.states().value) + allowReceiptRegistration.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + val tombstone = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", tombstone.method) + assertEquals(upload.headers["Idempotency-Key"], tombstone.headers["Idempotency-Key"]) + assertEquals(2, fixture.server.requestCount) + assertNull(fixture.server.takeRequest(200, TimeUnit.MILLISECONDS)) + } + } + + @Test + fun cancellationAfterReceiptLookupCompletesSendsAuthoritativeTombstone() = runBlocking { + val lookupCompleted = CountDownLatch(1) + val allowLookupResult = CountDownLatch(1) + testFixture( + afterReceiptLookup = { + lookupCompleted.countDown() + assertTrue(allowLookupResult.await(2, TimeUnit.SECONDS)) + }, + ).use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fixture.server.enqueue(MockResponse.Builder().code(404).build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + val reconciliation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("GET", reconciliation.method) + assertTrue(lookupCompleted.await(2, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + allowLookupResult.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertEquals(3, fixture.server.requestCount) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + } + } + + @Test + fun cancellationDoesNotCancelTheTombstoneStartedByResponseHandling() = runBlocking { + val uploadResponseCompleted = CountDownLatch(1) + val allowUploadResponseDisposition = CountDownLatch(1) + val cancellationIntentPublished = CountDownLatch(1) + val allowCancellationContinuation = CountDownLatch(1) + testFixture( + afterUploadResponse = { + uploadResponseCompleted.countDown() + assertTrue(allowUploadResponseDisposition.await(5, TimeUnit.SECONDS)) + }, + afterCancellationIntentPublished = { + cancellationIntentPublished.countDown() + assertTrue(allowCancellationContinuation.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + fixture.server.enqueue(MockResponse.Builder().code(503).build()) + fixture.server.enqueue( + MockResponse.Builder().code(204).headersDelay(2, TimeUnit.SECONDS).build(), + ) + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + assertEquals("POST", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + assertTrue(uploadResponseCompleted.await(5, TimeUnit.SECONDS)) + + val cancellation = launch(Dispatchers.Default) { + assertTrue(fixture.intake.cancel()) + } + assertTrue(cancellationIntentPublished.await(5, TimeUnit.SECONDS)) + allowUploadResponseDisposition.countDown() + val tombstone = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", tombstone.method) + assertEquals("/api/v1/receipts", tombstone.url.encodedPath) + allowCancellationContinuation.countDown() + + cancellation.join() + submission.join() + + assertIs(fixture.intake.states().value) + assertEquals(2, fixture.server.requestCount) + assertFalse(File(fixture.temporaryRoot, "pending.json").exists()) + } + } + + @Test + fun cancellationAfterReceiptIntentCheckContinuesWithAuthoritativeTombstone() = runBlocking { + val dispositionEntered = CountDownLatch(1) + val allowDisposition = CountDownLatch(1) + testFixture( + beforeReceiptResponseDisposition = { + dispositionEntered.countDown() + assertTrue(allowDisposition.await(2, TimeUnit.SECONDS)) + }, + ).use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fixture.server.enqueue(MockResponse.Builder().code(404).build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("GET", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + assertTrue(dispositionEntered.await(2, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + allowDisposition.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + } + } + + @Test + fun cancellationWinsAgainstPermanentUploadRejectionDisposition() = runBlocking { + val responseDispositionEntered = CountDownLatch(1) + val allowResponseDisposition = CountDownLatch(1) + testFixture( + beforeUploadResponseDisposition = { + responseDispositionEntered.countDown() + check(allowResponseDisposition.await(5, TimeUnit.SECONDS)) + }, + ).use { fixture -> + fixture.server.enqueue( + MockResponse.Builder().code(400).body( + """{"contractVersion":1,"code":"invalid_report","message":"Invalid."}""", + ).build(), + ) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertTrue(responseDispositionEntered.await(5, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + assertIs(fixture.intake.states().value) + allowResponseDisposition.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + val tombstone = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", tombstone.method) + assertEquals(upload.headers["Idempotency-Key"], tombstone.headers["Idempotency-Key"]) + assertFalse(File(fixture.temporaryRoot, "pending.json").exists()) + } + } + + @Test + fun cancellationAfterUploadResponseCompletesSendsAuthoritativeTombstone() = runBlocking { + val responseCompleted = CountDownLatch(1) + val allowResponseResult = CountDownLatch(1) + testFixture( + afterUploadResponse = { + responseCompleted.countDown() + assertTrue(allowResponseResult.await(2, TimeUnit.SECONDS)) + }, + ).use { fixture -> + fixture.server.enqueue(MockResponse.Builder().code(503).build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertTrue(responseCompleted.await(2, TimeUnit.SECONDS)) + + assertTrue(fixture.intake.cancel()) + allowResponseResult.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertEquals(2, fixture.server.requestCount) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + } + } + + @Test + fun cancellationAfterUploadIntentCheckContinuesWithAuthoritativeTombstone() = runBlocking { + val dispositionEntered = CountDownLatch(1) + val allowDisposition = CountDownLatch(1) + testFixture( + beforeUploadResponseDisposition = { + dispositionEntered.countDown() + assertTrue(allowDisposition.await(2, TimeUnit.SECONDS)) + }, + ).use { fixture -> fixture.server.enqueue(MockResponse.Builder().code(503).build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) - fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertTrue(dispositionEntered.await(2, TimeUnit.SECONDS)) - assertIs(fixture.intake.states().value) assertTrue(fixture.intake.cancel()) - assertIs(fixture.intake.states().value) - assertTrue(File(fixture.temporaryRoot, "pending.json").isFile) + allowDisposition.countDown() + submission.join() + + assertIs(fixture.intake.states().value) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + } + } + @Test + fun cancellationAfterReconciledReceiptAbsenceUsesAuthoritativeTombstone() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) fixture.server.enqueue(MockResponse.Builder().code(404).build()) - fixture.intake.retry() + + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) assertIs(fixture.intake.states().value) - assertTrue(File(fixture.temporaryRoot, "pending.json").isFile) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + val reconciliation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("GET", reconciliation.method) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) - fixture.intake.retry() + assertTrue(fixture.intake.cancel()) assertIs(fixture.intake.states().value) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertEquals(3, fixture.server.requestCount) assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) } } @Test - fun serverFailureRemainsAmbiguousUntilDiscardReconcilesAndDeletes() = runBlocking { + fun serverFailureRemainsAmbiguousUntilCancellationIsConfirmed() = runBlocking { testFixture().use { fixture -> fixture.server.enqueue(MockResponse.Builder().code(503).build()) @@ -1009,32 +1702,23 @@ class JvmSupportIntakeTest { val retryable = assertIs(fixture.intake.states().value) assertTrue(retryable.outcomeAmbiguous) val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) assertTrue(fixture.intake.cancel()) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) - - fixture.intake.retry() assertIs(fixture.intake.states().value) - val reconciliation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - val deletion = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - assertEquals(upload.headers["Idempotency-Key"], reconciliation.headers["Idempotency-Key"]) - assertEquals("GET", reconciliation.method) - assertEquals("DELETE", deletion.method) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) } } @Test - fun cancellationRechecksInitialAbsenceAndDeletesLateReceipt() = runBlocking { - testFixture( - cancellationReconcileWindowMillis = 1_000L, - cancellationReconcilePollMillis = 1L, - ).use { fixture -> + fun cancellationDoesNotPollReceiptAbsenceBeforeDiscardingRecovery() = runBlocking { + testFixture().use { fixture -> fixture.server.enqueue(receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build()) - fixture.server.enqueue(MockResponse.Builder().code(404).build()) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) val submission = launch(Dispatchers.Default) { fixture.intake.submit("A refresh failed.", "nightly", emptyList()) @@ -1044,41 +1728,151 @@ class JvmSupportIntakeTest { submission.join() assertIs(fixture.intake.states().value) - val firstReconcile = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - val secondReconcile = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - val deletion = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - assertEquals(upload.headers["Idempotency-Key"], firstReconcile.headers["Idempotency-Key"]) - assertEquals(upload.headers["Idempotency-Key"], secondReconcile.headers["Idempotency-Key"]) - assertEquals("GET", firstReconcile.method) - assertEquals("GET", secondReconcile.method) - assertEquals("DELETE", deletion.method) - assertTrue(deletion.url.encodedPath.startsWith("/api/v1/reports/")) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) + assertEquals(2, fixture.server.requestCount) assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) } } @Test - fun cancellationStillDeletesWhenReceiptPersistenceFails() = runBlocking { - var directorySyncs = 0 - testFixture( - directorySync = { - directorySyncs += 1 - if (directorySyncs == 5) throw IOException("Synthetic receipt persistence failure.") - }, - ).use { fixture -> - fixture.server.enqueue(receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build()) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + fun cancellationRetryWaitsForAuthoritativeTerminalResult() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue( + receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build(), + ) + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + val submission = launch(Dispatchers.Default) { fixture.intake.submit("A refresh failed.", "nightly", emptyList()) } requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertTrue(fixture.intake.cancel()) + submission.join() + + assertIs(fixture.intake.states().value) + val descriptor = File(fixture.temporaryRoot, "pending.json") + assertTrue(descriptor.isFile) + val firstCancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", firstCancellation.method) + fixture.intake.close() + val restored = fixture.newIntake() + assertIs(restored.states().value) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + restored.retry() + + assertIs(restored.states().value) + val retryCancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", retryCancellation.method) + assertEquals("/api/v1/receipts", retryCancellation.url.encodedPath) + assertEquals( + firstCancellation.headers["Idempotency-Key"], + retryCancellation.headers["Idempotency-Key"], + ) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + restored.close() + } + } + + @Test + fun restoredAmbiguousSubmissionAcceptsAuthoritativeCancellationTombstone() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue( + receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build(), + ) + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) assertTrue(fixture.intake.cancel()) submission.join() + requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + fixture.intake.close() - assertIs(fixture.intake.states().value) + val descriptor = File(fixture.temporaryRoot, "pending.json") + descriptor.writeText( + descriptor.readText().replace( + "\"cancellationPending\":true", + "\"cancellationPending\":false", + ), + ) + val restored = fixture.newIntake() + assertIs(restored.states().value) + fixture.server.enqueue(submissionCancelledResponse()) + + restored.retry() + + assertIs(restored.states().value) + val reconciliation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("GET", reconciliation.method) + assertEquals("/api/v1/receipts", reconciliation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], reconciliation.headers["Idempotency-Key"]) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + restored.close() + } + } + + @Test + fun unverifiedGoneUploadResponseRetainsRecovery() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue( + MockResponse.Builder().code(410).body( + """{"contractVersion":1,"code":"not_found","message":"Gone."}""", + ).build(), + ) + + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + + val retryable = assertIs( + fixture.intake.states().value, + ) + assertTrue(retryable.outcomeAmbiguous) + assertTrue(File(fixture.temporaryRoot, "pending.json").isFile) + assertEquals("POST", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + } + } + + @Test + fun unverifiedGoneReceiptResponseRetainsRecovery() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue(MockResponse.Builder().onResponseStart(SocketEffect.CloseSocket()).build()) + fixture.server.enqueue(MockResponse.Builder().code(410).body("gone").build()) + + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + + val retryable = assertIs( + fixture.intake.states().value, + ) + assertTrue(retryable.outcomeAmbiguous) + assertTrue(File(fixture.temporaryRoot, "pending.json").isFile) + assertEquals("POST", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) assertEquals("GET", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + } + } + + @Test + fun cancellationRetainsRecoveryUntilTerminalNoContent() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue(MockResponse.Builder().code(503).build()) + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + assertTrue(fixture.intake.cancel()) + + assertIs(fixture.intake.states().value) + assertTrue(File(fixture.temporaryRoot, "pending.json").isFile) + val nonTerminal = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", nonTerminal.method) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + fixture.intake.retry() + + assertIs(fixture.intake.states().value) assertEquals("DELETE", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) } @@ -1222,23 +2016,21 @@ class JvmSupportIntakeTest { submission.join() assertIs(fixture.intake.states().value) - val firstReconciliation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - assertEquals("GET", firstReconciliation.method) + val firstCancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", firstCancellation.method) fixture.intake.close() val restored = fixture.newIntake() assertIs(restored.states().value) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) restored.retry() assertIs(restored.states().value) - val retryReconciliation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - val deletion = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - assertEquals(upload.headers["Idempotency-Key"], retryReconciliation.headers["Idempotency-Key"]) - assertEquals("GET", retryReconciliation.method) - assertEquals("DELETE", deletion.method) + val retryCancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals(upload.headers["Idempotency-Key"], retryCancellation.headers["Idempotency-Key"]) + assertEquals("DELETE", retryCancellation.method) + assertEquals("/api/v1/receipts", retryCancellation.url.encodedPath) assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) } } @@ -1340,85 +2132,122 @@ class JvmSupportIntakeTest { } @Test - fun retriesPersistedDeletionCapabilityWithoutResubmitting() = runBlocking { + fun retriesPersistedCancellationWithoutResubmitting() = runBlocking { testFixture().use { fixture -> fixture.server.enqueue(receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build()) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) fixture.server.enqueue(MockResponse.Builder().code(503).build()) val submission = launch(Dispatchers.Default) { fixture.intake.submit("A refresh failed.", "nightly", emptyList()) } - requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) assertTrue(fixture.intake.cancel()) submission.join() assertIs(fixture.intake.states().value) - val reconciliation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - val failedDeletion = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - assertEquals("GET", reconciliation.method) - assertEquals("DELETE", failedDeletion.method) + val failedCancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", failedCancellation.method) + assertEquals("/api/v1/receipts", failedCancellation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], failedCancellation.headers["Idempotency-Key"]) fixture.intake.close() val restored = fixture.newIntake() - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) restored.retry() assertIs(restored.states().value) - val retriedDeletion = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - assertEquals("DELETE", retriedDeletion.method) - assertEquals(4, fixture.server.requestCount) + val retriedCancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", retriedCancellation.method) + assertEquals("/api/v1/receipts", retriedCancellation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], retriedCancellation.headers["Idempotency-Key"]) + assertEquals(3, fixture.server.requestCount) assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) } } @Test - fun acceptedDeletionKeepsReceiptUntilStatusConfirmsRemoval() = runBlocking { + fun failedCancellationRetainsOnlyMinimalRecoveryState() = runBlocking { testFixture().use { fixture -> - fixture.server.enqueue(receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build()) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) - fixture.server.enqueue(MockResponse.Builder().code(202).body("{}").build()) - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + fixture.server.enqueue(MockResponse.Builder().code(503).build()) + fixture.intake.submit( + "Private cancellation reproduction note.", + "nightly", + listOf(SupportDiagnosticFieldDraft("private_field", "Private value")), + ) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + fixture.server.enqueue(MockResponse.Builder().code(503).build()) - val submission = launch(Dispatchers.Default) { - fixture.intake.submit("A refresh failed.", "nightly", emptyList()) - } - requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) assertTrue(fixture.intake.cancel()) - submission.join() assertIs(fixture.intake.states().value) - assertTrue(File(fixture.temporaryRoot, "pending.json").isFile) - val reconciliation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - val acceptedDeletion = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - val statusCheck = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - assertEquals("GET", reconciliation.method) - assertEquals("DELETE", acceptedDeletion.method) - assertEquals("GET", statusCheck.method) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertFalse(fixture.temporaryRoot.listFiles().orEmpty().any { it.extension == "zip" }) + val descriptor = File(fixture.temporaryRoot, "pending.json").readText() + assertTrue(descriptor.contains(upload.headers["Idempotency-Key"].orEmpty())) + assertTrue(descriptor.contains("\"cancellationPending\":true")) + assertFalse(descriptor.contains("Private cancellation reproduction note.")) + assertFalse(descriptor.contains("private_field")) + assertFalse(descriptor.contains("Private value")) + } + } - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) - fixture.intake.retry() + @Test + fun restorationMinimizesPreviouslyPersistedCancellationIntent() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue(MockResponse.Builder().code(503).build()) + fixture.intake.submit( + "Private restored cancellation note.", + "nightly", + listOf(SupportDiagnosticFieldDraft("private_restored_field", "Private restored value")), + ) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + val descriptor = File(fixture.temporaryRoot, "pending.json") + assertTrue(descriptor.readText().contains("Private restored cancellation note.")) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().any { it.extension == "zip" }) + fixture.intake.close() + descriptor.writeText( + descriptor.readText().replace( + "\"cancellationPending\":false", + "\"cancellationPending\":true", + ), + ) - assertIs(fixture.intake.states().value) + val restored = fixture.newIntake() + + assertIs(restored.states().value) + assertFalse(fixture.temporaryRoot.listFiles().orEmpty().any { it.extension == "zip" }) + val minimized = descriptor.readText() + assertTrue(minimized.contains(upload.headers["Idempotency-Key"].orEmpty())) + assertTrue(minimized.contains("\"cancellationPending\":true")) + assertFalse(minimized.contains("Private restored cancellation note.")) + assertFalse(minimized.contains("private_restored_field")) + assertFalse(minimized.contains("Private restored value")) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + restored.retry() + + assertIs(restored.states().value) assertEquals("DELETE", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + restored.close() } } @Test - fun keepsDeletionCapabilityAfterTheLocalArchiveRetentionWindow() = runBlocking { + fun keepsCancellationKeyAfterTheLocalArchiveRetentionWindow() = runBlocking { testFixture().use { fixture -> fixture.server.enqueue(receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build()) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) fixture.server.enqueue(MockResponse.Builder().code(503).build()) val submission = launch(Dispatchers.Default) { fixture.intake.submit("A refresh failed.", "nightly", emptyList()) } - requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) assertTrue(fixture.intake.cancel()) submission.join() - requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertIs(fixture.intake.states().value) requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) fixture.intake.close() @@ -1435,29 +2264,69 @@ class JvmSupportIntakeTest { assertIs(restored.states().value) assertTrue(descriptor.isFile) assertFalse(fixture.temporaryRoot.listFiles().orEmpty().any { it.extension == "zip" }) - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) restored.retry() assertIs(restored.states().value) - assertEquals("DELETE", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) } } @Test - fun restoresDeletionCapabilityWhenTheWallClockMovesBackward() = runBlocking { + fun reconcilesRestoredCancellationBeforeApplyingRecoveryExpiry() = runBlocking { testFixture().use { fixture -> fixture.server.enqueue(receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build()) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) fixture.server.enqueue(MockResponse.Builder().code(503).build()) val submission = launch(Dispatchers.Default) { fixture.intake.submit("A refresh failed.", "nightly", emptyList()) } - requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) assertTrue(fixture.intake.cancel()) submission.join() + assertIs(fixture.intake.states().value) requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + fixture.intake.close() + + val descriptor = File(fixture.temporaryRoot, "pending.json") + val expiredCreatedAt = Instant.now().minus(31, ChronoUnit.DAYS).toEpochMilli() + descriptor.writeText( + descriptor.readText().replace( + Regex("\"createdAtEpochMillis\":\\d+"), + "\"createdAtEpochMillis\":$expiredCreatedAt", + ), + ) + val restored = fixture.newIntake() + assertIs(restored.states().value) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) + + restored.retry() + + assertIs(restored.states().value) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals("/api/v1/receipts", cancellation.url.encodedPath) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) + assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) + restored.close() + } + } + + @Test + fun restoresCancellationKeyWhenTheWallClockMovesBackward() = runBlocking { + testFixture().use { fixture -> + fixture.server.enqueue(receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build()) + fixture.server.enqueue(MockResponse.Builder().code(503).build()) + + val submission = launch(Dispatchers.Default) { + fixture.intake.submit("A refresh failed.", "nightly", emptyList()) + } + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertTrue(fixture.intake.cancel()) + submission.join() requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) fixture.intake.close() @@ -1474,20 +2343,21 @@ class JvmSupportIntakeTest { assertIs(restored.states().value) assertTrue(descriptor.isFile) - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) restored.retry() assertIs(restored.states().value) - assertEquals("DELETE", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) + assertEquals(upload.headers["Idempotency-Key"], cancellation.headers["Idempotency-Key"]) assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) } } @Test - fun preservesLastPersistedReceiptWhenRetryStateCannotBeRewritten() = runBlocking { + fun preservesLastCancellationRecordWhenRetryStateCannotBeRewritten() = runBlocking { testFixture().use { fixture -> fixture.server.enqueue(receiptResponse(fixture.statusUrl).newBuilder().headersDelay(10, TimeUnit.SECONDS).build()) - fixture.server.enqueue(receiptResponse(fixture.statusUrl)) fixture.server.enqueue( MockResponse.Builder().code(503).headersDelay(1, TimeUnit.SECONDS).build(), ) @@ -1495,10 +2365,10 @@ class JvmSupportIntakeTest { val submission = launch(Dispatchers.Default) { fixture.intake.submit("A refresh failed.", "nightly", emptyList()) } - requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + val upload = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) assertTrue(fixture.intake.cancel()) - requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) - requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + val cancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", cancellation.method) val retainedRoot = File(fixture.root, "submissions-retained") Files.move(fixture.temporaryRoot.toPath(), retainedRoot.toPath()) fixture.temporaryRoot.writeText("temporarily unavailable") @@ -1508,19 +2378,21 @@ class JvmSupportIntakeTest { assertTrue(state.message.contains("could not be stored")) val descriptor = File(retainedRoot, "pending.json") assertTrue(descriptor.isFile) - assertTrue(descriptor.readText().contains("OBI-ABCDE-23456")) + assertTrue(descriptor.readText().contains(upload.headers["Idempotency-Key"].orEmpty())) assertTrue(fixture.temporaryRoot.delete()) Files.move(retainedRoot.toPath(), fixture.temporaryRoot.toPath()) fixture.intake.close() val restored = fixture.newIntake() assertIs(restored.states().value) - fixture.server.enqueue(MockResponse.Builder().code(200).body("{}").build()) + fixture.server.enqueue(MockResponse.Builder().code(204).build()) restored.retry() assertIs(restored.states().value) - assertEquals("DELETE", requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)).method) + val retryCancellation = requireNotNull(fixture.server.takeRequest(2, TimeUnit.SECONDS)) + assertEquals("DELETE", retryCancellation.method) + assertEquals(upload.headers["Idempotency-Key"], retryCancellation.headers["Idempotency-Key"]) assertTrue(fixture.temporaryRoot.listFiles().orEmpty().isEmpty()) } } @@ -1677,18 +2549,31 @@ class JvmSupportIntakeTest { descriptorCleanupRetryMillis: Long = 60_000L, beforeCallRegistration: () -> Unit = {}, beforeSubmissionPreparation: () -> Unit = {}, + beforePendingSubmissionInstall: () -> Unit = {}, beforeBundlePackaging: () -> Unit = {}, afterBundlePackaging: () -> Unit = {}, + beforeArchivePromotion: () -> Unit = {}, + beforeRetryUploadTransition: () -> Unit = {}, + afterRetryUploadTransition: () -> Unit = {}, + beforeRecoveryExpiryDisposition: () -> Unit = {}, + beforeUploadMarker: () -> Unit = {}, + beforeUploadArchiveValidation: () -> Unit = {}, + beforeTransportGateFailureDisposition: () -> Unit = {}, + beforeReceiptCallRegistration: () -> Unit = {}, + afterCancellationIntentPublished: () -> Unit = {}, + afterUploadResponse: () -> Unit = {}, + beforeUploadResponseDisposition: () -> Unit = {}, + afterReceiptLookup: () -> Unit = {}, + beforeReceiptResponseDisposition: () -> Unit = {}, privateFileDelete: (File) -> Boolean = File::delete, pendingDescriptorRead: (File) -> String = { descriptor -> descriptor.readText() }, completedDescriptorRead: (File) -> String = { descriptor -> descriptor.readText() }, - cancellationReconcileWindowMillis: Long = 0L, - cancellationReconcilePollMillis: Long = 1L, submissionStorageBlocked: Boolean = false, pendingTemporaryBeforeInitialization: Boolean = false, archiveTemporaryBeforeInitialization: Boolean = false, invalidPendingBeforeInitialization: Boolean = false, invalidCompletedBeforeInitialization: Boolean = false, + currentTimeMillis: () -> Long = System::currentTimeMillis, ): Fixture { val root = createTempDirectory("support-intake-test").toFile() val diagnosticRoot = File(root, "diagnostics") @@ -1743,13 +2628,26 @@ class JvmSupportIntakeTest { descriptorCleanupRetryMillis = descriptorCleanupRetryMillis, beforeCallRegistration = beforeCallRegistration, beforeSubmissionPreparation = beforeSubmissionPreparation, + beforePendingSubmissionInstall = beforePendingSubmissionInstall, beforeBundlePackaging = beforeBundlePackaging, afterBundlePackaging = afterBundlePackaging, + beforeArchivePromotion = beforeArchivePromotion, + beforeRetryUploadTransition = beforeRetryUploadTransition, + afterRetryUploadTransition = afterRetryUploadTransition, + beforeRecoveryExpiryDisposition = beforeRecoveryExpiryDisposition, + beforeUploadMarker = beforeUploadMarker, + beforeUploadArchiveValidation = beforeUploadArchiveValidation, + beforeTransportGateFailureDisposition = beforeTransportGateFailureDisposition, + beforeReceiptCallRegistration = beforeReceiptCallRegistration, + afterCancellationIntentPublished = afterCancellationIntentPublished, + afterUploadResponse = afterUploadResponse, + beforeUploadResponseDisposition = beforeUploadResponseDisposition, + afterReceiptLookup = afterReceiptLookup, + beforeReceiptResponseDisposition = beforeReceiptResponseDisposition, privateFileDelete = privateFileDelete, pendingDescriptorRead = pendingDescriptorRead, completedDescriptorRead = completedDescriptorRead, - cancellationReconcileWindowMillis = cancellationReconcileWindowMillis, - cancellationReconcilePollMillis = cancellationReconcilePollMillis, + currentTimeMillis = currentTimeMillis, ) } @@ -1778,6 +2676,10 @@ class JvmSupportIntakeTest { ).build() } + private fun submissionCancelledResponse(): MockResponse = MockResponse.Builder().code(410).body( + """{"contractVersion":1,"code":"submission_cancelled","message":"Submission cancelled."}""", + ).build() + private data class Fixture( val root: File, val temporaryRoot: File, @@ -1789,13 +2691,26 @@ class JvmSupportIntakeTest { val descriptorCleanupRetryMillis: Long, val beforeCallRegistration: () -> Unit, val beforeSubmissionPreparation: () -> Unit, + val beforePendingSubmissionInstall: () -> Unit, val beforeBundlePackaging: () -> Unit, val afterBundlePackaging: () -> Unit, + val beforeArchivePromotion: () -> Unit, + val beforeRetryUploadTransition: () -> Unit, + val afterRetryUploadTransition: () -> Unit, + val beforeRecoveryExpiryDisposition: () -> Unit, + val beforeUploadMarker: () -> Unit, + val beforeUploadArchiveValidation: () -> Unit, + val beforeTransportGateFailureDisposition: () -> Unit, + val beforeReceiptCallRegistration: () -> Unit, + val afterCancellationIntentPublished: () -> Unit, + val afterUploadResponse: () -> Unit, + val beforeUploadResponseDisposition: () -> Unit, + val afterReceiptLookup: () -> Unit, + val beforeReceiptResponseDisposition: () -> Unit, val privateFileDelete: (File) -> Boolean, val pendingDescriptorRead: (File) -> String, val completedDescriptorRead: (File) -> String, - val cancellationReconcileWindowMillis: Long, - val cancellationReconcilePollMillis: Long, + val currentTimeMillis: () -> Long, ) : AutoCloseable { val intake = newIntake() val statusUrl: String get() = server.url("/r/abcdefghijklmnopqrstuvwxyzABCDEFGH_12345678").toString() @@ -1814,13 +2729,26 @@ class JvmSupportIntakeTest { descriptorCleanupRetryMillis = descriptorCleanupRetryMillis, beforeCallRegistration = beforeCallRegistration, beforeSubmissionPreparation = beforeSubmissionPreparation, + beforePendingSubmissionInstall = beforePendingSubmissionInstall, beforeBundlePackaging = beforeBundlePackaging, afterBundlePackaging = afterBundlePackaging, + beforeArchivePromotion = beforeArchivePromotion, + beforeRetryUploadTransition = beforeRetryUploadTransition, + afterRetryUploadTransition = afterRetryUploadTransition, + beforeRecoveryExpiryDisposition = beforeRecoveryExpiryDisposition, + beforeUploadMarker = beforeUploadMarker, + beforeUploadArchiveValidation = beforeUploadArchiveValidation, + beforeTransportGateFailureDisposition = beforeTransportGateFailureDisposition, + beforeReceiptCallRegistration = beforeReceiptCallRegistration, + afterCancellationIntentPublished = afterCancellationIntentPublished, + afterUploadResponse = afterUploadResponse, + beforeUploadResponseDisposition = beforeUploadResponseDisposition, + afterReceiptLookup = afterReceiptLookup, + beforeReceiptResponseDisposition = beforeReceiptResponseDisposition, privateFileDelete = privateFileDelete, pendingDescriptorRead = pendingDescriptorRead, completedDescriptorRead = completedDescriptorRead, - cancellationReconcileWindowMillis = cancellationReconcileWindowMillis, - cancellationReconcilePollMillis = cancellationReconcilePollMillis, + currentTimeMillis = currentTimeMillis, ).also { intake -> intake.setActiveAccountIdentity(TEST_ACCOUNT_IDENTITY) runBlocking { intake.awaitInitialization() } diff --git a/ui/src/jvmMain/kotlin/dev/obiente/nextcloudnative/app/JvmSupportIntake.kt b/ui/src/jvmMain/kotlin/dev/obiente/nextcloudnative/app/JvmSupportIntake.kt index 057543398..5209bab53 100644 --- a/ui/src/jvmMain/kotlin/dev/obiente/nextcloudnative/app/JvmSupportIntake.kt +++ b/ui/src/jvmMain/kotlin/dev/obiente/nextcloudnative/app/JvmSupportIntake.kt @@ -64,8 +64,22 @@ class JvmSupportIntake( private val descriptorCleanupRetryMillis: Long = SUPPORT_DESCRIPTOR_DELETE_RETRY_MILLIS, private val beforeCallRegistration: () -> Unit = {}, private val beforeSubmissionPreparation: () -> Unit = {}, + private val beforePendingSubmissionInstall: () -> Unit = {}, private val beforeBundlePackaging: () -> Unit = {}, private val afterBundlePackaging: () -> Unit = {}, + private val beforeArchivePromotion: () -> Unit = {}, + private val beforeRetryUploadTransition: () -> Unit = {}, + private val afterRetryUploadTransition: () -> Unit = {}, + private val beforeRecoveryExpiryDisposition: () -> Unit = {}, + private val beforeUploadMarker: () -> Unit = {}, + private val beforeUploadArchiveValidation: () -> Unit = {}, + private val beforeTransportGateFailureDisposition: () -> Unit = {}, + private val beforeReceiptCallRegistration: () -> Unit = {}, + private val afterCancellationIntentPublished: () -> Unit = {}, + private val afterUploadResponse: () -> Unit = {}, + private val beforeUploadResponseDisposition: () -> Unit = {}, + private val afterReceiptLookup: () -> Unit = {}, + private val beforeReceiptResponseDisposition: () -> Unit = {}, private val privateFileDelete: (File) -> Boolean = File::delete, private val pendingDescriptorRead: (File) -> String = { descriptor -> descriptor.readText(Charsets.UTF_8) @@ -73,8 +87,7 @@ class JvmSupportIntake( private val completedDescriptorRead: (File) -> String = { descriptor -> descriptor.readText(Charsets.UTF_8) }, - private val cancellationReconcileWindowMillis: Long = SUPPORT_CANCELLATION_RECONCILE_WINDOW_MILLIS, - private val cancellationReconcilePollMillis: Long = SUPPORT_CANCELLATION_RECONCILE_POLL_MILLIS, + private val currentTimeMillis: () -> Long = System::currentTimeMillis, ) : AutoCloseable { private val baseUrl = supportBaseUrl.toHttpUrl() private val client = client.newBuilder() @@ -94,6 +107,7 @@ class JvmSupportIntake( private val cancellationRequested = AtomicBoolean(false) private val shutdownRequested = AtomicBoolean(false) private val operationActive = AtomicBoolean(false) + private val cancellationContinuationRequested = AtomicBoolean(false) private val rejectedPendingDescriptorCleanup = AtomicBoolean(false) private val pendingDescriptorRestorePending = AtomicBoolean(false) private val completedDescriptorRestorePending = AtomicBoolean(false) @@ -109,8 +123,6 @@ class JvmSupportIntake( init { require(descriptorCleanupRetryMillis > 0L) - require(cancellationReconcileWindowMillis in 0L..MAX_CANCELLATION_RECONCILE_WINDOW_MILLIS) - require(cancellationReconcilePollMillis > 0L) scope.launch { try { val storageFailure = runCatching { preparePrivateStorage() }.exceptionOrNull() @@ -220,14 +232,28 @@ class JvmSupportIntake( }, ), idempotencyKey = secureIdempotencyKey(), - createdAtEpochMillis = System.currentTimeMillis().coerceAtLeast(0L), + createdAtEpochMillis = currentTimeMillis().coerceAtLeast(0L), originAccountIdentity = originAccountIdentity, - cancellationPending = cancellationRequested.get(), context = context, ) - synchronized(lock) { pending = submission } + beforePendingSubmissionInstall() + val submissionInstalled = synchronized(lock) { + if (cancellationRequested.get()) { + false + } else { + pending = submission + true + } + } + if (!submissionInstalled) { + publishState(SupportDiagnosticsSubmissionState.Cancelling, originAccountIdentity) + return@withContext + } if (!persistPendingSafely(submission)) { - finishRejected(submission, "The private support submission could not be retained safely on this device.") + finishRejectedIfCurrent( + submission, + "The private support submission could not be retained safely on this device.", + ) return@withContext } if (!packageSubmission(submission)) return@withContext @@ -268,24 +294,24 @@ class JvmSupportIntake( if (!submission.belongsTo(synchronized(lock) { activeAccountIdentity })) { return@withContext } - if (submission.recoveryExpired(System.currentTimeMillis())) { - finishRejected(submission, "The private report recovery capability expired and was removed from this device.") - return@withContext - } if (submission.cancellationPending) { cancellationRequested.set(true) - val receipt = submission.receipt - if (receipt == null) { - reconcileAfterAmbiguousResult( - submission, - IOException("Cancellation still needs to be reconciled."), - ) - } else { - deleteCancelledReceipt(submission, receipt) - } + publishState( + SupportDiagnosticsSubmissionState.Cancelling, + submission.originAccountIdentity, + ) + finishOrReconcileCancellation(submission) return@withContext } - val waitMillis = submission.retryNotBeforeEpochMillis?.minus(System.currentTimeMillis()) ?: 0L + if (submission.recoveryExpired(currentTimeMillis())) { + beforeRecoveryExpiryDisposition() + finishRejectedIfCurrent( + submission, + "The private report recovery capability expired and was removed from this device.", + ) + return@withContext + } + val waitMillis = submission.retryNotBeforeEpochMillis?.minus(currentTimeMillis()) ?: 0L if (waitMillis > 0L) { publishState(SupportDiagnosticsSubmissionState.RetryableFailure( "Obiente Support asked the app to wait before retrying. Try again shortly.", @@ -300,15 +326,35 @@ class JvmSupportIntake( ) return@withContext } - cancellationRequested.set(false) + beforeRetryUploadTransition() + val continueRetry = synchronized(lock) { + when { + pending !== submission -> false + submission.cancellationPending || cancellationRequested.get() -> false + else -> { + cancellationRequested.set(false) + true + } + } + } + if (!continueRetry) { + if (synchronized(lock) { pending === submission }) { + finishOrReconcileCancellation(submission) + } + return@withContext + } + afterRetryUploadTransition() submission.retryNotBeforeEpochMillis = null if (!persistPendingSafely(submission)) { - finishRejected(submission, "The private support submission could not be retained safely on this device.") + finishRejectedIfCurrent( + submission, + "The private support submission could not be retained safely on this device.", + ) return@withContext } if (submission.archive == null && !packageSubmission(submission)) return@withContext if (submission.archive?.isFile != true) { - finishRejected(submission, "The pending private report archive is unavailable.") + finishRejectedIfCurrent(submission, "The pending private report archive is unavailable.") return@withContext } upload(submission) @@ -324,6 +370,7 @@ class JvmSupportIntake( private fun supportMutationsAreAllowed(): Boolean = runCatching(supportMutationsAllowed).getOrDefault(false) private fun endOperation() { + var cancellationContinuation: PendingSubmission? = null synchronized(lock) { if (actualState is SupportDiagnosticsSubmissionState.Cancelling && pending == null) { publishStateLocked( @@ -332,18 +379,41 @@ class JvmSupportIntake( ) } operationActive.set(false) - refreshVisibleStateLocked() + val submission = pending + if ( + cancellationContinuationRequested.getAndSet(false) && + submission?.cancellationPending == true && + !shutdownRequested.get() && + operationActive.compareAndSet(false, true) + ) { + cancellationContinuation = submission + } else { + refreshVisibleStateLocked() + } + } + cancellationContinuation?.let { submission -> + try { + publishState( + SupportDiagnosticsSubmissionState.Cancelling, + submission.originAccountIdentity, + ) + cancelPendingSubmission(submission) + } finally { + endOperation() + } } } suspend fun cancel(): Boolean { awaitInitialization() + var callAtIntent: Call? = null // Serialize the terminal receipt decision with publication of the user's intent. If receipt // completion wins and clears pending first, cancellation is correctly reported as too late. val pendingCancellation: Boolean? = synchronized(lock) { val submission = pending when { submission?.belongsTo(activeAccountIdentity) == true -> { + callAtIntent = activeCall.get() cancellationRequested.set(true) true } @@ -358,7 +428,10 @@ class JvmSupportIntake( return when (pendingCancellation) { null -> false false -> true - true -> withContext(Dispatchers.IO) { cancelAfterIntentPublished() } + true -> { + afterCancellationIntentPublished() + withContext(Dispatchers.IO) { cancelAfterIntentPublished(callAtIntent) } + } } } @@ -395,17 +468,32 @@ class JvmSupportIntake( } } - private fun cancelAfterIntentPublished(): Boolean { - val submission = synchronized(lock) { pending } + private fun cancelAfterIntentPublished(callAtIntent: Call?): Boolean { + var localCancellationCommitted = false + val submission = synchronized(lock) { + pending?.also { current -> + if (!current.cancellationRequiresTombstone && activeCall.get() == null) { + // Claim the terminal decision while holding the same lock used by the upload + // marker transition. The upload cannot publish a tombstone requirement after + // cancellation has already removed this submission from the active state. + pending = null + localCancellationCommitted = true + } + } + } if (submission != null) { - if (!submission.outcomeAmbiguous && activeCall.get() == null) { + if (localCancellationCommitted) { finishCancelled(submission) return true } submission.cancellationPending = true submission.outcomeAmbiguous = true - val cancellationPersisted = persistPendingSafely(submission) - val call = activeCall.getAndSet(null) + val archive = submission.archive + val cancellationPersisted = persistMinimalCancellationSafely( + submission, + deleteArchiveAfterPersist = false, + ) + val call = callAtIntent?.takeIf { activeCall.compareAndSet(it, null) } if (!cancellationPersisted) { call?.cancel() publishState(SupportDiagnosticsSubmissionState.RetryableFailure( @@ -414,20 +502,39 @@ class JvmSupportIntake( )) return false } + call?.cancel() + deletePrivateFileOrRetry(archive) if (call != null) { publishState( SupportDiagnosticsSubmissionState.Cancelling, submission.originAccountIdentity, ) - call.cancel() return true } - } - if (submission != null) { - publishState(SupportDiagnosticsSubmissionState.RetryableFailure( - "Cancellation could not be confirmed. Retry safely to reconcile and delete the private report.", - outcomeAmbiguous = true, - )) + val startCancellation = synchronized(lock) { + if (operationActive.compareAndSet(false, true)) { + true + } else { + cancellationContinuationRequested.set(true) + false + } + } + if (startCancellation) { + try { + publishState( + SupportDiagnosticsSubmissionState.Cancelling, + submission.originAccountIdentity, + ) + cancelPendingSubmission(submission) + } finally { + endOperation() + } + } else { + publishState( + SupportDiagnosticsSubmissionState.Cancelling, + submission.originAccountIdentity, + ) + } return true } cancellationRequested.compareAndSet(true, false) @@ -586,7 +693,7 @@ class JvmSupportIntake( private suspend fun packageSubmission(submission: PendingSubmission): Boolean { if (cancellationRequested.get()) { - finishCancelled(submission) + finishOrReconcileCancellation(submission) return false } publishState(SupportDiagnosticsSubmissionState.Packaging) @@ -599,7 +706,7 @@ class JvmSupportIntake( } catch (cancellation: CancellationException) { deletePrivateFileOrRetry(destination) if (cancellationRequested.get() || synchronized(lock) { pending !== submission }) { - finishCancelled(submission) + finishOrReconcileCancellation(submission) } else { retainForRetry( submission, @@ -611,7 +718,7 @@ class JvmSupportIntake( } catch (_: Throwable) { deletePrivateFileOrRetry(destination) if (cancellationRequested.get() || synchronized(lock) { pending !== submission }) { - finishCancelled(submission) + finishOrReconcileCancellation(submission) } else { retainForRetry( submission, @@ -636,13 +743,28 @@ class JvmSupportIntake( deletePrivateFileOrRetry(prepared.archive) return false } - submission.archive = prepared.archive + beforeArchivePromotion() + val archivePromoted = synchronized(lock) { + if (pending !== submission || cancellationRequested.get() || submission.cancellationPending) { + false + } else { + submission.archive = prepared.archive + true + } + } + if (!archivePromoted) { + deletePrivateFileOrRetry(prepared.archive) + return false + } if (!persistPendingSafely(submission)) { - finishRejected(submission, "The private support submission could not be retained safely on this device.") + finishRejectedIfCurrent( + submission, + "The private support submission could not be retained safely on this device.", + ) return false } if (cancellationRequested.get()) { - finishCancelled(submission) + finishOrReconcileCancellation(submission) return false } return true @@ -650,26 +772,26 @@ class JvmSupportIntake( private suspend fun upload(submission: PendingSubmission) { if (cancellationRequested.get()) { - finishCancelled(submission) + finishOrReconcileCancellation(submission) return } val mutationAllowedBeforePreparation = supportMutationsAreAllowed() if (cancellationRequested.get() || synchronized(lock) { pending !== submission }) { - finishCancelled(submission) + finishOrReconcileCancellation(submission) return } if (!mutationAllowedBeforePreparation) { retainForRetry(submission, READ_ONLY_SUPPORT_MESSAGE, ambiguous = false) return } - submission.latestUploadAttemptAtEpochMillis = System.currentTimeMillis().coerceAtLeast(0L) - submission.outcomeAmbiguous = true - if (!persistPendingSafely(submission)) { - finishRejected(submission, "The private support submission could not be retained safely on this device.") + val archive = requireNotNull(submission.archive) { "The private support archive has not been prepared." } + beforeUploadArchiveValidation() + if (!archive.isFile || archive.length() !in 1L..MAX_SUPPORT_ARCHIVE_BYTES) { + if (synchronized(lock) { pending === submission }) { + finishRejectedIfCurrent(submission, "The pending private report archive is unavailable.") + } return } - val archive = requireNotNull(submission.archive) { "The private support archive has not been prepared." } - require(archive.isFile && archive.length() in 1L..MAX_SUPPORT_ARCHIVE_BYTES) val metadata = json.encodeToString(SupportIntakeMetadata.serializer(), submission.metadata) val progressBody = ProgressRequestBody( delegate = archive.asRequestBody(SUPPORT_ARCHIVE_MEDIA_TYPE), @@ -694,19 +816,78 @@ class JvmSupportIntake( .build() val mutationAllowedAtTransport = supportMutationsAreAllowed() if (cancellationRequested.get() || synchronized(lock) { pending !== submission }) { - finishCancelled(submission) + finishOrReconcileCancellation(submission) return } if (!mutationAllowedAtTransport) { + beforeTransportGateFailureDisposition() retainForRetry(submission, READ_ONLY_SUPPORT_MESSAGE, ambiguous = false) return } + beforeUploadMarker() + var previousLatestUploadAttemptAtEpochMillis: Long? = null + var previousOutcomeAmbiguous = false + var previousCancellationRequiresTombstone = false + val uploadMarked = synchronized(lock) { + if (pending !== submission || cancellationRequested.get() || submission.cancellationPending) { + false + } else { + previousLatestUploadAttemptAtEpochMillis = submission.latestUploadAttemptAtEpochMillis + previousOutcomeAmbiguous = submission.outcomeAmbiguous + previousCancellationRequiresTombstone = submission.cancellationRequiresTombstone + submission.latestUploadAttemptAtEpochMillis = currentTimeMillis().coerceAtLeast(0L) + submission.outcomeAmbiguous = true + submission.cancellationRequiresTombstone = true + true + } + } + if (!uploadMarked) { + if (synchronized(lock) { pending === submission }) { + finishOrReconcileCancellation(submission) + } + return + } + if (!persistPendingSafely(submission)) { + val stillPending = synchronized(lock) { + if (pending !== submission) { + false + } else { + submission.latestUploadAttemptAtEpochMillis = previousLatestUploadAttemptAtEpochMillis + submission.outcomeAmbiguous = previousOutcomeAmbiguous + submission.cancellationRequiresTombstone = previousCancellationRequiresTombstone + true + } + } + if (stillPending) { + if (previousCancellationRequiresTombstone) { + publishState(SupportDiagnosticsSubmissionState.RetryableFailure( + "The retry could not be stored safely. The earlier upload remains recoverable on this device.", + outcomeAmbiguous = previousOutcomeAmbiguous, + )) + } else { + finishRejectedIfCurrent( + submission, + "The private support submission could not be retained safely on this device.", + ) + } + } + return + } publishState(SupportDiagnosticsSubmissionState.Uploading(0f)) val call = client.newCall(request) beforeCallRegistration() if (!registerActiveCall(submission, call, allowCancellationRequested = false)) { call.cancel() + synchronized(lock) { + if (pending === submission) { + submission.latestUploadAttemptAtEpochMillis = previousLatestUploadAttemptAtEpochMillis + submission.outcomeAmbiguous = previousOutcomeAmbiguous + submission.cancellationRequiresTombstone = previousCancellationRequiresTombstone + } + } when { + cancellationRequested.get() && previousCancellationRequiresTombstone -> + finishOrReconcileCancellation(submission) cancellationRequested.get() -> finishCancelled(submission) synchronized(lock) { pending === submission } -> retainForRetry( submission, @@ -720,6 +901,12 @@ class JvmSupportIntake( call.execute().use { response -> val responseText = response.readBoundedText() activeCall.compareAndSet(call, null) + afterUploadResponse() + if (cancellationRequested.get() || submission.cancellationPending) { + cancelPendingSubmission(submission) + return + } + beforeUploadResponseDisposition() when { response.isSuccessful -> finishReceived(submission, decodeReceipt(responseText)) response.code == 408 -> reconcileAfterAmbiguousResult( @@ -732,7 +919,16 @@ class JvmSupportIntake( ambiguous = false, retryNotBeforeEpochMillis = response.retryNotBeforeEpochMillis(), ) - response.code in 400..499 -> finishRejected(submission, decodeProblem(responseText)) + response.code == 410 && isSubmissionCancelledProblem(responseText) -> finishCancelled(submission) + response.code == 410 -> retainForRetry( + submission, + "Obiente Support returned an unverified cancellation result. Retry safely to reconcile the report.", + ambiguous = true, + ) + response.code in 400..499 -> finishRejectedIfCurrent( + submission, + decodeProblem(responseText), + ) else -> retainForRetry( submission, "Obiente Support is temporarily unavailable.", @@ -766,90 +962,77 @@ class JvmSupportIntake( submission: PendingSubmission, uploadFailure: IOException, ) { + if (cancellationRequested.get() || submission.cancellationPending) { + cancelPendingSubmission(submission) + return + } val request = Request.Builder() .url(baseUrl.newBuilder().addPathSegments("api/v1/receipts").build()) .header("Accept", "application/json") .header("Idempotency-Key", submission.idempotencyKey) .get() .build() - var cancellationDeadlineNanos: Long? = null - while (true) { - val call = client.newCall(request) - if (!registerActiveCall(submission, call, allowCancellationRequested = true)) { - call.cancel() - if (synchronized(lock) { pending === submission }) { - retainForRetry( - submission, - "The upload result still needs to be reconciled. You can retry it safely.", - ambiguous = true, - ) - } - return + val call = client.newCall(request) + beforeReceiptCallRegistration() + if (!registerActiveCall(submission, call, allowCancellationRequested = false)) { + call.cancel() + retainForRetry( + submission, + "The upload result still needs to be reconciled. You can retry it safely.", + ambiguous = true, + ) + return + } + val responseResult = try { + call.execute().use { response -> + response.code to response.readBoundedText() } - val responseResult = try { - call.execute().use { response -> - response.code to response.readBoundedText() - } - } catch (_: IOException) { + } catch (_: IOException) { + if (cancellationRequested.get() || submission.cancellationPending) { + cancelPendingSubmission(submission) + } else { retainForRetry( submission, - if (cancellationRequested.get()) { - "Cancellation could not be confirmed. Reconcile the private submission before retrying." - } else uploadFailure.message?.filterSupportMetadata(MAX_SUPPORT_INTAKE_MESSAGE_LENGTH) + uploadFailure.message?.filterSupportMetadata(MAX_SUPPORT_INTAKE_MESSAGE_LENGTH) ?.takeIf(String::isNotBlank) ?: "The upload result is uncertain. Check your connection before retrying.", true, ) - return - } finally { - activeCall.compareAndSet(call, null) } - val (responseCode, responseText) = responseResult - when { - responseCode in 200..299 -> { - try { - finishReceived(submission, decodeReceipt(responseText)) - } catch (_: IOException) { - retainForRetry(submission, "Obiente Support returned an invalid receipt.", true) - } catch (_: IllegalArgumentException) { - retainForRetry(submission, "Obiente Support returned an invalid receipt.", true) - } - return - } - responseCode == 404 && cancellationRequested.get() -> { - val nowNanos = System.nanoTime() - val deadlineNanos = cancellationDeadlineNanos - ?: nowNanos.saturatingAdd(cancellationReconcileWindowMillis * NANOS_PER_MILLISECOND) - .also { cancellationDeadlineNanos = it } - val remainingNanos = deadlineNanos - nowNanos - if (remainingNanos > 0L) { - val delayMillis = minOf( - cancellationReconcilePollMillis, - (remainingNanos / NANOS_PER_MILLISECOND).coerceAtLeast(1L), - ) - delay(delayMillis) - continue - } - retainForRetry( - submission, - "Support has not confirmed receipt yet. Retry again to finish deleting the private report safely.", - ambiguous = true, - ) - return - } - responseCode == 404 -> { - retainForRetry(submission, "The upload did not complete. You can retry it safely.", false) - return - } - else -> { - retainForRetry( - submission, - "The upload result is uncertain. Check your connection before retrying.", - true, - ) - return + return + } finally { + activeCall.compareAndSet(call, null) + } + afterReceiptLookup() + if (cancellationRequested.get() || submission.cancellationPending) { + cancelPendingSubmission(submission) + return + } + beforeReceiptResponseDisposition() + val (responseCode, responseText) = responseResult + when { + responseCode in 200..299 -> { + try { + finishReceived(submission, decodeReceipt(responseText)) + } catch (_: IOException) { + retainForRetry(submission, "Obiente Support returned an invalid receipt.", true) + } catch (_: IllegalArgumentException) { + retainForRetry(submission, "Obiente Support returned an invalid receipt.", true) } } + responseCode == 404 -> + retainForRetry(submission, "The upload did not complete. You can retry it safely.", false) + responseCode == 410 && isSubmissionCancelledProblem(responseText) -> finishCancelled(submission) + responseCode == 410 -> retainForRetry( + submission, + "Obiente Support returned an unverified cancellation result. Retry safely to reconcile the report.", + ambiguous = true, + ) + else -> retainForRetry( + submission, + "The upload result is uncertain. Check your connection before retrying.", + true, + ) } } @@ -867,42 +1050,42 @@ class JvmSupportIntake( } } if (!submitReceivedReport) { - deleteCancelledReceipt(submission, receipt) + cancelPendingSubmission(submission, receipt) } else { finishSubmitted(submission, receipt) } } - private fun deleteCancelledReceipt(submission: PendingSubmission, receipt: SupportIntakeReceipt) { - val statusUrl = validateReceipt(receipt) - val deletionUrl = receipt.deletionUrl.toHttpUrl() - require( - deletionUrl.scheme == statusUrl.scheme && - deletionUrl.host == statusUrl.host && - deletionUrl.port == statusUrl.port && - deletionUrl.encodedPath == statusUrl.encodedPath && - deletionUrl.encodedQuery == null && - deletionUrl.fragment == null, - ) + private fun cancelPendingSubmission( + submission: PendingSubmission, + receipt: SupportIntakeReceipt? = submission.receipt, + ) { + cancellationContinuationRequested.set(false) submission.cancellationPending = true submission.outcomeAmbiguous = true submission.receipt = receipt - persistPendingSafely(submission) - val capability = statusUrl.pathSegments.last() + if (!persistMinimalCancellationSafely(submission)) { + publishState(SupportDiagnosticsSubmissionState.RetryableFailure( + "Cancellation was not sent because its recovery state could not be stored safely. Keep the app open and retry.", + outcomeAmbiguous = true, + )) + return + } val request = Request.Builder() - .url(baseUrl.newBuilder().addPathSegments("api/v1/reports").addPathSegment(capability).build()) + .url(baseUrl.newBuilder().addPathSegments("api/v1/receipts").build()) .header("Accept", "application/json") + .header("Idempotency-Key", submission.idempotencyKey) .delete() .build() if (!supportMutationsAreAllowed()) { - retainCancellationForRetry(submission, receipt, READ_ONLY_SUPPORT_MESSAGE) + retainCancellationForRetry(submission, READ_ONLY_SUPPORT_MESSAGE) return } val call = client.newCall(request) if (!registerActiveCall(submission, call, allowCancellationRequested = true)) { call.cancel() if (synchronized(lock) { pending === submission }) { - retainCancellationForRetry(submission, receipt, "Deletion still needs to be confirmed. Retry safely.") + retainCancellationForRetry(submission, "Cancellation still needs to be confirmed. Retry safely.") } return } @@ -911,75 +1094,62 @@ class JvmSupportIntake( response.readBoundedText() activeCall.compareAndSet(call, null) when { - response.code in TERMINAL_DELETION_STATUS_CODES || response.code == 404 -> - finishCancelled(submission) - response.isSuccessful -> verifyDeletionAfterAccepted( - submission, - receipt, - capability, - ) + response.code == 204 -> finishCancelled(submission) else -> retainCancellationForRetry( submission, - receipt, - "Deletion could not be confirmed. Retry safely to delete the private report.", + "Obiente Support did not confirm cancellation. Retry safely; the private report remains recoverable.", ) } } } catch (_: IOException) { retainCancellationForRetry( submission, - receipt, - "Deletion could not be confirmed. Check your connection, then retry safely.", + "Cancellation could not be confirmed. Check your connection, then retry safely.", ) } finally { activeCall.compareAndSet(call, null) } } - private fun verifyDeletionAfterAccepted( + private fun persistMinimalCancellationSafely( submission: PendingSubmission, - receipt: SupportIntakeReceipt, - capability: String, - ) { - val request = Request.Builder() - .url(baseUrl.newBuilder().addPathSegments("api/v1/reports").addPathSegment(capability).build()) - .header("Accept", "application/json") - .get() - .build() - val call = client.newCall(request) - if (!registerActiveCall(submission, call, allowCancellationRequested = true)) { - call.cancel() - if (synchronized(lock) { pending === submission }) { - retainCancellationForRetry( - submission, - receipt, - "Deletion verification was interrupted. Retry safely.", - ) - } - return - } - try { - call.execute().use { response -> - response.readBoundedText() - if (response.code == 404) { - finishCancelled(submission) - } else { - retainCancellationForRetry( - submission, - receipt, - "Deletion is still being processed. Retry safely to verify the private report was removed.", - ) - } - } - } catch (_: IOException) { - retainCancellationForRetry( - submission, - receipt, - "Deletion was accepted but could not be verified. Check your connection, then retry safely.", - ) - } finally { - activeCall.compareAndSet(call, null) - } + deleteArchiveAfterPersist: Boolean = true, + ): Boolean = synchronized(persistenceLock) { + val archive = submission.archive + val metadata = submission.metadata + val context = submission.context + val receipt = submission.receipt + val retryNotBeforeEpochMillis = submission.retryNotBeforeEpochMillis + stripPrivateCancellationPayload(submission) + if (!persistPendingSafely(submission)) { + submission.archive = archive + submission.metadata = metadata + submission.context = context + submission.receipt = receipt + submission.retryNotBeforeEpochMillis = retryNotBeforeEpochMillis + return@synchronized false + } + if (deleteArchiveAfterPersist) deletePrivateFileOrRetry(archive) + true + } + + private fun stripPrivateCancellationPayload(submission: PendingSubmission): File? { + val archive = submission.archive + submission.archive = null + submission.metadata = SupportIntakeMetadata( + title = "", + description = "", + release = SupportIntakeRelease("", "", "", "", ""), + ) + submission.context = PreparedSupportSubmissionContext( + sanitizedReproductionSteps = null, + featureState = emptyList(), + confirmedAtEpochMillis = 0L, + events = emptyList(), + ) + submission.receipt = null + submission.retryNotBeforeEpochMillis = null + return archive } private fun finishSubmitted(submission: PendingSubmission, receipt: SupportIntakeReceipt) { @@ -1124,6 +1294,22 @@ class JvmSupportIntake( publishState(SupportDiagnosticsSubmissionState.Rejected(message), submission.originAccountIdentity) } + private fun finishRejectedIfCurrent(submission: PendingSubmission, message: String) { + val rejectionCommitted = synchronized(lock) { + if ( + pending === submission && + !cancellationRequested.get() && + !submission.cancellationPending + ) { + pending = null + true + } else { + false + } + } + if (rejectionCommitted) finishRejected(submission, message) + } + private fun finishCancelled(submission: PendingSubmission) { finishTerminal(submission) publishState( @@ -1136,35 +1322,79 @@ class JvmSupportIntake( ) } + private fun finishOrReconcileCancellation(submission: PendingSubmission) { + if ( + synchronized(lock) { + pending === submission && submission.cancellationRequiresTombstone + } + ) { + cancelPendingSubmission(submission) + } else { + finishCancelled(submission) + } + } + private fun retainForRetry( submission: PendingSubmission, message: String, ambiguous: Boolean, retryNotBeforeEpochMillis: Long? = null, ) { - submission.outcomeAmbiguous = ambiguous - submission.retryNotBeforeEpochMillis = retryNotBeforeEpochMillis - synchronized(lock) { pending = submission } - if (persistPendingSafely(submission)) { - publishState(SupportDiagnosticsSubmissionState.RetryableFailure(message, ambiguous)) - } else { - // Atomic replacement keeps the descriptor from immediately before the request. That - // record retains the idempotency key and conservatively requires reconciliation. - publishState(SupportDiagnosticsSubmissionState.RetryableFailure( - "The updated retry state could not be stored. Keep the app open and retry safely to reconcile the report.", - outcomeAmbiguous = true, - )) + val retentionCommitted = synchronized(lock) { + if ( + pending !== submission || + cancellationRequested.get() || + submission.cancellationPending + ) { + false + } else { + submission.outcomeAmbiguous = ambiguous + submission.retryNotBeforeEpochMillis = retryNotBeforeEpochMillis + true + } + } + if (!retentionCommitted) { + if (synchronized(lock) { pending === submission }) { + finishOrReconcileCancellation(submission) + } + return + } + val persisted = persistPendingSafely(submission) + val retryStatePublished = synchronized(lock) { + if ( + pending === submission && + !cancellationRequested.get() && + !submission.cancellationPending + ) { + publishStateLocked( + if (persisted) { + SupportDiagnosticsSubmissionState.RetryableFailure(message, ambiguous) + } else { + // Atomic replacement keeps the descriptor from immediately before the + // request, including its idempotency key and recovery requirement. + SupportDiagnosticsSubmissionState.RetryableFailure( + "The updated retry state could not be stored. Keep the app open and retry safely to reconcile the report.", + outcomeAmbiguous = true, + ) + }, + submission.originAccountIdentity, + ) + true + } else { + false + } + } + if (!retryStatePublished && synchronized(lock) { pending === submission }) { + finishOrReconcileCancellation(submission) } } private fun retainCancellationForRetry( submission: PendingSubmission, - receipt: SupportIntakeReceipt, message: String, ) { submission.cancellationPending = true submission.outcomeAmbiguous = true - submission.receipt = receipt synchronized(lock) { pending = submission } if (persistPendingSafely(submission)) { publishState(SupportDiagnosticsSubmissionState.RetryableFailure(message, outcomeAmbiguous = true)) @@ -1172,7 +1402,7 @@ class JvmSupportIntake( // The receipt was persisted before deletion began. Atomic replacement leaves that last // valid recovery record in place when this newer retry-state write fails. publishState(SupportDiagnosticsSubmissionState.RetryableFailure( - "Deletion was not confirmed and its updated retry state could not be stored. Keep the app open and retry.", + "Cancellation was not confirmed and its updated retry state could not be stored. Keep the app open and retry.", outcomeAmbiguous = true, )) } @@ -1196,6 +1426,13 @@ class JvmSupportIntake( ?.takeIf(String::isNotBlank) ?: "Obiente Support rejected this diagnostic report." + private fun isSubmissionCancelledProblem(response: String): Boolean = runCatching { + json.decodeFromString(SupportIntakeProblem.serializer(), response).let { problem -> + problem.contractVersion == SUPPORT_INTAKE_CONTRACT_VERSION && + problem.code == SUPPORT_SUBMISSION_CANCELLED_CODE + } + }.getOrDefault(false) + private fun pruneTemporaryReports(retainedArchive: File?) { val cutoff = System.currentTimeMillis() - SUPPORT_TEMPORARY_MAX_AGE_MILLIS temporaryRoot.listFiles().orEmpty() @@ -1261,6 +1498,7 @@ class JvmSupportIntake( context = submission.context, cancellationPending = submission.cancellationPending, outcomeAmbiguous = submission.outcomeAmbiguous, + cancellationRequiresTombstone = submission.cancellationRequiresTombstone, latestUploadAttemptAtEpochMillis = submission.latestUploadAttemptAtEpochMillis, retryNotBeforeEpochMillis = submission.retryNotBeforeEpochMillis, receipt = submission.receipt, @@ -1381,13 +1619,13 @@ class JvmSupportIntake( } val recoveryDeadlineEpochMillis = persisted.receipt ?.let { receipt -> Instant.parse(receipt.retentionUntil).toEpochMilli() } - ?: if (persisted.outcomeAmbiguous) { + ?: if (persisted.outcomeAmbiguous || persisted.cancellationRequiresTombstone == true) { (persisted.latestUploadAttemptAtEpochMillis ?: persisted.createdAtEpochMillis) .saturatingAdd(SUPPORT_RECOVERY_MAX_AGE_MILLIS) } else { persisted.createdAtEpochMillis.saturatingAdd(SUPPORT_RECOVERY_MAX_AGE_MILLIS) } - require(nowEpochMillis <= recoveryDeadlineEpochMillis) + require(persisted.cancellationPending || nowEpochMillis <= recoveryDeadlineEpochMillis) val archiveAgeMillis = (nowEpochMillis - persisted.createdAtEpochMillis).coerceAtLeast(0L) val archiveIsRetained = archiveAgeMillis <= SUPPORT_TEMPORARY_MAX_AGE_MILLIS val archive = persisted.archiveName?.let { archiveName -> @@ -1411,7 +1649,7 @@ class JvmSupportIntake( } } pendingDescriptorRestorePending.set(false) - PendingSubmission( + val restored = PendingSubmission( archive = archive, metadata = persisted.metadata, idempotencyKey = persisted.idempotencyKey, @@ -1420,10 +1658,22 @@ class JvmSupportIntake( context = persisted.context, cancellationPending = persisted.cancellationPending, outcomeAmbiguous = persisted.outcomeAmbiguous, + cancellationRequiresTombstone = persisted.cancellationRequiresTombstone + ?: ( + persisted.latestUploadAttemptAtEpochMillis != null || + persisted.outcomeAmbiguous || + persisted.receipt != null + ), latestUploadAttemptAtEpochMillis = persisted.latestUploadAttemptAtEpochMillis, retryNotBeforeEpochMillis = retryNotBeforeEpochMillis, receipt = persisted.receipt, ) + if (restored.cancellationPending) { + val privateArchive = stripPrivateCancellationPayload(restored) + persistPending(restored) + deletePrivateFileOrRetry(privateArchive) + } + restored } catch (failure: Throwable) { if (failure is IOException || failure is SecurityException) { val retryWasNotScheduled = pendingDescriptorRestorePending.compareAndSet(false, true) @@ -1526,7 +1776,7 @@ class JvmSupportIntake( } ?: pendingSubmission?.let { submission -> SupportDiagnosticsSubmissionState.RetryableFailure( if (submission.cancellationPending) { - "Cancellation was interrupted. Retry safely to reconcile and delete the private report." + "Cancellation was interrupted. Retry safely to obtain terminal confirmation from Obiente Support." } else if (submission.archive == null) { "Private report preparation was interrupted. You can retry it safely." } else { @@ -1600,13 +1850,14 @@ class JvmSupportIntake( private data class PendingSubmission( var archive: File?, - val metadata: SupportIntakeMetadata, + var metadata: SupportIntakeMetadata, val idempotencyKey: String, val createdAtEpochMillis: Long, val originAccountIdentity: String, - val context: PreparedSupportSubmissionContext, + var context: PreparedSupportSubmissionContext, var cancellationPending: Boolean = false, var outcomeAmbiguous: Boolean = false, + var cancellationRequiresTombstone: Boolean = false, var latestUploadAttemptAtEpochMillis: Long? = null, var retryNotBeforeEpochMillis: Long? = null, var receipt: SupportIntakeReceipt? = null, @@ -1616,7 +1867,7 @@ class JvmSupportIntake( fun recoveryExpired(nowEpochMillis: Long): Boolean { val deadline = receipt ?.let { value -> runCatching { Instant.parse(value.retentionUntil).toEpochMilli() }.getOrNull() } - ?: if (outcomeAmbiguous) { + ?: if (outcomeAmbiguous || cancellationRequiresTombstone) { (latestUploadAttemptAtEpochMillis ?: createdAtEpochMillis) .saturatingAdd(SUPPORT_RECOVERY_MAX_AGE_MILLIS) } else { @@ -1637,6 +1888,13 @@ class JvmSupportIntake( fun isRetained(nowEpochMillis: Long): Boolean = nowEpochMillis <= retentionUntilEpochMillis } + @Serializable + private data class SupportIntakeProblem( + val contractVersion: Int, + val code: String, + val message: String, + ) + @Serializable private data class PersistedPendingSubmission( val archiveName: String?, @@ -1647,7 +1905,11 @@ class JvmSupportIntake( val context: PreparedSupportSubmissionContext, val cancellationPending: Boolean = false, val outcomeAmbiguous: Boolean = true, + val cancellationRequiresTombstone: Boolean? = null, val latestUploadAttemptAtEpochMillis: Long? = null, + // Read descriptors written by early PR #386 builds, but never use wall time to confirm + // cancellation. Only the server's idempotency-key tombstone is terminal. + val cancellationRequestedAtEpochMillis: Long? = null, val retryNotBeforeEpochMillis: Long? = null, val receipt: SupportIntakeReceipt? = null, ) @@ -1906,6 +2168,7 @@ private val SUPPORT_IDEMPOTENCY_PATTERN = Regex("[A-Za-z0-9_-]{43}") private val SUPPORT_ACCOUNT_IDENTITY_PATTERN = Regex("[0-9a-f]{32}(?:[0-9a-f]{32})?") private val RETRYABLE_CLIENT_STATUS_CODES = setOf(425, 429) private val TERMINAL_DELETION_STATUS_CODES = setOf(200, 204) +private const val SUPPORT_SUBMISSION_CANCELLED_CODE = "submission_cancelled" private const val MAX_SUPPORT_INTAKE_MESSAGE_LENGTH = 240 private const val MAX_SUPPORT_INTAKE_RESPONSE_BYTES = 64 * 1024 private const val MAX_SUPPORT_INTAKE_DESCRIPTION_BYTES = 8_000 @@ -1919,10 +2182,6 @@ private const val SUPPORT_RECOVERY_MAX_AGE_MILLIS = 30L * 24L * 60L * 60L * 1_00 private const val SUPPORT_SERVER_RETENTION_MAX_AGE_MILLIS = 30L * 24L * 60L * 60L * 1_000L private const val SUPPORT_RECEIPT_CLOCK_SKEW_MILLIS = 5L * 60L * 1_000L private const val SUPPORT_DESCRIPTOR_DELETE_RETRY_MILLIS = 60L * 1_000L -private const val SUPPORT_CANCELLATION_RECONCILE_WINDOW_MILLIS = 10L * 1_000L -private const val SUPPORT_CANCELLATION_RECONCILE_POLL_MILLIS = 500L -private const val MAX_CANCELLATION_RECONCILE_WINDOW_MILLIS = 60L * 1_000L -private const val NANOS_PER_MILLISECOND = 1_000_000L private const val SUPPORT_PENDING_RESTORE_MESSAGE = "Private support report recovery is temporarily unavailable. The app will retry automatically." private const val SUPPORT_COMPLETED_RESTORE_MESSAGE =