From ae2f7b68bd3b475908808dbe703d1a09a61cf6e9 Mon Sep 17 00:00:00 2001 From: Dylan Murphy Date: Thu, 1 Oct 2026 12:11:15 -0400 Subject: [PATCH] Release v9 bytes that no row owns, and reuse the v9 blob on a replace cancel(id) on an id with no row resolved and did nothing, so a v9 chunked upload that was still running at the upgrade, which gets no legacy row, kept its manifest and blob on disk for good. cancel(id) now deletes the id's directory when it has no entry file. A live entry and another id's directory are never touched. A same-id mutate() over a legacy row whose v9 manifest could not be read ignored the v9 blob, so it failed E_FILE_MISSING when the caller's source was already gone. It now runs over that blob on both platforms. A caller file that exists still wins. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 3 +- README.md | 6 +- .../backgroundupload/QueueController.kt | 19 ++++-- .../openspace/backgroundupload/QueueStore.kt | 14 +++- .../backgroundupload/QueueControllerTest.kt | 66 +++++++++++++++++++ ios/QueueCoordinator+Enqueue.swift | 6 +- ios/QueueCoordinator.swift | 8 ++- ios/QueueStore.swift | 15 ++++- ios/Tests/CoordinatorRelaunchTests.swift | 46 +++++++++++++ ios/Tests/CoordinatorSimpleTests.swift | 59 ++++++++++++++++- src/NativeRNFileUploader.ts | 8 +-- 11 files changed, 232 insertions(+), 18 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8f604e7d..288f45f8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -183,7 +183,8 @@ Native: is cancelled (Android WorkManager rows, iOS session tasks). v9 chunked manifests and their bytes stay: a same-id `mutate()` with the same parts resumes from the accepted parts, and one with different parts starts over - on the kept bytes. Wi-Fi only starts off. + on the kept bytes. `cancel(id)` on a v9 upload id with no row deletes its + kept files. Wi-Fi only starts off. - **Platform limits.** Android, API 31 and later: a WorkManager run started from the background usually cannot start its foreground service and then has JobScheduler's limit of about 10 minutes; a single body that does not diff --git a/README.md b/README.md index fc72fb7d..39d95296 100644 --- a/README.md +++ b/README.md @@ -390,7 +390,8 @@ kept. A paused entry past `expiresAt` settles `error` with ### `cancel(id): Promise` A live entry settles `cancelled` with reason `user` and is forgotten after its ack. A settled entry is forgotten now: row, bytes, and its -unacknowledged outcomes. An unknown id resolves and does nothing. If the +unacknowledged outcomes. An id with no row resolves. If its directory has +no entry file, such as a v9 upload's bytes, the directory is deleted. If the journal or the store cannot be written, `cancel()` rejects with `E_STORAGE`. The caller may call again. On Android, a cancel whose journal write landed but whose entry save failed is already in effect: the work stops and the @@ -475,7 +476,8 @@ What happens to work a v9 build left behind: - **v9 chunked uploads resume under the same id.** The v9 manifest and bytes stay on disk. A `mutate()` with the v9 upload id and the same parts resumes from the accepted parts; different parts start over on the kept bytes. An - id with no `mutate()` keeps its bytes until `cancel(id)`. + id with no `mutate()` keeps its bytes until `cancel(id)`, which deletes + them. Cancel every v9 upload id you will not send again. - **Wi-Fi only starts off.** Call `setWifiOnly(true)` again if the app had it on. diff --git a/android/src/main/java/ai/openspace/backgroundupload/QueueController.kt b/android/src/main/java/ai/openspace/backgroundupload/QueueController.kt index 98143e6f..448c16af 100644 --- a/android/src/main/java/ai/openspace/backgroundupload/QueueController.kt +++ b/android/src/main/java/ai/openspace/backgroundupload/QueueController.kt @@ -114,8 +114,13 @@ class QueueController( } EnqueueRules.Action.Replace -> { val old = existing!! - // Different parts: a present caller file wins; the old blob is the fallback. - val ownedBlob = if (old.body?.kind == StagedBody.CHUNKED) store.bodyFile(old) else null + // Different parts: a present caller file wins; the old blob is the + // fallback. A legacy row has no body, so its fallback is the v9 blob. + val ownedBlob = when { + old.legacy -> store.blobFile(p.id) + old.body?.kind == StagedBody.CHUNKED -> store.bodyFile(old) + else -> null + } val staged = preStaged(old.generation + 1) ?: stageOrThrow(p.descriptor, dir, old.generation + 1, ownedBlob) commit(EnqueueRules.replaced(old, p, staged, s.isPaused(p.key), s.headerGeneration, now)) @@ -205,7 +210,8 @@ class QueueController( /** * Live: journal a 'cancelled' (user) outcome, then forget after its ack. * Settled (or legacy): forget now, row, bytes, and its unacknowledged - * outcomes. Unknown: no-op. + * outcomes. No row: delete the id's directory when it has no entry file + * ([QueueStore.removeUnowned]). * * A failed journal write rejects E_STORAGE and changes nothing: the entry * goes on, and JS can call cancel() again. @@ -222,8 +228,11 @@ class QueueController( var saved: QueueEntry? = null try { store.locked { - stop = true // unknown or settled: stop any stray work, as before - val e = store.load(id) ?: return@locked + stop = true // no row or settled: stop any stray work + val e = store.load(id) ?: run { + store.removeUnowned(id) + return@locked + } if (!e.isLive || e.legacy) { journal.ackEntry(id) store.remove(id) diff --git a/android/src/main/java/ai/openspace/backgroundupload/QueueStore.kt b/android/src/main/java/ai/openspace/backgroundupload/QueueStore.kt index 3db76532..0f3af6e5 100644 --- a/android/src/main/java/ai/openspace/backgroundupload/QueueStore.kt +++ b/android/src/main/java/ai/openspace/backgroundupload/QueueStore.kt @@ -13,7 +13,7 @@ import java.util.Base64 * * A v10 directory holds `entry.json` and at most one staged body file (see * [StagedBody.fileName]). A v9 directory holds `manifest.json` and `blob` - * until a same-id enqueue adopts it. + * until a same-id enqueue adopts it or cancel deletes it. * * Every write is tmp + fsync + rename ([AtomicFiles]). Bodies are staged * before `entry.json` is saved, so an `entry.json` on disk means its body is @@ -116,6 +116,18 @@ class QueueStore(private val dir: File, private val index: RequestIndex = Reques index.remove(id) } + /** + * Deletes [id]'s directory when it holds no `entry.json`: v9 files that no + * row owns. A directory with an `entry.json`, readable or not, is kept. The + * empty id encodes to the store root, so it deletes nothing. + */ + @Synchronized + fun removeUnowned(id: String) { + if (id.isEmpty()) return + val d = entryDir(id) + if (d.isDirectory && !File(d, ENTRY_FILE).exists()) d.deleteRecursively() + } + /** Every v10 entry. A directory with only a v9 manifest is not a row. */ @Synchronized fun all(): List = diff --git a/android/src/test/java/ai/openspace/backgroundupload/QueueControllerTest.kt b/android/src/test/java/ai/openspace/backgroundupload/QueueControllerTest.kt index 88d631d6..889ae763 100644 --- a/android/src/test/java/ai/openspace/backgroundupload/QueueControllerTest.kt +++ b/android/src/test/java/ai/openspace/backgroundupload/QueueControllerTest.kt @@ -251,6 +251,32 @@ class QueueControllerTest { assertEquals(0, e.bytesSent) } + @Test + fun `a same-id enqueue over a legacy row with an unreadable manifest runs over the v9 blob`() { + val dir = store.entryDir("e1").apply { mkdirs() } + File(dir, "blob").writeBytes(ByteArray(20)) + File(dir, "manifest.json").writeText("""{"id":"e1","parts":[""") + store.save(LegacyImport.legacyRow(LegacyImport.V9Entry("x", "e1", "error", 5))!!) + controller.enqueue(parsed(descriptor = desc(url = null, method = "PUT", file = "/gone.bin", parts = listOf(part(0, 10), part(10, 20))))) + val e = store.load("e1")!! + assertFalse(e.legacy) + assertEquals(2, e.generation) + assertEquals("blob", e.body!!.fileName) + assertTrue(e.descriptor!!.parts!!.none { it.accepted }) + assertEquals(0, e.bytesSent) + assertEquals(setOf("entry.json", "blob"), dirFiles()) // the manifest is pruned + } + + @Test + fun `a same-id enqueue over a legacy row with no blob and a missing file rejects E_FILE_MISSING`() { + store.save(LegacyImport.legacyRow(LegacyImport.V9Entry("x", "e1", "error", 5))!!) + val e = assertThrows(QueueException::class.java) { + controller.enqueue(parsed(descriptor = desc(url = null, method = "PUT", file = "/gone.bin", parts = listOf(part(0, 20))))) + } + assertEquals(QueueException.E_FILE_MISSING, e.code) + assertTrue(store.load("e1")!!.legacy) + } + // MARK: - chunked replace (a present file wins over the old blob) private fun chunkedErrorEntry(): File { @@ -525,6 +551,46 @@ class QueueControllerTest { fun `cancel of an unknown id is a no-op`() { controller.cancel("nope") assertEquals(emptyList(), events.log) + assertFalse(store.entryDir("nope").exists()) + } + + @Test + fun `cancel of an id with only v9 files deletes its directory`() { + val dir = v9Dir("cap", v9Parts, 20) + controller.cancel("cap") + assertFalse(dir.exists()) + assertEquals(emptyList(), events.log) + assertTrue(journal.unacknowledged().isEmpty()) + } + + @Test + fun `cancel of an id with only a v9 blob deletes its directory`() { + val dir = store.entryDir("cap").apply { mkdirs() } + File(dir, "blob").writeBytes(ByteArray(20)) + controller.cancel("cap") + assertFalse(dir.exists()) + } + + @Test + fun `cancel keeps a directory whose entry file can not be read`() { + val dir = v9Dir("cap", v9Parts, 20) + File(dir, "entry.json").writeText("{") + controller.cancel("cap") + assertEquals(setOf("entry.json", "manifest.json", "blob"), dirFiles("cap")) + } + + @Test + fun `cancel of an empty id deletes nothing in the store`() { + controller.enqueue(parsed()) + events.log.clear() + val dir = store.entryDir("cap").apply { mkdirs() } + File(dir, "blob").writeBytes(ByteArray(20)) + controller.setWifiOnly(true) + controller.cancel("") + assertNotNull(store.load("e1")) + assertTrue(File(dir, "blob").exists()) + assertTrue(File(root, "settings.json").exists()) + assertEquals(emptyList(), events.log) } // MARK: - pause, resume, wifi, headers diff --git a/ios/QueueCoordinator+Enqueue.swift b/ios/QueueCoordinator+Enqueue.swift index b5234f0c..9c7bf586 100644 --- a/ios/QueueCoordinator+Enqueue.swift +++ b/ios/QueueCoordinator+Enqueue.swift @@ -146,11 +146,15 @@ extension QueueCoordinator { cancelTasks(existing.id, purpose: .superseded) chunked.stop(existing.id) // A chunked replace may keep the current blob when the caller already - // deleted its source (the part-404 recreate over moved bytes). + // deleted its source (the part-404 recreate over moved bytes). A legacy + // row with no body may still have the v9 blob on disk: its manifest + // could not be read at import. var fallback: String? if let path = existing.bodyPath, existing.isChunked || existing.legacy, path == ChunkedManifestV9.blobName || path.hasPrefix(BodyStaging.blobPrefix) { fallback = path + } else if existing.legacy, existing.bodyPath == nil { + fallback = ChunkedManifestV9.blobName } let staged = try mapStaging { try BodyStaging.stage(p.body, parts: p.parts, into: store.dir(p.id), fallbackBlob: fallback) diff --git a/ios/QueueCoordinator.swift b/ios/QueueCoordinator.swift index 6ee945ff..d7a66e61 100644 --- a/ios/QueueCoordinator.swift +++ b/ios/QueueCoordinator.swift @@ -213,13 +213,17 @@ final class QueueCoordinator { } /// Live entry: journal 'cancelled' (user); forgotten after its ack. - /// Settled entry: forgotten now, with its unacked outcomes. Unknown: no-op. + /// Settled entry: forgotten now, with its unacked outcomes. No row: the id + /// directory is deleted when it has no entry file. /// When the journal or the store cannot be written, rejects E_STORAGE and /// changes nothing: the entry keeps running (or stays settled), and JS may /// call again. func cancel(_ id: String, resolve: @escaping () -> Void, reject: @escaping (String, String) -> Void) { queue.async { - guard let e = self.index.entry(id) else { return resolve() } + guard let e = self.index.entry(id) else { + self.store.removeUnowned(id) + return resolve() + } if e.isLive && !e.legacy { guard self.settle(id, .cancelled(reason: "user"), requireJournal: true) else { return reject("E_STORAGE", "cancel: cannot journal the outcome of '\(id)'; nothing changed") diff --git a/ios/QueueStore.swift b/ios/QueueStore.swift index 3d3382eb..a16015ee 100644 --- a/ios/QueueStore.swift +++ b/ios/QueueStore.swift @@ -8,7 +8,8 @@ import Foundation /// - `entry.json`: the QueueEntry (v10). /// - `body-` or `blob-` / `blob`: the staged body. /// - `part-..-`: a chunked part while in flight. -/// - `manifest.json`: a v9 manifest, until a same-id enqueue adopts it. +/// - `manifest.json`: a v9 manifest, until a same-id enqueue adopts it or +/// cancel deletes it. /// /// Next to the directories: `settings.json` and the `v10-imported` marker. /// Writes are tmp + fsync + rename. A corrupt or half-written file reads as @@ -82,6 +83,18 @@ final class QueueStore { queue.sync { _ = try? FileManager.default.removeItem(at: dir(id)) } } + /// Deletes the id directory when it holds no `entry.json`: v9 files that + /// no row owns. A directory with an `entry.json`, readable or not, is kept. + /// The empty id names the store root, so it deletes nothing. + func removeUnowned(_ id: String) { + guard !id.isEmpty else { return } + queue.sync { + let d = dir(id) + guard FileIO.exists(d), !FileIO.exists(d.appendingPathComponent(Self.entryName)) else { return } + _ = try? FileManager.default.removeItem(at: d) + } + } + // MARK: - Forget in steps (cancel on a settled entry) /// Step one of a forget that must not half-happen: renames the id diff --git a/ios/Tests/CoordinatorRelaunchTests.swift b/ios/Tests/CoordinatorRelaunchTests.swift index 0fc54e3a..16e03ae6 100644 --- a/ios/Tests/CoordinatorRelaunchTests.swift +++ b/ios/Tests/CoordinatorRelaunchTests.swift @@ -486,6 +486,52 @@ final class CoordinatorRelaunchTests: XCTestCase { XCTAssertEqual(l.transport.live.count, 3) } + /// v9 state with a journaled error for a chunked id whose manifest is + /// truncated, and a blob when `blobBytes` is set. + private func legacyErrorWithUnreadableManifest(blobBytes: Int?) -> Harness { + let root = makeTempDir() + let v9Store = QueueStore(root: root.appendingPathComponent("queue")) + let v9Journal = EventJournal(root: root.appendingPathComponent("events")) + V9Journal.write(V9Journal.error(eventId: "ev1", id: "cap-1", timestamp: 100), eventId: "ev1", + into: v9Journal.root) + writeFile(v9Store.fileURL("cap-1", QueueStore.manifestName), #"{"id":"cap-1","parts":["#) + if let blobBytes { writeFile(v9Store.fileURL("cap-1", "blob"), bytes: blobBytes) } + let fresh = Harness(root: root) + fresh.boot() + return fresh + } + + func testLegacyErrorRowWithAnUnreadableManifestRunsOverTheV9Blob() throws { + let l = legacyErrorWithUnreadableManifest(blobBytes: 30) + defer { try? FileManager.default.removeItem(at: l.root) } + let legacy = try XCTUnwrap(l.entry("cap-1")) + XCTAssertTrue(legacy.legacy) + XCTAssertEqual(legacy.state, .error) + XCTAssertNil(legacy.bodyPath, "the import could not read the manifest") + + let gone = l.root.appendingPathComponent("gone") + _ = try l.enqueue(l.chunkedRaw(id: "cap-1", size: 30, parts: 3, source: gone)).get() + let e = try XCTUnwrap(l.entry("cap-1")) + XCTAssertFalse(e.legacy) + XCTAssertEqual(e.generation, 2) + XCTAssertEqual(e.bodyPath, "blob") + XCTAssertTrue(e.parts.allSatisfy { !$0.accepted }, "no manifest to resume from: every part is sent") + XCTAssertEqual(l.row("cap-1")?["bytesSent"] as? Int64, 0) + XCTAssertTrue(FileIO.exists(l.store.fileURL("cap-1", "blob"))) + XCTAssertFalse(FileIO.exists(l.store.fileURL("cap-1", QueueStore.manifestName)), "the bad manifest is deleted") + XCTAssertEqual(l.transport.live.count, 3) + } + + func testLegacyErrorRowWithNoBlobAndAMissingFileRejectsFileMissing() throws { + let l = legacyErrorWithUnreadableManifest(blobBytes: nil) + defer { try? FileManager.default.removeItem(at: l.root) } + let gone = l.root.appendingPathComponent("gone") + guard case .failure(let missing) = l.enqueue(l.chunkedRaw(id: "cap-1", size: 30, parts: 3, source: gone)) + else { return XCTFail("expected E_FILE_MISSING") } + XCTAssertEqual(missing.code, "E_FILE_MISSING") + XCTAssertEqual(l.entry("cap-1")?.legacy, true, "the legacy row stays") + } + func testV9JournalFilesOfEachKindImportAsLegacyRows() throws { let root = makeTempDir() let events = root.appendingPathComponent("events") diff --git a/ios/Tests/CoordinatorSimpleTests.swift b/ios/Tests/CoordinatorSimpleTests.swift index ea1f014d..7915cbc2 100644 --- a/ios/Tests/CoordinatorSimpleTests.swift +++ b/ios/Tests/CoordinatorSimpleTests.swift @@ -248,7 +248,64 @@ final class CoordinatorSimpleTests: XCTestCase { XCTAssertNil(h.row("a")) XCTAssertFalse(FileIO.exists(h.store.dir("a"))) XCTAssertTrue(h.journal.unacknowledged().isEmpty) - h.cancel("unknown") // no-op + XCTAssertTrue(cancelResolves("unknown")) + XCTAssertFalse(FileIO.exists(h.store.dir("unknown"))) + } + + /// True when cancel resolved; a rejection or no answer is false. + private func cancelResolves(_ id: String) -> Bool { + var resolved = false + h.coordinator.cancel(id, resolve: { resolved = true }, reject: { _, _ in }) + h.drain() + return resolved + } + + private func writeV9Upload(_ id: String) throws { + let manifest = ChunkedManifestV9(id: id, parts: [ + .init(url: "https://s3.test/part1", start: 0, end: 10, accepted: true), + .init(url: "https://s3.test/part2", start: 10, end: 20, accepted: false), + ], expiresAt: h.clock + 60_000, incarnation: "v9inc") + writeFile(h.store.fileURL(id, QueueStore.manifestName), + String(data: try JSONEncoder().encode(manifest), encoding: .utf8)!) + writeFile(h.store.fileURL(id, ChunkedManifestV9.blobName), bytes: 20) + } + + func testCancelOfAnIdWithOnlyV9FilesDeletesItsDirectory() throws { + try writeV9Upload("cap") + XCTAssertTrue(cancelResolves("cap")) + XCTAssertFalse(FileIO.exists(h.store.dir("cap"))) + XCTAssertTrue(h.sink.states.isEmpty) + XCTAssertTrue(h.sink.settled.isEmpty) + } + + func testCancelOfAnIdWithOnlyABlobDeletesItsDirectory() { + writeFile(h.store.fileURL("cap", ChunkedManifestV9.blobName), bytes: 20) + XCTAssertTrue(cancelResolves("cap")) + XCTAssertFalse(FileIO.exists(h.store.dir("cap"))) + } + + func testCancelKeepsADirectoryWhoseEntryFileCannotBeRead() throws { + try writeV9Upload("cap") + writeFile(h.store.fileURL("cap", QueueStore.entryName), "{") + XCTAssertTrue(cancelResolves("cap")) + XCTAssertTrue(FileIO.exists(h.store.fileURL("cap", QueueStore.entryName))) + XCTAssertTrue(FileIO.exists(h.store.fileURL("cap", QueueStore.manifestName))) + XCTAssertTrue(FileIO.exists(h.store.fileURL("cap", ChunkedManifestV9.blobName))) + } + + func testCancelOfAnEmptyIdDeletesNothingInTheStore() throws { + _ = try h.enqueue(h.dataRaw(id: "a")).get() + writeFile(h.store.fileURL("cap", ChunkedManifestV9.blobName), bytes: 20) + h.setWifiOnly(true) + XCTAssertTrue(h.store.isImported()) + let states = h.sink.states.count + XCTAssertTrue(cancelResolves("")) + XCTAssertNotNil(h.entry("a")) + XCTAssertNotNil(h.store.load("a")) + XCTAssertTrue(FileIO.exists(h.store.fileURL("cap", ChunkedManifestV9.blobName))) + XCTAssertTrue(FileIO.exists(h.store.root.appendingPathComponent("settings.json"))) + XCTAssertTrue(h.store.isImported()) + XCTAssertEqual(h.sink.states.count, states) } func testCancelWhoseJournalWriteFailsRejectsStorageAndChangesNothing() throws { diff --git a/src/NativeRNFileUploader.ts b/src/NativeRNFileUploader.ts index d9614d9a..1c84804a 100644 --- a/src/NativeRNFileUploader.ts +++ b/src/NativeRNFileUploader.ts @@ -68,10 +68,10 @@ export interface Spec extends TurboModule { pause(scope: CodegenTypes.UnsafeObject): Promise; resume(scope: CodegenTypes.UnsafeObject): Promise; // Live entry: journal 'cancelled' (user), forget after its ack. Settled - // entry: forget now: row, bytes, and its unacknowledged outcomes. Unknown - // id: resolve, no-op. If the journal or the store cannot be written, - // cancel rejects with E_STORAGE and changes nothing; the caller may call - // again. + // entry: forget now: row, bytes, and its unacknowledged outcomes. No row: + // delete the files left under the id, then resolve. If the journal or the + // store cannot be written, cancel rejects with E_STORAGE and changes + // nothing; the caller may call again. cancel(id: string): Promise; // The queue's Wi-Fi setting. Persisted natively. Applies to queued and // future entries whose descriptor has no wifiOnly. An entry with