diff --git a/CHANGELOG.md b/CHANGELOG.md index 30e40e87..4816febb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,8 +21,9 @@ Breaking: - **`cancelUpload` and `removeUpload` fold into `cancel(id)`.** A live entry settles `cancelled` and is forgotten after its ack; a settled entry is forgotten now, row and bytes. -- **Per-upload `wifiOnly` becomes `setWifiOnly(enabled)`** on the queue, - persisted natively. +- **`wifiOnly` moves from the upload options to the request descriptor.** + `setWifiOnly(enabled)` on the queue, persisted natively, is the value for + entries that omit it. - **`progress` carries `{ id, bytesSent, totalBytes }`** instead of a percentage. - **`configure()` must be called at boot, after every `define()`.** It starts @@ -78,11 +79,24 @@ Added: a named error and warns when the native write has not settled by then, so a native bug cannot hang a caller in silence), and the v9 `android` notification options. -- **Queue control**: `pause()` / `resume()` for the whole queue, `cancel(id)`, +- **Queue control**: `pause(scope?)` / `resume(scope?)`, `cancel(id)`, `setWifiOnly(enabled)`, and `updateHeaders(patch)` to re-auth entries parked on 401 or 403. Each `updateHeaders()` call bumps a header generation, so a 401 from an attempt sent under older headers re-issues at once instead of parking. +- **Key-scoped pause.** `pause()` and `resume()` take an optional + `{ keys?: string[] }`. No scope is the whole queue; `{ keys }` pauses or + resumes the entries of those definition keys, queued and future. An entry + is paused while the whole queue or its key is paused, so resuming one + scope does not resume an entry the other still pauses. Both states persist + natively. A misspelled scope field, `keys: undefined`, or an empty key + rejects, so a mistake cannot pause the whole queue. The native `pause` and + `resume` now take a scope object, `{}` for the whole queue, because + codegen has no optional arguments. +- **Per-request `wifiOnly`.** `RequestDescriptor.wifiOnly?: boolean` + overrides `setWifiOnly()` for that entry and is persisted with it. An + entry that omits it follows `setWifiOnly()`, including later toggles. + Both are checked before each attempt. - **`getRequests(filter?)`**: synchronous, from native's in-memory index. Returns every entry native has not yet forgotten, so completed and cancelled rows appear until their ack and `error` rows until `cancel()` or diff --git a/README.md b/README.md index e76194b6..fc72fb7d 100644 --- a/README.md +++ b/README.md @@ -128,6 +128,7 @@ TypeScript does not flag a misspelled key on an inferred arrow return. | `accept` | Non-2xx responses to treat as success: `[{ status, bodyIncludes? }]`. | | `expiresAt` | Epoch ms. Default now + `lifetimeMs` (14 days). Past it: `error` with `errorKind: 'expired'`. | | `retry` | Per-request override of the `configure()` retry defaults. | +| `wifiOnly` | `true` waits for Wi-Fi before each attempt; `false` never waits. Overrides `setWifiOnly()` for this entry. Omit it to follow `setWifiOnly()`, including later toggles. | | `android` | `{ noNotification?: boolean }`. See Silent uploads. | `vars` and `data` cross to native as JSON strings, and native parses them. @@ -256,10 +257,10 @@ while the app was dead from being sent twice. live. A handler that keeps throwing sees the number grow. The library never gives up on its own; the app decides a poison policy from that number. -3. **One outcome per settle cycle.** `pause()` produces none. A paused entry - past `expiresAt` settles `error` with `errorKind: 'expired'` at - `resume()`. A same-id `mutate()` on a settled entry reopens it, and it - settles once more. +3. **One outcome per settle cycle.** `pause()` produces none, for the whole + queue or for a set of keys. A paused entry past `expiresAt` settles + `error` with `errorKind: 'expired'` when it is resumed. A same-id + `mutate()` on a settled entry reopens it, and it settles once more. 4. **Never before `mutate()` resolves.** Delivery for an id waits for the caller's promise. 5. **Replay starts after `configure()`.** Outcomes journaled by a dead session @@ -364,10 +365,27 @@ outcomes. A second call updates the settings and does not replay again. | `enqueueTimeoutMs` | Default 10 s. `mutate()` rejects and warns when native enqueue has not settled by then. Enqueue includes the time to stage a copy of a `file` body and of form `path` parts, so a large file takes longer. A timeout means native did not answer, not that the request failed. A watchdog for a native bug, not a tuning knob. | | `android` | Notification text and identity: `notificationId/Title/TitleNoWifi/TitleNoInternet/Channel`. Persisted natively. | -### `pause(): Promise` and `resume(): Promise` -Whole-queue pause. No outcome is produced; live rows show `paused`. A paused -entry past `expiresAt` settles `error` with `errorKind: 'expired'` at -`resume()`. +### `pause(scope?): Promise` and `resume(scope?): Promise` +`scope` is `{ keys?: string[] }`. No scope pauses or resumes the whole queue. +`{ keys }` pauses or resumes the entries of those definition keys, queued and +future. An empty `keys` list changes nothing. A scope field other than `keys`, +`keys: undefined`, or a key that is not a non-empty string rejects, so a +mistake cannot pause the whole queue. + +An entry is paused while the whole queue is paused or its key is paused. A +resume of one scope does not resume an entry that the other still pauses: + +```ts +await Upload.pause({ keys: ['capture.upload'] }); // captures wait, notes run +await Upload.pause(); // everything waits +await Upload.resume(); // notes run, captures still wait +await Upload.resume({ keys: ['capture.upload'] }); // captures run +``` + +Both states are persisted natively. Each entry that moves emits one `state` +event (`paused`, then `queued`). No outcome is produced and the bytes are +kept. A paused entry past `expiresAt` settles `error` with +`errorKind: 'expired'` when it is resumed. ### `cancel(id): Promise` A live entry settles `cancelled` with reason `user` and is forgotten after @@ -380,7 +398,10 @@ whose entry save failed is already in effect: the work stops and the finishes it. ### `setWifiOnly(enabled): Promise` -Persisted natively. Applies to queued and future entries. +The queue's Wi-Fi setting. Persisted natively. Applies to queued and future +entries whose descriptor does not set `wifiOnly`. An entry that sets it keeps +its own value. Both are checked before each attempt, so a toggle moves the +queued entries that follow the setting. ### `updateHeaders(patch): Promise` Merges the patch into the headers of every entry not yet forgotten and @@ -439,7 +460,8 @@ The CHANGELOG lists every removed v9 export with its replacement. In short: `startUpload` becomes a `define()` plus `mutate()`; `getAllUploads` becomes `getRequests()`; `cancelUpload` and `removeUpload` become `cancel(id)`; the terminal event names become the definition's handlers; per-upload `wifiOnly` -becomes `setWifiOnly()`. +stays on the descriptor, and `setWifiOnly()` sets the value for entries that +omit it. What happens to work a v9 build left behind: diff --git a/android/src/main/java/ai/openspace/backgroundupload/EnqueueRules.kt b/android/src/main/java/ai/openspace/backgroundupload/EnqueueRules.kt index 4331797e..5fd1ac15 100644 --- a/android/src/main/java/ai/openspace/backgroundupload/EnqueueRules.kt +++ b/android/src/main/java/ai/openspace/backgroundupload/EnqueueRules.kt @@ -121,11 +121,11 @@ object EnqueueRules { ) /** - * Same body. New headers, expiresAt, vars, accept, retry, and notification - * flag replace the stored ones; the body and accepted parts stay. A settled - * entry reopens with a fresh generation and attempts 0 (attempts count the - * current generation). A running one stays running (the worker reads the - * new headers before its next attempt). + * Same body. New headers, expiresAt, vars, accept, retry, wifiOnly, and + * notification flag replace the stored ones; the body and accepted parts + * stay. A settled entry reopens with a fresh generation and attempts 0 + * (attempts count the current generation). A running one stays running + * (the worker reads the new headers before its next attempt). */ fun resumed( existing: QueueEntry, @@ -148,6 +148,7 @@ object EnqueueRules { accept = p.descriptor.accept, retry = p.descriptor.retry, noNotification = p.descriptor.noNotification, + wifiOnly = p.descriptor.wifiOnly, ), state = if (running) EntryState.RUNNING else initialState(paused), attempts = if (reopen) 0 else existing.attempts, diff --git a/android/src/main/java/ai/openspace/backgroundupload/EntryParsing.kt b/android/src/main/java/ai/openspace/backgroundupload/EntryParsing.kt index f3c77ac6..5992de45 100644 --- a/android/src/main/java/ai/openspace/backgroundupload/EntryParsing.kt +++ b/android/src/main/java/ai/openspace/backgroundupload/EntryParsing.kt @@ -65,6 +65,7 @@ object EntryParsing { val headers = parseHeaderMap(d.map("headers")) requireValidHeaders(headers, "headers") + if (d.isSet("wifiOnly") && d.bool("wifiOnly") == null) invalid("wifiOnly must be a boolean") return Descriptor( url = url, @@ -77,9 +78,25 @@ object EntryParsing { accept = parseAcceptRules(d.array("accept")), retry = d.map("retry")?.let { parseRetry(it) }, noNotification = d.map("android")?.bool("noNotification") ?: false, + wifiOnly = d.bool("wifiOnly"), ) } + /** + * pause/resume scope `{ keys?: string[] }`. A null scope or a scope with + * no keys is the whole queue (null comes from a caller that sent no + * argument). A keys field that is not a list of non-empty strings is + * refused: dropping it would widen the scope to the whole queue. + */ + fun scopeKeys(scope: ReadableMap?): List? { + if (scope == null || !scope.hasKey("keys")) return null + val arr = scope.array("keys") ?: invalid("scope.keys must be an array of strings") + return (0 until arr.size()).map { i -> + if (arr.getType(i) != ReadableType.String) invalid("scope.keys[$i] must be a string") + arr.getString(i)?.takeIf { it.isNotEmpty() } ?: invalid("scope.keys[$i] must be a non-empty string") + } + } + /** updateHeaders(patch): the header map, checked the same way as a descriptor's. */ fun headerPatch(patch: ReadableMap): Map = parseHeaderMap(patch).also { requireValidHeaders(it, "updateHeaders") } diff --git a/android/src/main/java/ai/openspace/backgroundupload/EntryRun.kt b/android/src/main/java/ai/openspace/backgroundupload/EntryRun.kt index d9b5ffbe..aa130380 100644 --- a/android/src/main/java/ai/openspace/backgroundupload/EntryRun.kt +++ b/android/src/main/java/ai/openspace/backgroundupload/EntryRun.kt @@ -32,8 +32,8 @@ internal class EntryRun( class ExpiredException : Exception("expired before completion") - /** The queue was paused between the module's pause and the work cancel reaching us. */ - class PausedException : Exception("queue paused") + /** The entry was paused (whole queue or its key) between the module's pause and the work cancel reaching us. */ + class PausedException : Exception("paused") /** How one attempt ended, after the retry table. */ sealed class AttemptResult { @@ -96,7 +96,7 @@ internal class EntryRun( return } if (initial.state != EntryState.QUEUED && initial.state != EntryState.RUNNING) return - if (ops.settings().paused) return + if (ops.settings().isPaused(initial.key)) return if (!waitUntilDue(initial)) return val entry = ops.begin(entryId) ?: return generation = entry.generation @@ -220,16 +220,17 @@ internal class EntryRun( // MARK: - helpers for the transfers /** - * Waits until the network fits the queue's wifi-only setting. Re-reads - * the settings and the entry at every poll. + * Waits until the network fits the entry's Wi-Fi rule: its own wifiOnly, + * else the queue's setting. Re-reads the settings and the entry at every + * poll, so a setWifiOnly() toggle reaches an entry that follows it. */ suspend fun waitForNetwork() { while (true) { val s = ops.settings() - if (s.paused) throw PausedException() val entry = ops.latest(entryId, generation) + if (s.isPaused(entry.key)) throw PausedException() if (RetryClassifier.isExpired(host.now(), entry.expiresAt)) throw ExpiredException() - if (host.connectivity(s.wifiOnly) == Connectivity.Ok) return + if (host.connectivity(s.wifiOnlyFor(entry)) == Connectivity.Ok) return host.sleep(CONNECTIVITY_POLL_MS) } } diff --git a/android/src/main/java/ai/openspace/backgroundupload/QueueController.kt b/android/src/main/java/ai/openspace/backgroundupload/QueueController.kt index c72cb0fe..98143e6f 100644 --- a/android/src/main/java/ai/openspace/backgroundupload/QueueController.kt +++ b/android/src/main/java/ai/openspace/backgroundupload/QueueController.kt @@ -96,13 +96,13 @@ class QueueController( // With no listener yet, the next drain delivers it. is EnqueueRules.Action.ReEmit -> Enqueued(null, journal.redeliver(action.eventId)) EnqueueRules.Action.Resume -> { - val next = EnqueueRules.resumed(existing!!, p, s.paused, s.headerGeneration, now) + val next = EnqueueRules.resumed(existing!!, p, s.isPaused(p.key), s.headerGeneration, now) saveOrThrow(next) Enqueued(next, null) } EnqueueRules.Action.Create -> { val staged = preStaged(1) ?: stageOrThrow(p.descriptor, dir, 1) - commit(EnqueueRules.created(p, staged, null, s.paused, s.headerGeneration, now)) + commit(EnqueueRules.created(p, staged, null, s.isPaused(p.key), s.headerGeneration, now)) } is EnqueueRules.Action.AdoptV9 -> { val incoming = p.descriptor.parts @@ -110,7 +110,7 @@ class QueueController( // The same parts resume over the v9 blob, as a same-body enqueue does. val keepOwned = incoming != null && ChunkedParts.sameParts(action.manifest.parts, incoming) val staged = stageOrThrow(p.descriptor.copy(parts = parts), dir, action.generation, store.blobFile(p.id), keepOwned) - commit(EnqueueRules.adopted(p, staged, parts, existing, action.generation, s.paused, s.headerGeneration, now)) + commit(EnqueueRules.adopted(p, staged, parts, existing, action.generation, s.isPaused(p.key), s.headerGeneration, now)) } EnqueueRules.Action.Replace -> { val old = existing!! @@ -118,7 +118,7 @@ class QueueController( val ownedBlob = if (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.paused, s.headerGeneration, now)) + commit(EnqueueRules.replaced(old, p, staged, s.isPaused(p.key), s.headerGeneration, now)) } } } @@ -164,22 +164,38 @@ class QueueController( // MARK: - queue control - /** Whole-queue pause. Live rows move to paused and their work stops. No outcome. */ - fun pause() { - settings.update { it.copy(paused = true) } + /** + * [keys] null is the whole-queue gate; a list adds those keys to the + * paused set, and an empty list changes nothing. Each live row that is now + * paused (gate on, or its key in the set) moves to paused and its work + * stops. No outcome. + */ + fun pause(keys: List? = null) { + if (keys != null && keys.isEmpty()) return + val s = settings.update { + if (keys == null) it.copy(paused = true) else it.copy(pausedKeys = it.pausedKeys + keys) + } val now = clock() val paused = transformAll { e -> - if (e.isLive && e.state != EntryState.PAUSED) EntryTransitions.toPaused(e, now) else e + if (e.isLive && e.state != EntryState.PAUSED && s.isPaused(e.key)) EntryTransitions.toPaused(e, now) else e } paused.forEach { scheduler.cancel(it.id) } paused.forEach { events.state(it.toRow()) } } - fun resume() { - val s = settings.update { it.copy(paused = false) } + /** + * Undoes the pause of the same scope: [keys] null turns the gate off, a + * list removes those keys from the set. A row comes back only when no + * scope still pauses it, so the gate on + its key resumed stays paused. + */ + fun resume(keys: List? = null) { + if (keys != null && keys.isEmpty()) return + val s = settings.update { + if (keys == null) it.copy(paused = false) else it.copy(pausedKeys = it.pausedKeys - keys.toSet()) + } val now = clock() val resumed = transformAll { e -> - if (e.state == EntryState.PAUSED) EntryTransitions.toResumed(e, s.headerGeneration, now) else e + if (e.state == EntryState.PAUSED && !s.isPaused(e.key)) EntryTransitions.toResumed(e, s.headerGeneration, now) else e } resumed.forEach { events.state(it.toRow()) } // Every queued entry, not only the resumed ones: a run is idempotent. @@ -290,7 +306,7 @@ class QueueController( val patched = e.withHeadersPatched(patch, s.headerGeneration) if (patched.state != EntryState.AWAITING_AUTH) return@transformAll patched unparkedIds += e.id - EntryTransitions.toUnparked(patched, s.paused, now) + EntryTransitions.toUnparked(patched, s.isPaused(patched.key), now) } changed.filter { it.id in unparkedIds }.forEach { entry -> scheduleRun(entry) @@ -353,8 +369,8 @@ class QueueController( * 1. A live entry with a journal record of its own generation: the process * died between the journal append and the store transition. Apply it. * 2. A running entry with no worker: a process death mid-run. Queue it. - * A live row that disagrees with the queue's paused setting (a death - * partway through pause() or resume()): make it agree. + * A live row that disagrees with the paused settings, the gate or its + * key (a death partway through pause() or resume()): make it agree. * 3. A settled entry with records of its generation other than its own: * orphans from a cancel race. Ack them. * 4. A completed or cancelled entry whose own record is gone: the ack @@ -366,7 +382,7 @@ class QueueController( val changed = mutableListOf() val toForget = mutableListOf() val toSchedule = mutableListOf() - val paused = settings.load().paused + val s = settings.load() store.locked { val records = journal.unacknowledged().groupBy { it.id } for (e in store.all()) { @@ -386,11 +402,12 @@ class QueueController( cur = EntryTransitions.toStopped(e, now) } // A process death partway through pause() or resume() leaves rows - // that disagree with the queue setting. + // that disagree with the settings (the gate or the paused keys). + val paused = s.isPaused(e.key) if (paused && cur.state != EntryState.PAUSED) { cur = EntryTransitions.toPaused(cur, now) } else if (!paused && cur.state == EntryState.PAUSED) { - cur = EntryTransitions.toResumed(cur, settings.load().headerGeneration, now) + cur = EntryTransitions.toResumed(cur, s.headerGeneration, now) } if (cur !== e) { if (trySave(cur)) changed += cur else continue diff --git a/android/src/main/java/ai/openspace/backgroundupload/QueueEntry.kt b/android/src/main/java/ai/openspace/backgroundupload/QueueEntry.kt index 7a5e7f2f..09f1b364 100644 --- a/android/src/main/java/ai/openspace/backgroundupload/QueueEntry.kt +++ b/android/src/main/java/ai/openspace/backgroundupload/QueueEntry.kt @@ -49,6 +49,8 @@ data class Descriptor( val accept: List, val retry: RetryOverride?, val noNotification: Boolean, + /** The request's own Wi-Fi rule. Null follows the queue's setWifiOnly(). */ + val wifiOnly: Boolean? = null, ) { val bodyKind: String get() = when { diff --git a/android/src/main/java/ai/openspace/backgroundupload/QueueSettings.kt b/android/src/main/java/ai/openspace/backgroundupload/QueueSettings.kt index 1a364526..1adda7fe 100644 --- a/android/src/main/java/ai/openspace/backgroundupload/QueueSettings.kt +++ b/android/src/main/java/ai/openspace/backgroundupload/QueueSettings.kt @@ -14,12 +14,22 @@ data class RetryDefaults( /** Queue-wide settings. They live next to the entries, in `settings.json`. */ data class QueueSettings( + /** The queue's Wi-Fi setting. An entry whose descriptor sets wifiOnly ignores it. */ val wifiOnly: Boolean = false, + /** The whole-queue pause gate. */ val paused: Boolean = false, + /** The keys paused by pause({ keys }). An entry is paused under the gate or its key. */ + val pausedKeys: Set = emptySet(), /** +1 per updateHeaders(). A 401 from an attempt sent under an older value re-issues at once. */ val headerGeneration: Int = 0, val retry: RetryDefaults = RetryDefaults(), -) +) { + /** Whether an entry of [key] is paused: the gate is on, or its key is in the set. */ + fun isPaused(key: String) = paused || key in pausedKeys + + /** The Wi-Fi rule of one attempt: the entry's own value, else the queue's. */ + fun wifiOnlyFor(entry: QueueEntry) = entry.descriptor?.wifiOnly ?: wifiOnly +} /** * Reads and writes [QueueSettings]. The value is cached after the first read. @@ -77,7 +87,7 @@ class QueueSettingsStore(private val file: File) { } // Gson does not run constructors, so absent fields read as null or 0. - @Suppress("SENSELESS_COMPARISON", "USELESS_ELVIS") + @Suppress("SENSELESS_COMPARISON", "USELESS_ELVIS", "UNNECESSARY_SAFE_CALL") private fun validated(s: QueueSettings?): QueueSettings { if (s == null) return QueueSettings() val d = RetryDefaults() @@ -85,6 +95,7 @@ class QueueSettingsStore(private val file: File) { return QueueSettings( wifiOnly = s.wifiOnly, paused = s.paused, + pausedKeys = s.pausedKeys?.toSet() ?: emptySet(), headerGeneration = s.headerGeneration, retry = if (r == null) d else RetryDefaults( baseMs = if (r.baseMs > 0) r.baseMs else d.baseMs, diff --git a/android/src/main/java/ai/openspace/backgroundupload/UploaderModule.kt b/android/src/main/java/ai/openspace/backgroundupload/UploaderModule.kt index 83ad4505..1b6e61a5 100644 --- a/android/src/main/java/ai/openspace/backgroundupload/UploaderModule.kt +++ b/android/src/main/java/ai/openspace/backgroundupload/UploaderModule.kt @@ -164,9 +164,28 @@ class UploaderModule(context: ReactApplicationContext) : onQueue(promise) { controller.enqueue(parsed) } } - override fun pause(promise: Promise) = onQueue(promise) { controller.pause(); null } + // scope is { keys?: string[] }; {} is the whole queue. Null (a stale JS + // bundle that sends no argument) is also the whole queue, not a crash. + // Parse on the calling thread, as enqueue does. + override fun pause(scope: ReadableMap?, promise: Promise) { + val keys = try { + EntryParsing.scopeKeys(scope) + } catch (e: EntryParsing.InvalidEntryException) { + promise.reject(QueueException.E_INVALID, e.message, e) + return + } + onQueue(promise) { controller.pause(keys); null } + } - override fun resume(promise: Promise) = onQueue(promise) { controller.resume(); null } + override fun resume(scope: ReadableMap?, promise: Promise) { + val keys = try { + EntryParsing.scopeKeys(scope) + } catch (e: EntryParsing.InvalidEntryException) { + promise.reject(QueueException.E_INVALID, e.message, e) + return + } + onQueue(promise) { controller.resume(keys); null } + } override fun cancel(id: String, promise: Promise) = onQueue(promise) { controller.cancel(id); null } diff --git a/android/src/test/java/ai/openspace/backgroundupload/EntryParsingTest.kt b/android/src/test/java/ai/openspace/backgroundupload/EntryParsingTest.kt index 5ed1c534..900a94db 100644 --- a/android/src/test/java/ai/openspace/backgroundupload/EntryParsingTest.kt +++ b/android/src/test/java/ai/openspace/backgroundupload/EntryParsingTest.kt @@ -116,6 +116,35 @@ class EntryParsingTest { assertThrows(EntryParsing.InvalidEntryException::class.java) { EntryParsing.parse(entryMap(both)) } } + @Test + fun `wifiOnly parses true, false, or absent`() { + assertEquals(true, EntryParsing.parse(entryMap(base("wifiOnly", true))).descriptor.wifiOnly) + assertEquals(false, EntryParsing.parse(entryMap(base("wifiOnly", false))).descriptor.wifiOnly) + assertNull(EntryParsing.parse(entryMap(base())).descriptor.wifiOnly) + assertThrows(EntryParsing.InvalidEntryException::class.java) { + EntryParsing.parse(entryMap(base("wifiOnly", "yes"))) + } + } + + @Test + fun `a pause scope with no keys is the whole queue, and keys must be non-empty strings`() { + assertNull(EntryParsing.scopeKeys(JavaOnlyMap())) + // Codegen passes null when a stale bundle calls pause() with no argument. + assertNull(EntryParsing.scopeKeys(null)) + assertEquals(listOf("capture", "video"), EntryParsing.scopeKeys(JavaOnlyMap.of("keys", JavaOnlyArray.of("capture", "video")))) + assertEquals(emptyList(), EntryParsing.scopeKeys(JavaOnlyMap.of("keys", JavaOnlyArray()))) + // A keys field that native drops would widen the scope to the whole queue. + listOf( + JavaOnlyMap.of("keys", null), + JavaOnlyMap.of("keys", "capture"), + JavaOnlyMap.of("keys", JavaOnlyArray.of("capture", 1.0)), + // JS and the README refuse empty keys; native agrees. + JavaOnlyMap.of("keys", JavaOnlyArray.of("capture", "")), + ).forEach { m -> + assertThrows("$m", EntryParsing.InvalidEntryException::class.java) { EntryParsing.scopeKeys(m) } + } + } + @Test fun `what native can not run is rejected`() { val cases = listOf( diff --git a/android/src/test/java/ai/openspace/backgroundupload/EntryRunTest.kt b/android/src/test/java/ai/openspace/backgroundupload/EntryRunTest.kt index 1e7e0511..4471e889 100644 --- a/android/src/test/java/ai/openspace/backgroundupload/EntryRunTest.kt +++ b/android/src/test/java/ai/openspace/backgroundupload/EntryRunTest.kt @@ -25,6 +25,8 @@ internal class FakeHost(var clock: Long = 10_000L) : TransferHost { var network = ArrayDeque() var timeout = false var foregroundError: Throwable? = null + /** The wifiOnly value of every connectivity check, in order. */ + val wifiChecks = mutableListOf() override fun now() = clock @@ -39,7 +41,10 @@ internal class FakeHost(var clock: Long = 10_000L) : TransferHost { return handler(request, onProgress) } - override fun connectivity(wifiOnly: Boolean) = network.removeFirstOrNull() ?: Connectivity.Ok + override fun connectivity(wifiOnly: Boolean): Connectivity { + wifiChecks += wifiOnly + return network.removeFirstOrNull() ?: Connectivity.Ok + } override suspend fun foreground(entry: QueueEntry) { foregroundError?.let { throw it } @@ -263,6 +268,48 @@ class EntryRunTest { assertEquals(EntryState.ERROR, store.load("e1")!!.state) } + @Test + fun `the connectivity check uses the entry's wifiOnly, else the queue setting at each poll`() { + controller.setWifiOnly(true) + controller.enqueue(parsed(id = "cell", descriptor = desc(dataJson = """{"a":1}""", wifiOnly = false))) + respond(200) + run("cell") + assertEquals(listOf(false), host.wifiChecks) // its own false beats the queue's true + + host.wifiChecks.clear() + controller.enqueue(parsed(id = "follow")) + host.network = ArrayDeque(listOf(Connectivity.NoWifi)) + host.onSleep = { controller.setWifiOnly(false) } // a toggle while it waits + respond(200) + run("follow") + assertEquals(listOf(true, false), host.wifiChecks) + assertEquals(EntryState.COMPLETED, store.load("follow")!!.state) + } + + @Test + fun `a run under a key pause does not start, and other keys run`() { + controller.enqueue(parsed(id = "c", key = "capture")) + controller.enqueue(parsed(id = "n", key = "note")) + controller.pause(listOf("capture")) + respond(200) + run("c") + run("n") + assertEquals(EntryState.PAUSED, store.load("c")!!.state) + assertEquals(EntryState.COMPLETED, store.load("n")!!.state) + assertEquals(1, host.requests.size) + } + + @Test + fun `a key pause while waiting for the network stops the run with no outcome`() { + controller.enqueue(parsed(id = "c", key = "capture")) + host.network = ArrayDeque(listOf(Connectivity.NoInternet)) + host.onSleep = { controller.pause(listOf("capture")) } + run("c") + assertEquals(EntryState.PAUSED, store.load("c")!!.state) + assertEquals(emptyList(), host.requests) + assertEquals(emptyList(), journal.unacknowledged()) + } + @Test fun `a failure that lands after pause is not an outcome`() { controller.enqueue(parsed()) diff --git a/android/src/test/java/ai/openspace/backgroundupload/EntryTransitionsTest.kt b/android/src/test/java/ai/openspace/backgroundupload/EntryTransitionsTest.kt index a40a156f..e6121262 100644 --- a/android/src/test/java/ai/openspace/backgroundupload/EntryTransitionsTest.kt +++ b/android/src/test/java/ai/openspace/backgroundupload/EntryTransitionsTest.kt @@ -153,6 +153,14 @@ class EnqueueRulesTest { assertEquals(1, next.generation) // a live entry keeps its life } + @Test + fun `resume takes the incoming wifiOnly, including its removal`() { + val stored = entry(descriptor = desc(dataJson = """{"a":1}""", wifiOnly = true)) + val off = parsed(descriptor = desc(dataJson = """{"a":1}""", wifiOnly = false)) + assertEquals(false, EnqueueRules.resumed(stored, off, false, 0, 9).descriptor!!.wifiOnly) + assertNull(EnqueueRules.resumed(stored, parsed(), false, 0, 9).descriptor!!.wifiOnly) + } + @Test fun `resume of a settled entry reopens it with a fresh generation and attempts 0`() { val settled = entry(state = EntryState.ERROR, settledEventId = "ev", generation = 2, attempts = 3) diff --git a/android/src/test/java/ai/openspace/backgroundupload/QueueControllerTest.kt b/android/src/test/java/ai/openspace/backgroundupload/QueueControllerTest.kt index 6cfc408f..88d631d6 100644 --- a/android/src/test/java/ai/openspace/backgroundupload/QueueControllerTest.kt +++ b/android/src/test/java/ai/openspace/backgroundupload/QueueControllerTest.kt @@ -554,6 +554,99 @@ class QueueControllerTest { assertEquals(listOf("a" to FAR_FUTURE), scheduler.wakes) // the parked one waits for its expiry } + @Test + fun `a key pause moves only that key's live rows, one state event each, and no outcome`() { + store.save(entry(id = "c1", key = "capture", state = EntryState.QUEUED)) + store.save(entry(id = "c2", key = "capture", state = EntryState.RUNNING)) + store.save(entry(id = "c3", key = "capture", state = EntryState.ERROR)) + store.save(entry(id = "n1", key = "note", state = EntryState.QUEUED)) + controller.pause(listOf("capture")) + assertFalse(settings.load().paused) + assertEquals(setOf("capture"), settings.load().pausedKeys) + assertEquals(listOf("paused", "paused", "error", "queued"), listOf("c1", "c2", "c3", "n1").map { store.load(it)!!.state.wire }) + assertEquals(setOf("c1", "c2"), scheduler.cancelled.toSet()) + assertEquals(listOf("state:c1:paused", "state:c2:paused"), events.log.sorted()) + assertEquals(emptyList(), journal.unacknowledged()) + + events.log.clear() + controller.resume(listOf("capture")) + assertEquals(emptySet(), settings.load().pausedKeys) + assertEquals(listOf("queued", "queued", "error", "queued"), listOf("c1", "c2", "c3", "n1").map { store.load(it)!!.state.wire }) + assertEquals(listOf("state:c1:queued", "state:c2:queued"), events.log.sorted()) + assertTrue(scheduler.scheduled.containsAll(listOf("c1", "c2"))) + } + + @Test + fun `a key resume under the global gate stays paused, and a global resume under a key pause too`() { + store.save(entry(id = "c", key = "capture")) + store.save(entry(id = "n", key = "note")) + controller.pause(listOf("capture")) + controller.pause() + events.log.clear() + controller.resume(listOf("capture")) + assertEquals(EntryState.PAUSED, store.load("c")!!.state) // the gate still pauses it + assertEquals(emptyList(), events.log) + + controller.pause(listOf("capture")) + controller.resume() + assertEquals(EntryState.PAUSED, store.load("c")!!.state) // its key still pauses it + assertEquals(EntryState.QUEUED, store.load("n")!!.state) + assertEquals(listOf("state:n:queued"), events.log) + assertFalse("c" in scheduler.scheduled) + } + + @Test + fun `a key pause of an already paused entry emits nothing`() { + store.save(entry(id = "c", key = "capture")) + controller.pause() + events.log.clear() + controller.pause(listOf("capture")) + assertEquals(emptyList(), events.log) + assertEquals(setOf("capture"), settings.load().pausedKeys) + } + + @Test + fun `an empty key list changes nothing`() { + store.save(entry(id = "c", key = "capture")) + controller.pause(emptyList()) + assertEquals(QueueSettings(), settings.load()) + assertEquals(EntryState.QUEUED, store.load("c")!!.state) + controller.pause(listOf("capture")) + controller.resume(emptyList()) + assertEquals(setOf("capture"), settings.load().pausedKeys) + assertEquals(EntryState.PAUSED, store.load("c")!!.state) + } + + @Test + fun `enqueue under a key pause starts paused, and other keys run`() { + controller.pause(listOf("capture")) + controller.enqueue(parsed(id = "c", key = "capture")) + controller.enqueue(parsed(id = "n", key = "note")) + assertEquals(EntryState.PAUSED, store.load("c")!!.state) + assertEquals(EntryState.QUEUED, store.load("n")!!.state) + assertEquals(listOf("n"), scheduler.scheduled) + } + + @Test + fun `updateHeaders unparks an entry under a key pause to paused`() { + store.save(entry(id = "c", key = "capture", state = EntryState.AWAITING_AUTH, parkedGeneration = 0)) + settings.update { it.copy(pausedKeys = setOf("capture")) } + controller.updateHeaders(mapOf("Authorization" to "Bearer new")) + assertEquals(EntryState.PAUSED, store.load("c")!!.state) + assertFalse("c" in scheduler.scheduled) + } + + @Test + fun `the per-request wifiOnly is persisted with the entry`() { + controller.enqueue(parsed(id = "w", descriptor = desc(dataJson = """{"a":1}""", wifiOnly = true))) + controller.enqueue(parsed(id = "c", descriptor = desc(dataJson = """{"a":1}""", wifiOnly = false))) + controller.enqueue(parsed(id = "f")) + val reread = QueueStore(root, RequestIndex()) + assertEquals(true, reread.load("w")!!.descriptor!!.wifiOnly) + assertEquals(false, reread.load("c")!!.descriptor!!.wifiOnly) + assertNull(reread.load("f")!!.descriptor!!.wifiOnly) + } + @Test fun `setWifiOnly persists`() { controller.setWifiOnly(true) @@ -774,4 +867,19 @@ class QueueControllerTest { assertEquals(EntryState.PAUSED, store.load("r")!!.state) assertEquals(listOf("p"), scheduler.scheduled) } + + @Test + fun `sweep finishes a key pause or resume that a process death cut short`() { + // pause({ keys }) saved the set, then died before the rows. + settings.update { it.copy(pausedKeys = setOf("capture")) } + store.save(entry(id = "c", key = "capture", state = EntryState.QUEUED)) + store.save(entry(id = "n", key = "note", state = EntryState.QUEUED)) + // resume({ keys }) of "video" saved the set, then died before the rows. + store.save(entry(id = "v", key = "video", state = EntryState.PAUSED)) + controller.sweep() + assertEquals(EntryState.PAUSED, store.load("c")!!.state) + assertEquals(EntryState.QUEUED, store.load("n")!!.state) + assertEquals(EntryState.QUEUED, store.load("v")!!.state) + assertEquals(setOf("n", "v"), scheduler.scheduled.toSet()) + } } diff --git a/android/src/test/java/ai/openspace/backgroundupload/QueueSettingsTest.kt b/android/src/test/java/ai/openspace/backgroundupload/QueueSettingsTest.kt index 6dbcfff7..d21668ee 100644 --- a/android/src/test/java/ai/openspace/backgroundupload/QueueSettingsTest.kt +++ b/android/src/test/java/ai/openspace/backgroundupload/QueueSettingsTest.kt @@ -28,6 +28,40 @@ class QueueSettingsTest { assertTrue(reloaded.paused) } + @Test + fun `paused keys persist and reload in a new instance`() { + val file = File(tmp.newFolder(), "settings.json") + QueueSettingsStore(file).update { it.copy(pausedKeys = setOf("capture", "video")) } + assertEquals(setOf("capture", "video"), QueueSettingsStore(file).load().pausedKeys) + } + + @Test + fun `a null paused keys field reads an empty set`() { + val file = File(tmp.newFolder(), "settings.json").apply { writeText("""{"paused":true,"pausedKeys":null}""") } + val s = QueueSettingsStore(file).load() + assertTrue(s.paused) + assertEquals(emptySet(), s.pausedKeys) + } + + @Test + fun `an entry is paused under the gate or its key`() { + assertFalse(QueueSettings().isPaused("note")) + assertTrue(QueueSettings(paused = true).isPaused("note")) + val keyed = QueueSettings(pausedKeys = setOf("capture")) + assertTrue(keyed.isPaused("capture")) + assertFalse(keyed.isPaused("note")) + } + + @Test + fun `the entry's own wifiOnly wins, and absent follows the queue setting`() { + val on = QueueSettings(wifiOnly = true) + val off = QueueSettings(wifiOnly = false) + assertTrue(on.wifiOnlyFor(entry(descriptor = desc()))) + assertFalse(off.wifiOnlyFor(entry(descriptor = desc()))) + assertFalse(on.wifiOnlyFor(entry(descriptor = desc(wifiOnly = false)))) + assertTrue(off.wifiOnlyFor(entry(descriptor = desc(wifiOnly = true)))) + } + @Test fun `a corrupt file reads as the defaults`() { val file = File(tmp.newFolder(), "settings.json").apply { writeText("{not json") } @@ -40,6 +74,7 @@ class QueueSettingsTest { val s = QueueSettingsStore(file).load() assertTrue(s.wifiOnly) assertFalse(s.paused) + assertEquals(emptySet(), s.pausedKeys) assertEquals(RetryDefaults(), s.retry) } diff --git a/android/src/test/java/ai/openspace/backgroundupload/TestSupport.kt b/android/src/test/java/ai/openspace/backgroundupload/TestSupport.kt index 1effda61..c90665d4 100644 --- a/android/src/test/java/ai/openspace/backgroundupload/TestSupport.kt +++ b/android/src/test/java/ai/openspace/backgroundupload/TestSupport.kt @@ -15,7 +15,8 @@ internal fun desc( accept: List = emptyList(), retry: RetryOverride? = null, noNotification: Boolean = false, -) = Descriptor(url, method, headers, dataJson, form, file, parts, accept, retry, noNotification) + wifiOnly: Boolean? = null, +) = Descriptor(url, method, headers, dataJson, form, file, parts, accept, retry, noNotification, wifiOnly) internal fun part(start: Long, end: Long, accepted: Boolean = false, url: String? = null) = Part( url = url ?: "https://example.com/part?start=$start", diff --git a/example/RNBGUExample/README.md b/example/RNBGUExample/README.md index d72c86fd..ac14efc8 100644 --- a/example/RNBGUExample/README.md +++ b/example/RNBGUExample/README.md @@ -218,6 +218,13 @@ back to `Bearer bad` before a step that needs a 401. `code=E_FILE_MISSING`. Press "GET with a body". Pass: rejected. The JS layer rejects it before native, so there is no code. +15. **Dispatch latency.** Mode `ok`, app in the foreground. Press "JSON POST" + and note the clock on the log line the press writes. Note `at=` on the + first attempt line for that id. The difference is the time from `mutate()` + to the first send. Repeat five times and report the median. This number + decides whether a follow-up adds an in-process fast path for small + requests: a median under about 500 ms means no. + ### iOS Set Settings > Developer > Network Link Conditioner > "3G" on the phone, so an @@ -295,6 +302,9 @@ step says otherwise. outcomes are `ok`. If one retries without end with a `network` error, write down the error message: a bodiless GET needs a download task. +18. **Dispatch latency.** Same as Android step 15. iOS creates a background + session task per attempt, so this is the number to watch. + ## Tests `yarn test` in this folder runs a Jest smoke test of the screen against a stub diff --git a/ios/ChunkedCoordinator.swift b/ios/ChunkedCoordinator.swift index 1aca7410..92b98c3a 100644 --- a/ios/ChunkedCoordinator.swift +++ b/ios/ChunkedCoordinator.swift @@ -39,7 +39,7 @@ final class ChunkedCoordinator { /// queued -> running, then fill the window. func start(_ id: String) { guard var e = q.index.entry(id), e.isChunked, e.state == .queued || e.state == .running, - !q.settings.paused else { return } + !q.settings.isPaused(e.key) else { return } if e.state == .queued { e.state = .running e.nextAttemptAt = nil @@ -221,7 +221,7 @@ final class ChunkedCoordinator { let ownedOrUnknown = !q.ready || inFlight[id]?[part] == key guard let e = q.index.entry(id), e.isChunked, incarnation == e.incarnation, e.parts.indices.contains(part), !e.parts[part].accepted, e.state == .running, - !q.settings.paused, ownedOrUnknown, let url = URL(string: e.parts[part].url) else { + !q.settings.isPaused(e.key), ownedOrUnknown, let url = URL(string: e.parts[part].url) else { q.taskMap.setPurpose(.superseded, forKey: key, id: id) q.liveTasks[key] = nil if inFlight[id]?[part] == key { @@ -253,8 +253,8 @@ final class ChunkedCoordinator { /// background-wake refill that keeps the upload moving while the app is /// dead), at start, and at the end of every reconcile. func refill(_ id: String) { - guard q.ready, !q.settings.paused, let e = q.index.entry(id), e.isChunked, - e.state == .running else { return } + guard q.ready, let e = q.index.entry(id), e.isChunked, e.state == .running, + !q.settings.isPaused(e.key) else { return } if e.allAccepted { q.settle(id, .completed(RawResponseRecord(bodyTruncated: false))) return @@ -320,7 +320,7 @@ final class ChunkedCoordinator { generation: e.generation, purpose: .attempt) let task = q.transport.upload( q.buildRequest(e, url: url, requestId: requestId, partHeaders: part.headers), - fromFile: file, wifiOnly: q.settings.wifiOnly, + fromFile: file, wifiOnly: q.settings.wifiOnly(e.wifiOnly), description: ChunkedEngine.taskDescription(id: id, part: index, incarnation: e.incarnation), beginAt: delayMs.map { Date(timeIntervalSince1970: (q.now() + Double($0)) / 1000) }, beforeResume: { key in self.q.taskMap.set(meta, forKey: key) }) diff --git a/ios/EnqueueParser.swift b/ios/EnqueueParser.swift index 6e5474d5..0b206d52 100644 --- a/ios/EnqueueParser.swift +++ b/ios/EnqueueParser.swift @@ -30,6 +30,8 @@ struct ParsedEnqueue { let accept: [UploadOutcome.AcceptRule] let expiresAt: Double let retry: RetryOverride? + /// nil when the descriptor has none: the entry follows the queue setting. + let wifiOnly: Bool? /// Body identity for the same-id rules. let fingerprint: String } @@ -65,6 +67,11 @@ enum EnqueueParser { let url = present(d["url"]) as? String if let url { try requireURL(url, "url") } + var wifiOnly: Bool? + if let raw = present(d["wifiOnly"]) { + guard let b = raw as? Bool else { throw ParseError(message: "'wifiOnly' must be a boolean") } + wifiOnly = b + } var kinds: [ParsedEnqueue.Body] = [] if let text = present(d["dataJson"]) { @@ -92,7 +99,7 @@ enum EnqueueParser { id: id, key: key, varsJSON: varsJSON, url: url, method: method, headers: try headers(present(d["headers"])), body: body, parts: parts, accept: UploadOutcome.parseAcceptRules(present(d["accept"])), expiresAt: expiresAt, - retry: RetryOverride.parse(present(d["retry"])), + retry: RetryOverride.parse(present(d["retry"])), wifiOnly: wifiOnly, fingerprint: fingerprint(body, parts: parts, url: url, method: method)) } diff --git a/ios/QueueCoordinator+Enqueue.swift b/ios/QueueCoordinator+Enqueue.swift index 3a852eaf..b5234f0c 100644 --- a/ios/QueueCoordinator+Enqueue.swift +++ b/ios/QueueCoordinator+Enqueue.swift @@ -59,7 +59,7 @@ extension QueueCoordinator { } var e = QueueEntry.created(from: p, staged: staged, headerGeneration: settings.headerGeneration, - paused: settings.paused, now: now()) + paused: settings.isPaused(p.key), now: now()) if let manifest = adopted { e.incarnation = manifest.incarnation for i in e.parts.indices { e.parts[i].accepted = manifest.parts[i].accepted } @@ -70,12 +70,13 @@ extension QueueCoordinator { store.sweep(e) publish(e) armExpiry(e) - return (p.id, !settings.paused) + return (p.id, !settings.isPaused(p.key)) } private func enqueueExisting(_ existing: QueueEntry, _ p: ParsedEnqueue) throws -> (id: String, issue: Bool) { - let paused = settings.paused + // By the incoming key: a same-id enqueue may move the entry to another key. + let paused = settings.isPaused(p.key) if existing.bodyFingerprint == p.fingerprint && !existing.legacy { switch existing.state { case .completed: @@ -115,11 +116,19 @@ extension QueueCoordinator { case .queued, .running, .paused: // The in-flight task keeps its request. A retry waiting in the daemon - // picks up the new headers in willBeginDelayedRequest. - let n = existing.resumed(with: p, resetBudget: false, now: now()) + // picks up the new headers in willBeginDelayedRequest. A paused entry + // parked on auth leaves the parking spot: fresh headers came with the + // call, as in the awaiting-auth case. + var n = existing.resumed(with: p, resetBudget: false, now: now()) + n.authParked = false try saveOrThrow(n) publish(n) armExpiry(n) + if paused != (n.state == .paused) { + // A new key that another scope pauses, or no longer pauses. + applyPauseGate(n) + return (p.id, false) + } if ready, !paused, n.state == .queued, !n.isChunked, let at = n.nextAttemptAt, at > now() { // A simple retry waiting out its backoff: the caller asks again, // so retry now. The waiting attempt never ran: keep its ordinal. diff --git a/ios/QueueCoordinator+Reconcile.swift b/ios/QueueCoordinator+Reconcile.swift index 206ed7b4..310da76a 100644 --- a/ios/QueueCoordinator+Reconcile.swift +++ b/ios/QueueCoordinator+Reconcile.swift @@ -113,6 +113,9 @@ extension QueueCoordinator { withoutTask(e) } } + // A crash between the settings save and the entry saves of a pause or + // resume leaves rows that do not match the gates. Match them now. + applyPauseGates() } /// A queued or running entry with no live task. Either the daemon finished diff --git a/ios/QueueCoordinator+Simple.swift b/ios/QueueCoordinator+Simple.swift index 511a25ae..7659f1e0 100644 --- a/ios/QueueCoordinator+Simple.swift +++ b/ios/QueueCoordinator+Simple.swift @@ -15,7 +15,8 @@ extension QueueCoordinator { /// the ordinal and request id. It applies only while the entry holds such /// an attempt (nextAttemptAt set); otherwise a new attempt is minted. func issue(_ id: String, delayMs: Int? = nil, advanceAttempt: Bool = true) { - guard var e = index.entry(id), !e.legacy, e.state == .queued, !settings.paused else { return } + guard var e = index.entry(id), !e.legacy, e.state == .queued, + !settings.isPaused(e.key) else { return } let t = now() if t >= e.expiresAt { settle(id, .expired) @@ -76,7 +77,8 @@ extension QueueCoordinator { id: id, attempt: e.attempts, requestId: requestId, headerGeneration: settings.headerGeneration, generation: e.generation, purpose: .attempt) let task = transport.upload( - buildRequest(e, url: url, requestId: requestId), fromFile: body, wifiOnly: settings.wifiOnly, + buildRequest(e, url: url, requestId: requestId), fromFile: body, + wifiOnly: settings.wifiOnly(e.wifiOnly), description: ChunkedEngine.taskDescription(id: id, attempt: e.attempts, generation: e.generation), beginAt: beginAt.map { Date(timeIntervalSince1970: $0 / 1000) }, beforeResume: { key in self.taskMap.set(meta, forKey: key) }) @@ -199,7 +201,7 @@ extension QueueCoordinator { guard meta?.purpose != .superseded, case .request(let id, let generation, let attempt) = owner, var e = index.entry(id), !e.isChunked, e.generation == generation, - e.attempts == attempt, e.state == .queued || e.state == .running, !settings.paused, + e.attempts == attempt, e.state == .queued || e.state == .running, !settings.isPaused(e.key), let url = e.url.flatMap(URL.init(string:)) else { taskMap.setPurpose(.superseded, forKey: key, id: owner.id) liveTasks[key] = nil diff --git a/ios/QueueCoordinator.swift b/ios/QueueCoordinator.swift index 93360c71..6ee945ff 100644 --- a/ios/QueueCoordinator.swift +++ b/ios/QueueCoordinator.swift @@ -132,51 +132,83 @@ final class QueueCoordinator { } } - /// Whole-queue pause. Cancels every task with purpose "pause", so its - /// NSURLErrorCancelled produces no outcome and no attempt event. A single - /// body restarts from byte 0 on resume; a chunked upload keeps its - /// accepted parts. - func pause(resolve: @escaping () -> Void, reject: @escaping (String, String) -> Void) { + /// pause(scope). `{}` turns the global gate on; `{ keys }` adds the keys + /// to the paused set. Each live entry that is now paused stops: every task + /// is cancelled with purpose "pause", so its NSURLErrorCancelled produces + /// no outcome and no attempt event. A single body restarts from byte 0 on + /// resume; a chunked upload keeps its accepted parts. + func pause(_ scope: [String: Any], resolve: @escaping () -> Void, + reject: @escaping (String, String) -> Void) { + changeGate(scope, method: "pause", resolve: resolve, reject: reject) { next, keys in + if let keys { next.pausedKeys.formUnion(keys) } else { next.paused = true } + } + } + + /// resume(scope). `{}` turns the global gate off; `{ keys }` removes the + /// keys from the paused set. An entry another scope still pauses stays + /// paused. + func resume(_ scope: [String: Any], resolve: @escaping () -> Void, + reject: @escaping (String, String) -> Void) { + changeGate(scope, method: "resume", resolve: resolve, reject: reject) { next, keys in + if let keys { next.pausedKeys.subtract(keys) } else { next.paused = false } + } + } + + /// Parses the scope, saves the changed settings, then moves every entry + /// whose paused state changed. `change` gets nil keys for the whole queue. + private func changeGate(_ scope: [String: Any], method: String, resolve: @escaping () -> Void, + reject: @escaping (String, String) -> Void, + change: @escaping (inout QueueSettings, Set?) -> Void) { queue.async { + var keys: Set? + if let raw = scope["keys"], !(raw is NSNull) { + guard let list = raw as? [String] else { + reject("E_INVALID", "\(method): 'keys' must be a list of strings") + return + } + keys = Set(list) + } var next = self.settings - next.paused = true - guard self.saveSettings(next) else { - reject("E_STORAGE", "pause: cannot save the queue settings") + change(&next, keys) + guard next == self.settings || self.saveSettings(next) else { + reject("E_STORAGE", "\(method): cannot save the queue settings") return } - for e in self.index.entries() where !e.legacy - && [.queued, .running, .awaitingAuth].contains(e.state) { - self.cancelTasks(e.id, purpose: .pause) - self.chunked.stop(e.id) - var n = e - n.state = .paused - n.nextAttemptAt = nil - self.commit(n) - } + self.applyPauseGates() resolve() } } - func resume(resolve: @escaping () -> Void, reject: @escaping (String, String) -> Void) { - queue.async { - var next = self.settings - next.paused = false - guard self.saveSettings(next) else { - reject("E_STORAGE", "resume: cannot save the queue settings") + /// Moves every live entry to match the gates. See applyPauseGate. + func applyPauseGates() { + for e in index.entries() { applyPauseGate(e) } + } + + /// An entry is paused when the global gate is on or its key is paused. + /// A live entry that should be paused stops and moves to 'paused'. A + /// paused entry that should not be moves back: past expiresAt it settles + /// expired (a pause does not expire an entry; the resume does), else to + /// awaiting-auth when it was parked, else to queued and issues. + func applyPauseGate(_ e: QueueEntry) { + guard e.isLive, !e.legacy else { return } + let gated = settings.isPaused(e.key) + if gated, e.state != .paused { + cancelTasks(e.id, purpose: .pause) + chunked.stop(e.id) + var n = e + n.state = .paused + n.nextAttemptAt = nil + commit(n) + } else if !gated, e.state == .paused { + if now() >= e.expiresAt { + settle(e.id, .expired) return } - for e in self.index.entries() where e.state == .paused && !e.legacy { - // A pause does not expire an entry; the resume does. - if self.now() >= e.expiresAt { - self.settle(e.id, .expired) - continue - } - var n = e - n.state = e.authParked ? .awaitingAuth : .queued - self.commit(n) - if n.state == .queued { self.issue(n.id) } - } - resolve() + var n = e + n.state = e.authParked ? .awaitingAuth : .queued + commit(n) + armExpiry(n) + if n.state == .queued { issue(n.id) } } } @@ -203,9 +235,10 @@ final class QueueCoordinator { } } - /// Persisted. Queued and future entries use the session it picks. A queued - /// entry whose retry waits in the daemon moves to the new session; a - /// running task finishes where it started (session config is fixed). + /// Persisted. Queued and future entries with no wifiOnly of their own use + /// the session it picks. Such a queued entry whose retry waits in the + /// daemon moves to the new session; a running task finishes where it + /// started (session config is fixed). An entry that pins wifiOnly stays. func setWifiOnly(_ enabled: Bool, resolve: @escaping () -> Void, reject: @escaping (String, String) -> Void) { queue.async { @@ -217,7 +250,8 @@ final class QueueCoordinator { return } if changed && self.ready { - for e in self.index.entries() where e.state == .queued && !e.isChunked && !e.legacy { + for e in self.index.entries() where e.state == .queued && !e.isChunked && !e.legacy + && e.wifiOnly == nil { let remaining = e.nextAttemptAt.map { Int($0 - self.now()) } self.cancelTasks(e.id, purpose: .superseded) // The waiting attempt never ran: keep its ordinal and request id. @@ -257,7 +291,7 @@ final class QueueCoordinator { n.headerGeneration = self.settings.headerGeneration if n.state == .awaitingAuth { n.authParked = false - n.state = self.settings.paused ? .paused : .queued + n.state = self.settings.isPaused(n.key) ? .paused : .queued self.commit(n) self.issue(n.id) } else { diff --git a/ios/QueueEntry.swift b/ios/QueueEntry.swift index 7ebc2437..984d089d 100644 --- a/ios/QueueEntry.swift +++ b/ios/QueueEntry.swift @@ -36,6 +36,9 @@ struct QueueEntry: Codable, Equatable { var method: String var accept: [UploadOutcome.AcceptRule] var retry: RetryOverride? + /// descriptor.wifiOnly. nil follows the queue setting (setWifiOnly), so a + /// toggle moves this entry; true or false pins it. + var wifiOnly: Bool? var bodyKind: BodyKind /// File name inside the entry directory: "body-" or a blob name. /// nil for a legacy row with no bytes. @@ -127,7 +130,7 @@ extension QueueEntry { paused: Bool, now: Double, createdAt: Double? = nil) -> QueueEntry { QueueEntry( id: p.id, key: p.key, varsJSON: p.varsJSON, url: p.url, method: p.method, accept: p.accept, retry: p.retry, - bodyKind: staged.kind, bodyPath: staged.relativePath, + wifiOnly: p.wifiOnly, bodyKind: staged.kind, bodyPath: staged.relativePath, bodyContentType: staged.contentType, forceContentType: staged.forceContentType, bodyFingerprint: p.fingerprint, parts: p.parts, incarnation: UUID().uuidString, headers: p.headers, headerGeneration: headerGeneration, @@ -150,6 +153,7 @@ extension QueueEntry { next.expiresAt = p.expiresAt next.accept = p.accept next.retry = p.retry + next.wifiOnly = p.wifiOnly // Same parts by definition of "same body"; the incoming ones carry the // new per-part headers. if isChunked, p.parts.count == parts.count { diff --git a/ios/QueueSettings.swift b/ios/QueueSettings.swift index 6249af68..eecb96b5 100644 --- a/ios/QueueSettings.swift +++ b/ios/QueueSettings.swift @@ -48,7 +48,11 @@ struct RetryOverride: Codable, Equatable { /// Queue-wide settings, persisted as `settings.json` in the queue directory. struct QueueSettings: Codable, Equatable { var wifiOnly = false + /// The global pause gate: pause() / resume() with no keys. var paused = false + /// Keys paused by pause({ keys }). An entry is paused when the gate is on + /// or its key is in this set. + var pausedKeys: Set = [] /// Bumped by every updateHeaders(). An attempt records the value it was /// issued under; a 401/403 from an older value re-issues instead of parking. var headerGeneration = 0 @@ -62,10 +66,17 @@ struct QueueSettings: Codable, Equatable { let c = try decoder.container(keyedBy: CodingKeys.self) wifiOnly = try c.decodeIfPresent(Bool.self, forKey: .wifiOnly) ?? false paused = try c.decodeIfPresent(Bool.self, forKey: .paused) ?? false + pausedKeys = try c.decodeIfPresent(Set.self, forKey: .pausedKeys) ?? [] headerGeneration = try c.decodeIfPresent(Int.self, forKey: .headerGeneration) ?? 0 retry = try c.decodeIfPresent(RetryOverride.self, forKey: .retry) } + /// true when either scope pauses entries of `key`. + func isPaused(_ key: String) -> Bool { paused || pausedKeys.contains(key) } + + /// The queue setting unless the entry pins its own. + func wifiOnly(_ entryWifiOnly: Bool?) -> Bool { entryWifiOnly ?? wifiOnly } + /// configure(options). iOS reads `retry` only: JS turns lifetimeMs into /// each entry's expiresAt, and the Android notification keys are Android's. mutating func apply(configure options: [String: Any]) { diff --git a/ios/RNBackgroundUpload.swift b/ios/RNBackgroundUpload.swift index eb8fa4a2..6e63f890 100644 --- a/ios/RNBackgroundUpload.swift +++ b/ios/RNBackgroundUpload.swift @@ -139,14 +139,17 @@ public class RNBackgroundUpload: NSObject, URLSessionDataDelegate { coordinator.enqueue(entry, resolve: { resolve($0) }, reject: { reject($0, $1, nil) }) } - @objc(pause:reject:) - public func pause(_ resolve: @escaping RCTPromiseResolveBlock, reject: @escaping RCTPromiseRejectBlock) { - coordinator.pause(resolve: { resolve(nil) }, reject: { reject($0, $1, nil) }) + // scope is { keys?: [String] }; [:] is the whole queue. + @objc(pause:resolve:reject:) + public func pause(_ scope: [String: Any], resolve: @escaping RCTPromiseResolveBlock, + reject: @escaping RCTPromiseRejectBlock) { + coordinator.pause(scope, resolve: { resolve(nil) }, reject: { reject($0, $1, nil) }) } - @objc(resume:reject:) - public func resume(_ resolve: @escaping RCTPromiseResolveBlock, reject: @escaping RCTPromiseRejectBlock) { - coordinator.resume(resolve: { resolve(nil) }, reject: { reject($0, $1, nil) }) + @objc(resume:resolve:reject:) + public func resume(_ scope: [String: Any], resolve: @escaping RCTPromiseResolveBlock, + reject: @escaping RCTPromiseRejectBlock) { + coordinator.resume(scope, resolve: { resolve(nil) }, reject: { reject($0, $1, nil) }) } @objc(cancel:resolve:reject:) diff --git a/ios/RNFileUploader.mm b/ios/RNFileUploader.mm index 2a2395c3..f2659478 100644 --- a/ios/RNFileUploader.mm +++ b/ios/RNFileUploader.mm @@ -82,16 +82,18 @@ - (void)enqueue:(NSDictionary *)entry [RNBackgroundUpload.shared enqueue:entry resolve:resolve reject:reject]; } -- (void)pause:(RCTPromiseResolveBlock)resolve +- (void)pause:(NSDictionary *)scope + resolve:(RCTPromiseResolveBlock)resolve reject:(RCTPromiseRejectBlock)reject { - [RNBackgroundUpload.shared pause:resolve reject:reject]; + [RNBackgroundUpload.shared pause:scope resolve:resolve reject:reject]; } -- (void)resume:(RCTPromiseResolveBlock)resolve +- (void)resume:(NSDictionary *)scope + resolve:(RCTPromiseResolveBlock)resolve reject:(RCTPromiseRejectBlock)reject { - [RNBackgroundUpload.shared resume:resolve reject:reject]; + [RNBackgroundUpload.shared resume:scope resolve:resolve reject:reject]; } - (void)cancel:(NSString *)id diff --git a/ios/Tests/CoordinatorChunkedTests.swift b/ios/Tests/CoordinatorChunkedTests.swift index fc6b635e..2e00e0e5 100644 --- a/ios/Tests/CoordinatorChunkedTests.swift +++ b/ios/Tests/CoordinatorChunkedTests.swift @@ -266,6 +266,32 @@ final class CoordinatorChunkedTests: XCTestCase { XCTAssertNil(partTask(0)) } + func testKeyPauseStopsOnlyThatKeysPartsAndResumeRefills() throws { + _ = try h.enqueue(h.chunkedRaw(id: "cap", key: "capture", size: 50, parts: 5)).get() + _ = try h.enqueue(h.dataRaw(id: "note", key: "fieldNote")).get() + h.complete(try XCTUnwrap(partTask(0))) + h.pause(keys: ["capture"]) + XCTAssertEqual(h.entry("cap")?.state, .paused) + XCTAssertNil(partTask(1)) + XCTAssertEqual(h.entry("note")?.state, .running) + XCTAssertEqual(h.entry("cap")?.parts[0].accepted, true, "accepted parts are kept") + h.resume(keys: ["capture"]) + XCTAssertEqual(h.entry("cap")?.state, .running) + XCTAssertNil(partTask(0)) + XCTAssertNotNil(partTask(1)) + } + + func testEntryWifiOnlyPicksThePartSession() throws { + _ = try h.enqueue(h.chunkedRaw(id: "cap", extra: ["wifiOnly": true])).get() + let parts = h.transport.live.filter { ChunkedEngine.parsePartDescription($0.taskDescription) != nil } + XCTAssertEqual(parts.count, 3) + XCTAssertTrue(parts.allSatisfy(\.wifiOnly)) + h.setWifiOnly(true) + h.setWifiOnly(false) + h.complete(try XCTUnwrap(partTask(0)), status: 503) + XCTAssertEqual(partTask(0)?.wifiOnly, true, "a part retry after a toggle keeps the entry's own setting") + } + func testCancelLiveChunked() throws { _ = try h.enqueue(h.chunkedRaw(id: "cap", size: 50, parts: 5)).get() h.complete(try XCTUnwrap(partTask(0))) diff --git a/ios/Tests/CoordinatorSimpleTests.swift b/ios/Tests/CoordinatorSimpleTests.swift index 7bca024a..4b0a3d31 100644 --- a/ios/Tests/CoordinatorSimpleTests.swift +++ b/ios/Tests/CoordinatorSimpleTests.swift @@ -603,6 +603,209 @@ final class CoordinatorSimpleTests: XCTestCase { XCTAssertEqual(h.entry("a")?.state, .queued) } + // MARK: - Per-entry wifiOnly + + private func liveTask(_ id: String) -> FakeTask? { + h.transport.live.first { $0.taskDescription?.contains("\"\(id)\"") == true } + } + + func testEntryWifiOnlyOverridesTheQueueSetting() throws { + _ = try h.enqueue(h.dataRaw(id: "cap", extra: ["wifiOnly": true])).get() + XCTAssertEqual(liveTask("cap")?.wifiOnly, true) + h.setWifiOnly(true) + _ = try h.enqueue(h.dataRaw(id: "note", extra: ["wifiOnly": false])).get() + XCTAssertEqual(liveTask("note")?.wifiOnly, false, "false pins cellular too") + _ = try h.enqueue(h.dataRaw(id: "plain")).get() + XCTAssertEqual(liveTask("plain")?.wifiOnly, true, "no field follows the queue") + } + + func testEntryWifiOnlyIsPersistedAndUsedForEachRetry() throws { + _ = try h.enqueue(h.dataRaw(id: "cap", extra: ["wifiOnly": true])).get() + XCTAssertEqual(h.store.load("cap")?.wifiOnly, true) + _ = try h.enqueue(h.dataRaw(id: "b")).get() + XCTAssertNil(h.store.load("b")?.wifiOnly) + h.relaunch() + h.boot() + h.complete(try XCTUnwrap(liveTask("cap")), status: 503) + XCTAssertEqual(liveTask("cap")?.wifiOnly, true, "the retry after a relaunch keeps it") + } + + func testSetWifiOnlyMovesOnlyEntriesThatFollowIt() throws { + _ = try h.enqueue(h.dataRaw(id: "pinned", extra: ["wifiOnly": false])).get() + _ = try h.enqueue(h.dataRaw(id: "follows")).get() + h.complete(try XCTUnwrap(liveTask("pinned")), status: 503) + h.complete(try XCTUnwrap(liveTask("follows")), status: 503) + let pinnedWait = try XCTUnwrap(liveTask("pinned")) + let followsWait = try XCTUnwrap(liveTask("follows")) + h.setWifiOnly(true) + XCTAssertFalse(pinnedWait.cancelled, "a pinned entry keeps its task") + XCTAssertTrue(followsWait.cancelled) + XCTAssertEqual(liveTask("follows")?.wifiOnly, true) + XCTAssertEqual(liveTask("pinned")?.key, pinnedWait.key) + } + + func testSameIdEnqueueReplacesTheEntryWifiOnly() throws { + _ = try h.enqueue(h.dataRaw(id: "a", extra: ["wifiOnly": true])).get() + h.complete(try XCTUnwrap(liveTask("a")), status: 503) + _ = try h.enqueue(h.dataRaw(id: "a")).get() + XCTAssertNil(h.entry("a")?.wifiOnly) + XCTAssertEqual(liveTask("a")?.wifiOnly, false, "the retry re-issued under the queue setting") + } + + // MARK: - Key-scoped pause + + func testKeyPauseStopsOnlyThatKeyAndItsFutureEntries() throws { + _ = try h.enqueue(h.dataRaw(id: "cap", key: "capture")).get() + _ = try h.enqueue(h.dataRaw(id: "note", key: "fieldNote")).get() + let capTask = try XCTUnwrap(liveTask("cap")) + h.sink.states = [] + h.pause(keys: ["capture"]) + XCTAssertTrue(capTask.cancelled) + XCTAssertEqual(h.map.meta(forKey: capTask.key)?.purpose, .pause) + XCTAssertEqual(h.entry("cap")?.state, .paused) + XCTAssertEqual(h.entry("note")?.state, .running) + XCTAssertEqual(h.sink.states.map { $0["id"] as? String }, ["cap"], "one state event, for the paused row") + h.deliverCancel(capTask) + XCTAssertTrue(h.sink.settled.isEmpty) + XCTAssertTrue(h.sink.attempts.isEmpty) + + _ = try h.enqueue(h.dataRaw(id: "cap2", key: "capture")).get() + XCTAssertEqual(h.entry("cap2")?.state, .paused, "a future entry of the key starts paused") + XCTAssertNil(liveTask("cap2")) + _ = try h.enqueue(h.dataRaw(id: "note2", key: "fieldNote")).get() + XCTAssertNotNil(liveTask("note2")) + + h.sink.states = [] + h.resume(keys: ["capture"]) + XCTAssertEqual(h.entry("cap")?.state, .running) + XCTAssertEqual(h.entry("cap2")?.state, .running) + XCTAssertEqual(h.entry("cap")?.attempts, 2) + XCTAssertEqual(Set(h.sink.states.compactMap { $0["id"] as? String }), ["cap", "cap2"]) + } + + func testAKeyResumeKeepsEntriesTheGlobalGatePauses() throws { + _ = try h.enqueue(h.dataRaw(id: "cap", key: "capture")).get() + _ = try h.enqueue(h.dataRaw(id: "note", key: "fieldNote")).get() + h.pause(keys: ["capture"]) + h.pause() + h.sink.states = [] + h.resume(keys: ["capture"]) + XCTAssertEqual(h.entry("cap")?.state, .paused, "the global gate still pauses it") + XCTAssertTrue(h.sink.states.isEmpty) + XCTAssertTrue(h.transport.live.isEmpty) + h.resume() + XCTAssertEqual(h.entry("cap")?.state, .running) + XCTAssertEqual(h.entry("note")?.state, .running) + } + + func testAGlobalResumeKeepsEntriesAKeyPauses() throws { + _ = try h.enqueue(h.dataRaw(id: "cap", key: "capture")).get() + _ = try h.enqueue(h.dataRaw(id: "note", key: "fieldNote")).get() + h.pause() + h.pause(keys: ["capture"]) + h.resume() + XCTAssertEqual(h.entry("cap")?.state, .paused, "the key still pauses it") + XCTAssertEqual(h.entry("note")?.state, .running) + XCTAssertNil(liveTask("cap")) + } + + func testPausedKeysSurviveRelaunch() throws { + h.pause(keys: ["capture"]) + XCTAssertEqual(h.store.loadSettings().pausedKeys, ["capture"]) + h.relaunch() + h.boot() + _ = try h.enqueue(h.dataRaw(id: "cap", key: "capture")).get() + _ = try h.enqueue(h.dataRaw(id: "note", key: "fieldNote")).get() + XCTAssertEqual(h.entry("cap")?.state, .paused) + XCTAssertEqual(h.entry("note")?.state, .running) + } + + func testEmptyKeysChangeNothing() throws { + _ = try h.enqueue(h.dataRaw(id: "a")).get() + h.pause(keys: []) + XCTAssertEqual(h.entry("a")?.state, .running, "an empty list never means the whole queue") + XCTAssertEqual(h.store.loadSettings(), QueueSettings()) + } + + func testKeysThatAreNotAStringListReject() throws { + var code: String? + h.coordinator.pause(["keys": "capture"], resolve: { XCTFail("resolved") }, reject: { c, _ in code = c }) + h.drain() + XCTAssertEqual(code, "E_INVALID") + XCTAssertEqual(h.store.loadSettings(), QueueSettings()) + } + + func testKeyPausedEntryCrossingExpiresAtSettlesAtItsResume() throws { + _ = try h.enqueue(h.dataRaw(id: "cap", key: "capture", extra: ["expiresAt": h.clock + 60_000])).get() + h.pause(keys: ["capture"]) + h.advance(60_200) + XCTAssertEqual(h.entry("cap")?.state, .paused) + h.resume(keys: ["capture"]) + XCTAssertEqual(h.entry("cap")?.state, .error) + XCTAssertTrue(h.transport.live.isEmpty) + } + + func testKeyPauseKeepsAuthParkingAcrossResume() throws { + _ = try h.enqueue(h.dataRaw(id: "a", key: "capture")).get() + h.complete(try XCTUnwrap(liveTask("a")), status: 403) + h.pause(keys: ["capture"]) + XCTAssertEqual(h.entry("a")?.state, .paused) + h.resume(keys: ["capture"]) + XCTAssertEqual(h.entry("a")?.state, .awaitingAuth) + XCTAssertTrue(h.transport.live.isEmpty) + } + + func testSameIdEnqueueUnderAPausedKeyPausesTheEntry() throws { + h.pause(keys: ["capture"]) + _ = try h.enqueue(h.dataRaw(id: "a", key: "fieldNote")).get() + let task = try XCTUnwrap(liveTask("a")) + _ = try h.enqueue(h.dataRaw(id: "a", key: "capture")).get() + XCTAssertTrue(task.cancelled) + XCTAssertEqual(h.entry("a")?.state, .paused) + _ = try h.enqueue(h.dataRaw(id: "a", key: "fieldNote")).get() + XCTAssertEqual(h.entry("a")?.state, .running, "back under a key nothing pauses") + XCTAssertNotNil(liveTask("a")) + } + + func testSameIdEnqueueOffAPausedKeyLeavesAuthParking() throws { + _ = try h.enqueue(h.dataRaw(id: "a", key: "capture")).get() + h.complete(try XCTUnwrap(liveTask("a")), status: 403) + h.pause(keys: ["capture"]) + XCTAssertEqual(h.entry("a")?.state, .paused) + XCTAssertEqual(h.entry("a")?.authParked, true) + _ = try h.enqueue(h.dataRaw(id: "a", key: "fieldNote")).get() + XCTAssertEqual(h.entry("a")?.state, .running, "the fresh headers are tried, not parked") + XCTAssertEqual(h.entry("a")?.authParked, false) + XCTAssertEqual(h.store.load("a")?.authParked, false) + XCTAssertNotNil(liveTask("a")) + } + + func testSameIdEnqueueOnAGloballyPausedParkedEntryResumesToQueued() throws { + _ = try h.enqueue(h.dataRaw(id: "a")).get() + h.complete(try XCTUnwrap(liveTask("a")), status: 401) + h.pause() + _ = try h.enqueue(h.dataRaw(id: "a")).get() + XCTAssertEqual(h.entry("a")?.state, .paused, "the gate still pauses it") + XCTAssertNil(liveTask("a")) + h.resume() + XCTAssertEqual(h.entry("a")?.state, .running, "the fresh headers are tried at the resume") + XCTAssertNotNil(liveTask("a")) + } + + func testRelaunchMatchesRowsThatAPauseDidNotReach() throws { + _ = try h.enqueue(h.dataRaw(id: "cap", key: "capture")).get() + let task = try XCTUnwrap(liveTask("cap")) + // A crash after the settings save, before the entry save. + var s = h.store.loadSettings() + s.pausedKeys = ["capture"] + try h.store.saveSettings(s) + h.relaunch() + h.boot() + XCTAssertEqual(h.entry("cap")?.state, .paused) + XCTAssertTrue(task.cancelled) + XCTAssertTrue(h.transport.live.isEmpty) + } + // MARK: - Write-ahead failures func testFailedAttemptSaveCreatesNoTaskAndIssuesAfterABackoff() throws { diff --git a/ios/Tests/SupportComponentTests.swift b/ios/Tests/SupportComponentTests.swift index 4780d67d..635c2c1e 100644 --- a/ios/Tests/SupportComponentTests.swift +++ b/ios/Tests/SupportComponentTests.swift @@ -317,3 +317,56 @@ final class AttemptEventTests: XCTestCase { XCTAssertEqual(e["responseBodyTruncated"] as? Bool, true) } } + +final class QueueSettingsTests: XCTestCase { + func testPausedDerivesFromTheGateOrTheKey() { + var s = QueueSettings() + XCTAssertFalse(s.isPaused("capture")) + s.pausedKeys = ["capture"] + XCTAssertTrue(s.isPaused("capture")) + XCTAssertFalse(s.isPaused("fieldNote")) + s.paused = true + XCTAssertTrue(s.isPaused("fieldNote")) + } + + func testEntryWifiOnlyWinsOverTheQueueSetting() { + var s = QueueSettings() + s.wifiOnly = true + XCTAssertTrue(s.wifiOnly(nil)) + XCTAssertFalse(s.wifiOnly(false)) + s.wifiOnly = false + XCTAssertTrue(s.wifiOnly(true)) + } + + func testPausedKeysRoundTripAndAnOlderFileLoads() throws { + var s = QueueSettings() + s.paused = true + s.pausedKeys = ["capture", "video"] + let back = try JSONDecoder().decode(QueueSettings.self, from: try JSONEncoder().encode(s)) + XCTAssertEqual(back, s) + let old = try JSONDecoder().decode(QueueSettings.self, from: Data(#"{"wifiOnly":true,"paused":true}"#.utf8)) + XCTAssertEqual(old.pausedKeys, []) + XCTAssertTrue(old.paused) + } + + func testAnEntryFromAnOlderBuildHasNoWifiOnly() throws { + var json = try JSONSerialization.jsonObject(with: try QueueStore.encode(sampleEntry("a", createdAt: 1))) + as! [String: Any] + json["wifiOnly"] = nil + let e = try JSONDecoder().decode(QueueEntry.self, from: try JSONSerialization.data(withJSONObject: json)) + XCTAssertNil(e.wifiOnly, "follows the queue setting") + } + + func testParserReadsWifiOnly() throws { + func parse(_ extra: [String: Any]) throws -> ParsedEnqueue { + var d: [String: Any] = ["url": "https://a.test", "expiresAt": 1.0] + for (k, v) in extra { d[k] = v } + return try EnqueueParser.parse(["id": "a", "key": "k", "varsJson": "{}", "descriptor": d]) + } + XCTAssertNil(try parse([:]).wifiOnly) + XCTAssertNil(try parse(["wifiOnly": NSNull()]).wifiOnly) + XCTAssertEqual(try parse(["wifiOnly": true]).wifiOnly, true) + XCTAssertEqual(try parse(["wifiOnly": false]).wifiOnly, false) + XCTAssertThrowsError(try parse(["wifiOnly": "yes"])) + } +} diff --git a/ios/Tests/TestSupport.swift b/ios/Tests/TestSupport.swift index c9e21d58..66fb6742 100644 --- a/ios/Tests/TestSupport.swift +++ b/ios/Tests/TestSupport.swift @@ -212,13 +212,16 @@ final class Harness { drain() } - func pause() { - coordinator.pause(resolve: {}, reject: { _, _ in XCTFail("pause rejected") }) + /// No keys: the whole queue. `keys`: those definition keys. + func pause(keys: [String]? = nil) { + coordinator.pause(keys.map { ["keys": $0] } ?? [:], resolve: {}, + reject: { _, _ in XCTFail("pause rejected") }) drain() } - func resume() { - coordinator.resume(resolve: {}, reject: { _, _ in XCTFail("resume rejected") }) + func resume(keys: [String]? = nil) { + coordinator.resume(keys.map { ["keys": $0] } ?? [:], resolve: {}, + reject: { _, _ in XCTFail("resume rejected") }) drain() } @@ -268,16 +271,16 @@ final class Harness { } /// `data` crosses as its JSON text, as the JS layer sends it. - func dataRaw(id: String, data: Any = ["x": 1], url: String = "https://api.test/x", + func dataRaw(id: String, key: String = "k", data: Any = ["x": 1], url: String = "https://api.test/x", headers: [String: Any] = [:], extra: [String: Any] = [:]) -> [String: Any] { var d: [String: Any] = ["url": url, "dataJson": jsonText(data), "headers": headers] for (k, v) in extra { d[k] = v } - return raw(id: id, descriptor: d) + return raw(id: id, key: key, descriptor: d) } /// A chunked descriptor over a fresh source file of `size` bytes cut into /// `parts` equal parts. - func chunkedRaw(id: String, size: Int = 30, parts: Int = 3, source: URL? = nil, + func chunkedRaw(id: String, key: String = "k", size: Int = 30, parts: Int = 3, source: URL? = nil, urlPrefix: String = "https://s3.test/part", extra: [String: Any] = [:]) -> [String: Any] { let file = source ?? root.appendingPathComponent("src-\(UUID().uuidString)") if source == nil { writeFile(file, bytes: size) } @@ -291,6 +294,6 @@ final class Harness { var d: [String: Any] = ["file": file.path, "parts": list, "method": "PUT", "headers": ["Content-Type": "video/mp4"]] for (k, v) in extra { d[k] = v } - return raw(id: id, descriptor: d) + return raw(id: id, key: key, descriptor: d) } } diff --git a/src/NativeRNFileUploader.ts b/src/NativeRNFileUploader.ts index bd97302b..d9614d9a 100644 --- a/src/NativeRNFileUploader.ts +++ b/src/NativeRNFileUploader.ts @@ -17,7 +17,9 @@ export interface Spec extends TurboModule { // relaunch (no JS) can read it. Each call replaces the full configuration. configure(options: CodegenTypes.UnsafeObject): void; // Persists { id, key, varsJson, descriptor } and schedules it. varsJson is - // JSON.stringify(vars). descriptor.dataJson is JSON.stringify(data) and + // JSON.stringify(vars). descriptor.wifiOnly, when present, is a boolean + // persisted with the entry; see setWifiOnly. An entry enqueued while its + // scope is paused starts 'paused'; see pause. descriptor.dataJson is JSON.stringify(data) and // replaces data; a bodiless request omits it. Both cross as strings // because React Native on iOS drops object keys whose value is null, so // { status: null } would arrive as {}. Native parses them. dataJson @@ -48,17 +50,33 @@ export interface Spec extends TurboModule { // // The resolved value is the entry's id. The JS layer does not read it. enqueue(entry: CodegenTypes.UnsafeObject): Promise; - // Whole-queue pause. No outcome is produced; live rows move to 'paused'. - // A paused entry past expiresAt settles error/expired at resume. - pause(): Promise; - resume(): Promise; + // scope is { keys?: string[] }. JS always passes an object, because + // codegen has no optional arguments. + // - {} (no keys): the global gate. pause turns it on, resume turns it off. + // - { keys }: pause adds the keys to a paused-keys set, resume removes + // them. An empty list changes nothing. keys never mean the whole queue. + // An entry is paused when the global gate is on OR its key is in the set, + // so a resume of one scope does not resume an entry that another scope + // still pauses (gate on + key resumed stays paused). The gate and the set + // are persisted, and apply to queued and future entries. + // + // Each live entry that becomes paused moves to 'paused', and each that + // stops being paused moves back to 'queued', with one 'state' event per + // entry. A running attempt stops. No outcome is produced and no attempt + // event is emitted. A paused entry past expiresAt settles error/expired + // when it stops being paused. + 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. cancel(id: string): Promise; - // Persisted natively. Applies to queued and future entries. + // The queue's Wi-Fi setting. Persisted natively. Applies to queued and + // future entries whose descriptor has no wifiOnly. An entry with + // descriptor.wifiOnly set (true or false) ignores it. Both are evaluated + // per attempt, so a toggle moves queued entries that follow the setting. setWifiOnly(enabled: boolean): Promise; // Merges the patch into every entry not yet forgotten and bumps a header // generation. The patch also replaces same-named headers a part carries. A 401/403 from an attempt issued under an older generation diff --git a/src/__tests__/client.test.ts b/src/__tests__/client.test.ts index 1a2d453d..1e2f14fe 100644 --- a/src/__tests__/client.test.ts +++ b/src/__tests__/client.test.ts @@ -246,14 +246,48 @@ describe('configure', () => { describe('queue control forwards', () => { const client = createUploadClient(); - it('pause', async () => { + it('pause with no scope forwards {} (the whole queue)', async () => { await client.pause(); - expect(native.pause).toHaveBeenCalledTimes(1); + expect(native.pause).toHaveBeenCalledWith({}); + await client.pause({}); + expect(native.pause).toHaveBeenLastCalledWith({}); }); - it('resume', async () => { + it('resume with no scope forwards {} (the whole queue)', async () => { await client.resume(); - expect(native.resume).toHaveBeenCalledTimes(1); + expect(native.resume).toHaveBeenCalledWith({}); + }); + + it('forwards a keys scope as a copy', async () => { + const keys = ['capture.upload']; + await client.pause({ keys }); + expect(native.pause).toHaveBeenCalledWith({ keys: ['capture.upload'] }); + expect(native.pause.mock.calls[0][0].keys).not.toBe(keys); + await client.resume({ keys: ['capture.upload', 'capture.meta'] }); + expect(native.resume).toHaveBeenCalledWith({ + keys: ['capture.upload', 'capture.meta'], + }); + }); + + it('forwards an empty keys list, which changes nothing natively', async () => { + await client.pause({ keys: [] }); + expect(native.pause).toHaveBeenCalledWith({ keys: [] }); + }); + + it.each([ + ['a non-object scope', 'capture', /scope must be an object/], + ['an array scope', ['capture'], /scope must be an object/], + ['a null scope', null, /scope must be an object/], + ['a misspelled field', { key: ['capture'] }, /unknown scope field "key"/], + ['keys that is not an array', { keys: 'capture' }, /keys must be an array/], + ['keys: undefined', { keys: undefined }, /keys must be an array/], + ['an empty key', { keys: ['a', ''] }, /non-empty strings/], + ['a non-string key', { keys: [1] }, /non-empty strings/], + ])('rejects %s and never reaches native', async (_name, scope, message) => { + await expect(client.pause(scope as any)).rejects.toThrow(message); + await expect(client.resume(scope as any)).rejects.toThrow(message); + expect(native.pause).not.toHaveBeenCalled(); + expect(native.resume).not.toHaveBeenCalled(); }); it('cancel', async () => { diff --git a/src/__tests__/registry.test.ts b/src/__tests__/registry.test.ts index 89cf1512..d44c7aee 100644 --- a/src/__tests__/registry.test.ts +++ b/src/__tests__/registry.test.ts @@ -566,6 +566,33 @@ describe('mutate', () => { ).resolves.toBeDefined(); }); + it('forwards wifiOnly to native when set, and omits it when absent', async () => { + const on = mutateWith({ url: 'https://x', wifiOnly: true }); + await on.promise; + expect(on.enqueue.mock.calls[0][0].descriptor.wifiOnly).toBe(true); + const off = mutateWith({ url: 'https://x', wifiOnly: false }); + await off.promise; + expect(off.enqueue.mock.calls[0][0].descriptor.wifiOnly).toBe(false); + // Absent means the entry follows setWifiOnly(), so nothing crosses. + const follow = mutateWith({ url: 'https://x' }); + await follow.promise; + expect(follow.enqueue.mock.calls[0][0].descriptor).not.toHaveProperty( + 'wifiOnly', + ); + }); + + it('rejects a wifiOnly that is not a boolean, and a misspelling', async () => { + await expect( + mutateWith({ url: 'https://x', wifiOnly: 'yes' }).promise, + ).rejects.toThrow(/wifiOnly must be a boolean/); + await expect( + mutateWith({ url: 'https://x', wifiOnly: null }).promise, + ).rejects.toThrow(/wifiOnly must be a boolean/); + await expect( + mutateWith({ url: 'https://x', wifionly: true }).promise, + ).rejects.toThrow(/unknown descriptor field "wifionly".*"wifiOnly"/); + }); + it('never reaches native on a rejected descriptor', async () => { const { promise, enqueue } = mutateWith({ data: {} }); await expect(promise).rejects.toThrow(); diff --git a/src/__tests__/testing.test.ts b/src/__tests__/testing.test.ts index 64f96416..9352fb72 100644 --- a/src/__tests__/testing.test.ts +++ b/src/__tests__/testing.test.ts @@ -115,18 +115,95 @@ describe('queue controls', () => { const fake = createFakeNative(); await fake.enqueue(entry()); fake.seedRows([row()]); - await fake.pause(); + await fake.pause({}); expect(fake.getRequests().map((r) => (r as RequestRow).state)).toEqual([ 'paused', 'error', ]); - await fake.resume(); + await fake.resume({}); expect(fake.getRequests().map((r) => (r as RequestRow).state)).toEqual([ 'queued', 'error', ]); - expect(fake.calls.pause).toBe(1); - expect(fake.calls.resume).toBe(1); + expect(fake.calls.pause).toEqual([{}]); + expect(fake.calls.resume).toEqual([{}]); + }); + + describe('key-scoped pause', () => { + const states = (fake: ReturnType) => + Object.fromEntries( + fake.getRequests().map((r) => [ + (r as RequestRow).id, + (r as RequestRow).state, + ]), + ); + + const twoKeys = async () => { + const fake = createFakeNative(); + await fake.enqueue(entry({ id: 'cap', key: 'capture' })); + await fake.enqueue(entry({ id: 'note', key: 'note' })); + return fake; + }; + + it('pauses only the entries of the given keys, with one state event each', async () => { + const fake = await twoKeys(); + const events: RequestRow[] = []; + fake.onState((r) => { + events.push(r as RequestRow); + }); + await fake.pause({ keys: ['capture'] }); + expect(states(fake)).toEqual({ cap: 'paused', note: 'queued' }); + expect(events.map((e) => [e.id, e.state])).toEqual([['cap', 'paused']]); + await fake.resume({ keys: ['capture'] }); + expect(states(fake)).toEqual({ cap: 'queued', note: 'queued' }); + expect(fake.calls.pause).toEqual([{ keys: ['capture'] }]); + expect(fake.calls.resume).toEqual([{ keys: ['capture'] }]); + }); + + it('keeps a key-paused entry paused when the whole queue resumes', async () => { + const fake = await twoKeys(); + await fake.pause({ keys: ['capture'] }); + await fake.pause({}); + expect(states(fake)).toEqual({ cap: 'paused', note: 'paused' }); + await fake.resume({}); + expect(states(fake)).toEqual({ cap: 'paused', note: 'queued' }); + }); + + it('keeps an entry paused when its key resumes but the gate is on', async () => { + const fake = await twoKeys(); + await fake.pause({ keys: ['capture'] }); + await fake.pause({}); + await fake.resume({ keys: ['capture'] }); + expect(states(fake)).toEqual({ cap: 'paused', note: 'paused' }); + await fake.resume({}); + expect(states(fake)).toEqual({ cap: 'queued', note: 'queued' }); + }); + + it('starts an entry paused when its scope is paused at enqueue', async () => { + const fake = createFakeNative(); + await fake.pause({ keys: ['capture'] }); + await fake.enqueue(entry({ id: 'cap', key: 'capture' })); + await fake.enqueue(entry({ id: 'note', key: 'note' })); + expect(states(fake)).toEqual({ cap: 'paused', note: 'queued' }); + await fake.pause({}); + await fake.enqueue(entry({ id: 'late', key: 'note' })); + expect(states(fake).late).toBe('paused'); + }); + + it('an empty keys list changes nothing', async () => { + const fake = await twoKeys(); + await fake.pause({ keys: [] }); + expect(states(fake)).toEqual({ cap: 'queued', note: 'queued' }); + }); + + it('reset() clears the pause state', async () => { + const fake = createFakeNative(); + await fake.pause({}); + await fake.pause({ keys: ['capture'] }); + fake.reset(); + await fake.enqueue(entry({ id: 'cap', key: 'capture' })); + expect(states(fake)).toEqual({ cap: 'queued' }); + }); }); it('settles a live entry cancelled on cancel(), as native does', async () => { @@ -381,7 +458,7 @@ describe('reset', () => { fake.onSettled(settled); await fake.enqueue(entry()); fake.configure({}); - await fake.pause(); + await fake.pause({}); const before = fake.buildSettled('a', { kind: 'completed' }).eventId; fake.failNext('cancel', 'E_STORAGE'); fake.reset(); @@ -391,8 +468,8 @@ describe('reset', () => { expect(fake.ackedEventIds).toEqual([]); expect(fake.calls).toEqual({ configure: [], - pause: 0, - resume: 0, + pause: [], + resume: [], cancel: [], setWifiOnly: [], updateHeaders: [], diff --git a/src/__typetests__/define.ts b/src/__typetests__/define.ts index 2240a258..dc8b63eb 100644 --- a/src/__typetests__/define.ts +++ b/src/__typetests__/define.ts @@ -277,3 +277,23 @@ client.addListener('attempt', (e) => { // @ts-expect-error an attempt event carries no cancelReason void e.cancelReason; }); + +// pause and resume take an optional scope of keys. +void client.pause(); +void client.pause({ keys: ['capture.upload'] }); +void client.resume({}); +// @ts-expect-error keys is a list of strings +void client.pause({ keys: 'capture.upload' }); +// @ts-expect-error the scope field is keys, not key +void client.resume({ key: ['capture.upload'] }); + +// wifiOnly on a descriptor is a boolean. +client.define({ + key: 'wifi.only', + request: (_vars: null) => ({ url: 'https://x', wifiOnly: true }), +}); +client.define({ + key: 'wifi.only.bad', + // @ts-expect-error wifiOnly is a boolean + request: (_vars: null) => ({ url: 'https://x', wifiOnly: 'yes' }), +}); diff --git a/src/index.ts b/src/index.ts index 1733deec..e7e8d279 100644 --- a/src/index.ts +++ b/src/index.ts @@ -19,6 +19,7 @@ import { import type { AddListener, ConfigureOptions, + PauseScope, RequestRow, StateEvent, UploadClient, @@ -38,6 +39,42 @@ export type UploadClientOptions = { native?: UploadNative; }; +/** + * The scope as it crosses to native. Codegen has no optional arguments, so + * the whole queue is `{}`. A misspelled field throws, because dropping it + * would pause the whole queue. Throws inside the async caller, so the + * promise rejects. + */ +const toNativeScope = ( + method: 'pause' | 'resume', + scope: PauseScope | undefined, +): PauseScope => { + if (scope === undefined) { + return {}; + } + if (typeof scope !== 'object' || scope === null || Array.isArray(scope)) { + throw new Error(`${method}: scope must be an object when given`); + } + Object.keys(scope).forEach((field) => { + if (field !== 'keys') { + throw new Error(`${method}: unknown scope field "${field}"`); + } + }); + if (!('keys' in scope)) { + return {}; + } + // A present but undefined list is refused too: it is more often a + // computed list gone missing than a request for the whole queue. + const { keys } = scope; + if ( + !Array.isArray(keys) || + !keys.every((key) => typeof key === 'string' && key.length > 0) + ) { + throw new Error(`${method}: keys must be an array of non-empty strings`); + } + return { keys: [...keys] }; +}; + /** * Builds one client over the native queue. Each client has its own * definitions and settings. An app needs one; the default export is one. @@ -123,13 +160,20 @@ export const createUploadClient = ({ }; /** - * Pauses the whole queue. No outcome is produced; live rows show 'paused'. - * A paused entry past its expiresAt settles error/expired at resume. + * No scope pauses the whole queue. `{ keys }` pauses the entries of those + * keys, queued and future. No outcome is produced; live rows show + * 'paused'. A paused entry past its expiresAt settles error/expired when + * it resumes. */ - const pause = (): Promise => native.pause(); + const pause = async (scope?: PauseScope): Promise => + native.pause(toNativeScope('pause', scope)); - /** Resumes a paused queue. */ - const resume = (): Promise => native.resume(); + /** + * Undoes the pause of the same scope. An entry stays paused while another + * scope still pauses it: the whole-queue pause, or its key. + */ + const resume = async (scope?: PauseScope): Promise => + native.resume(toNativeScope('resume', scope)); /** * On a live entry: settles it 'cancelled' with reason 'user', then forgets @@ -138,7 +182,10 @@ export const createUploadClient = ({ */ const cancel = (id: string): Promise => native.cancel(id); - /** Persisted natively. Applies to queued and future entries. */ + /** + * Persisted natively. Applies to queued and future entries whose + * descriptor does not set `wifiOnly`. + */ const setWifiOnly = (enabled: boolean): Promise => native.setWifiOnly(enabled); diff --git a/src/registry.ts b/src/registry.ts index 405be906..ceb5471d 100644 --- a/src/registry.ts +++ b/src/registry.ts @@ -39,6 +39,7 @@ const DESCRIPTOR_KEYS = [ 'accept', 'expiresAt', 'retry', + 'wifiOnly', 'android', ]; const PART_KEYS = ['url', 'headers', 'range']; @@ -81,7 +82,8 @@ export const withTimeout = ( /** * The descriptor as it crosses to native. `dataJson` is `JSON.stringify(data)` - * and replaces `data`. It is absent for a bodiless request. + * and replaces `data`. It is absent for a bodiless request. `wifiOnly` crosses + * only when the definition set it; absent means "follow setWifiOnly()". */ export type NativeDescriptor = Omit & { dataJson?: string; @@ -478,6 +480,9 @@ export const validateDescriptor = (descriptor: unknown): NativeDescriptor => { if (d.android !== undefined) { validateAndroid(d.android); } + if (d.wifiOnly !== undefined && typeof d.wifiOnly !== 'boolean') { + throw new Error('mutate: wifiOnly must be a boolean when present'); + } if ( d.expiresAt !== undefined && (!Number.isFinite(d.expiresAt) || d.expiresAt <= 0) diff --git a/src/testing.ts b/src/testing.ts index 5b53fedd..fe255db9 100644 --- a/src/testing.ts +++ b/src/testing.ts @@ -11,6 +11,9 @@ * settle() journals an outcome and emits it, and ackEvents() forgets a * completed or cancelled row. It emits a 'state' row for each of these * transitions, as native does. It never sends HTTP and never retries. + * pause() and resume() keep the whole-queue gate and the paused-keys set, + * and move rows between 'paused' and 'queued' as native does. It records + * setWifiOnly() and a descriptor's wifiOnly but has no network to gate. * * The fake does not model the same-id rules of enqueue() (resume, replace, * E_RUNNING, re-emit). A second enqueue() on an id only overwrites its row. @@ -27,6 +30,7 @@ import type { ErrorKind, Json, OutcomeError, + PauseScope, ProgressEvent, RawResponse, RequestDescriptor, @@ -98,8 +102,9 @@ export type FakeNative = Spec & { /** The arguments of the calls that have no other record. */ readonly calls: { configure: object[]; - pause: number; - resume: number; + /** The scope of each pause() call, as it crossed to native. */ + pause: PauseScope[]; + resume: PauseScope[]; cancel: string[]; setWifiOnly: boolean[]; updateHeaders: Record[]; @@ -154,7 +159,8 @@ export type FakeNative = Spec & { */ failNext: (method: FailableMethod, code: string, message?: string) => void; /** - * Clears entries, rows, journal, calls, acks and queued failures. A + * Clears entries, rows, journal, calls, acks, the pause state and queued + * failures. A * settle() that still waits for its ack resolves now, so it cannot reject * later inside another test. Keeps the subscriptions, because a client subscribes one time, at its first * configure(). New eventIds never repeat old ones, so delivery's dedupe @@ -246,12 +252,17 @@ export const createFakeNative = ( const ackedEventIds: string[] = []; const calls: FakeNative['calls'] = { configure: [], - pause: 0, - resume: 0, + pause: [], + resume: [], cancel: [], setWifiOnly: [], updateHeaders: [], }; + // The pause model: the whole-queue gate and the paused-keys set. + let pausedAll = false; + const pausedKeys = new Set(); + const isPaused = (key: string): boolean => + pausedAll || pausedKeys.has(key); const failures = new Map(); const ackWaiters = new Map void>(); @@ -413,7 +424,7 @@ export const createFakeNative = ( id: raw.id, key: raw.key, vars, - state: 'queued', + state: isPaused(raw.key) ? 'paused' : 'queued', bytesSent: 0, totalBytes: prior?.totalBytes ?? 0, attempts: 0, @@ -422,21 +433,38 @@ export const createFakeNative = ( return raw.id; }), - pause: () => + pause: (input) => run('pause', () => { - calls.pause += 1; + const scope = input as PauseScope; + calls.pause.push({ ...scope }); + if (scope.keys === undefined) { + pausedAll = true; + } else { + scope.keys.forEach((key) => pausedKeys.add(key)); + } rows.forEach((row) => { - if (LIVE_STATES.includes(row.state) && row.state !== 'paused') { + if ( + LIVE_STATES.includes(row.state) && + row.state !== 'paused' && + isPaused(row.key) + ) { setRow({ ...row, state: 'paused', updatedAt: Date.now() }); } }); }), - resume: () => + // A row stays paused while the other scope still pauses it. + resume: (input) => run('resume', () => { - calls.resume += 1; + const scope = input as PauseScope; + calls.resume.push({ ...scope }); + if (scope.keys === undefined) { + pausedAll = false; + } else { + scope.keys.forEach((key) => pausedKeys.delete(key)); + } rows.forEach((row) => { - if (row.state === 'paused') { + if (row.state === 'paused' && !isPaused(row.key)) { setRow({ ...row, state: 'queued', updatedAt: Date.now() }); } }); @@ -566,8 +594,10 @@ export const createFakeNative = ( journal.length = 0; ackedEventIds.length = 0; calls.configure.length = 0; - calls.pause = 0; - calls.resume = 0; + calls.pause.length = 0; + calls.resume.length = 0; + pausedAll = false; + pausedKeys.clear(); calls.cancel.length = 0; calls.setWifiOnly.length = 0; calls.updateHeaders.length = 0; diff --git a/src/types.ts b/src/types.ts index 045d4b59..870f3bd8 100644 --- a/src/types.ts +++ b/src/types.ts @@ -99,6 +99,13 @@ export type RequestDescriptor = { /** Epoch ms. Default now + `lifetimeMs`. */ expiresAt?: number; retry?: Partial; + /** + * Wait for Wi-Fi before each attempt. When set, it overrides + * `setWifiOnly()` for this entry. When absent, the entry follows + * `setWifiOnly()`, so a later toggle moves it too. Persisted with the + * entry. + */ + wifiOnly?: boolean; android?: { noNotification?: boolean }; }; @@ -333,11 +340,19 @@ export interface AddListener { (event: 'attempt', listener: (e: AttemptEvent) => void): EventSubscription; } +/** + * What `pause()` and `resume()` act on. No `keys` field means the whole + * queue. `keys` means the entries of those definition keys, queued and + * future; an empty list changes nothing. `keys: undefined` is refused, so a + * missing list cannot pause the whole queue by accident. + */ +export type PauseScope = { keys?: string[] }; + export type UploadClient = { configure: (options: ConfigureOptions) => void; define: Define; - pause: () => Promise; - resume: () => Promise; + pause: (scope?: PauseScope) => Promise; + resume: (scope?: PauseScope) => Promise; cancel: (id: string) => Promise; setWifiOnly: (enabled: boolean) => Promise; updateHeaders: (patch: Record) => Promise;