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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 17 additions & 3 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
42 changes: 32 additions & 10 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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<void>` and `resume(): Promise<void>`
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<void>` and `resume(scope?): Promise<void>`
`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<void>`
A live entry settles `cancelled` with reason `user` and is forgotten after
Expand All @@ -380,7 +398,10 @@ whose entry save failed is already in effect: the work stops and the
finishes it.

### `setWifiOnly(enabled): Promise<void>`
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<void>`
Merges the patch into the headers of every entry not yet forgotten and
Expand Down Expand Up @@ -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:

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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<String>? {
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<String, String> =
parseHeaderMap(patch).also { requireValidHeaders(it, "updateHeaders") }
Expand Down
15 changes: 8 additions & 7 deletions android/src/main/java/ai/openspace/backgroundupload/EntryRun.kt
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -96,29 +96,29 @@ 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
val parts = incoming?.let { EnqueueRules.adoptedParts(action.manifest, it) }
// 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!!
// Different parts: a present caller file wins; the old blob is the fallback.
val ownedBlob = if (old.body?.kind == StagedBody.CHUNKED) store.bodyFile(old) else null
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))
}
}
}
Expand Down Expand Up @@ -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<String>? = 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<String>? = 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.
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand All @@ -366,7 +382,7 @@ class QueueController(
val changed = mutableListOf<QueueEntry>()
val toForget = mutableListOf<String>()
val toSchedule = mutableListOf<QueueEntry>()
val paused = settings.load().paused
val s = settings.load()
store.locked {
val records = journal.unacknowledged().groupBy { it.id }
for (e in store.all()) {
Expand 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,8 @@ data class Descriptor(
val accept: List<UploadOutcome.AcceptRule>,
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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String> = 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.
Expand Down Expand Up @@ -77,14 +87,15 @@ 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()
val r = s.retry
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,
Expand Down
Loading
Loading