diff --git a/CHANGELOG.md b/CHANGELOG.md index d5df7d5..bfe4c4d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,125 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +Fixes 3 issues found by a full `lib/` audit (2026-09-23), reviewed in detail +in the linked commits. Ships as **1.8.3** (1.8.2 is already live on pub.dev +and immutable). + +### Fixed + +- **`moveToSharedStorage()`'s `subDir` had no path-traversal check on either + platform** — the one genuine cross-platform bypass in this batch. Neither + native implementation ever checked it: Android does + `File(publicDir, config.subDir)` directly, iOS does + `docsURL.appendingPathComponent(subDir)` directly, which does **not** + resolve `..` safely. `subDir` feeds `MediaStore.RELATIVE_PATH` on Android + and a Documents subfolder on iOS, so + `moveToSharedStorage(sourcePath: p, subDir: '../../OtherApp/Camera')` + reached native and would have actually escaped the sandbox on both + platforms. Everything else validated in the same commit + (`multiUpload()`'s url/files, `webSocket()`'s url/storeResponseAt, + `ParallelHttpUploadWorker`'s constructor) is Dart-side defense-in-depth, + not a closed system-level hole — native's `SecurityValidator` + (Kotlin/Swift, at parity) already independently validated those; see the + commit for the exact per-worker breakdown. +- **A `DartWorker` placed in a `TaskGraph` node or a `RemoteTriggerRule` + mapping never reached native with a resolved callback handle** — + `enqueue()` and task chains converted `DartWorker` to `DartWorkerInternal` + (resolving the native callback handle) before sending it; `TaskNode` and + `RemoteTriggerRule` never did, so the task was enqueued but its callback + could never be resolved and the task silently never ran. Confirmed + reachable on both platforms before fixing (neither native worker factory + special-cases `DartCallbackWorker`; both require `callbackHandle` and fail + cleanly without it). + - A chain-step or `TaskGraph`-node `DartWorker` now gets the same + `isHeavyTask: true` promotion plain `enqueue()` applies on iOS, instead of + silently skipping it. **Correction, checked after first writing this + entry:** this is not currently an observable behavior change. `isHeavyTask` + only ever affects scheduling on iOS for the `periodic` and `windowed` + trigger types (both routed through `BGTaskScheduler`); chain steps and + graph nodes have no `TaskTrigger` of their own and always execute inline + via a plain `Task {}`, never through `BGTaskScheduler`, so the promoted + value is stored correctly but not read anywhere today. Kept for + consistency with `enqueue()`'s existing behavior and so a future fix to + chain/graph iOS scheduling doesn't need its own audit of this — not + claimed as a fix for OS-kill risk on either mechanism today. + - **Behavior change:** an unregistered `DartWorker` (`callbackId` not + passed to `initialize(dartWorkers:)`) used in a `TaskGraph` node or a + `RemoteTriggerRule` mapping now throws a `StateError` at the point of + `enqueueGraph()`/`registerRemoteTrigger()`, instead of silently reaching + native and never running. + - **Bonus fix, found while re-checking this change:** the same unregistered + `callbackId` used in a task **chain** used to throw too, but with the + wrong message — chains never checked registration at all, so they fell + straight into the "should never happen" internal-error branch meant for + a genuinely impossible state, printing `INTERNAL ERROR: Callback handle + not found... Please report this bug.` for what is actually a completely + ordinary mistake. Chains now get the same clear "not registered, here's + how to fix it" message `enqueue()` has always given. +- **iOS's foreground `handleEnqueue` never read `existingPolicy` at all** — + found while investigating the item above. Every repeat `enqueue()` call + for a reused `taskId` silently started a second, fully independent + concurrent execution, regardless of what policy the caller asked for. + Confirmed on a simulator: two `DartWorker` executions of one `taskId`, + 600ms apart, both ran to full completion independently. + `existingPolicy: .replace` (the default) now actually cancels the outgoing + execution before starting the new one; `.keep` now actually leaves the + running execution alone and ignores the new request — both matching + Android's WorkManager semantics. + - Fixing this alone reproduced **issue #72's exact bug shape on iOS**: the + replacement execution resolves almost instantly (same running engine, no + boot delay), sees the outgoing execution's cancellation mark, and its own + cleanup — previously keyed by bare `taskId` — cleared that mark before + the outgoing execution's next poll could observe it. So this also ports + issue #72's fix to iOS: `DartTaskCancellationRegistry` is now keyed by a + fresh per-execution id (minted in `executeDartWorkerViaMethodChannel`, + covering both the foreground and headless paths since they share that + function), not by bare `taskId`, mirroring the Android fix. `isTaskCancelled()` + now binds this id into a Zone on the foreground/simulator path too + (`method_channel.dart`'s `_executeDartCallback`), matching what the + headless isolate's `_callbackDispatcher` already did. + - Known residual limitation, traced through but deliberately not closed + (closing it means threading a pre-minted execution id through every + caller of `executeDartWorkerViaMethodChannel` — direct enqueue, chains, + `TaskGraph`, `BGTaskScheduler`-resumed tasks, offline queue — more surface + area than this PR's blast radius should grow to without its own device + verification pass): the per-execution id is minted lazily, inside + `executeDartWorkerViaMethodChannel`, once the replacement `Task` actually + starts running — not synchronously when `handleEnqueue` decides to + replace. In the narrow window between a `.replace` swapping + `activeTasks[taskId]` to the new `Task` and that `Task` reaching its + first `beginExecution` call, `DartTaskCancellationRegistry`'s + `currentExecutionId[taskId]` still points at the OLD (already-replaced) + execution. A `cancel(taskId)` — or another `.replace` — landing in that + window marks the wrong (stale) execution id; the new one starts moments + later unaffected by that mark, so it does **not** stop when the caller + thought it just told it to. Pure double/triple-replace with no + intervening `cancel()` was traced through and does not corrupt state — + only an explicit cancel landing in that specific gap does. The window is + on the order of the time from `Task { }` construction to its first + `await` inside `executeDartWorkerViaMethodChannel` — real, but requires a + caller to `enqueue()`-then-immediately-`cancel()`/`enqueue()` again the + same `taskId` back-to-back, not a pattern normal usage hits. + - Device-verified on an iOS simulator (both `.replace` and `.keep`); no + Android changes were needed (Android's `existingPolicy` handling and + `DartTaskCancellationRegistry` were already correct — that's what issue + #72 fixed). + - **Follow-up found in second-pass review, before this ever shipped**: the + `existingPolicy` fix above read `activeTasks[taskId]` as its "is this + still running" signal, but that dictionary was never cleared when a + direct one-time task finished *naturally* (only explicit cancel ever + removed an entry) — a leftover from before anything read it as a + liveness signal. Confirmed on a simulator: re-enqueuing a `taskId` whose + task had already completed, with `existingPolicy: .keep`, was silently + dropped forever, because the stale entry made `.keep` think something + was still running. Fixed with a per-enqueue generation id that lets a + task's own completion clear its entry — but only if nothing has replaced + it in the meantime, the same guard pattern used by + `DartTaskCancellationRegistry`. Verified fixed on the same simulator, and + the full `Cancellation` device-test group (8 tests) still passes. + ## [1.8.2] - 2026-09-23 ### Added diff --git a/example/integration_test/device_integration_test.dart b/example/integration_test/device_integration_test.dart index 1e63d2e..d5db719 100644 --- a/example/integration_test/device_integration_test.dart +++ b/example/integration_test/device_integration_test.dart @@ -2070,6 +2070,207 @@ void main() { }, ); + testWidgets( + 'lib_audit_3: existingPolicy.replace cancels the outgoing execution ' + 'precisely and the replacement runs to completion, not inheriting its ' + 'stale cancel mark (iOS)', + (tester) async { + // https://github.com/brewkits/native_workmanager — found by the + // 2026-09-23 lib/ audit while investigating the iOS analogue of + // issue_72 above. Two things were wrong before this fix, discovered + // in order: + // 1. iOS's foreground handleEnqueue never read existingPolicy at + // all — every repeat enqueue() of one taskId silently started a + // second, fully independent concurrent execution, confirmed by + // running exactly this test's shape pre-fix: both old and new + // ran to full completion (50/50), regardless of policy. + // 2. Fixing #1 alone (mark-and-replace on the outgoing execution) + // reproduced issue_72's "direction 1" bug on iOS: the new + // execution resolves almost instantly (same running engine, no + // boot delay), sees the mark, and its own cleanup — keyed by + // bare taskId — cleared it before the outgoing execution's next + // poll could observe it. Confirmed the same way: the new + // execution died at iteration 1, the old one ran all 50. + // The real fix needed both: existingPolicy.replace AND per-execution + // (executionId-keyed) cancellation tracking, mirroring issue #72's + // Android fix, ported to iOS's DartTaskCancellationRegistry. + if (!Platform.isIOS) { + markTestSkipped('iOS foreground existingPolicy path'); + return; + } + + final id = _id('lib_audit_3_replace'); + final oldCounterFile = + File('${tmpDir.path}/lib_audit_3_replace_old.txt'); + final newCounterFile = + File('${tmpDir.path}/lib_audit_3_replace_new.txt'); + + await NativeWorkManager.enqueue( + taskId: id, + trigger: const TaskTrigger.oneTime(), + worker: DartWorker( + callbackId: 'dit_cancel_poll', + input: {'counterFile': oldCounterFile.path}, + ), + ); + + await Future.delayed(const Duration(milliseconds: 600)); + + // Default policy is replace. + await NativeWorkManager.enqueue( + taskId: id, + trigger: const TaskTrigger.oneTime(), + worker: DartWorker( + callbackId: 'dit_cancel_poll', + input: {'counterFile': newCounterFile.path}, + ), + ); + + await Future.delayed(const Duration(seconds: 12)); + + expect( + oldCounterFile.existsSync(), + isTrue, + reason: 'lib_audit_3: the replaced execution must have started', + ); + final oldIterations = + int.parse(oldCounterFile.readAsStringSync().trim()); + expect( + oldIterations, + lessThan(50), + reason: 'lib_audit_3: the replaced execution must have observed ' + 'cancellation and stopped', + ); + + expect( + newCounterFile.existsSync(), + isTrue, + reason: 'lib_audit_3: the replacement execution must have started', + ); + final newIterations = + int.parse(newCounterFile.readAsStringSync().trim()); + expect( + newIterations, + equals(50), + reason: 'lib_audit_3: the replacement must run to completion — a ' + 'lower count means it inherited the replaced execution\'s ' + 'stale cancellation mark and self-aborted', + ); + }, + ); + + testWidgets( + 'lib_audit_3: existingPolicy.keep leaves the running execution alone ' + 'and ignores the new request (iOS)', + (tester) async { + if (!Platform.isIOS) { + markTestSkipped('iOS foreground existingPolicy path'); + return; + } + + final id = _id('lib_audit_3_keep'); + final oldCounterFile = File('${tmpDir.path}/lib_audit_3_keep_old.txt'); + final newCounterFile = File('${tmpDir.path}/lib_audit_3_keep_new.txt'); + + await NativeWorkManager.enqueue( + taskId: id, + trigger: const TaskTrigger.oneTime(), + worker: DartWorker( + callbackId: 'dit_cancel_poll', + input: {'counterFile': oldCounterFile.path}, + ), + ); + + await Future.delayed(const Duration(milliseconds: 600)); + + await NativeWorkManager.enqueue( + taskId: id, + trigger: const TaskTrigger.oneTime(), + worker: DartWorker( + callbackId: 'dit_cancel_poll', + input: {'counterFile': newCounterFile.path}, + ), + existingPolicy: ExistingTaskPolicy.keep, + ); + + await Future.delayed(const Duration(seconds: 12)); + + expect( + oldCounterFile.existsSync(), + isTrue, + reason: 'lib_audit_3: keep must leave the running execution alone', + ); + expect( + int.parse(oldCounterFile.readAsStringSync().trim()), + equals(50), + reason: 'lib_audit_3: keep must not cancel the running execution', + ); + expect( + newCounterFile.existsSync(), + isFalse, + reason: 'lib_audit_3: keep must ignore the new request entirely — ' + 'no second execution should ever have started', + ); + }, + ); + + testWidgets( + 'lib_audit_3: re-enqueuing a taskId whose previous execution already ' + 'completed runs the new one, regardless of existingPolicy (iOS)', + (tester) async { + // Found while re-verifying the two tests above: iOS's activeTasks + // dict is never cleared when a direct one-time task finishes + // NATURALLY (only explicit cancel paths ever call removeValue). + // Before existingPolicy read that dict as a liveness signal this was + // a harmless leak; the first version of the existingPolicy fix + // treated a long-finished taskId as "still running" forever, so + // ANY later re-enqueue with existingPolicy.keep was silently + // dropped — confirmed with this exact test shape before the second + // fix (a per-enqueue generation id, cleared only by the Task that + // is still the current occupant of activeTasks[taskId] when it + // finishes — mirrors DartTaskCancellationRegistry's "clear only if + // still current" guard). + if (!Platform.isIOS) { + markTestSkipped('iOS foreground existingPolicy path'); + return; + } + + final id = _id('lib_audit_3_keep_after_completion'); + + final firstEvent = _waitEvent(id, timeout: const Duration(seconds: 15)); + await NativeWorkManager.enqueue( + taskId: id, + trigger: const TaskTrigger.oneTime(), + worker: DartWorker(callbackId: 'dit_pass'), + ); + final first = await firstEvent; + expect(first?.success, isTrue, + reason: 'lib_audit_3: the first execution must complete'); + + // Give the natural-completion cleanup a moment to run before + // re-enqueuing, so this genuinely exercises the "already finished, + // not just finishing" case. + await Future.delayed(const Duration(seconds: 2)); + + final secondEvent = + _waitEvent(id, timeout: const Duration(seconds: 15)); + await NativeWorkManager.enqueue( + taskId: id, + trigger: const TaskTrigger.oneTime(), + worker: DartWorker(callbackId: 'dit_pass'), + existingPolicy: ExistingTaskPolicy.keep, + ); + final second = await secondEvent; + expect( + second?.success, + isTrue, + reason: 'lib_audit_3: a taskId reused after its previous execution ' + 'already completed must run the new request — keep must not ' + 'mistake a long-finished task for one still running', + ); + }, + ); + testWidgets( 'issue_69: cancelling a background-session download actually aborts the transfer (iOS)', (tester) async { diff --git a/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+Cancel.swift b/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+Cancel.swift index a4426a6..dc0221d 100644 --- a/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+Cancel.swift +++ b/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+Cancel.swift @@ -16,6 +16,7 @@ extension NativeWorkmanagerPlugin { activeTasks.keys.forEach { DartTaskCancellationRegistry.shared.markCancelled($0) } activeTasks.values.forEach { $0.cancel() } activeTasks.removeAll() + activeTaskGenerations.removeAll() // see its doc comment on NativeWorkmanagerPlugin taskStates.removeAll() taskTags.removeAll() workers.values.forEach { $0.stop() } @@ -53,6 +54,7 @@ extension NativeWorkmanagerPlugin { DartTaskCancellationRegistry.shared.markCancelled(taskId) // issue #66 activeTasks[taskId]?.cancel() activeTasks.removeValue(forKey: taskId) + activeTaskGenerations.removeValue(forKey: taskId) // see its doc comment taskStates[taskId] = .cancelled taskTags.removeValue(forKey: taskId) workers[taskId]?.stop() @@ -74,6 +76,7 @@ extension NativeWorkmanagerPlugin { stateQueue.async(flags: .barrier) { self.activeTasks[taskId]?.cancel() self.activeTasks.removeValue(forKey: taskId) + self.activeTaskGenerations.removeValue(forKey: taskId) // see its doc comment self.taskStates[taskId] = .cancelled self.taskTags.removeValue(forKey: taskId) self.workers[taskId]?.stop() diff --git a/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+Execution.swift b/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+Execution.swift index c73df79..3917d75 100644 --- a/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+Execution.swift +++ b/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+Execution.swift @@ -578,24 +578,36 @@ extension NativeWorkmanagerPlugin { workerConfig: [String: Any], taskId: String ) async -> WorkerResult { - // Issue #66: whichever branch below runs, always drop this taskId's - // cancellation-registry entry once execution is done — otherwise a - // cancelled taskId (or, worse, a reused one on a later run) leaks or - // misreports "cancelled" forever. - defer { DartTaskCancellationRegistry.shared.clear(taskId) } - // Issue #75: clear() above already drops the stop notifier, but be - // explicit — an early `return` on a config error below must not leave a - // notifier behind that a later cancel of the same taskId would fire. - defer { DartTaskCancellationRegistry.shared.clearStopNotifier(taskId) } + // Issue #72: mint a fresh executionId for THIS invocation, distinct + // from taskId, mirroring Android's DartCallbackWorker (which mints + // one per doWork() call). handleEnqueue's existingPolicy: .replace + // can cancel an outgoing execution and start a new one under the SAME + // taskId — a bare taskId-keyed cancellation mark cannot tell the two + // apart, so this is what lets isTaskCancelled() resolve against the + // specific execution polling it rather than whichever one happens to + // share its taskId. Registered immediately so a replace racing in + // concurrently always has something to resolve to. + let executionId = UUID().uuidString + DartTaskCancellationRegistry.shared.beginExecution(executionId, taskId: taskId) + // Issue #66/#72: whichever branch below runs, always drop this + // execution's cancellation-registry entry once it's done — otherwise + // a cancelled execution (or, worse, a reused taskId on a later run) + // leaks or misreports "cancelled" forever. Scoped to executionId, not + // taskId, so a slow-finishing outgoing execution's cleanup can never + // wipe a replacing execution's tracking or stop notifier out from + // under it (endExecution only clears taskId-level state if it still + // points at this executionId) — see DartTaskCancellationRegistry. + defer { DartTaskCancellationRegistry.shared.endExecution(executionId, taskId: taskId) } guard let callbackId = workerConfig["callbackId"] as? String else { return WorkerResult.failure(message: "DartCallbackWorker: missing callbackId in config") } - // Inject __taskId into the input JSON so the Dart callback can call - // NativeWorkManager.reportDartWorkerProgress(). The Dart side only receives - // the inner "input" string — mirror Android's DartCallbackWorker, which merges - // the outer __taskId into that inner object before forwarding to Dart. - let input = Self.mergeTaskId(into: workerConfig["input"] as? String, taskId: taskId) + // Inject __taskId (progress reporting) and __executionId (issue #72 + // cancellation precision) into the input JSON. The Dart side only + // receives the inner "input" string — mirrors Android's + // DartCallbackWorker, which merges both into that inner object before + // forwarding to Dart. + let input = Self.mergeTaskId(into: workerConfig["input"] as? String, taskId: taskId, executionId: executionId) // Honor user-configured DartWorker.timeoutMs across both execution paths // (foreground main-channel and killed-app FlutterEngineManager fallback). @@ -709,11 +721,12 @@ extension NativeWorkmanagerPlugin { } } - /// Merge `__taskId` into a DartWorker input JSON string so the callback can - /// report progress. Returns a JSON object string. Mirrors Android's - /// DartCallbackWorker input enrichment; falls back to the original string if - /// the input is a non-object JSON that cannot carry the key. - static func mergeTaskId(into inputJson: String?, taskId: String) -> String? { + /// Merge `__taskId` (progress reporting) and, when given, `__executionId` + /// (issue #72 cancellation precision) into a DartWorker input JSON + /// string. Returns a JSON object string. Mirrors Android's + /// DartCallbackWorker input enrichment; falls back to the original string + /// if the input is a non-object JSON that cannot carry the keys. + static func mergeTaskId(into inputJson: String?, taskId: String, executionId: String? = nil) -> String? { var obj: [String: Any] = [:] if let json = inputJson, !json.isEmpty, json != "null" { if let data = json.data(using: .utf8), @@ -725,6 +738,9 @@ extension NativeWorkmanagerPlugin { } } obj["__taskId"] = taskId + if let executionId { + obj["__executionId"] = executionId + } guard let merged = try? JSONSerialization.data(withJSONObject: obj), let mergedString = String(data: merged, encoding: .utf8) else { return inputJson diff --git a/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+StreamHandlers.swift b/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+StreamHandlers.swift index 0d65742..45cf09e 100644 --- a/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+StreamHandlers.swift +++ b/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin+StreamHandlers.swift @@ -98,6 +98,7 @@ extension NativeWorkmanagerPlugin: UNUserNotificationCenterDelegate { stateQueue.async(flags: .barrier) { self.activeTasks[taskId]?.cancel() self.activeTasks.removeValue(forKey: taskId) + self.activeTaskGenerations.removeValue(forKey: taskId) // see its doc comment self.taskStates[taskId] = .paused } taskStore?.updateStatus(taskId: taskId, status: "paused") @@ -113,6 +114,7 @@ extension NativeWorkmanagerPlugin: UNUserNotificationCenterDelegate { stateQueue.async(flags: .barrier) { self.activeTasks[taskId]?.cancel() self.activeTasks.removeValue(forKey: taskId) + self.activeTaskGenerations.removeValue(forKey: taskId) // see its doc comment self.taskStates[taskId] = .cancelled self.taskNotifTitles.removeValue(forKey: taskId) self.taskAllowPause.removeValue(forKey: taskId) diff --git a/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin.swift b/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin.swift index 0ebd8f3..9508596 100644 --- a/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin.swift +++ b/ios/native_workmanager/Sources/native_workmanager/NativeWorkmanagerPlugin.swift @@ -83,6 +83,25 @@ public class NativeWorkmanagerPlugin: NSObject, FlutterPlugin { var activeTasks: [String: Task] = [:] var workers: [String: IosWorker] = [:] + /// taskId -> a fresh id minted each time `handleEnqueue`'s direct + /// (one-time) path stores a new entry in `activeTasks`. `activeTasks` + /// itself is never cleared when a direct task finishes NATURALLY + /// (success, failure, or timeout) — only explicit cancel paths + /// (handleCancel/cancelAll/cancelByTag/notification-cancel) ever call + /// `removeValue`. Before existingPolicy was implemented that was a + /// harmless leak (a `Task` value is cheap and nothing read the dict as a + /// liveness signal). It is not harmless now: `existingPolicy` reads + /// `activeTasks[taskId] != nil` to decide whether a taskId is "still + /// running". Confirmed on a simulator (2026-09-23 lib/ audit follow-up): + /// re-enqueuing a taskId whose task had already completed, with + /// `existingPolicy: .keep`, was silently dropped forever — `.keep` saw + /// the stale entry and concluded something was still running. This id + /// lets the completing Task's own cleanup remove its `activeTasks` entry + /// exactly once it is truly done, while guarding against a replacing + /// execution's entry being wiped out from under it (same "clear only if + /// still current" pattern as `DartTaskCancellationRegistry.endExecution`). + var activeTaskGenerations: [String: UUID] = [:] + @available(iOS 13.0, *) var chainStateManager: ChainStateManager { ChainStateManager.shared } @@ -134,7 +153,15 @@ public class NativeWorkmanagerPlugin: NSObject, FlutterPlugin { // headless FlutterEngineManager engine — see DartTaskCancellationRegistry). case "isTaskCancelled": let taskId = args?["taskId"] as? String ?? "" - result(DartTaskCancellationRegistry.shared.isCancelled(taskId)) + // Issue #72: precise per-execution check when the Dart side's + // Zone had an executionId to send (see method_channel.dart's + // _executeDartCallback); falls back to the coarse taskId + // check otherwise. + if let executionId = args?["executionId"] as? String { + result(DartTaskCancellationRegistry.shared.isCancelled(executionId: executionId)) + } else { + result(DartTaskCancellationRegistry.shared.isCancelled(taskId)) + } default: result(FlutterMethodNotImplemented) } @@ -362,21 +389,76 @@ public class NativeWorkmanagerPlugin: NSObject, FlutterPlugin { let directQos = (directConstraintsMap?["qos"] as? String) ?? "background" let directRetryConfig = RetryConfig.from(constraintsMap: directConstraintsMap) - let task = Task { [weak self] in - guard let self else { return } - if initialDelayMs > 0 { - try? await Task.sleep(nanoseconds: UInt64(initialDelayMs) * 1_000_000) + // existingPolicy was accepted from Dart but never read here — every repeat + // enqueue() of the same taskId silently started a second, fully independent + // concurrent Task, regardless of what policy the caller asked for, because this + // dictionary write always just clobbered whatever was there. Confirmed on a + // simulator (2026-09-23 lib/ audit): two DartWorker executions of one taskId, + // 600ms apart, both ran to full completion independently. "replace" (the + // default, matching Android and NativeWorkManager.enqueue's own default) now + // stops the outgoing execution the same way handleCancel does before starting + // the new one; "keep" leaves the running execution alone and ignores the new + // request, matching WorkManager's ExistingWorkPolicy.KEEP on Android, which also + // always reports the enqueue as accepted regardless of whether it was a no-op. + // + // This whole check-decide-store sequence is one atomic stateQueue block so two + // overlapping handleEnqueue calls for the same taskId can't both see "nothing + // running yet" and both proceed. + let existingPolicyStr = (args["existingPolicy"] as? String)?.lowercased() ?? "replace" + var skippedForKeep = false + stateQueue.sync(flags: .barrier) { + if let existingTask = self.activeTasks[taskId] { + if existingPolicyStr == "keep" { + skippedForKeep = true + return + } + // Cancelling the Swift Task only unblocks whatever it's synchronously + // awaiting (irrelevant for a DartCallbackWorker, which awaits a method + // channel round-trip, not a cancellable operation). The registry mark is + // what a running DartWorker's isTaskCancelled() poll actually sees — + // same two calls handleCancel makes for an explicit user cancel(). + existingTask.cancel() + DartTaskCancellationRegistry.shared.markCancelled(taskId) + self.workers[taskId]?.stop() } - guard !Task.isCancelled else { return } - await self.executeWorkerSync( - taskId: taskId, - workerClassName: workerClassName, - workerConfig: workerConfig, - qos: directQos, - retryConfig: directRetryConfig - ) + + // See activeTaskGenerations' doc comment: this id is what lets the + // Task below tell, once IT finishes, whether it is still the + // current occupant of activeTasks[taskId] — a naturally-completing + // task must remove its own entry so a later enqueue() doesn't + // mistake a long-finished taskId for one still running. + let generationId = UUID() + let task = Task { [weak self] in + guard let self else { return } + defer { + self.stateQueue.sync(flags: .barrier) { + // Only clear if nothing has replaced us in the meantime — + // a .replace enqueue() racing in after we started but + // before we finish must not have its brand-new entry + // wiped out by our own late cleanup. + guard self.activeTaskGenerations[taskId] == generationId else { return } + self.activeTasks.removeValue(forKey: taskId) + self.activeTaskGenerations.removeValue(forKey: taskId) + } + } + if initialDelayMs > 0 { + try? await Task.sleep(nanoseconds: UInt64(initialDelayMs) * 1_000_000) + } + guard !Task.isCancelled else { return } + await self.executeWorkerSync( + taskId: taskId, + workerClassName: workerClassName, + workerConfig: workerConfig, + qos: directQos, + retryConfig: directRetryConfig + ) + } + self.activeTasks[taskId] = task + self.activeTaskGenerations[taskId] = generationId + } + if skippedForKeep { + NativeLogger.d("handleEnqueue: '\(taskId)' already running, existingPolicy=keep — new request ignored") } - stateQueue.sync(flags: .barrier) { self.activeTasks[taskId] = task } result("ACCEPTED") } @@ -400,6 +482,7 @@ public class NativeWorkmanagerPlugin: NSObject, FlutterPlugin { task.cancel() } activeTasks.removeAll() + activeTaskGenerations.removeAll() // see its doc comment } NativeLogger.w("⚠️ OS Expiration: Stopped all active workers") } @@ -452,6 +535,7 @@ public class NativeWorkmanagerPlugin: NSObject, FlutterPlugin { stateQueue.async(flags: .barrier) { self.activeTasks[taskId]?.cancel() self.activeTasks.removeValue(forKey: taskId) + self.activeTaskGenerations.removeValue(forKey: taskId) // see its doc comment self.taskStates[taskId] = .cancelled self.workers[taskId]?.stop() } diff --git a/ios/native_workmanager/Sources/native_workmanager/engine/DartTaskCancellationRegistry.swift b/ios/native_workmanager/Sources/native_workmanager/engine/DartTaskCancellationRegistry.swift index d5e654d..053a979 100644 --- a/ios/native_workmanager/Sources/native_workmanager/engine/DartTaskCancellationRegistry.swift +++ b/ios/native_workmanager/Sources/native_workmanager/engine/DartTaskCancellationRegistry.swift @@ -1,38 +1,99 @@ import Foundation -/// Tracks which DartWorker task IDs have been cancelled, so a Dart callback -/// running in either Flutter engine (the main isolate, or the headless -/// engine spun up by `FlutterEngineManager`) can ask +/// Tracks which DartWorker **executions** have been cancelled, so a Dart +/// callback running in either Flutter engine (the main isolate, or the +/// headless engine spun up by `FlutterEngineManager`) can ask /// `NativeWorkManager.isTaskCancelled(taskId)` (issue #66) and get a real /// answer instead of always `false`. /// +/// Keyed by **executionId**, not by the app-level `taskId` (issue #72, +/// ported from Android's `DartTaskCancellationRegistry.kt`). A single +/// `taskId` can have more than one execution alive at once: `handleEnqueue`'s +/// `existingPolicy: .replace` (the `enqueue()` default) now cancels the +/// outgoing execution and immediately starts a new one under the SAME +/// `taskId`. A bare `taskId`-keyed mark/clear cannot tell those two +/// executions apart — confirmed on a simulator (2026-09-23 lib/ audit): +/// without this, the new execution's own `clear` wiped the mark meant for the +/// outgoing one, and the outgoing one ran to full completion never having +/// noticed it should stop. +/// /// Marked from every place that cancels a running task — /// `NativeWorkmanagerPlugin`'s `cancel`/`cancelAll`/`cancelByTag` handlers, -/// notification-driven cancel, and `BGTaskSchedulerManager`'s expiration -/// handler. Read from the `dev.brewkits/dart_worker_channel` handlers in -/// both `NativeWorkmanagerPlugin` (main isolate) and `FlutterEngineManager` +/// notification-driven cancel, `BGTaskSchedulerManager`'s expiration handler, +/// and `handleEnqueue`'s replace path — via the taskId-only `markCancelled`, +/// which resolves to whichever executionId is currently live for that taskId +/// (none of those call sites know a specific executionId; only +/// `executeDartWorkerViaMethodChannel`, which mints one per invocation, does). +/// Read from the `dev.brewkits/dart_worker_channel` handlers in both +/// `NativeWorkmanagerPlugin` (main isolate) and `FlutterEngineManager` /// (headless isolate). /// -/// This is **cooperative only**: marking a taskId here does not interrupt -/// whatever the Dart isolate is currently `await`-ing — it only lets a -/// polling callback see the request and return early. +/// This is **cooperative only**: marking an execution here does not +/// interrupt whatever the Dart isolate is currently `await`-ing — it only +/// lets a polling callback see the request and return early. final class DartTaskCancellationRegistry { static let shared = DartTaskCancellationRegistry() private init() {} private let lock = NSLock() - private var cancelled: Set = [] - /// Issue #75: per-task stop notifiers, keyed by taskId. + /// executionId -> taskId. The value is only needed for the coarse + /// taskId-only `isCancelled` fallback (callers with no executionId). + private var cancelled: [String: String] = [:] + + /// taskId -> the executionId of whichever execution is currently "the" + /// live one for that taskId. This is what a bare `markCancelled(taskId)` + /// call (from every call site except `executeDartWorkerViaMethodChannel` + /// itself, none of which know a specific executionId) actually targets. + private var currentExecutionId: [String: String] = [:] + + /// Issue #75: per-task stop notifiers, keyed by taskId (not per-execution + /// — a known, deliberately out-of-scope simplification; see the 2026-09-23 + /// lib/ audit notes. Two overlapping executions of one taskId could still + /// clobber each other's notifier registration, same as before this file's + /// issue #72 rework). /// /// Registered by whichever execution path is actually running the callback /// (main channel vs headless engine), because only that path knows which /// channel to notify on. Hanging this off the registry means every existing /// `markCancelled` call site — explicit cancel, cancelByTag, cancelAll, - /// notification-driven cancel, BGTask expiration — fires the notification - /// for free, instead of seven sites each having to remember to. + /// notification-driven cancel, BGTask expiration, replace — fires the + /// notification for free, instead of each having to remember to. private var stopNotifiers: [String: (Int64?) -> Void] = [:] + // MARK: - Per-execution lifecycle (issue #72) + + /// Record that `executionId` is now the live execution for `taskId`. + /// Called once, right after `executeDartWorkerViaMethodChannel` mints an + /// executionId for a fresh invocation — before the Dart callback starts, + /// so a `markCancelled(taskId)` racing in concurrently always has + /// something to resolve to. + func beginExecution(_ executionId: String, taskId: String) { + lock.lock() + defer { lock.unlock() } + currentExecutionId[taskId] = executionId + } + + /// Drop `executionId`'s tracking once it has finished, successfully or + /// not — otherwise every cancelled execution leaks in `cancelled` forever. + /// + /// Only clears `currentExecutionId[taskId]` (and `taskId`'s stop + /// notifier) if it still points at THIS executionId. Without that guard, + /// a slow-finishing outgoing execution's cleanup could run AFTER a + /// replacing execution has already begun and wipe the newer execution's + /// tracking and notifier out from under it. + func endExecution(_ executionId: String, taskId: String) { + lock.lock() + cancelled.removeValue(forKey: executionId) + if currentExecutionId[taskId] == executionId { + currentExecutionId.removeValue(forKey: taskId) + stopNotifiers.removeValue(forKey: taskId) + } + lock.unlock() + } + + // MARK: - Stop notifiers (issue #75) + /// Register the stop notifier for a running DartWorker execution (issue #75). /// /// `cancelGraceMs` is handed to the notifier rather than stored here: the @@ -54,15 +115,20 @@ final class DartTaskCancellationRegistry { stopNotifiers.removeValue(forKey: taskId) } - /// Record that `taskId` has been cancelled/stopped. + // MARK: - Marking / reading cancellation + + /// Mark `taskId`'s CURRENT execution cancelled/stopped. Used by every + /// cancellation source that only knows a taskId — `cancel`/`cancelAll`/ + /// `cancelByTag`, notification-driven cancel, and BGTask expiration. /// /// Issue #75: also fires that task's stop notifier, exactly once — the - /// notifier is removed as it is taken, so a `cancelAll` that sweeps the same - /// taskId twice, or an explicit cancel racing a BGTask expiration, cannot - /// notify the Dart handler twice. + /// notifier is removed as it is taken, so sweeping the same taskId twice + /// (e.g. an explicit cancel racing a BGTask expiration) cannot notify the + /// Dart handler twice. func markCancelled(_ taskId: String) { lock.lock() - cancelled.insert(taskId) + let executionId = currentExecutionId[taskId] ?? taskId + cancelled[executionId] = taskId let notifier = stopNotifiers.removeValue(forKey: taskId) lock.unlock() @@ -73,19 +139,23 @@ final class DartTaskCancellationRegistry { notifier?(nil) } - /// Whether `taskId` has been marked cancelled. - func isCancelled(_ taskId: String) -> Bool { + /// Whether the execution identified by `executionId` is cancelled. + /// Precise — use this whenever an executionId is available (from + /// `isTaskCancelled`'s Zone-bound `executionId` argument). + func isCancelled(executionId: String) -> Bool { lock.lock() defer { lock.unlock() } - return cancelled.contains(taskId) + return cancelled[executionId] != nil } - /// Remove `taskId`'s entry once its execution has finished, successfully - /// or not — otherwise every cancelled taskId leaks in this set forever. - func clear(_ taskId: String) { + /// Coarse: whether `taskId`'s CURRENT execution is cancelled. Fallback + /// for callers with no executionId in scope. Do not use this when an + /// executionId is available — it cannot distinguish a replaced + /// generation of the same taskId from its replacement. + func isCancelled(_ taskId: String) -> Bool { lock.lock() defer { lock.unlock() } - cancelled.remove(taskId) - stopNotifiers.removeValue(forKey: taskId) + let executionId = currentExecutionId[taskId] ?? taskId + return cancelled[executionId] != nil } } diff --git a/ios/native_workmanager/Sources/native_workmanager/engine/FlutterEngineManager.swift b/ios/native_workmanager/Sources/native_workmanager/engine/FlutterEngineManager.swift index adaa030..0ca973f 100644 --- a/ios/native_workmanager/Sources/native_workmanager/engine/FlutterEngineManager.swift +++ b/ios/native_workmanager/Sources/native_workmanager/engine/FlutterEngineManager.swift @@ -495,8 +495,15 @@ class FlutterEngineManager { } else if call.method == "isTaskCancelled" { // Issue #66: cooperative cancellation poll from inside a // running DartWorker callback. See DartTaskCancellationRegistry. - let taskId = (call.arguments as? [String: Any])?["taskId"] as? String ?? "" - result(DartTaskCancellationRegistry.shared.isCancelled(taskId)) + // Issue #72: precise per-execution check when provided (see + // NativeWorkmanagerPlugin.swift's identical handler). + let args = call.arguments as? [String: Any] + let taskId = args?["taskId"] as? String ?? "" + if let executionId = args?["executionId"] as? String { + result(DartTaskCancellationRegistry.shared.isCancelled(executionId: executionId)) + } else { + result(DartTaskCancellationRegistry.shared.isCancelled(taskId)) + } } else { result(FlutterMethodNotImplemented) } diff --git a/lib/native_workmanager.dart b/lib/native_workmanager.dart index 3e8e1bb..c0fe048 100644 --- a/lib/native_workmanager.dart +++ b/lib/native_workmanager.dart @@ -41,7 +41,7 @@ export 'src/task_id.dart'; export 'src/enqueue_request.dart'; export 'src/events.dart'; export 'src/ios_live_activity_bridge.dart'; -export 'src/native_work_manager.dart'; +export 'src/native_work_manager.dart' hide executionIdZoneKey; export 'src/observability.dart'; export 'src/offline_queue.dart'; export 'src/middleware.dart'; diff --git a/lib/src/method_channel.dart b/lib/src/method_channel.dart index 4e2d647..9018b6e 100644 --- a/lib/src/method_channel.dart +++ b/lib/src/method_channel.dart @@ -9,7 +9,7 @@ import 'battery_restriction.dart'; import 'constraints.dart'; import 'events.dart'; import 'native_work_manager.dart' - show resolveDispatcherTimeout, resolveStopHandlerBudget; + show executionIdZoneKey, resolveDispatcherTimeout, resolveStopHandlerBudget; import 'platform_interface.dart'; import 'remote_trigger.dart'; import 'task_trigger.dart'; @@ -206,17 +206,34 @@ class MethodChannelNativeWorkManager extends NativeWorkManagerPlatform { // native-side BGTask deadline (release on a real device) — in tests it // ran to completion regardless of timeoutMs. Mirrors the dispatcher. final timeoutDuration = resolveDispatcherTimeout(args); - return _callbackExecutor!(callbackId, input).timeout( - timeoutDuration, - onTimeout: () { - developer.log( - '[NativeWorkManager] DartWorker callback "$callbackId" timed out ' - 'after ${timeoutDuration.inSeconds} s on the main method channel. ' - 'Increase DartWorker.timeoutMs or split the work.', - level: 900, - ); - return false; - }, + + // Issue #72 (iOS foreground/simulator — found by the 2026-09-23 lib/ + // audit): bind this invocation's executionId into a Zone around the + // callback call, exactly like the headless isolate's + // `_callbackDispatcher` does. Without this, isTaskCancelled() falls back + // to its coarse by-taskId check on this path, which is exactly the + // instance-identity bug issue #72 fixed for the headless path — a + // DartWorker replaced via `existingPolicy: .replace` (the enqueue() + // default) could have its brand-new, legitimate execution mistaken for + // the outgoing one it replaced, since both share one taskId. iOS's + // `executeDartWorkerViaMethodChannel` mints and injects `__executionId` + // into `input` for this exact reason. + final executionId = input?['__executionId'] as String?; + + return runZoned( + () => _callbackExecutor!(callbackId, input).timeout( + timeoutDuration, + onTimeout: () { + developer.log( + '[NativeWorkManager] DartWorker callback "$callbackId" timed out ' + 'after ${timeoutDuration.inSeconds} s on the main method channel. ' + 'Increase DartWorker.timeoutMs or split the work.', + level: 900, + ); + return false; + }, + ), + zoneValues: {executionIdZoneKey: executionId}, ); } diff --git a/lib/src/native_work_manager.dart b/lib/src/native_work_manager.dart index 1739b52..e173d33 100644 --- a/lib/src/native_work_manager.dart +++ b/lib/src/native_work_manager.dart @@ -137,16 +137,14 @@ Duration resolveStopHandlerBudget(Map args) { /// Issue #72: Zone key `isTaskCancelled` uses to find the executionId of /// whichever DartWorker invocation is currently running, without requiring -/// the app's callback to pass it explicitly. See `_callbackDispatcher`. -const Object _executionIdZoneKey = #nativeWorkmanagerExecutionId; - -/// Test-only accessor for [_executionIdZoneKey] — lets unit tests simulate -/// being inside a dispatched DartWorker callback's Zone without needing to -/// go through the full `_callbackDispatcher` (which is private and driven by -/// a native-supplied callback handle, not something a unit test can invoke -/// directly). Not part of the public API. -@visibleForTesting -const Object executionIdZoneKeyForTesting = _executionIdZoneKey; +/// the app's callback to pass it explicitly. Bound by both the headless +/// isolate's `_callbackDispatcher` and, since the 2026-09-23 lib/ audit's +/// iOS-foreground fix, `method_channel.dart`'s `_executeDartCallback` — not +/// private, so that separate library can use the same key. Also used by unit +/// tests to simulate being inside a dispatched DartWorker callback's Zone +/// without going through the full dispatcher. Not part of the public API. +@internal +const Object executionIdZoneKey = #nativeWorkmanagerExecutionId; /// Top-level callback dispatcher for background Dart execution. /// @@ -280,7 +278,7 @@ Future _callbackDispatcher() async { return false; }, ), - zoneValues: {_executionIdZoneKey: executionId}, + zoneValues: {executionIdZoneKey: executionId}, ); // Return execution result to native side @@ -352,7 +350,7 @@ Future _callbackDispatcher() async { ); }), zoneValues: { - _executionIdZoneKey: input?['__executionId'] as String?, + executionIdZoneKey: input?['__executionId'] as String?, }, ); } catch (e, stackTrace) { @@ -766,6 +764,79 @@ class NativeWorkManager { return handle; } + /// Resolves a [Worker] for the platform channel. + /// + /// If [worker] is a plain [DartWorker], this validates that its callback is + /// registered and returns the equivalent [DartWorkerInternal] carrying the + /// resolved native callback handle — plus, when [constraints] is given and + /// the platform is iOS, constraints promoted to a heavy task (a DartWorker + /// spins up a Flutter Engine, so it needs BGProcessingTask's larger budget + /// instead of BGAppRefreshTask's ~30s, the same promotion [enqueue] always + /// applied). Any other worker — including an already-resolved + /// [DartWorkerInternal] — passes through unchanged. + /// + /// This is the single choke point for every Dart→native path that can carry + /// a [DartWorker]: [enqueue], task chains ([_enqueueChain]), [TaskGraph] + /// nodes (via [enqueueTaskGraph]), and [registerRemoteTrigger]. Before this + /// existed, only [enqueue] and task chains did the conversion — a + /// DartWorker placed in a TaskGraph node or a RemoteTriggerRule mapping + /// reached native with no callbackHandle at all, so its callback could + /// never be resolved (found by the 2026-09-23 lib/ audit). Task chains also + /// silently skipped the iOS heavy-task promotion that plain enqueue() + /// applied, which this fixes as a side effect of sharing one implementation. + @internal + static (Worker, Constraints?) resolveWorkerForWire( + Worker worker, [ + Constraints? constraints, + ]) { + if (worker is! DartWorker) return (worker, constraints); + + if (!_dartWorkers.containsKey(worker.callbackId)) { + throw StateError( + 'Dart worker "${worker.callbackId}" not registered.\n' + 'Register it in NativeWorkManager.initialize():\n' + ' await NativeWorkManager.initialize(\n' + ' dartWorkers: {\n' + ' "${worker.callbackId}": (input) async { ... },\n' + ' },\n' + ' );', + ); + } + + var resolvedConstraints = constraints; + if (resolvedConstraints != null && + defaultTargetPlatform == TargetPlatform.iOS && + !resolvedConstraints.isHeavyTask) { + resolvedConstraints = resolvedConstraints.copyWith(isHeavyTask: true); + developer.log( + 'NativeWorkManager: DartWorker on iOS detected. Promoting to heavy task ' + '(isHeavyTask=true) to prevent OS termination.', + name: 'NativeWorkManager', + ); + } + + final callbackHandle = _callbackHandles[worker.callbackId]; + if (callbackHandle == null) { + throw StateError( + 'INTERNAL ERROR: Callback handle not found for "${worker.callbackId}". ' + 'This should never happen. Please report this bug.', + ); + } + + final resolvedWorker = DartWorkerInternal( + callbackId: worker.callbackId, + callbackHandle: callbackHandle, + input: worker.input, + autoDispose: worker.autoDispose, + timeoutMs: worker.timeoutMs, + onStoppedId: worker.onStoppedId, + onStoppedHandle: _resolveOnStoppedHandle(worker.onStoppedId), + cancelGraceMs: worker.cancelGrace?.inMilliseconds, + ); + + return (resolvedWorker, resolvedConstraints); + } + // ═══════════════════════════════════════════════════════════════════════════ // TASK SCHEDULING // ═══════════════════════════════════════════════════════════════════════════ @@ -983,56 +1054,9 @@ class NativeWorkManager { } // Validate DartWorker registration and prepare worker data - Worker workerToEnqueue = worker; - Constraints finalConstraints = constraints; - - if (worker is DartWorker) { - if (!_dartWorkers.containsKey(worker.callbackId)) { - throw StateError( - 'Dart worker "${worker.callbackId}" not registered.\n' - 'Register it in NativeWorkManager.initialize():\n' - ' await NativeWorkManager.initialize(\n' - ' dartWorkers: {\n' - ' "${worker.callbackId}": (input) async { ... },\n' - ' },\n' - ' );', - ); - } - - // iOS safety: DartWorkers are heavy by definition because they spin up - // a Flutter Engine. Force BGProcessingTask (60s+) instead of - // BGAppRefreshTask (30s) to prevent immediate OS kills. - if (defaultTargetPlatform == TargetPlatform.iOS && - !constraints.isHeavyTask) { - finalConstraints = constraints.copyWith(isHeavyTask: true); - developer.log( - 'NativeWorkManager: DartWorker on iOS detected. Promoting to heavy task ' - '(isHeavyTask=true) to prevent OS termination.', - name: 'NativeWorkManager', - ); - } - - // Get the callback handle for this worker - final callbackHandle = _callbackHandles[worker.callbackId]; - if (callbackHandle == null) { - throw StateError( - 'INTERNAL ERROR: Callback handle not found for "${worker.callbackId}". ' - 'This should never happen. Please report this bug.', - ); - } - - // Create enhanced DartWorker with callback handle - workerToEnqueue = DartWorkerInternal( - callbackId: worker.callbackId, - callbackHandle: callbackHandle, - input: worker.input, - autoDispose: worker.autoDispose, - timeoutMs: worker.timeoutMs, - onStoppedId: worker.onStoppedId, - onStoppedHandle: _resolveOnStoppedHandle(worker.onStoppedId), - cancelGraceMs: worker.cancelGrace?.inMilliseconds, - ); - } + final (workerToEnqueue, resolvedConstraints) = + resolveWorkerForWire(worker, constraints); + final finalConstraints = resolvedConstraints ?? constraints; final scheduleResult = await NativeWorkManagerPlatform.instance.enqueue( taskId: taskId, @@ -1348,7 +1372,7 @@ class NativeWorkManager { // other. Falls back to the coarse by-taskId check when called from // outside that Zone (e.g. main-isolate code with no dispatched // execution in scope). - final executionId = Zone.current[_executionIdZoneKey] as String?; + final executionId = Zone.current[executionIdZoneKey] as String?; final result = await channel.invokeMethod('isTaskCancelled', { 'taskId': taskId, @@ -2090,44 +2114,19 @@ class NativeWorkManager { static Future _enqueueChain(TaskChainBuilder chain) { // Convert DartWorker to DartWorkerInternal for all tasks in the chain + // (and, on iOS, apply the same heavy-task promotion enqueue() does — + // this used to be skipped here, leaving a chain-step DartWorker on the + // short BGAppRefreshTask budget instead of BGProcessingTask). final convertedSteps = chain.steps.map((step) { return step.map((task) { - final worker = task.worker; - - // Check if worker is DartWorker and needs conversion - if (worker is DartWorker) { - // Get the callback handle for this worker - final callbackHandle = _callbackHandles[worker.callbackId]; - if (callbackHandle == null) { - throw StateError( - 'INTERNAL ERROR: Callback handle not found for "${worker.callbackId}". ' - 'This should never happen. Please report this bug.', - ); - } - - // Convert DartWorker to DartWorkerInternal - final convertedWorker = DartWorkerInternal( - callbackId: worker.callbackId, - callbackHandle: callbackHandle, - input: worker.input, - autoDispose: worker.autoDispose, - timeoutMs: worker.timeoutMs, - onStoppedId: worker.onStoppedId, - onStoppedHandle: _resolveOnStoppedHandle(worker.onStoppedId), - cancelGraceMs: worker.cancelGrace?.inMilliseconds, - ); - - // Return modified task map with converted worker - return { - 'id': task.id, - 'workerClassName': convertedWorker.workerClassName, - 'workerConfig': convertedWorker.toMap(), - 'constraints': task.constraints.toMap(), - }; - } - - // For non-DartWorker tasks, use original toMap() - return task.toMap(); + final (resolvedWorker, resolvedConstraints) = + resolveWorkerForWire(task.worker, task.constraints); + return { + 'id': task.id, + 'workerClassName': resolvedWorker.workerClassName, + 'workerConfig': resolvedWorker.toMap(), + 'constraints': (resolvedConstraints ?? task.constraints).toMap(), + }; }).toList(); }).toList(); @@ -2554,9 +2553,22 @@ class NativeWorkManager { required RemoteTriggerRule rule, }) async { _checkInitialized(); + // Resolve any DartWorker in workerMappings before it reaches native — + // RemoteTriggerRule.toMap() alone never did this, so a DartWorker mapped + // here reached native with no callbackHandle and could never be resolved + // (found by the 2026-09-23 lib/ audit). No Constraints exist on a + // mapping, so only the worker half of resolveWorkerForWire's result is + // used. + final resolvedRule = RemoteTriggerRule( + payloadKey: rule.payloadKey, + workerMappings: rule.workerMappings.map( + (key, worker) => MapEntry(key, resolveWorkerForWire(worker).$1), + ), + secretKey: rule.secretKey, + ); return NativeWorkManagerPlatform.instance.registerRemoteTrigger( source: source, - rule: rule, + rule: resolvedRule, ); } diff --git a/lib/src/task_graph.dart b/lib/src/task_graph.dart index b072e28..d9116a4 100644 --- a/lib/src/task_graph.dart +++ b/lib/src/task_graph.dart @@ -47,6 +47,28 @@ class TaskNode { 'constraints': constraints.toMap(), }; } + + /// Same shape as [toMap], but with [worker]/[constraints] run through + /// [NativeWorkManager.resolveWorkerForWire] first — so a [DartWorker] node + /// carries a resolved callback handle instead of reaching native with none + /// at all (found by the 2026-09-23 lib/ audit: [toMap] alone never did this + /// conversion, unlike [NativeWorkManager.enqueue] and task chains). + /// + /// Kept separate from [toMap] so a plain `toMap()` call — used in tests, + /// docs examples, debug printing — never requires + /// `NativeWorkManager.initialize()` to have run. Only [enqueueTaskGraph] + /// calls this, at the moment a graph is actually sent to native. + Map _toResolvedMap() { + final (resolvedWorker, resolvedConstraints) = + NativeWorkManager.resolveWorkerForWire(worker, constraints); + return { + 'id': id, + 'workerClassName': resolvedWorker.workerClassName, + 'workerConfig': resolvedWorker.toMap(), + 'dependsOn': dependsOn, + 'constraints': (resolvedConstraints ?? constraints).toMap(), + }; + } } /// A directed acyclic graph (DAG) of background tasks. @@ -216,6 +238,16 @@ class TaskGraph { 'nodes': _nodes.map((n) => n.toMap()).toList(), }; } + + /// Same shape as [toMap], but every node goes through + /// [TaskNode._toResolvedMap] first. See that method for why this is + /// separate from [toMap]. + Map _toResolvedMap() { + return { + 'id': id, + 'nodes': _nodes.map((n) => n._toResolvedMap()).toList(), + }; + } } /// Result of a [TaskGraph] execution. @@ -462,7 +494,7 @@ Future enqueueTaskGraph(TaskGraph graph) async { // 1. Send graph to native for persistent orchestration. // This ensures the graph continues even if the app is killed. - await NativeWorkManagerPlatform.instance.enqueueGraph(graph.toMap()); + await NativeWorkManagerPlatform.instance.enqueueGraph(graph._toResolvedMap()); // 2. Start the Dart-side listener so we can resolve the result future // if the app stays alive. diff --git a/lib/src/worker.dart b/lib/src/worker.dart index b72f229..2ea7e02 100644 --- a/lib/src/worker.dart +++ b/lib/src/worker.dart @@ -59,27 +59,55 @@ class NativeWorker { NativeWorker._(); /// Validate URL format and throw helpful error if invalid. - static void _validateUrl(String url) { + static void _validateUrl(String url) => _validateUrlWithSchemes( + url, + insecureScheme: 'http', + secureScheme: 'https', + example: 'https://api.example.com/endpoint', + ); + + /// Validate a WebSocket URL. Same checks as [_validateUrl] (null-byte / + /// injection / HTTPS-equivalent enforcement / private-IP blocking), just + /// against `ws`/`wss` instead of `http`/`https`. + /// + /// Added by the 2026-09-23 lib/ audit: `NativeWorker.webSocket()` used to + /// accept any string with a `ws://`/`wss://` prefix and nothing else — + /// `enforceHttps(true)` had no effect on it, and it never blocked private + /// IPs, unlike every HTTP-based worker. + static void _validateWebSocketUrl(String url) => _validateUrlWithSchemes( + url, + insecureScheme: 'ws', + secureScheme: 'wss', + example: 'wss://api.example.com/socket', + ); + + static void _validateUrlWithSchemes( + String url, { + required String insecureScheme, + required String secureScheme, + required String example, + }) { _validateInput(url, 'URL', isUrl: true); if (url.isEmpty) { throw ArgumentError( 'URL cannot be empty.\n' - 'Provide a valid HTTP/HTTPS URL like "https://api.example.com/endpoint"', + 'Provide a valid $insecureScheme:// or $secureScheme:// URL like "$example"', ); } final uri = Uri.tryParse(url); if (uri == null || - (!uri.hasScheme || (uri.scheme != 'http' && uri.scheme != 'https'))) { + (!uri.hasScheme || + (uri.scheme != insecureScheme && uri.scheme != secureScheme))) { throw ArgumentError( 'Invalid URL format: "$url"\n' - 'URL must start with http:// or https://\n' - 'Example: "https://api.example.com/endpoint"', + 'URL must start with $insecureScheme:// or $secureScheme://\n' + 'Example: "$example"', ); } - // SECURITY: Enforce HTTPS if configured - if (NativeWorkManager.enforceHttps && uri.scheme == 'http') { + // SECURITY: Enforce the secure scheme if configured + if (NativeWorkManager.enforceHttps && uri.scheme == insecureScheme) { throw ArgumentError( 'Insecure URL blocked: "$url"\n' 'HTTPS is enforced by NativeWorkManager.initialize(enforceHttps: true)', @@ -165,6 +193,35 @@ class NativeWorker { } } + /// Validate a relative sub-path segment (e.g. [moveToSharedStorage]'s + /// `subDir`) for path traversal and injection, without requiring it to be + /// absolute — unlike [_validatePath], a relative segment is exactly what + /// these fields are meant to carry (Android's `MediaStore.RELATIVE_PATH` + /// / iOS's app-Documents subfolder). + static void _validateRelativeSegment(String value, String label) { + _validateInput(value, label, isUrl: false); + final normalized = value.toLowerCase(); + if (normalized.contains('..') || normalized.contains('%2e%2e')) { + throw ArgumentError( + '$label cannot contain ".." or encoded dot-segments (path traversal attempt blocked).'); + } + } + + /// Validates [url] the same way every built-in HTTP worker's + /// `NativeWorker.*` factory does. Exposed (not private) so a package + /// `Worker` subclass with no `NativeWorker.*` factory of its own — e.g. + /// [ParallelHttpUploadWorker], whose constructor is the only entry point — + /// can still call the same check from its own file. Not part of the public + /// API surface. + @internal + static void validateUrlForWorkerConstructor(String url) => _validateUrl(url); + + /// Validates [path] the same way every built-in file-based worker's + /// `NativeWorker.*` factory does. See [validateUrlForWorkerConstructor]. + @internal + static void validateFilePathForWorkerConstructor(String path, String label) => + _validateFilePath(path, label); + /// Validate file path and throw helpful error if invalid. static void _validateFilePath(String path, String parameterName) { if (path.isEmpty) { diff --git a/lib/src/workers/native_worker_http.dart b/lib/src/workers/native_worker_http.dart index 6a26007..c1f7929 100644 --- a/lib/src/workers/native_worker_http.dart +++ b/lib/src/workers/native_worker_http.dart @@ -459,12 +459,20 @@ MultiUploadWorker _buildMultiUpload({ Duration timeout = const Duration(minutes: 10), bool useBackgroundSession = false, }) { + // These two validators were missing entirely until the 2026-09-23 lib/ + // audit — every other HTTP worker in this file calls them, but this one + // didn't, so multiUpload() bypassed HTTPS enforcement, SSRF/private-IP + // blocking, and path-traversal checks. + NativeWorker._validateUrl(url); if (files.isEmpty) { throw ArgumentError('files must not be empty'); } if (files.length > 50) { throw ArgumentError('Maximum 50 files per upload request'); } + for (final file in files) { + NativeWorker._validateFilePath(file.filePath, 'files[].filePath'); + } return MultiUploadWorker( url: url, files: files, @@ -496,6 +504,13 @@ MoveToSharedStorageWorker _buildMoveToSharedStorage({ String? mimeType, String? subDir, }) { + // Missing entirely until the 2026-09-23 lib/ audit: sourcePath reached + // native with no path-traversal check, and subDir (which feeds + // MediaStore.RELATIVE_PATH on Android) wasn't checked at all. + NativeWorker._validateFilePath(sourcePath, 'sourcePath'); + if (subDir != null) { + NativeWorker._validateRelativeSegment(subDir, 'subDir'); + } return MoveToSharedStorageWorker( sourcePath: sourcePath, storageType: storageType, diff --git a/lib/src/workers/native_worker_websocket.dart b/lib/src/workers/native_worker_websocket.dart index ce391ea..90cc698 100644 --- a/lib/src/workers/native_worker_websocket.dart +++ b/lib/src/workers/native_worker_websocket.dart @@ -15,14 +15,12 @@ Worker _buildWebSocket({ 'Use a DartWorker with dart:io WebSocket for cross-platform WebSocket support.', ); } - if (url.isEmpty) { - throw ArgumentError('url cannot be empty for webSocket'); - } - final uri = Uri.tryParse(url); - if (uri == null || (uri.scheme != 'ws' && uri.scheme != 'wss')) { - throw ArgumentError( - 'Invalid WebSocket URL: "$url". Must start with ws:// or wss://', - ); + // Was a bare scheme-prefix check until the 2026-09-23 lib/ audit — url + // content (injection chars, null bytes) went unchecked, enforceHttps(true) + // had no effect on ws:// vs wss://, and blockPrivateIPs didn't apply here. + NativeWorker._validateWebSocketUrl(url); + if (storeResponseAt != null) { + NativeWorker._validateFilePath(storeResponseAt, 'storeResponseAt'); } if (timeoutSeconds <= 0) { throw ArgumentError('timeoutSeconds must be > 0, got $timeoutSeconds'); diff --git a/lib/src/workers/parallel_http_upload_worker.dart b/lib/src/workers/parallel_http_upload_worker.dart index a155af3..93859dd 100644 --- a/lib/src/workers/parallel_http_upload_worker.dart +++ b/lib/src/workers/parallel_http_upload_worker.dart @@ -68,9 +68,18 @@ final class ParallelHttpUploadWorker extends Worker { this.notificationBody, this.certificatePinning, }) { + // Missing entirely until the 2026-09-23 lib/ audit: unlike every other + // HTTP worker, this class has no NativeWorker.* factory to gate + // construction, so its own constructor is the only place these checks + // can live. + NativeWorker.validateUrlForWorkerConstructor(url); if (files.isEmpty) { throw ArgumentError.value(files, 'files', 'must not be empty'); } + for (final file in files) { + NativeWorker.validateFilePathForWorkerConstructor( + file.filePath, 'files[].filePath'); + } if (maxConcurrent < 1 || maxConcurrent > 16) { throw RangeError.range(maxConcurrent, 1, 16, 'maxConcurrent'); } diff --git a/test/security/issue_lib_audit_1_missing_validators_test.dart b/test/security/issue_lib_audit_1_missing_validators_test.dart new file mode 100644 index 0000000..7d38765 --- /dev/null +++ b/test/security/issue_lib_audit_1_missing_validators_test.dart @@ -0,0 +1,152 @@ +import 'package:flutter_test/flutter_test.dart'; +import 'package:flutter/services.dart'; +import 'package:native_workmanager/native_workmanager.dart'; + +/// Guards a real gap found by the 2026-09-23 lib/ audit: `multiUpload`, +/// `moveToSharedStorage`, `webSocket`, and `ParallelHttpUploadWorker`'s +/// constructor (which has no `NativeWorker.*` factory) all reached native +/// with none of the URL/path validation every sibling HTTP/file worker +/// enforces — no HTTPS enforcement, no SSRF/private-IP blocking, no +/// path-traversal blocking. +void main() { + TestWidgetsFlutterBinding.ensureInitialized(); + + const MethodChannel channel = + MethodChannel('dev.brewkits/native_workmanager'); + + setUp(() { + TestDefaultBinaryMessengerBinding.instance.defaultBinaryMessenger + .setMockMethodCallHandler(channel, (MethodCall methodCall) async { + switch (methodCall.method) { + case 'initialize': + return null; + default: + return null; + } + }); + }); + + tearDown(() { + NativeWorkManager.resetSecurityFlags(); + TestDefaultBinaryMessengerBinding.instance.defaultBinaryMessenger + .setMockMethodCallHandler(channel, null); + }); + + group('multiUpload() validators', () { + test('rejects a private-IP URL when blockPrivateIPs is set', () async { + NativeWorkManager.resetInitializedState(); + try { + await NativeWorkManager.initialize(blockPrivateIPs: true); + } catch (_) {} + + expect( + () => NativeWorker.multiUpload( + url: 'https://10.0.0.5/upload', + files: const [UploadFile(filePath: '/tmp/a.jpg')], + ), + throwsArgumentError, + ); + }); + + test('rejects a file path with ".." traversal', () { + expect( + () => NativeWorker.multiUpload( + url: 'https://upload.example.com/batch', + files: const [UploadFile(filePath: '/tmp/../../etc/passwd')], + ), + throwsArgumentError, + ); + }); + + test('accepts a normal https URL and absolute file paths', () { + final w = NativeWorker.multiUpload( + url: 'https://upload.example.com/batch', + files: const [UploadFile(filePath: '/tmp/a.jpg')], + ); + expect(w.toMap()['url'], 'https://upload.example.com/batch'); + }); + }); + + group('moveToSharedStorage() validators', () { + test('rejects a sourcePath with ".." traversal', () { + expect( + () => NativeWorker.moveToSharedStorage( + sourcePath: '/tmp/../../etc/passwd', + storageType: SharedStorageType.downloads, + ), + throwsArgumentError, + ); + }); + + test('rejects a subDir with ".." traversal', () { + expect( + () => NativeWorker.moveToSharedStorage( + sourcePath: '/tmp/photo.jpg', + storageType: SharedStorageType.photos, + subDir: '../../OtherApp/Camera', + ), + throwsArgumentError, + ); + }); + + test('accepts a normal relative subDir', () { + final w = NativeWorker.moveToSharedStorage( + sourcePath: '/tmp/photo.jpg', + storageType: SharedStorageType.photos, + subDir: 'Holidays', + ); + expect(w.toMap()['subDir'], 'Holidays'); + }); + }); + + group('webSocket() validators', () { + test('rejects ws:// when enforceHttps is set', () async { + NativeWorkManager.resetInitializedState(); + try { + await NativeWorkManager.initialize(enforceHttps: true); + } catch (_) {} + + expect( + () => NativeWorker.webSocket(url: 'ws://insecure.example.com'), + throwsArgumentError, + ); + }); + + test('rejects a storeResponseAt path with ".." traversal', () { + expect( + () => NativeWorker.webSocket( + url: 'wss://api.example.com', + storeResponseAt: '../../etc/hosts', + ), + throwsArgumentError, + ); + }); + }); + + group('ParallelHttpUploadWorker constructor validators', () { + test('rejects a private-IP URL when blockPrivateIPs is set', () async { + NativeWorkManager.resetInitializedState(); + try { + await NativeWorkManager.initialize(blockPrivateIPs: true); + } catch (_) {} + + expect( + () => ParallelHttpUploadWorker( + url: 'https://192.168.1.1/upload', + files: const [UploadFile(filePath: '/tmp/a.jpg')], + ), + throwsArgumentError, + ); + }); + + test('rejects a file path with ".." traversal', () { + expect( + () => ParallelHttpUploadWorker( + url: 'https://upload.example.com', + files: const [UploadFile(filePath: '/tmp/../../etc/passwd')], + ), + throwsArgumentError, + ); + }); + }); +} diff --git a/test/unit/issue_66_dart_worker_cancellation_test.dart b/test/unit/issue_66_dart_worker_cancellation_test.dart index f55b15a..b239871 100644 --- a/test/unit/issue_66_dart_worker_cancellation_test.dart +++ b/test/unit/issue_66_dart_worker_cancellation_test.dart @@ -3,6 +3,11 @@ import 'dart:async'; import 'package:flutter/services.dart'; import 'package:flutter_test/flutter_test.dart'; import 'package:native_workmanager/native_workmanager.dart'; +// executionIdZoneKey is @internal (not part of the public API, so hidden +// from the barrel above) — reachable directly since this test lives inside +// the same package. +import 'package:native_workmanager/src/native_work_manager.dart' + show executionIdZoneKey; /// Issue #66: https://github.com/brewkits/native_workmanager/discussions/66 /// @@ -87,13 +92,13 @@ void main() { // forward it transparently and native can answer per-execution instead // of per-taskId. This test cannot drive the real dispatcher (private, // driven by a native-supplied callback handle) but simulates being - // inside its Zone via the test-only [executionIdZoneKeyForTesting] hook. + // inside its Zone via the test-only [executionIdZoneKey] hook. test( 'forwards the current execution\'s executionId from the dispatcher Zone, when present', () async { final result = await runZoned( () => NativeWorkManager.isTaskCancelled('shared-task'), - zoneValues: {executionIdZoneKeyForTesting: 'exec-123'}, + zoneValues: {executionIdZoneKey: 'exec-123'}, ); expect(result, isFalse); diff --git a/test/unit/issue_lib_audit_2_dartworker_wire_resolution_test.dart b/test/unit/issue_lib_audit_2_dartworker_wire_resolution_test.dart new file mode 100644 index 0000000..3c4b63b --- /dev/null +++ b/test/unit/issue_lib_audit_2_dartworker_wire_resolution_test.dart @@ -0,0 +1,243 @@ +import 'package:flutter/foundation.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:native_workmanager/native_workmanager.dart'; +import 'package:native_workmanager/src/platform_interface.dart'; +import 'package:plugin_platform_interface/plugin_platform_interface.dart'; + +/// Guards a real gap found by the 2026-09-23 lib/ audit. +/// +/// [NativeWorkManager.enqueue] and task chains convert a [DartWorker] to a +/// [DartWorkerInternal] (resolving its native callback handle) before +/// sending it to native. [TaskGraph] nodes and [RemoteTriggerRule] +/// `workerMappings` never did — they called the worker's own `toMap()` +/// directly, which for a plain `DartWorker` never includes `callbackHandle` +/// at all (only `DartWorkerInternal.toMap()` does). Native's +/// `DartCallbackWorker` requires `callbackHandle`, so the task was enqueued +/// but the callback could never resolve — a silent no-op, discoverable only +/// by the task never running. +class _MockPlatform extends NativeWorkManagerPlatform + with MockPlatformInterfaceMixin { + Map? capturedGraphMap; + RemoteTriggerRule? capturedRule; + Map? capturedChainMap; + + @override + Future initialize({ + int? callbackHandle, + bool debugMode = false, + int maxConcurrentTasks = 4, + int diskSpaceBufferMB = 20, + int cleanupAfterDays = 30, + bool enforceHttps = false, + bool blockPrivateIPs = false, + bool registerPlugins = false, + }) async {} + + @override + void setCallbackExecutor( + Future Function(String callbackId, Map? input) + executor) {} + + @override + Future enqueueGraph(Map graphMap) async { + capturedGraphMap = graphMap; + return 'accepted'; + } + + // enqueueTaskGraph() subscribes to this right after enqueueGraph() — an + // empty stream is enough since these tests only assert on the outgoing + // payload, not on graph completion. + @override + Stream get events => const Stream.empty(); + + @override + Future registerRemoteTrigger({ + required RemoteTriggerSource source, + required RemoteTriggerRule rule, + }) async { + capturedRule = rule; + } + + @override + Future enqueueChain(Map chainMap) async { + capturedChainMap = chainMap; + return ScheduleResult.accepted; + } +} + +// Top-level function: PluginUtilities.getCallbackHandle requires this. +Future _testCallback(Map? input) async => true; + +void main() { + late _MockPlatform mockPlatform; + + setUp(() { + mockPlatform = _MockPlatform(); + NativeWorkManagerPlatform.instance = mockPlatform; + NativeWorkManager.initialize( + dartWorkers: {'test-worker': _testCallback}, + ); + }); + + group('TaskGraph node DartWorker resolution', () { + test('a DartWorker node carries a resolved callbackHandle', () async { + final graph = TaskGraph(id: 'g1') + ..add(TaskNode( + id: 'a', + worker: DartWorker(callbackId: 'test-worker'), + )); + + await NativeWorkManager.enqueueGraph(graph); + + final nodes = mockPlatform.capturedGraphMap!['nodes'] as List; + final nodeConfig = (nodes.first as Map)['workerConfig'] as Map; + expect(nodeConfig['callbackHandle'], isNotNull, + reason: 'Without resolution, a plain DartWorker.toMap() never ' + 'includes callbackHandle — native could never resolve it.'); + }); + + test('a non-DartWorker node is unaffected', () async { + final graph = TaskGraph(id: 'g2') + ..add(TaskNode( + id: 'a', + worker: NativeWorker.httpSync(url: 'https://example.com'), + )); + + await NativeWorkManager.enqueueGraph(graph); + + final nodes = mockPlatform.capturedGraphMap!['nodes'] as List; + expect((nodes.first as Map)['workerClassName'], 'HttpSyncWorker'); + }); + + test('a DartWorker node is promoted to isHeavyTask on iOS', () async { + debugDefaultTargetPlatformOverride = TargetPlatform.iOS; + addTearDown(() => debugDefaultTargetPlatformOverride = null); + + final graph = TaskGraph(id: 'g3') + ..add(TaskNode( + id: 'a', + worker: DartWorker(callbackId: 'test-worker'), + constraints: const Constraints(isHeavyTask: false), + )); + + await NativeWorkManager.enqueueGraph(graph); + + final nodes = mockPlatform.capturedGraphMap!['nodes'] as List; + final constraints = (nodes.first as Map)['constraints'] as Map; + expect(constraints['isHeavyTask'], isTrue); + }); + }); + + group('RemoteTriggerRule workerMappings DartWorker resolution', () { + test('a mapped DartWorker carries a resolved callbackHandle', () async { + await NativeWorkManager.registerRemoteTrigger( + source: RemoteTriggerSource.fcm, + rule: RemoteTriggerRule( + payloadKey: 'action', + workerMappings: { + 'sync': DartWorker(callbackId: 'test-worker'), + }, + ), + ); + + final mapped = mockPlatform.capturedRule!.workerMappings['sync']!; + expect(mapped, isA()); + expect((mapped as DartWorkerInternal).callbackHandle, isNonZero); + }); + + test('a mapped NativeWorker is unaffected', () async { + await NativeWorkManager.registerRemoteTrigger( + source: RemoteTriggerSource.fcm, + rule: RemoteTriggerRule( + payloadKey: 'action', + workerMappings: { + 'sync': NativeWorker.httpSync(url: 'https://example.com'), + }, + ), + ); + + final mapped = mockPlatform.capturedRule!.workerMappings['sync']!; + expect(mapped, isNot(isA())); + }); + }); + + group('Task chain DartWorker resolution (regression guard)', () { + test('a chain-step DartWorker is promoted to isHeavyTask on iOS', () async { + debugDefaultTargetPlatformOverride = TargetPlatform.iOS; + addTearDown(() => debugDefaultTargetPlatformOverride = null); + + await NativeWorkManager.beginWith( + TaskRequest(id: 'step1', worker: DartWorker(callbackId: 'test-worker')), + ).enqueue(); + + final steps = mockPlatform.capturedChainMap!['steps'] as List; + final firstStep = (steps.first as List).first as Map; + final constraints = firstStep['constraints'] as Map; + expect(constraints['isHeavyTask'], isTrue, + reason: 'This was skipped before resolveWorkerForWire unified the ' + 'chain and enqueue() code paths.'); + }); + }); + + group('Chain-step unregistered DartWorker error message', () { + test( + 'is the same helpful message enqueue() gives, not a misleading ' + '"INTERNAL ERROR" — pre-fix, chains skipped the registration check ' + 'and fell straight into the internal-error branch meant for a truly ' + 'impossible state', () async { + expect( + () => NativeWorkManager.beginWith( + TaskRequest( + id: 'step1', + worker: DartWorker(callbackId: 'never-registered'), + ), + ).enqueue(), + throwsA( + isA().having( + (e) => e.message, + 'message', + contains('not registered'), + ), + ), + ); + }); + }); + + group('Unregistered DartWorker fails loudly instead of silently', () { + test( + 'enqueueGraph() rejects with a StateError, not a crash or a silent ' + 'no-op native call', () async { + final graph = TaskGraph(id: 'g4') + ..add(TaskNode( + id: 'a', + worker: DartWorker(callbackId: 'never-registered'), + )); + + await expectLater( + NativeWorkManager.enqueueGraph(graph), + throwsA(isA()), + ); + // The whole point of failing before the native call: it must never + // have been reached with an unresolvable worker. + expect(mockPlatform.capturedGraphMap, isNull); + }); + + test( + 'registerRemoteTrigger() rejects with a StateError, not a crash or ' + 'a silent no-op native call', () async { + await expectLater( + NativeWorkManager.registerRemoteTrigger( + source: RemoteTriggerSource.fcm, + rule: RemoteTriggerRule( + payloadKey: 'action', + workerMappings: { + 'sync': DartWorker(callbackId: 'never-registered'), + }, + ), + ), + throwsA(isA()), + ); + expect(mockPlatform.capturedRule, isNull); + }); + }); +}