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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 4 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -390,7 +390,8 @@ kept. A paused entry past `expiresAt` settles `error` with
### `cancel(id): Promise<void>`
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
Expand Down Expand Up @@ -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.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down Expand Up @@ -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.
Expand All @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<QueueEntry> =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -525,6 +551,46 @@ class QueueControllerTest {
fun `cancel of an unknown id is a no-op`() {
controller.cancel("nope")
assertEquals(emptyList<String>(), 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<String>(), 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<String>(), events.log)
}

// MARK: - pause, resume, wifi, headers
Expand Down
6 changes: 5 additions & 1 deletion ios/QueueCoordinator+Enqueue.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
8 changes: 6 additions & 2 deletions ios/QueueCoordinator.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
15 changes: 14 additions & 1 deletion ios/QueueStore.swift
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,8 @@ import Foundation
/// - `entry.json`: the QueueEntry (v10).
/// - `body-<uuid>` or `blob-<uuid>` / `blob`: the staged body.
/// - `part-<i>.<incarnation>.<start>-<end>`: 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
Expand Down Expand Up @@ -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
Expand Down
46 changes: 46 additions & 0 deletions ios/Tests/CoordinatorRelaunchTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
59 changes: 58 additions & 1 deletion ios/Tests/CoordinatorSimpleTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
8 changes: 4 additions & 4 deletions src/NativeRNFileUploader.ts
Original file line number Diff line number Diff line change
Expand Up @@ -68,10 +68,10 @@ export interface Spec extends TurboModule {
pause(scope: CodegenTypes.UnsafeObject): Promise<void>;
resume(scope: CodegenTypes.UnsafeObject): Promise<void>;
// 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<void>;
// The queue's Wi-Fi setting. Persisted natively. Applies to queued and
// future entries whose descriptor has no wifiOnly. An entry with
Expand Down
Loading