diff --git a/Resources/analytics-events.psv b/Resources/analytics-events.psv index eeb47332..a9df3bf0 100644 --- a/Resources/analytics-events.psv +++ b/Resources/analytics-events.psv @@ -22,7 +22,7 @@ app_launched| app_session_stall_detected|app_version,build_channel,build_revision,build_version,duration_bucket,format_ready,heartbeat_age_bucket,last_event,os_major,previous_clean_shutdown,reason,recovering,session_active,session_duration_bucket,session_kind,session_stage,stall_kind,stall_stage,trigger app_unclean_shutdown_detected|app_version,build_channel,build_revision,build_version,duration_bucket,format_ready,heartbeat_age_bucket,last_event,os_major,previous_clean_shutdown,reason,recovering,session_active,session_duration_bucket,session_kind,session_stage,stall_kind,stall_stage,trigger dictation_artifact_saved|delivery,duration_bucket,save_outcome,surface,trigger,word_count_bucket -dictation_audio_route_changed|default_input_class,default_output_class,format_ready,hfp_suspected,input_channels,input_device_class,input_rate_hz,output_channels,output_device_class,output_rate_hz,recovering,recovery_latency_bucket,route_shape,sample_flow_started,selected_input_class,selection_overrode_default,selection_reason,was_recording +dictation_audio_route_changed|default_input_class,default_output_class,format_ready,hfp_suspected,input_device_class,output_device_class,recovering,route_shape,sample_flow_started,selected_input_class,selection_overrode_default,selection_reason,was_recording dictation_audio_route_recovery_finished|default_input_class,default_output_class,format_ready,hfp_suspected,input_channels,input_device_class,input_rate_hz,outcome,output_channels,output_device_class,output_rate_hz,recovering,recovery_latency_bucket,route_shape,sample_flow_started,selected_input_class,selection_overrode_default,selection_reason,was_recording dictation_audio_route_recovery_timeout|default_input_class,default_output_class,format_ready,hfp_suspected,input_channels,input_device_class,input_rate_hz,output_channels,output_device_class,output_rate_hz,recovering,recovery_latency_bucket,route_shape,sample_flow_started,selected_input_class,selection_overrode_default,selection_reason,was_recording dictation_cancelled|default_input_class,default_output_class,duration_bucket,format_ready,hfp_suspected,input_channels,input_device_class,input_rate_hz,output_channels,output_device_class,output_rate_hz,recovering,recovery_latency_bucket,route_shape,sample_flow_started,selected_input_class,selection_overrode_default,selection_reason,trigger,was_recording @@ -33,6 +33,7 @@ dictation_recording_too_short|default_input_class,default_output_class,duration_ dictation_start_failed|default_input_class,default_output_class,failure_kind,format_ready,hfp_suspected,input_channels,input_device_class,input_rate_hz,output_channels,output_device_class,output_rate_hz,recovering,recovery_latency_bucket,route_shape,sample_flow_started,selected_input_class,selection_overrode_default,selection_reason,start_attempt_bucket,trigger,was_recording dictation_started|default_input_class,default_output_class,format_ready,hfp_suspected,input_channels,input_device_class,input_rate_hz,output_channels,output_device_class,output_rate_hz,recovering,recovery_latency_bucket,route_shape,sample_flow_started,selected_input_class,selection_overrode_default,selection_reason,trigger,was_recording dictation_stop_latency_measured|auto_enter_bucket,auto_send,auto_send_block_reason,auto_send_expected,auto_send_key,cleanup_bucket,cleanup_changed,cleanup_enabled,copy_reason,decode_bucket,default_input_class,default_output_class,delivery,format_ready,hfp_suspected,input_channels,input_device_class,input_rate_hz,mic_stop_bucket,model_wait_bucket,outcome,output_channels,output_device_class,output_rate_hz,paste_bucket,recovering,recovery_latency_bucket,route_shape,sample_flow_started,save_bucket,save_outcome,selected_input_class,selection_overrode_default,selection_reason,stop_to_done_bucket,stop_to_paste_bucket,target_confirmation_mode,trigger,was_recording,word_count_bucket +dictation_zombie_recovery_finished|failure_kind,hfp_suspected,input_device_class,output_device_class,result,route_shape,stage local_summary_cancelled|duration_bucket,provider,result,runtime,setup_ready,stage,summary_action local_summary_failed|duration_bucket,failure_kind,provider,result,runtime,setup_ready,stage,summary_action local_summary_finished|chunk_count_bucket,duration_bucket,provider,result,runtime,setup_ready,stage,summary_action diff --git a/Sources/Observability/CLAUDE.md b/Sources/Observability/CLAUDE.md index 990ad7e5..ae1b78ec 100644 --- a/Sources/Observability/CLAUDE.md +++ b/Sources/Observability/CLAUDE.md @@ -52,6 +52,8 @@ anonymous analytics, and Sparkle update plumbing. - PostHog config is read from `Info.plist` (`TranscriptedPostHogAPIKey`, `TranscriptedPostHogHost`) or process environment (`POSTHOG_API_KEY`, `POSTHOG_HOST`), and anonymous analytics must stay event-allowlisted and bucketed rather than sending raw payloads - New analytics events should be added to `Resources/analytics-events.psv`; new reviewed non-bucket property names should be added to `Resources/analytics-reviewed-properties.psv`; run `python3 scripts/ops/normalize-analytics-taxonomy.py --check` after edits or union merges. - Activation analytics should route through `ActivationTelemetry` so artifact action, agent prompt/setup, and saved-recent artifact return-proxy events keep stable names, targets, result enums, and coarse age/window buckets. +- Timeline analytics should route through `TimelineAnalyticsTelemetry` so future Dayflow events keep screen-derived data local and only send coarse enum/bucket payloads. +- `dictation_zombie_recovery_finished` is the single PostHog terminal event for each zombie-engine attempt. Keep its trigger, stage, result, and route fields categorical; raw device labels and exact sample counts are forbidden. - Non-fatal error forwarding to Sentry is allowlisted. New `.error` events should not automatically assume they are safe to send off-device. - `RuntimeDiagnostics` writes only coarse runtime state under app-owned state. Keep it free of transcript text, raw audio, file paths, device names, meeting titles, and speaker names. - `ReliabilityPacketRecorder` derives packets from already-reviewed observability events (`EventReporter.capture` calls `ReliabilityPacketRecorder.record` directly — see the sink map in `docs/observability.md`). Keep its context allowlist coarse and bucketed; do not add raw error text, transcript text, raw audio, file paths, device names, meeting titles, speaker names, emails, tokens, or source app names. diff --git a/Sources/Observability/SentryEventPolicy.swift b/Sources/Observability/SentryEventPolicy.swift index e13f4907..0171c2a5 100644 --- a/Sources/Observability/SentryEventPolicy.swift +++ b/Sources/Observability/SentryEventPolicy.swift @@ -77,6 +77,7 @@ struct SentryEventPolicy: Equatable { "readiness_refreshes", "recovering", "recovery_start_attempts", + "result", "route_change_count_bucket", "route_stability_warning", "route_shape", @@ -91,6 +92,7 @@ struct SentryEventPolicy: Equatable { "session_stage", "stall_kind", "stall_stage", + "stage", "start_attempts", "stt_model", "stop_timed_out", diff --git a/Sources/Speech/CLAUDE.md b/Sources/Speech/CLAUDE.md index 84753d2a..ebda6d26 100644 --- a/Sources/Speech/CLAUDE.md +++ b/Sources/Speech/CLAUDE.md @@ -34,6 +34,8 @@ - `ParakeetEngine` mirrors `ParakeetRecoveryState` flags into `@Published var isRecovering` and `@Published var inputFormatReady` (forwarded through `STTRouter`). `DictationSessionController` uses `DictationReadinessWaitPolicy` in its wait-for-ready loop so it can distinguish between "still recovering", "refresh input readiness", and "safe to start" instead of blindly retrying. AirPods Hands-Free Profile (24kHz hw / 48kHz output bus) is supported: the tap is installed with `format: nil` so buffers arrive at the output rate, and `nativeSampleRate` tracks the output rate so downstream resampling is correct. - `ParakeetEngine` consults `DictationInputDeviceSelectionPolicy` through `ParakeetAudioDeviceLookup` before recording so Bluetooth headset input does not unnecessarily hijack playback when a built-in mic fallback is available. - `ParakeetEngine` consults `ParakeetStartRecordingFailurePolicy` when startup fails so format-reset, engine rebuild, and prewarm-retry behavior stay consistent across direct starts and recovery attempts. +- Route-change analytics are committed only after the config-change debounce resolves to a new categorical route. Notification churn and A -> B -> A oscillation stay quiet, but the underlying device-recovery state machine still runs. +- Zombie-engine recovery owns a separate generation-gated task, replaces the stale `AVAudioEngine` through a timed reset, and retries once. Do not fold that task back into the startup watchdog or reuse the detected zombie graph. - `ParakeetEngine` stays `@MainActor` for app state, published UI state, and event reporting, but all `AVAudioEngine` graph work runs through its private serial audio-engine queue. Keep recording start/stop/readiness APIs async so callers do not block the main actor while CoreAudio settles, starts, stops, or rebuilds. - `ParakeetEngine` reports model-init failures with `ParakeetModelInitDiagnostics.failureContext(...)`, which keeps diagnostics useful for packaging/download/debugging issues without shipping raw transcript or device content. - Dictation intentionally exposes only final transcription. The abandoned provisional-text/EOU path was removed rather than kept behind a false feature flag. A future live dictation experience must add a real streaming engine and end-to-end tests instead of reviving dormant audio-tap branches. diff --git a/Sources/Speech/ParakeetDeviceRecovery.swift b/Sources/Speech/ParakeetDeviceRecovery.swift index 01c564c4..b27c003f 100644 --- a/Sources/Speech/ParakeetDeviceRecovery.swift +++ b/Sources/Speech/ParakeetDeviceRecovery.swift @@ -83,16 +83,8 @@ extension ParakeetEngine { } private func handleDefaultInputDeviceChange(selection: DictationInputDeviceSelection) { - cachedInputDeviceName = selection.selectedInput.name - cachedInputDeviceSelection = selection - AppLogger.transcription.info("PARAKEET | default input changed → \(selection.defaultInput.name); dictation input → \(selection.selectedInput.name)") - EventReporter.shared.capture( - level: .info, - engine: "parakeet", - event: "default_input_device_changed", - message: "Default input device changed", - context: inputSelectionContext(selection) - ) + routeTransitionDebounceState.observe(categoricalAudioRoute(for: selection)) + updateCachedInputDeviceSelection(selection) Task { @MainActor [weak self] in await self?.handleAudioConfigChange() } @@ -129,16 +121,6 @@ extension ParakeetEngine { cancelConfigRecoveryTimeout() let recoveryGeneration = recoveryState.beginConfigChange() publishRecoveryState() - // Load the new route off the main actor — the enumeration is blocking - // coreaudiod IPC — then refresh the analytics cache and report. - let wasRecordingForAnalytics = configChangeWasRecording - Task.detached(priority: .utility) { [weak self] in - let selection = Self.loadDictationInputDeviceSelection() - await self?.recordRouteChangeAnalytics( - selection: selection, - wasRecording: wasRecordingForAnalytics - ) - } scheduleConfigRecoveryTimeout( generation: recoveryGeneration, wasRecording: configChangeWasRecording @@ -152,17 +134,21 @@ extension ParakeetEngine { cancelAudioWatchdog() prewarmRetryTask?.cancel() prewarmRetryTask = nil + let configCleanupOwner = currentAudioEngineQueueOwnerToken() if isRecording { preserveCurrentRecordingBuffersForRecovery() await removeRecordingTap() + guard ownsAudioEngineQueue(configCleanupOwner) else { return } isRecording = false audioLevel = 0 } await stopAudioEngine() + guard ownsAudioEngineQueue(configCleanupOwner) else { return } isEnginePrewarmed = false - await rebuildAudioEngine(reason: "configuration_change") + guard let rebuiltOwner = await rebuildAudioEngine(reason: "configuration_change"), + ownsAudioGraph(rebuiltOwner) else { return } // Cancel any in-flight recovery — the latest device change wins. // Bluetooth disconnect/reconnect fires multiple notifications over @@ -176,25 +162,50 @@ extension ParakeetEngine { // short enough that dictation recovery feels responsive. try? await Task.sleep(nanoseconds: TranscriptedConstants.audioConfigChangeDebounceDelay) guard !Task.isCancelled, let self = self else { return } + let wasRecordingForAnalytics = self.configChangeWasRecording + Task.detached(priority: .utility) { [weak self] in + let selection = Self.loadDictationInputDeviceSelection() + await self?.recordStableRouteChangeAnalytics( + selection: selection, + wasRecording: wasRecordingForAnalytics, + recoveryGeneration: recoveryGeneration + ) + } + // Telemetry coalescing must never suppress the real recovery state + // transition, including an A -> B -> A notification burst. self.attemptDeviceRecovery() } } - private func recordRouteChangeAnalytics( + private func recordStableRouteChangeAnalytics( selection: DictationInputDeviceSelection?, - wasRecording: Bool + wasRecording: Bool, + recoveryGeneration: UInt64 ) { - if let selection { - updateCachedInputDeviceSelection(selection) + guard !recoveryState.isStale(generation: recoveryGeneration) else { return } + guard let selection else { + routeTransitionDebounceState.discardPendingRoute() + return } + routeTransitionDebounceState.observe(categoricalAudioRoute(for: selection)) + updateCachedInputDeviceSelection(selection) + guard let stableRoute = routeTransitionDebounceState.commitPendingRoute() else { return } + + AppLogger.transcription.info("PARAKEET | stable input route changed → \(stableRoute.routeShape)") + let context = dictationRouteAnalyticsContext( + selection: selection, + extra: ["was_recording": "\(wasRecording)"] + ) + EventReporter.shared.capture( + level: .info, + engine: "parakeet", + event: "default_input_device_changed", + message: "Stable categorical input route changed", + context: context + ) AnalyticsReporter.track( "dictation_audio_route_changed", - properties: dictationRouteAnalyticsContext( - selection: selection, - extra: [ - "was_recording": "\(wasRecording)" - ] - ) + properties: context ) } @@ -274,11 +285,13 @@ extension ParakeetEngine { return } + var lastSnapshotOwner: ParakeetAudioEngineQueueOwnerToken? do { var recoveryAttempt = 0 var readySnapshot: ParakeetAudioInputSnapshot? while readySnapshot == nil { recoveryAttempt += 1 + lastSnapshotOwner = self.currentAudioEngineQueueOwnerToken() let snapshot = try await self.audioInputSnapshot( operation: recoveryAttempt == 1 ? "device_recovery" : "device_recovery_retry", recoveryGeneration: myGeneration @@ -344,7 +357,10 @@ extension ParakeetEngine { for attempt in 1...TranscriptedConstants.recordingRestartAttempts { guard !Task.isCancelled else { return } guard !self.recoveryState.isStale(generation: myGeneration) else { return } - if await self.startRecording() { + let startSucceeded = await self.startRecording() + guard !Task.isCancelled else { return } + guard !self.recoveryState.isStale(generation: myGeneration) else { return } + if startSucceeded { restarted = true AppLogger.transcription.info("PARAKEET | recording recovered on new device (\(self.inputDeviceName)) after \(attempt) attempt(s)") EventReporter.shared.capture(level: .info, engine: "parakeet", @@ -387,6 +403,10 @@ extension ParakeetEngine { // / Bluetooth route-switch hang). Rebuilding on that same queue would // never run, so fail safe by abandoning the blocked graph instead. let audioEngineQueueBlocked = error is ParakeetAudioEngineWorkError + if audioEngineQueueBlocked { + guard let lastSnapshotOwner, + self.ownsAudioEngineQueue(lastSnapshotOwner) else { return } + } let failureAction = ParakeetDeviceRecoveryFailurePolicy.action(wasRecording: shouldRestartRecording) if self.recoveryState.finishRecovery(success: false, generation: myGeneration) { self.cancelConfigRecoveryTimeout() @@ -446,9 +466,13 @@ extension ParakeetEngine { audioEngineQueueBlocked: audioEngineQueueBlocked ) { case .queuedOnAudioEngineQueue: - await self.rebuildAudioEngine(reason: "device_change_rewarm_failed") + guard await self.rebuildAudioEngine(reason: "device_change_rewarm_failed") != nil else { return } case .abandonBlockedAudioGraph: - self.abandonBlockedAudioEngine(reason: "device_change_rewarm_failed") + guard let lastSnapshotOwner, + self.abandonBlockedAudioEngine( + reason: "device_change_rewarm_failed", + expectedOwner: lastSnapshotOwner + ) else { return } } if failureAction.schedulePrewarmRetry { self.prewarmRetryCount = 0 @@ -528,11 +552,15 @@ extension ParakeetEngine { ) ) } + let timeoutOwner = self.currentAudioEngineQueueOwnerToken() switch timeoutAction.rebuildStrategy { case .queuedOnAudioEngineQueue: - await self.rebuildAudioEngine(reason: "device_change_recovery_timeout") + guard await self.rebuildAudioEngine(reason: "device_change_recovery_timeout") != nil else { return } case .abandonBlockedAudioGraph: - self.abandonBlockedAudioEngine(reason: "device_change_recovery_timeout") + guard self.abandonBlockedAudioEngine( + reason: "device_change_recovery_timeout", + expectedOwner: timeoutOwner + ) else { return } } if failureAction.schedulePrewarmRetry { self.prewarmRetryCount = 0 diff --git a/Sources/Speech/ParakeetEngine.swift b/Sources/Speech/ParakeetEngine.swift index dce12a06..772ccb9e 100644 --- a/Sources/Speech/ParakeetEngine.swift +++ b/Sources/Speech/ParakeetEngine.swift @@ -11,6 +11,11 @@ import FluidAudio import Foundation import TranscriptedCore +private struct ParakeetSystemInputRestoreTarget: Equatable { + let temporaryInput: AudioDeviceID + let previousInput: AudioDeviceID +} + @MainActor class ParakeetEngine: ObservableObject { @Published var isRecording = false @@ -28,8 +33,12 @@ class ParakeetEngine: ObservableObject { var audioEngine = AVAudioEngine() private var audioEngineQueue = ParakeetEngine.makeAudioEngineQueue() + private static let systemInputWorkCoordinator = ParakeetSerialSystemInputWorkCoordinator( + label: "com.transcripted.parakeet.system-input" + ) var audioGraphGeneration = 0 - var audioStartInProgress = false + private var audioStartAdmission = ParakeetAudioStartAdmissionState() + var audioStartInProgress: Bool { audioStartAdmission.isInProgress } private var inputTapInstalled = false var sharedMeetingMicRecording = false private nonisolated let sharedMeetingMicRecorder = SharedMeetingMicRecorder() @@ -53,6 +62,7 @@ class ParakeetEngine: ObservableObject { var configChangeDebounceTask: Task? var configRecoveryTask: Task? var configRecoveryTimeoutTask: Task? + var routeTransitionDebounceState = ParakeetRouteTransitionDebounceState() /// Tracks whether a recording was active when the first config change in a /// burst arrived. Subsequent changes during recovery inherit this flag so /// the final recovery attempt knows to restart recording. @@ -74,7 +84,11 @@ class ParakeetEngine: ObservableObject { var modelFilePrefetchTask: Task? var prefetchedModelPath: URL? private var audioWatchdogTask: Task? - private var zombieRecoveryRestartPending = false + private var zombieRecoveryTask: Task? + private var zombieRecoveryState = ParakeetZombieRecoveryState() + private let zombieEngineWorkOwnership = ParakeetTimedAudioEngineWorkOwnership() + private var zombieRecoveryStartGeneration: UInt64? + private var zombieRecoveryRestartPending: Bool { zombieRecoveryState.isActive } private var asrInferenceActivity = ParakeetASRInferenceActivityState() private var asrInferenceHandoffCount = 0 private var asrInferenceWaiters: [CheckedContinuation] = [] @@ -91,7 +105,7 @@ class ParakeetEngine: ObservableObject { private var lastAudioStartFailureReportAt: TimeInterval? private var lastInputSelectionReportKey: String? var ignoreInputSelectionConfigChangesUntil: CFAbsoluteTime = 0 - private var pendingSystemInputRestore: (temporaryInput: AudioDeviceID, previousInput: AudioDeviceID)? + private var pendingSystemInputRestore = ParakeetOwnerBoundPendingState() var isModelLoaded: Bool { asrManagerReady } var inputDeviceName: String { cachedInputDeviceName } @@ -322,7 +336,7 @@ class ParakeetEngine: ObservableObject { } } - private static func cleanUpLateAudioStart(on audioEngine: AVAudioEngine) { + private nonisolated static func cleanUpLateAudioStart(on audioEngine: AVAudioEngine) { safelyRemoveInputTap(on: audioEngine) audioEngine.reset() } @@ -344,6 +358,20 @@ class ParakeetEngine: ObservableObject { func updateCachedInputDeviceSelection(_ selection: DictationInputDeviceSelection) { cachedInputDeviceName = selection.selectedInput.name cachedInputDeviceSelection = selection + routeTransitionDebounceState.seedStableRouteIfNeeded( + categoricalAudioRoute(for: selection) + ) + } + + func categoricalAudioRoute( + for selection: DictationInputDeviceSelection + ) -> ParakeetCategoricalAudioRoute { + let context = dictationRouteAnalyticsContext(selection: selection) + return ParakeetCategoricalAudioRoute( + inputDeviceClass: context["input_device_class"] ?? "unknown", + outputDeviceClass: context["output_device_class"] ?? "unknown", + routeShape: context["route_shape"] ?? "unknown" + ) } // MARK: - Input readiness @@ -358,8 +386,8 @@ class ParakeetEngine: ObservableObject { await releaseIdleAudioHardware(removeTap: false) guard !Task.isCancelled else { return } - let prewarmGeneration = audioGraphGeneration - guard canContinuePrewarm(generation: prewarmGeneration) else { return } + let prewarmOwner = currentAudioEngineQueueOwnerToken() + guard canContinuePrewarm(owner: prewarmOwner) else { return } let microphoneStatus = AVCaptureDevice.authorizationStatus(for: .audio) switch ParakeetPrewarmPolicy.decision(for: microphoneStatus) { @@ -386,7 +414,7 @@ class ParakeetEngine: ObservableObject { let prewarmSelection = await Task.detached(priority: .utility) { Self.loadDictationInputDeviceSelection() }.value - guard canContinuePrewarm(generation: prewarmGeneration) else { return } + guard canContinuePrewarm(owner: prewarmOwner) else { return } if ParakeetPrewarmPolicy.shouldDeferHardwarePrewarm(for: prewarmSelection) { if let prewarmSelection { updateCachedInputDeviceSelection(prewarmSelection) @@ -407,7 +435,7 @@ class ParakeetEngine: ObservableObject { do { snapshot = try await audioInputSnapshot(operation: "prewarm") } catch { - guard canContinuePrewarm(generation: prewarmGeneration) else { return } + guard canContinuePrewarm(owner: prewarmOwner) else { return } EventReporter.shared.capture( level: .warning, engine: "parakeet", @@ -418,7 +446,7 @@ class ParakeetEngine: ObservableObject { schedulePrewarmRetry() return } - guard canContinuePrewarm(generation: prewarmGeneration) else { return } + guard canContinuePrewarm(owner: prewarmOwner) else { return } // Validate both formats. AirPods on macOS run input in Hands-Free Profile // (24kHz hw / 48kHz output bus); CoreAudio's internal converter handles @@ -460,11 +488,11 @@ class ParakeetEngine: ObservableObject { AppLogger.transcription.info("PARAKEET | input ready (\(inputDeviceName), \(safeNativeSampleRate())Hz)") } - private func canContinuePrewarm(generation: Int) -> Bool { + private func canContinuePrewarm(owner: ParakeetAudioEngineQueueOwnerToken) -> Bool { !isShuttingDown && !isRecording && !audioStartInProgress - && audioGraphGeneration == generation + && ownsAudioEngineQueue(owner) } func schedulePrewarmRetry() { @@ -731,12 +759,12 @@ class ParakeetEngine: ObservableObject { ) async { ignoreInputSelectionConfigChangesUntil = CFAbsoluteTimeGetCurrent() + TranscriptedConstants.selfInducedConfigChangeIgnoreWindow - let restoreError = await Task.detached(priority: .utility) { + let restoreError = await Self.systemInputWorkCoordinator.run { Self.restoreSystemInputDeviceIfStillTemporary( temporaryInput: temporaryInput, previousInput: previousInput ) - }.value + } if let restoreError { EventReporter.shared.capture( level: .warning, @@ -749,13 +777,18 @@ class ParakeetEngine: ObservableObject { ] ) } else { - DictationPersistentInputPreferences.setTemporaryRecoveryMarker(nil) + if !pendingSystemInputRestore.hasPendingValue { + DictationPersistentInputPreferences.setTemporaryRecoveryMarker(nil) + } } } - private func restorePendingSystemInputAfterRecording(operation: String) async { - guard let restoreTarget = pendingSystemInputRestore else { return } - pendingSystemInputRestore = nil + private func restorePendingSystemInputAfterRecording( + ownedBy owner: ParakeetAudioGraphOwnerToken?, + operation: String + ) async { + guard let owner, + let restoreTarget = pendingSystemInputRestore.take(ownedBy: owner) else { return } await restoreSystemInputIfStillTemporary( temporaryInput: restoreTarget.temporaryInput, previousInput: restoreTarget.previousInput, @@ -763,18 +796,21 @@ class ParakeetEngine: ObservableObject { ) } - private func schedulePendingSystemInputRestore(operation: String) { - guard let restoreTarget = pendingSystemInputRestore else { return } - pendingSystemInputRestore = nil + private func schedulePendingSystemInputRestore( + ownedBy owner: ParakeetAudioGraphOwnerToken?, + operation: String + ) { + guard let owner, + let restoreTarget = pendingSystemInputRestore.take(ownedBy: owner) else { return } ignoreInputSelectionConfigChangesUntil = CFAbsoluteTimeGetCurrent() + TranscriptedConstants.selfInducedConfigChangeIgnoreWindow - Task.detached(priority: .utility) { + Self.systemInputWorkCoordinator.schedule { [weak self] in let restoreError = Self.restoreSystemInputDeviceIfStillTemporary( temporaryInput: restoreTarget.temporaryInput, previousInput: restoreTarget.previousInput ) - if let restoreError { - await MainActor.run { + Task { @MainActor [weak self] in + if let restoreError { EventReporter.shared.capture( level: .warning, engine: "parakeet", @@ -785,9 +821,8 @@ class ParakeetEngine: ObservableObject { "error": restoreError, ] ) - } - } else { - await MainActor.run { + } else { + guard let self, !self.pendingSystemInputRestore.hasPendingValue else { return } DictationPersistentInputPreferences.setTemporaryRecoveryMarker(nil) } } @@ -799,13 +834,15 @@ class ParakeetEngine: ObservableObject { recoveryGeneration: UInt64? = nil, allowsBuiltInBluetoothFallback: Bool = true ) async throws -> ParakeetAudioInputSnapshot { + let operationOwner = currentAudioEngineQueueOwnerToken() let snapshotStartedAt = CFAbsoluteTimeGetCurrent() let selectionStartedAt = CFAbsoluteTimeGetCurrent() - let selection = await Task.detached(priority: .utility) { + let selection = await Self.systemInputWorkCoordinator.run { Self.loadDictationInputDeviceSelection( allowsBuiltInBluetoothFallback: allowsBuiltInBluetoothFallback ) - }.value + } + guard ownsAudioEngineQueue(operationOwner) else { throw CancellationError() } var stageTimings = [ "audio_input_selection_load_ms": Self.elapsedMilliseconds(since: selectionStartedAt) ] @@ -832,16 +869,43 @@ class ParakeetEngine: ObservableObject { if let recoveryMarker { DictationPersistentInputPreferences.setTemporaryRecoveryMarker(recoveryMarker) } + let systemInputOverrideOwner = operationOwner.graphOwner + let systemInputRestoreTarget = selection.flatMap { selection -> ParakeetSystemInputRestoreTarget? in + guard selection.didOverrideDefault, + selection.reason == .preferredBuiltInForBluetoothHeadset else { return nil } + return ParakeetSystemInputRestoreTarget( + temporaryInput: selection.selectedInput.id, + previousInput: selection.defaultInput.id + ) + } + if shouldRestoreSystemInputOnStop, let systemInputRestoreTarget { + pendingSystemInputRestore.replace( + systemInputRestoreTarget, + ownedBy: systemInputOverrideOwner + ) + } let systemInputOverrideStartedAt = CFAbsoluteTimeGetCurrent() - let systemInputOverrideError = await Task.detached(priority: .utility) { + let systemInputOverrideError = await Self.systemInputWorkCoordinator.run { Self.applyPreferredSystemInputDevice(for: selection) - }.value - var systemInputRestoreTarget: (temporaryInput: AudioDeviceID, previousInput: AudioDeviceID)? + } + func restoreSystemInputAfterOwnershipLoss(stage: String) async { + pendingSystemInputRestore.clear(ownedBy: systemInputOverrideOwner) + guard systemInputOverrideError == nil, let systemInputRestoreTarget else { return } + await restoreSystemInputIfStillTemporary( + temporaryInput: systemInputRestoreTarget.temporaryInput, + previousInput: systemInputRestoreTarget.previousInput, + operation: "\(operation)_system_input_stale_\(stage)" + ) + } + guard ownsAudioEngineQueue(operationOwner) else { + await restoreSystemInputAfterOwnershipLoss(stage: "override") + throw CancellationError() + } stageTimings["audio_input_system_override_ms"] = Self.elapsedMilliseconds(since: systemInputOverrideStartedAt) if let selection, selection.didOverrideDefault { var context = inputSelectionContext(selection, operation: "\(operation)_system_input") if let systemInputOverrideError { - pendingSystemInputRestore = nil + pendingSystemInputRestore.clear(ownedBy: systemInputOverrideOwner) if recoveryMarker != nil { DictationPersistentInputPreferences.setTemporaryRecoveryMarker(nil) } @@ -854,14 +918,6 @@ class ParakeetEngine: ObservableObject { context: context ) } else if selection.reason == .preferredBuiltInForBluetoothHeadset { - let restoreTarget = ( - temporaryInput: selection.selectedInput.id, - previousInput: selection.defaultInput.id - ) - if shouldRestoreSystemInputOnStop { - pendingSystemInputRestore = restoreTarget - } - systemInputRestoreTarget = restoreTarget EventReporter.shared.capture( level: .info, engine: "parakeet", @@ -882,7 +938,6 @@ class ParakeetEngine: ObservableObject { } throw CancellationError() } - let snapshotGraphGeneration = audioGraphGeneration let snapshotReadStartedAt = CFAbsoluteTimeGetCurrent() let snapshotResult: ( outputFormat: ParakeetAudioFormatSummary, @@ -902,6 +957,10 @@ class ParakeetEngine: ObservableObject { ) } } catch { + guard ownsAudioEngineQueue(operationOwner) else { + await restoreSystemInputAfterOwnershipLoss(stage: "snapshot_failure") + throw CancellationError() + } if !shouldRestoreSystemInputOnStop, let systemInputRestoreTarget { await restoreSystemInputIfStillTemporary( temporaryInput: systemInputRestoreTarget.temporaryInput, @@ -911,6 +970,10 @@ class ParakeetEngine: ObservableObject { } throw error } + guard ownsAudioEngineQueue(operationOwner) else { + await restoreSystemInputAfterOwnershipLoss(stage: "snapshot_success") + throw CancellationError() + } if !shouldRestoreSystemInputOnStop, let systemInputRestoreTarget { await restoreSystemInputIfStillTemporary( temporaryInput: systemInputRestoreTarget.temporaryInput, @@ -931,9 +994,8 @@ class ParakeetEngine: ObservableObject { if let recoveryGeneration, recoveryState.isStale(generation: recoveryGeneration) { throw CancellationError() } - if snapshotGraphGeneration == audioGraphGeneration { - recordInputSelection(snapshot.selectionApplication, operation: operation) - } + guard ownsAudioEngineQueue(operationOwner) else { throw CancellationError() } + recordInputSelection(snapshot.selectionApplication, operation: operation) guard snapshot.selectionApplication?.didApplyOverride == true else { return snapshot @@ -954,11 +1016,11 @@ class ParakeetEngine: ObservableObject { let settleSleepStartedAt = CFAbsoluteTimeGetCurrent() try? await Task.sleep(nanoseconds: overrideSettleDelay) stageTimings["audio_input_override_settle_sleep_ms"] = Self.elapsedMilliseconds(since: settleSleepStartedAt) + guard ownsAudioEngineQueue(operationOwner) else { throw CancellationError() } if let recoveryGeneration, recoveryState.isStale(generation: recoveryGeneration) { throw CancellationError() } - let settledGraphGeneration = audioGraphGeneration let settledSnapshotStartedAt = CFAbsoluteTimeGetCurrent() let settledSnapshotResult = try await runTimedAudioEngineWork(operation: "\(operation)_settled_snapshot") { audioEngine in let inputNode = audioEngine.inputNode @@ -978,31 +1040,30 @@ class ParakeetEngine: ObservableObject { engineWasRunning: settledSnapshotResult.engineWasRunning, stageTimings: stageTimings ) + guard ownsAudioEngineQueue(operationOwner) else { throw CancellationError() } if let recoveryGeneration, recoveryState.isStale(generation: recoveryGeneration) { throw CancellationError() } - if settledGraphGeneration == audioGraphGeneration { - let readiness = audioFormatReadiness( + let readiness = audioFormatReadiness( + outputFormat: settledSnapshot.outputFormat, + hwFormat: settledSnapshot.hwFormat, + selection: settledSnapshot.selection + ) + EventReporter.shared.capture( + level: .info, + engine: "parakeet", + event: "dictation_input_device_override_settled", + message: "Dictation input override settled before reading microphone format", + context: audioFormatContext( outputFormat: settledSnapshot.outputFormat, hwFormat: settledSnapshot.hwFormat, - selection: settledSnapshot.selection + selection: settledSnapshot.selection, + readiness: readiness + ).merging( + ["operation": operation], + uniquingKeysWith: { current, _ in current } ) - EventReporter.shared.capture( - level: .info, - engine: "parakeet", - event: "dictation_input_device_override_settled", - message: "Dictation input override settled before reading microphone format", - context: audioFormatContext( - outputFormat: settledSnapshot.outputFormat, - hwFormat: settledSnapshot.hwFormat, - selection: settledSnapshot.selection, - readiness: readiness - ).merging( - ["operation": operation], - uniquingKeysWith: { current, _ in current } - ) - ) - } + ) return settledSnapshot } @@ -1121,6 +1182,7 @@ class ParakeetEngine: ObservableObject { func removeRecordingTap(force: Bool = false) async { guard force || inputTapInstalled else { return } + let tapOwner = currentAudioGraphOwnerToken() await runAudioEngineWork { audioEngine in // Stop + drain before removing the tap; the canonical stop path // (`removeRecordingTap()` then `stopAudioEngine()`) otherwise removes @@ -1128,6 +1190,7 @@ class ParakeetEngine: ObservableObject { // audio IO thread with `isSink || tap != nullptr`. Self.safelyRemoveInputTap(on: audioEngine) } + guard ownsAudioGraph(tapOwner) else { return } inputTapInstalled = false } @@ -1172,10 +1235,13 @@ class ParakeetEngine: ObservableObject { } } - private func resetAudioGraphAfterStartFailure(reason: String, rebuildEngine: Bool) async { + private func resetAudioGraphAfterStartFailure( + reason: String, + rebuildEngine: Bool + ) async -> ParakeetAudioGraphOwnerToken? { // Keep runtime/UI state coherent when startRecording fails before we ever // transition to a stable recording session. - cancelAudioWatchdog() + cancelAudioWatchdogForRecordingStart() isRecording = false audioLevel = 0 didReceiveAudioSamples = false @@ -1183,16 +1249,18 @@ class ParakeetEngine: ObservableObject { recordingStartedOnLikelyBluetoothHandsFreeRoute = false if rebuildEngine { - await rebuildAudioEngine(reason: reason) - return + return await rebuildAudioEngine(reason: reason) } audioGraphGeneration += 1 + let resetOwner = currentAudioGraphOwnerToken() await runAudioEngineWork { audioEngine in Self.safelyRemoveInputTap(on: audioEngine) audioEngine.reset() } + guard ownsAudioGraph(resetOwner) else { return nil } inputTapInstalled = false isEnginePrewarmed = false + return resetOwner } /// Tracks rebuild frequency and reports once if rebuilds are churning — @@ -1223,16 +1291,18 @@ class ParakeetEngine: ObservableObject { ) } - func rebuildAudioEngine(reason: String) async { + @discardableResult + func rebuildAudioEngine(reason: String) async -> ParakeetAudioGraphOwnerToken? { trackAudioEngineRebuildChurn(reason: reason) audioGraphGeneration += 1 + let rebuildOwner = currentAudioGraphOwnerToken() removeAudioEngineConfigObserver() let retiredEngine = audioEngine await runAudioEngineWork { audioEngine in Self.safelyRemoveInputTap(on: audioEngine) audioEngine.reset() } - guard audioEngine === retiredEngine else { return } + guard ownsAudioGraph(rebuildOwner) else { return nil } audioEngine = AVAudioEngine() ParakeetRetiredAudioEngineStore.shared.retire(retiredEngine, reason: reason) inputTapInstalled = false @@ -1255,10 +1325,22 @@ class ParakeetEngine: ObservableObject { "generation": "\(recoveryState.generation)" ] ) + return currentAudioGraphOwnerToken() } - func abandonBlockedAudioEngine(reason: String) { + @discardableResult + func abandonBlockedAudioEngine( + reason: String, + expectedOwner: ParakeetAudioEngineQueueOwnerToken? = nil + ) -> Bool { + if let expectedOwner, !ownsAudioEngineQueue(expectedOwner) { + return false + } trackAudioEngineRebuildChurn(reason: reason) + _ = zombieEngineWorkOwnership.claimPendingWorkForSuccessor( + currentEngine: audioEngine, + currentQueue: audioEngineQueue + ) audioGraphGeneration += 1 removeAudioEngineConfigObserver() let retiredEngine = audioEngine @@ -1286,6 +1368,7 @@ class ParakeetEngine: ObservableObject { "generation": "\(recoveryState.generation)" ] ) + return true } private func handleSystemWake() async { @@ -1297,14 +1380,18 @@ class ParakeetEngine: ObservableObject { audioGraphGeneration += 1 let wasRecording = isRecording cancelAudioWatchdog() + audioStartAdmission.cancel() + let wakeCleanupOwner = currentAudioEngineQueueOwnerToken() if isRecording { preserveCurrentRecordingBuffersForRecovery() await removeRecordingTap() + guard ownsAudioEngineQueue(wakeCleanupOwner) else { return } isRecording = false audioLevel = 0 } await stopAudioEngine() + guard ownsAudioEngineQueue(wakeCleanupOwner) else { return } isEnginePrewarmed = false if wasRecording { @@ -1517,12 +1604,14 @@ class ParakeetEngine: ObservableObject { ) return false } - audioStartInProgress = true audioStartReferenceTime = CFAbsoluteTimeGetCurrent() audioGraphGeneration += 1 - var startGeneration = audioGraphGeneration + var startOwner = currentAudioEngineQueueOwnerToken() + guard audioStartAdmission.begin(owner: startOwner) else { return false } + var startEngine = audioEngine + var startQueue = audioEngineQueue defer { - audioStartInProgress = false + audioStartAdmission.finish(owner: startOwner) } func failAudioStart() async -> Bool { // Keep the temporary built-in input selected across the controller's @@ -1562,7 +1651,7 @@ class ParakeetEngine: ObservableObject { didReceiveAudioSamples = false didReceiveNonZeroAudioSamples = false recordingStartedOnLikelyBluetoothHandsFreeRoute = false - cancelAudioWatchdog() + cancelAudioWatchdogForRecordingStart() if !isRecoveryAttempt && !preservingRecordingAcrossRecovery { recoveredRecordingTimeline.removeAll(keepingCapacity: true) } @@ -1575,7 +1664,10 @@ class ParakeetEngine: ObservableObject { let maxAttempts = isRecoveryAttempt ? 1 : 1 + TranscriptedConstants.audioStartRecoveryAttempts for attempt in 1...maxAttempts { - guard startGeneration == audioGraphGeneration else { + let attemptOwner = startOwner + let attemptEngine = startEngine + let attemptQueue = startQueue + guard ownsAudioEngineQueue(attemptOwner) else { EventReporter.shared.capture( level: .warning, engine: "parakeet", @@ -1593,6 +1685,7 @@ class ParakeetEngine: ObservableObject { allowsBuiltInBluetoothFallback: !isRecoveryAttempt ) } catch { + guard ownsAudioEngineQueue(attemptOwner) else { return await failAudioStart() } let operationTimedOut = error is ParakeetAudioEngineWorkError EventReporter.shared.capture( level: operationTimedOut ? .error : .warning, @@ -1608,15 +1701,21 @@ class ParakeetEngine: ObservableObject { ] ) if operationTimedOut { - abandonBlockedAudioEngine(reason: "audio_format_read_timeout") + guard abandonBlockedAudioEngine( + reason: "audio_format_read_timeout", + expectedOwner: attemptOwner + ) else { return await failAudioStart() } } else { - await resetAudioGraphAfterStartFailure(reason: "audio_format_read_failed", rebuildEngine: true) + guard await resetAudioGraphAfterStartFailure( + reason: "audio_format_read_failed", + rebuildEngine: true + ) != nil else { return await failAudioStart() } } markFormatUnreadyAndPublish() schedulePrewarmRetry() return await failAudioStart() } - guard startGeneration == audioGraphGeneration else { + guard ownsAudioEngineQueue(attemptOwner) else { EventReporter.shared.capture( level: .warning, engine: "parakeet", @@ -1660,10 +1759,10 @@ class ParakeetEngine: ObservableObject { message: "Audio hardware format not ready while starting dictation", context: context ) - await resetAudioGraphAfterStartFailure( + guard await resetAudioGraphAfterStartFailure( reason: readiness == .routeNotSettled ? "audio_route_not_settled" : "invalid_audio_format", rebuildEngine: startFailureAction.rebuildAudioEngine - ) + ) != nil else { return await failAudioStart() } if startFailureAction.markFormatUnready { markFormatUnreadyAndPublish() } @@ -1682,11 +1781,31 @@ class ParakeetEngine: ObservableObject { ) reserveNativeSampleBufferCapacity() + let zombieStartLeaseOwner: ParakeetAudioEngineQueueOwnerToken? + if isRecoveryAttempt, + let generation = zombieRecoveryStartGeneration, + zombieRecoveryState.canContinue(generation: generation) { + zombieStartLeaseOwner = attemptOwner + zombieEngineWorkOwnership.begin( + owner: attemptOwner, + phase: .zombieRecoveryStart + ) + } else { + zombieStartLeaseOwner = nil + } + do { let startSnapshot = try await installTapAndStartEngine(isRecoveryAttempt: isRecoveryAttempt) - guard startGeneration == audioGraphGeneration else { - await removeRecordingTap(force: true) - await stopAudioEngine() + if let zombieStartLeaseOwner { + zombieEngineWorkOwnership.finish( + owner: zombieStartLeaseOwner, + phase: .zombieRecoveryStart + ) + } + guard ownsAudioEngineQueue(attemptOwner) else { + attemptQueue.async { + Self.cleanUpLateAudioStart(on: attemptEngine) + } EventReporter.shared.capture( level: .warning, engine: "parakeet", @@ -1729,6 +1848,13 @@ class ParakeetEngine: ObservableObject { ]) } } catch { + if let zombieStartLeaseOwner { + zombieEngineWorkOwnership.finish( + owner: zombieStartLeaseOwner, + phase: .zombieRecoveryStart + ) + } + guard ownsAudioEngineQueue(attemptOwner) else { return await failAudioStart() } let operationTimedOut = error is ParakeetAudioEngineWorkError var context = audioStartContext( attempt: attempt, @@ -1755,16 +1881,25 @@ class ParakeetEngine: ObservableObject { failedAttempts: attempt ) if operationTimedOut { - abandonBlockedAudioEngine(reason: "audio_engine_start_timeout") + guard abandonBlockedAudioEngine( + reason: "audio_engine_start_timeout", + expectedOwner: attemptOwner + ) else { return await failAudioStart() } } else { - await resetAudioGraphAfterStartFailure( + guard await resetAudioGraphAfterStartFailure( reason: failureReason == .audioRouteNotSettled ? "audio_route_not_settled" : "audio_engine_start_failed", rebuildEngine: startFailureAction.rebuildAudioEngine - ) + ) != nil else { return await failAudioStart() } } if shouldRetry { - startGeneration = audioGraphGeneration + let retryOwner = currentAudioEngineQueueOwnerToken() + guard audioStartAdmission.transfer(from: startOwner, to: retryOwner) else { + return await failAudioStart() + } + startOwner = retryOwner + startEngine = audioEngine + startQueue = audioEngineQueue AppLogger.transcription.warning("PARAKEET | audio engine start failed, resetting graph and retrying once: \(error.localizedDescription)") EventReporter.shared.capture( level: .warning, @@ -1905,13 +2040,20 @@ class ParakeetEngine: ObservableObject { let started = await startRecording(isRecoveryAttempt: true) guard sharedMeetingMicTransition.finishResume(token: transitionToken) else { if started { + let pendingRestoreOwner = pendingSystemInputRestore.owner audioGraphGeneration += 1 cancelAudioWatchdog() + let staleResumeOwner = currentAudioEngineQueueOwnerToken() await removeRecordingTap() + guard ownsAudioEngineQueue(staleResumeOwner) else { return } await stopAudioEngine() + guard ownsAudioEngineQueue(staleResumeOwner) else { return } isRecording = false audioLevel = 0 - await restorePendingSystemInputAfterRecording(operation: "stale_shared_meeting_mic_resume") + await restorePendingSystemInputAfterRecording( + ownedBy: pendingRestoreOwner, + operation: "stale_shared_meeting_mic_resume" + ) } return } @@ -1985,7 +2127,8 @@ class ParakeetEngine: ObservableObject { /// Watchdog that detects zombie audio engines — running but producing no usable signal. /// After sleep/wake, CoreAudio may report the engine as running but the hardware graph - /// is disconnected. If no samples arrive within 2 seconds, tear down and retry once. + /// is disconnected. If no samples arrive within 2 seconds, replace the stale engine + /// through a bounded reset and retry once. /// If the user stops dictation during the recovery delay, the pending retry is cleared /// so the watchdog does not revive a recording the user already ended. private func startAudioWatchdog() { @@ -2007,51 +2150,269 @@ class ParakeetEngine: ObservableObject { ) guard shouldReset else { return } + let failureKind = sampleCount == 0 ? "no_sample_callbacks" : "silent_hfp_callbacks" EventReporter.shared.capture(level: .warning, engine: "parakeet", event: "zombie_engine_detected", message: sampleCount == 0 ? "No audio samples received after recording start — resetting engine" : "Only silent audio samples received after recording start — resetting engine", - context: [ - "audio_device": self.inputDeviceName, - "hfp_suspected": "\(self.recordingStartedOnLikelyBluetoothHandsFreeRoute)", - "sample_count": "\(sampleCount)", - "sample_signal_started": "\(self.didReceiveNonZeroAudioSamples)", - ]) + context: self.zombieRecoveryTelemetryContext( + failureKind: failureKind, + stage: .detected, + result: nil + )) - self.pendingSamplesLock.withLock { - self.pendingSamples.removeAll(keepingCapacity: true) - self.didReportPendingSampleTruncation = false - } - self.zombieRecoveryRestartPending = true - self.isRecording = false - self.audioLevel = 0 - self.configChangeWasRecording = false - self.ignoreInputSelectionConfigChangesUntil = CFAbsoluteTimeGetCurrent() - + TranscriptedConstants.selfInducedConfigChangeIgnoreWindow - await self.removeRecordingTap() - await self.stopAudioEngine() - self.isEnginePrewarmed = false + // Detection and recovery use separate task lifetimes. Otherwise the + // recovery's call into startRecording cancels the watchdog task that + // is currently executing, skipping cancellation-aware settle work. + self.audioWatchdogTask = nil + self.startZombieEngineRecovery(failureKind: failureKind) + } + } - try? await Task.sleep(nanoseconds: TranscriptedConstants.audioRecoveryDelay) - guard !Task.isCancelled, self.zombieRecoveryRestartPending else { - self.zombieRecoveryRestartPending = false - return + private func startZombieEngineRecovery(failureKind: String) { + guard !zombieRecoveryState.isActive else { return } + let generation = zombieRecoveryState.begin(failureKind: failureKind) + zombieRecoveryTask = Task { @MainActor [weak self] in + await self?.runZombieEngineRecovery(generation: generation) + } + } + + private func runZombieEngineRecovery(generation: UInt64) async { + defer { + clearZombieRecoveryStartGeneration(ifMatching: generation) + if zombieRecoveryState.canContinue(generation: generation) { + finishZombieEngineRecovery( + generation: generation, + result: Task.isCancelled ? .cancelled : .failed + ) } - self.zombieRecoveryRestartPending = false + } - // Retry once — isRecoveryAttempt prevents another watchdog - if await self.startRecording(isRecoveryAttempt: true) { - AppLogger.transcription.info("PARAKEET | zombie engine recovered — recording restarted") - EventReporter.shared.capture(level: .info, engine: "parakeet", event: "zombie_engine_recovered", - message: "Audio engine recovered after reset") - } else { - AppLogger.transcription.error("PARAKEET | zombie engine recovery failed") - EventReporter.shared.capture(level: .error, engine: "parakeet", event: "zombie_engine_recovery_failed", - message: "Audio engine could not recover after reset", - context: ["audio_device": self.inputDeviceName]) - self.interruptRecordingPreservingRecoveredTimeline() + guard zombieRecoveryState.advance(to: .reset, generation: generation) else { return } + let recoveryGraphOwner = currentAudioGraphOwnerToken() + pendingSamplesLock.withLock { + pendingSamples.removeAll(keepingCapacity: true) + didReportPendingSampleTruncation = false + } + isRecording = false + audioLevel = 0 + configChangeWasRecording = false + ignoreInputSelectionConfigChangesUntil = CFAbsoluteTimeGetCurrent() + + TranscriptedConstants.selfInducedConfigChangeIgnoreWindow + + // Stop/config-change cancellation takes ownership of graph cleanup. The + // superseded zombie task must not enter recreation after this suspension. + guard canContinueZombieEngineRecovery( + generation: generation, + expectedOwner: recoveryGraphOwner + ) else { return } + guard await recreateAudioEngineForZombieRecovery( + generation: generation, + expectedOwner: recoveryGraphOwner + ) else { return } + guard !Task.isCancelled, zombieRecoveryState.canContinue(generation: generation) else { return } + + guard zombieRecoveryState.advance(to: .settle, generation: generation) else { return } + do { + try await Task.sleep(nanoseconds: TranscriptedConstants.audioRecoveryDelay) + } catch { + return + } + guard !Task.isCancelled, zombieRecoveryState.canContinue(generation: generation) else { return } + + guard zombieRecoveryState.advance(to: .restart, generation: generation) else { return } + zombieRecoveryStartGeneration = generation + let started = await startRecording(isRecoveryAttempt: true) + clearZombieRecoveryStartGeneration(ifMatching: generation) + guard zombieRecoveryState.canContinue(generation: generation) else { return } + + if started { + AppLogger.transcription.info("PARAKEET | zombie engine recovered — recording restarted") + finishZombieEngineRecovery(generation: generation, result: .succeeded) + } else { + AppLogger.transcription.error("PARAKEET | zombie engine recovery failed") + interruptRecordingPreservingRecoveredTimeline() + finishZombieEngineRecovery(generation: generation, result: .failed) + } + } + + /// A detected zombie is evidence that the current AVAudioEngine graph is stale. + /// Replace that instance rather than stopping and starting it again. Queue work + /// is bounded; if CoreAudio does not return, abandon the old graph and queue. + private func recreateAudioEngineForZombieRecovery( + generation: UInt64, + expectedOwner: ParakeetAudioGraphOwnerToken + ) async -> Bool { + guard canContinueZombieEngineRecovery( + generation: generation, + expectedOwner: expectedOwner + ) else { return false } + + trackAudioEngineRebuildChurn(reason: "zombie_engine_recovery") + audioGraphGeneration += 1 + let resetOwner = currentAudioGraphOwnerToken() + let resetQueueOwner = currentAudioEngineQueueOwnerToken() + removeAudioEngineConfigObserver() + let retiredEngine = audioEngine + zombieEngineWorkOwnership.begin(owner: resetQueueOwner, phase: .zombieReset) + + do { + try await runTimedAudioEngineWork(operation: "zombie_engine_reset") { [zombieEngineWorkOwnership] audioEngine in + defer { + zombieEngineWorkOwnership.finish( + owner: resetQueueOwner, + phase: .zombieReset + ) + } + Self.safelyRemoveInputTap(on: audioEngine) + audioEngine.reset() } + } catch { + zombieEngineWorkOwnership.finish(owner: resetQueueOwner, phase: .zombieReset) + guard error is ParakeetAudioEngineWorkError else { return false } + guard canContinueZombieEngineRecovery( + generation: generation, + expectedOwner: resetOwner + ) else { return false } + + // Only the exact generation+engine owner may abandon a timed-out + // queue; a newer graph may reuse the same engine instance. + return abandonBlockedAudioEngine( + reason: "zombie_engine_reset_timeout", + expectedOwner: resetQueueOwner + ) } + + guard canContinueZombieEngineRecovery( + generation: generation, + expectedOwner: resetOwner + ) else { return false } + + inputTapInstalled = false + isEnginePrewarmed = false + didReceiveAudioSamples = false + didReceiveNonZeroAudioSamples = false + recordingStartedOnLikelyBluetoothHandsFreeRoute = false + audioEngine = AVAudioEngine() + ParakeetRetiredAudioEngineStore.shared.retire(retiredEngine, reason: "zombie_engine_recovery") + if !isShuttingDown { + installAudioEngineConfigObserverIfNeeded() + } + EventReporter.shared.capture( + level: .warning, + engine: "parakeet", + event: "audio_engine_rebuilt", + message: "Audio engine replaced after zombie-state detection", + context: ["reason": "zombie_engine_recovery"] + ) + return true + } + + func currentAudioGraphOwnerToken() -> ParakeetAudioGraphOwnerToken { + ParakeetAudioGraphOwnerToken(generation: audioGraphGeneration, engine: audioEngine) + } + + func ownsAudioGraph(_ owner: ParakeetAudioGraphOwnerToken) -> Bool { + owner.matches(generation: audioGraphGeneration, engine: audioEngine) + } + + func currentAudioEngineQueueOwnerToken() -> ParakeetAudioEngineQueueOwnerToken { + ParakeetAudioEngineQueueOwnerToken( + generation: audioGraphGeneration, + engine: audioEngine, + queue: audioEngineQueue + ) + } + + func ownsAudioEngineQueue(_ owner: ParakeetAudioEngineQueueOwnerToken) -> Bool { + owner.matches( + generation: audioGraphGeneration, + engine: audioEngine, + queue: audioEngineQueue + ) + } + + private func canContinueZombieEngineRecovery( + generation: UInt64, + expectedOwner: ParakeetAudioGraphOwnerToken + ) -> Bool { + ParakeetZombieRecoveryOwnershipPolicy.canContinue( + taskIsCancelled: Task.isCancelled, + recoveryIsCurrent: zombieRecoveryState.canContinue(generation: generation), + expectedOwner: expectedOwner, + currentGraphGeneration: audioGraphGeneration, + currentEngine: audioEngine + ) + } + + private func clearZombieRecoveryStartGeneration(ifMatching generation: UInt64) { + guard zombieRecoveryStartGeneration == generation else { return } + zombieRecoveryStartGeneration = nil + } + + private func finishZombieEngineRecovery( + generation: UInt64, + result: ParakeetZombieRecoveryResult + ) { + guard let terminal = zombieRecoveryState.finish(result: result, generation: generation) else { return } + zombieRecoveryTask = nil + reportZombieEngineRecoveryTerminal(terminal) + } + + private func reportZombieEngineRecoveryTerminal(_ terminal: ParakeetZombieRecoveryTerminal) { + let context = zombieRecoveryTelemetryContext( + failureKind: terminal.failureKind, + stage: terminal.stage, + result: terminal.result + ) + AnalyticsReporter.track("dictation_zombie_recovery_finished", properties: context) + + switch terminal.result { + case .succeeded: + EventReporter.shared.capture( + level: .info, + engine: "parakeet", + event: "zombie_engine_recovered", + message: "Audio engine recovered after bounded replacement", + context: context + ) + case .failed: + EventReporter.shared.capture( + level: .error, + engine: "parakeet", + event: "zombie_engine_recovery_failed", + message: "Audio engine could not recover after bounded replacement", + context: context + ) + case .cancelled: + EventReporter.shared.capture( + level: .info, + engine: "parakeet", + event: "zombie_engine_recovery_cancelled", + message: "Audio engine recovery was cancelled", + context: context + ) + } + } + + private func zombieRecoveryTelemetryContext( + failureKind: String, + stage: ParakeetZombieRecoveryStage, + result: ParakeetZombieRecoveryResult? + ) -> [String: String] { + let route = dictationRouteAnalyticsContext(selection: cachedInputDeviceSelection) + var context: [String: String] = [ + "failure_kind": failureKind, + "hfp_suspected": route["hfp_suspected"] ?? "false", + "input_device_class": route["input_device_class"] ?? "unknown", + "output_device_class": route["output_device_class"] ?? "unknown", + "route_shape": route["route_shape"] ?? "unknown", + "stage": stage.rawValue, + ] + if let result { + context["result"] = result.rawValue + } + return context } func stopRecording() async { @@ -2069,6 +2430,7 @@ class ParakeetEngine: ObservableObject { return } + let pendingRestoreOwner = pendingSystemInputRestore.owner guard isRecording else { // Genuinely preserved/recovered audio (e.g. real pre-sleep audio held // across a wake-recovery gap) must win over a merely-pending zombie @@ -2076,7 +2438,10 @@ class ParakeetEngine: ObservableObject { // zombie retry drains real audio instead of discarding it. if preservingRecordingAcrossRecovery || !recoveredRecordingTimeline.isEmpty { cancelPendingRecordingRecovery() - await restorePendingSystemInputAfterRecording(operation: "stop_recording_preserved_recovery") + await restorePendingSystemInputAfterRecording( + ownedBy: pendingRestoreOwner, + operation: "stop_recording_preserved_recovery" + ) return } // A zombie reset marks recording idle while it waits to retry, with @@ -2084,24 +2449,44 @@ class ParakeetEngine: ObservableObject { // as cancellation of the pending restart. if zombieRecoveryRestartPending { audioGraphGeneration += 1 + let stopGraphGeneration = audioGraphGeneration cancelAudioWatchdog() + audioStartAdmission.cancel() clearRecoveredRecordingTimeline(keepingCapacity: true) - await restorePendingSystemInputAfterRecording(operation: "stop_recording_zombie_restart") + await releaseIdleAudioHardware( + removeTap: true, + expectedGeneration: stopGraphGeneration + ) + await restorePendingSystemInputAfterRecording( + ownedBy: pendingRestoreOwner, + operation: "stop_recording_zombie_restart" + ) return } if audioStartInProgress { + audioStartAdmission.cancel() audioGraphGeneration += 1 } else { clearRecoveredRecordingTimeline(keepingCapacity: true) } - await restorePendingSystemInputAfterRecording(operation: "stop_recording_idle") + await restorePendingSystemInputAfterRecording( + ownedBy: pendingRestoreOwner, + operation: "stop_recording_idle" + ) return } audioGraphGeneration += 1 cancelAudioWatchdog() + let stopOwner = currentAudioEngineQueueOwnerToken() await removeRecordingTap() + guard ownsAudioEngineQueue(stopOwner) else { return } await stopAudioEngine() - await restorePendingSystemInputAfterRecording(operation: "stop_recording") + guard ownsAudioEngineQueue(stopOwner) else { return } + await restorePendingSystemInputAfterRecording( + ownedBy: pendingRestoreOwner, + operation: "stop_recording" + ) + guard ownsAudioEngineQueue(stopOwner) else { return } isEnginePrewarmed = false drainPendingSamplesIntoSampleBuffer() isRecording = false @@ -2155,6 +2540,7 @@ class ParakeetEngine: ObservableObject { private func cancelPendingRecordingRecovery() { audioGraphGeneration += 1 cancelAudioWatchdog() + audioStartAdmission.cancel() prewarmRetryTask?.cancel() prewarmRetryTask = nil configChangeDebounceTask?.cancel() @@ -2638,10 +3024,12 @@ class ParakeetEngine: ObservableObject { // MARK: - Cleanup func resetAfterFailedRecordingStart() async { + let pendingRestoreOwner = pendingSystemInputRestore.owner sharedMeetingMicTransition.invalidate() sharedMeetingMicRecorder.cancel() sharedMeetingMicRecording = false cancelAudioWatchdog() + audioStartAdmission.cancel() prewarmRetryTask?.cancel() prewarmRetryTask = nil configChangeDebounceTask?.cancel() @@ -2656,6 +3044,9 @@ class ParakeetEngine: ObservableObject { pendingSamplesLock.withLock { pendingSamples.removeAll(keepingCapacity: true) } + audioGraphGeneration += 1 + let failedStartCleanupOwner = currentAudioEngineQueueOwnerToken() + guard ownsAudioEngineQueue(failedStartCleanupOwner) else { return } isRecording = false isTranscribing = false audioLevel = 0 @@ -2664,17 +3055,23 @@ class ParakeetEngine: ObservableObject { recordingStartedOnLikelyBluetoothHandsFreeRoute = false sampleBuffer.removeAll(keepingCapacity: true) clearRecoveredRecordingTimeline(keepingCapacity: true) - audioGraphGeneration += 1 - let cleanupGeneration = audioGraphGeneration - await releaseIdleAudioHardware(removeTap: true, expectedGeneration: cleanupGeneration) - await restorePendingSystemInputAfterRecording(operation: "reset_after_failed_recording_start") + guard await releaseIdleAudioHardware( + removeTap: true, + expectedGeneration: failedStartCleanupOwner.graphOwner.generation + ) != nil else { return } + await restorePendingSystemInputAfterRecording( + ownedBy: pendingRestoreOwner, + operation: "reset_after_failed_recording_start" + ) } func abandonBlockedRecordingStart(reason: String) { + let pendingRestoreOwner = pendingSystemInputRestore.owner sharedMeetingMicTransition.invalidate() sharedMeetingMicRecorder.cancel() sharedMeetingMicRecording = false cancelAudioWatchdog() + audioStartAdmission.cancel() prewarmRetryTask?.cancel() prewarmRetryTask = nil configChangeDebounceTask?.cancel() @@ -2691,22 +3088,26 @@ class ParakeetEngine: ObservableObject { } isRecording = false isTranscribing = false - audioStartInProgress = false audioLevel = 0 didReceiveAudioSamples = false didReceiveNonZeroAudioSamples = false recordingStartedOnLikelyBluetoothHandsFreeRoute = false sampleBuffer.removeAll(keepingCapacity: true) clearRecoveredRecordingTimeline(keepingCapacity: true) - schedulePendingSystemInputRestore(operation: "abandon_blocked_recording_start") + schedulePendingSystemInputRestore( + ownedBy: pendingRestoreOwner, + operation: "abandon_blocked_recording_start" + ) abandonBlockedAudioEngine(reason: reason) } func cancel() { + let pendingRestoreOwner = pendingSystemInputRestore.owner sharedMeetingMicTransition.invalidate() sharedMeetingMicRecorder.cancel() sharedMeetingMicRecording = false cancelAudioWatchdog() + audioStartAdmission.cancel() prewarmRetryTask?.cancel() prewarmRetryTask = nil configChangeDebounceTask?.cancel() @@ -2726,7 +3127,7 @@ class ParakeetEngine: ObservableObject { } audioGraphGeneration += 1 let cleanupGeneration = audioGraphGeneration - schedulePendingSystemInputRestore(operation: "cancel") + schedulePendingSystemInputRestore(ownedBy: pendingRestoreOwner, operation: "cancel") Task { @MainActor [weak self] in await self?.releaseIdleAudioHardware(removeTap: true, expectedGeneration: cleanupGeneration) } @@ -2735,38 +3136,74 @@ class ParakeetEngine: ObservableObject { isTranscribing = false } - private func releaseIdleAudioHardware(removeTap: Bool, expectedGeneration: Int? = nil) async { + @discardableResult + private func releaseIdleAudioHardware( + removeTap: Bool, + expectedGeneration: Int? = nil + ) async -> ParakeetAudioEngineQueueOwnerToken? { if let expectedGeneration, expectedGeneration != audioGraphGeneration { - return + return nil } audioGraphGeneration += 1 - let cleanupGeneration = audioGraphGeneration + let idleCleanupOwner = currentAudioEngineQueueOwnerToken() if removeTap { await removeRecordingTap(force: true) } - if expectedGeneration != nil, cleanupGeneration != audioGraphGeneration { - return - } + guard ownsAudioEngineQueue(idleCleanupOwner) else { return nil } await stopAudioEngine() - if expectedGeneration != nil, cleanupGeneration != audioGraphGeneration { + guard ownsAudioEngineQueue(idleCleanupOwner) else { return nil } + isEnginePrewarmed = false + return idleCleanupOwner + } + + private func cancelAudioWatchdogForRecordingStart() { + audioWatchdogTask?.cancel() + audioWatchdogTask = nil + guard let zombieRecoveryStartGeneration, + zombieRecoveryState.canContinue(generation: zombieRecoveryStartGeneration) else { + cancelZombieEngineRecovery() return } - isEnginePrewarmed = false + } + + private func cancelZombieEngineRecovery() { + zombieRecoveryTask?.cancel() + if let blockedLease = zombieEngineWorkOwnership.claimPendingWorkForSuccessor( + currentEngine: audioEngine, + currentQueue: audioEngineQueue + ) { + // Cancellation advances logical ownership before this method runs. + // If reset or restart work still owns these exact resources, replace + // both the engine and queue before successor cleanup can enqueue. + let reason = blockedLease.phase == .zombieReset + ? "zombie_engine_reset_cancelled" + : "zombie_engine_recovery_start_cancelled" + if blockedLease.phase == .zombieRecoveryStart { + audioStartAdmission.finish(owner: blockedLease.owner) + } + abandonBlockedAudioEngine(reason: reason) + } + zombieRecoveryTask = nil + zombieRecoveryStartGeneration = nil + guard let terminal = zombieRecoveryState.cancelActiveAttempt() else { return } + reportZombieEngineRecoveryTerminal(terminal) } func cancelAudioWatchdog() { audioWatchdogTask?.cancel() audioWatchdogTask = nil - zombieRecoveryRestartPending = false + cancelZombieEngineRecovery() } func cleanup() { + let pendingRestoreOwner = pendingSystemInputRestore.owner sharedMeetingMicTransition.invalidate() sharedMeetingMicRecorder.cancel() sharedMeetingMicRecording = false isShuttingDown = true cancelModelWork() cancelAudioWatchdog() + audioStartAdmission.cancel() prewarmRetryTask?.cancel() prewarmRetryTask = nil configChangeDebounceTask?.cancel() @@ -2776,7 +3213,7 @@ class ParakeetEngine: ObservableObject { cancelConfigRecoveryTimeout() audioGraphGeneration += 1 let cleanupGeneration = audioGraphGeneration - schedulePendingSystemInputRestore(operation: "cleanup") + schedulePendingSystemInputRestore(ownedBy: pendingRestoreOwner, operation: "cleanup") Task { @MainActor [weak self] in await self?.releaseIdleAudioHardware(removeTap: true, expectedGeneration: cleanupGeneration) } diff --git a/Sources/Speech/ParakeetRecoveryState.swift b/Sources/Speech/ParakeetRecoveryState.swift index 57a84104..9c3539d2 100644 --- a/Sources/Speech/ParakeetRecoveryState.swift +++ b/Sources/Speech/ParakeetRecoveryState.swift @@ -58,6 +58,338 @@ struct ParakeetRecoveryState: Equatable { } } +struct ParakeetCategoricalAudioRoute: Equatable { + let inputDeviceClass: String + let outputDeviceClass: String + let routeShape: String +} + +/// Coalesces noisy CoreAudio notifications into one categorical route transition. +/// The side-effecting recovery path still runs for every debounced config-change +/// burst; this state only decides whether that burst is new analytics signal. +struct ParakeetRouteTransitionDebounceState: Equatable { + private(set) var stableRoute: ParakeetCategoricalAudioRoute? + private(set) var pendingRoute: ParakeetCategoricalAudioRoute? + + mutating func seedStableRouteIfNeeded(_ route: ParakeetCategoricalAudioRoute) { + guard stableRoute == nil else { return } + stableRoute = route + } + + mutating func observe(_ route: ParakeetCategoricalAudioRoute) { + pendingRoute = route + } + + mutating func commitPendingRoute() -> ParakeetCategoricalAudioRoute? { + guard let pendingRoute else { return nil } + self.pendingRoute = nil + + guard let stableRoute else { + self.stableRoute = pendingRoute + return nil + } + guard pendingRoute != stableRoute else { return nil } + + self.stableRoute = pendingRoute + return pendingRoute + } + + mutating func discardPendingRoute() { + pendingRoute = nil + } +} + +struct ParakeetAudioGraphOwnerToken: Equatable, Sendable { + let generation: Int + let engineIdentity: ObjectIdentifier + + init(generation: Int, engine: AnyObject) { + self.generation = generation + engineIdentity = ObjectIdentifier(engine) + } + + func matches(generation: Int, engine: AnyObject) -> Bool { + self.generation == generation && engineIdentity == ObjectIdentifier(engine) + } +} + +struct ParakeetAudioEngineQueueOwnerToken: Equatable, Sendable { + let graphOwner: ParakeetAudioGraphOwnerToken + let queueIdentity: ObjectIdentifier + + init(generation: Int, engine: AnyObject, queue: AnyObject) { + graphOwner = ParakeetAudioGraphOwnerToken(generation: generation, engine: engine) + queueIdentity = ObjectIdentifier(queue) + } + + func matches(generation: Int, engine: AnyObject, queue: AnyObject) -> Bool { + graphOwner.matches(generation: generation, engine: engine) + && queueIdentity == ObjectIdentifier(queue) + } + + func matchesResources(engine: AnyObject, queue: AnyObject) -> Bool { + graphOwner.engineIdentity == ObjectIdentifier(engine) + && queueIdentity == ObjectIdentifier(queue) + } +} + +/// Owns the single admitted audio-start task. Finishing or cancelling an older +/// start cannot clear a successor that already owns a replacement graph. +struct ParakeetAudioStartAdmissionState: Equatable { + private(set) var owner: ParakeetAudioEngineQueueOwnerToken? + + var isInProgress: Bool { + owner != nil + } + + mutating func begin(owner: ParakeetAudioEngineQueueOwnerToken) -> Bool { + guard self.owner == nil else { return false } + self.owner = owner + return true + } + + mutating func transfer( + from previousOwner: ParakeetAudioEngineQueueOwnerToken, + to nextOwner: ParakeetAudioEngineQueueOwnerToken + ) -> Bool { + guard owner == previousOwner else { return false } + owner = nextOwner + return true + } + + @discardableResult + mutating func finish(owner: ParakeetAudioEngineQueueOwnerToken) -> Bool { + guard self.owner == owner else { return false } + self.owner = nil + return true + } + + @discardableResult + mutating func cancel() -> ParakeetAudioEngineQueueOwnerToken? { + defer { owner = nil } + return owner + } +} + +enum ParakeetTimedAudioEngineWorkPhase: String, Equatable, Sendable { + case zombieReset + case zombieRecoveryStart +} + +struct ParakeetTimedAudioEngineWorkLease: Equatable, Sendable { + let owner: ParakeetAudioEngineQueueOwnerToken + let phase: ParakeetTimedAudioEngineWorkPhase +} + +/// Thread-safe ownership for one bounded zombie-recovery engine operation. +/// Completion clears only its exact lease; a newer MainActor owner can claim +/// still-pending work by engine+queue identity after advancing the generation. +final class ParakeetTimedAudioEngineWorkOwnership: @unchecked Sendable { + private let lock = NSLock() + private var pendingLease: ParakeetTimedAudioEngineWorkLease? + + func begin( + owner: ParakeetAudioEngineQueueOwnerToken, + phase: ParakeetTimedAudioEngineWorkPhase + ) { + lock.lock() + pendingLease = ParakeetTimedAudioEngineWorkLease(owner: owner, phase: phase) + lock.unlock() + } + + @discardableResult + func finish( + owner: ParakeetAudioEngineQueueOwnerToken, + phase: ParakeetTimedAudioEngineWorkPhase + ) -> Bool { + lock.lock() + defer { lock.unlock() } + let lease = ParakeetTimedAudioEngineWorkLease(owner: owner, phase: phase) + guard pendingLease == lease else { return false } + pendingLease = nil + return true + } + + func claimPendingWorkForSuccessor( + currentEngine: AnyObject, + currentQueue: AnyObject + ) -> ParakeetTimedAudioEngineWorkLease? { + lock.lock() + defer { lock.unlock() } + guard let pendingLease, + pendingLease.owner.matchesResources(engine: currentEngine, queue: currentQueue) else { + return nil + } + self.pendingLease = nil + return pendingLease + } +} + +/// Serializes system-input selection, apply, and restore work. A replacement +/// recording therefore observes and applies its route only after any restore +/// already taken by an older graph owner has completed. +final class ParakeetSerialSystemInputWorkCoordinator: @unchecked Sendable { + private let queue: DispatchQueue + + init(label: String) { + queue = DispatchQueue(label: label, qos: .utility) + } + + func run(_ work: @escaping () -> T) async -> T { + await withCheckedContinuation { continuation in + queue.async { + continuation.resume(returning: work()) + } + } + } + + func schedule(_ work: @escaping () -> Void) { + queue.async { + work() + } + } +} + +/// Pending route state is consumed by the graph owner that captured it. A new +/// recording replaces the entry with a new owner, making delayed cleanup a +/// no-op instead of restoring the replacement recording's system input. +struct ParakeetOwnerBoundPendingState: Equatable { + private struct Entry: Equatable { + let owner: ParakeetAudioGraphOwnerToken + let value: Value + } + + private var entry: Entry? + + var owner: ParakeetAudioGraphOwnerToken? { + entry?.owner + } + + var hasPendingValue: Bool { + entry != nil + } + + mutating func replace(_ value: Value, ownedBy owner: ParakeetAudioGraphOwnerToken) { + entry = Entry(owner: owner, value: value) + } + + mutating func clear() { + entry = nil + } + + @discardableResult + mutating func clear(ownedBy owner: ParakeetAudioGraphOwnerToken) -> Bool { + guard entry?.owner == owner else { return false } + entry = nil + return true + } + + mutating func take(ownedBy owner: ParakeetAudioGraphOwnerToken) -> Value? { + guard entry?.owner == owner else { return nil } + defer { entry = nil } + return entry?.value + } + + func value(ownedBy owner: ParakeetAudioGraphOwnerToken) -> Value? { + guard entry?.owner == owner else { return nil } + return entry?.value + } +} + +enum ParakeetZombieRecoveryOwnershipPolicy { + static func canContinue( + taskIsCancelled: Bool, + recoveryIsCurrent: Bool, + expectedOwner: ParakeetAudioGraphOwnerToken, + currentGraphGeneration: Int, + currentEngine: AnyObject + ) -> Bool { + !taskIsCancelled + && recoveryIsCurrent + && expectedOwner.matches(generation: currentGraphGeneration, engine: currentEngine) + } +} + +enum ParakeetZombieRecoveryStage: String, Equatable { + case detected + case reset + case settle + case restart +} + +enum ParakeetZombieRecoveryResult: String, Equatable { + case succeeded + case failed + case cancelled +} + +struct ParakeetZombieRecoveryTerminal: Equatable { + let generation: UInt64 + let stage: ParakeetZombieRecoveryStage + let result: ParakeetZombieRecoveryResult + let failureKind: String +} + +/// Generation-gated lifecycle for the single bounded zombie-engine retry. +/// `finish` consumes the active attempt so every generation can produce at most +/// one terminal telemetry result, even when cancellation races a late callback. +struct ParakeetZombieRecoveryState: Equatable { + private struct Attempt: Equatable { + let generation: UInt64 + let failureKind: String + var stage: ParakeetZombieRecoveryStage + } + + private(set) var generation: UInt64 = 0 + private var activeAttempt: Attempt? + + var isActive: Bool { + activeAttempt != nil + } + + mutating func begin(failureKind: String) -> UInt64 { + if let activeAttempt { + return activeAttempt.generation + } + generation &+= 1 + activeAttempt = Attempt( + generation: generation, + failureKind: failureKind, + stage: .detected + ) + return generation + } + + mutating func advance(to stage: ParakeetZombieRecoveryStage, generation: UInt64) -> Bool { + guard activeAttempt?.generation == generation else { return false } + activeAttempt?.stage = stage + return true + } + + func canContinue(generation: UInt64) -> Bool { + activeAttempt?.generation == generation + } + + mutating func finish( + result: ParakeetZombieRecoveryResult, + generation: UInt64 + ) -> ParakeetZombieRecoveryTerminal? { + guard let attempt = activeAttempt, attempt.generation == generation else { return nil } + activeAttempt = nil + return ParakeetZombieRecoveryTerminal( + generation: generation, + stage: attempt.stage, + result: result, + failureKind: attempt.failureKind + ) + } + + mutating func cancelActiveAttempt() -> ParakeetZombieRecoveryTerminal? { + guard let attempt = activeAttempt else { return nil } + return finish(result: .cancelled, generation: attempt.generation) + } +} + struct ParakeetAudioStartRecoveryPolicy: Equatable { static func shouldRetryStartFailure( isRecoveryAttempt: Bool, diff --git a/Tests/AnalyticsEventPolicyTests.swift b/Tests/AnalyticsEventPolicyTests.swift index d2fb9f6b..9db7fe0b 100644 --- a/Tests/AnalyticsEventPolicyTests.swift +++ b/Tests/AnalyticsEventPolicyTests.swift @@ -1079,8 +1079,15 @@ func testAnalyticsEventPolicy() { let finished = AnalyticsEventPolicy.policy(forEvent: "dictation_audio_route_recovery_finished") let timeout = AnalyticsEventPolicy.policy(forEvent: "dictation_audio_route_recovery_timeout") + assertEqual( + changed?.allowedProperties ?? [], + ["default_input_class", "default_output_class", "format_ready", "hfp_suspected", "input_device_class", "output_device_class", "recovering", "route_shape", "sample_flow_started", "selected_input_class", "selection_overrode_default", "selection_reason", "was_recording"], + "stable route changes should permit only categorical route and state fields" + ) assertEqual(changed?.allowedProperties.contains("was_recording"), true, "route change should preserve whether an active recording was interrupted") assertEqual(changed?.allowedProperties.contains("selected_input_class"), true, "route change should preserve selected input class") + assertEqual(changed?.allowedProperties.contains("input_rate_hz"), false, "stable route changes should not expose exact input rates") + assertEqual(changed?.allowedProperties.contains("output_rate_hz"), false, "stable route changes should not expose exact output rates") assertEqual(finished?.allowedProperties.contains("outcome"), true, "route recovery should preserve success/failure") assertEqual(finished?.allowedProperties.contains("recovery_latency_bucket"), true, "route recovery should preserve latency as a bucket") assertEqual(timeout?.allowedProperties.contains("hfp_suspected"), true, "route timeout should preserve Bluetooth HFP suspicion only as a boolean") @@ -1100,6 +1107,36 @@ func testAnalyticsEventPolicy() { assertEqual(sanitized["was_recording"], "true", "recording interruption state should survive sanitization") } + runSuite("AnalyticsEventPolicy keeps zombie recovery terminal telemetry categorical") { + let terminal = AnalyticsEventPolicy.policy(forEvent: "dictation_zombie_recovery_finished") + + assertEqual( + terminal?.allowedProperties ?? [], + ["failure_kind", "hfp_suspected", "input_device_class", "output_device_class", "result", "route_shape", "stage"], + "zombie recovery should expose one categorical terminal payload" + ) + + let sanitized = AnalyticsPayloadSanitizer.sanitizeProperties( + [ + "audio_device": "Private microphone name", + "failure_kind": "no_sample_callbacks", + "hfp_suspected": "false", + "input_device_class": "built_in", + "output_device_class": "built_in", + "result": "failed", + "route_shape": "built_in_input_to_built_in_output", + "sample_count": "0", + "stage": "restart", + ], + allowedKeys: terminal?.allowedProperties ?? [] + ) + + assertEqual(sanitized["result"], "failed", "terminal result should survive") + assertEqual(sanitized["stage"], "restart", "coarse recovery stage should survive") + assertNil(sanitized["audio_device"], "raw device labels must stay out of zombie telemetry") + assertNil(sanitized["sample_count"], "exact callback counts must stay out of zombie telemetry") + } + runSuite("AnalyticsEventPolicy only permits reviewed analytics events") { let dictationStartFailed = AnalyticsEventPolicy.policy(forEvent: "dictation_start_failed") let dictationCompleted = AnalyticsEventPolicy.policy(forEvent: "dictation_completed") diff --git a/Tests/BluetoothRouteContractTests.swift b/Tests/BluetoothRouteContractTests.swift index ce44d97c..40cb0f9d 100644 --- a/Tests/BluetoothRouteContractTests.swift +++ b/Tests/BluetoothRouteContractTests.swift @@ -364,9 +364,10 @@ func testBluetoothRouteContract() { } let snapshotBody = String(source[snapshotStart.lowerBound.. Bool"), "failed starts should share one bounded-retry path" @@ -418,28 +434,37 @@ func testBluetoothRouteContract() { "retryable start failures should keep the temporary built-in input stable instead of restoring and reapplying it" ) assertTrue( - stopBody.contains("await restorePendingSystemInputAfterRecording(operation: \"stop_recording\")"), - "normal stop should restore the prior system input if Transcripted still owns the temporary input" + stopBody.contains("operation: \"stop_recording\"") + && stopBody.contains("ownedBy: pendingRestoreOwner"), + "normal stop should restore the prior system input only for its captured owner" ) assertTrue( - stopBody.contains("await restorePendingSystemInputAfterRecording(operation: \"stop_recording_idle\")"), - "canceled or interrupted start paths should not leave the temporary system input behind" + stopBody.contains("operation: \"stop_recording_idle\"") + && stopBody.contains("ownedBy: pendingRestoreOwner"), + "canceled or interrupted start paths should restore only their captured temporary input" ) assertTrue( - cleanupBody.contains("schedulePendingSystemInputRestore(operation: \"cancel\")"), - "explicit cancellation should restore the temporary system input even though cancel() is synchronous" + cleanupBody.contains("schedulePendingSystemInputRestore(ownedBy: pendingRestoreOwner, operation: \"cancel\")"), + "explicit cancellation should restore its owned temporary input even though cancel() is synchronous" ) assertTrue( - cleanupBody.contains("schedulePendingSystemInputRestore(operation: \"cleanup\")"), - "quit cleanup should restore the temporary system input even though cleanup() is synchronous" + cleanupBody.contains("schedulePendingSystemInputRestore(ownedBy: pendingRestoreOwner, operation: \"cleanup\")"), + "quit cleanup should restore its owned temporary input even though cleanup() is synchronous" ) assertTrue( - cleanupBody.contains("schedulePendingSystemInputRestore(operation: \"abandon_blocked_recording_start\")"), - "blocked-start abandonment should restore the temporary system input" + source.contains("let restoreError = await Self.systemInputWorkCoordinator.run") + && source.contains("Self.systemInputWorkCoordinator.schedule"), + "awaited and scheduled restores should share the replacement start's serial route coordinator" ) assertTrue( - cleanupBody.contains("await restorePendingSystemInputAfterRecording(operation: \"reset_after_failed_recording_start\")"), - "final failed-start cleanup should restore the temporary system input after retries are exhausted" + cleanupBody.contains("operation: \"abandon_blocked_recording_start\"") + && cleanupBody.contains("ownedBy: pendingRestoreOwner"), + "blocked-start abandonment should restore only the temporary input owned by that start" + ) + assertTrue( + cleanupBody.contains("operation: \"reset_after_failed_recording_start\"") + && cleanupBody.contains("ownedBy: pendingRestoreOwner"), + "final failed-start cleanup should restore its owned temporary input after retries are exhausted" ) assertTrue( source.contains("DictationPersistentInputPreferences.setTemporaryRecoveryMarker(recoveryMarker)"), @@ -456,7 +481,7 @@ func testBluetoothRouteContract() { // (codebase audit 2026-07-08 wave 2). let source = readBluetoothRouteContractFile("Sources/Speech/ParakeetDeviceRecovery.swift") guard let handlerStart = source.range(of: "private func handleAudioConfigChange() async"), - let handlerEnd = source.range(of: "private func recordRouteChangeAnalytics", range: handlerStart.upperBound.. B -> A notification churn is not a stable route change") + assertEqual(state.stableRoute, builtIn, "oscillation should preserve the original stable route") + } + + runSuite("ParakeetRouteTransitionDebounceState treats the first known route as a baseline") { + let builtIn = categoricalRoute(input: "built_in", output: "built_in", shape: "built_in_to_built_in") + var state = ParakeetRouteTransitionDebounceState() + + state.observe(builtIn) + + assertEqual(state.commitPendingRoute(), nil, "initial discovery should seed a baseline instead of claiming a transition") + assertEqual(state.stableRoute, builtIn, "initial discovery should become the stable baseline") + } + + runSuite("ParakeetZombieRecoveryOwnershipPolicy accepts only the exact active graph owner") { + let engine = NSObject() + let owner = ParakeetAudioGraphOwnerToken(generation: 7, engine: engine) + + assertTrue( + ParakeetZombieRecoveryOwnershipPolicy.canContinue( + taskIsCancelled: false, + recoveryIsCurrent: true, + expectedOwner: owner, + currentGraphGeneration: 7, + currentEngine: engine + ), + "the active task should mutate only its exact graph generation and engine" + ) + } + + runSuite("ParakeetZombieRecoveryOwnershipPolicy rejects stale reset interleavings") { + let staleEngine = NSObject() + let healthyReplacement = NSObject() + let owner = ParakeetAudioGraphOwnerToken(generation: 11, engine: staleEngine) + + assertFalse( + ParakeetZombieRecoveryOwnershipPolicy.canContinue( + taskIsCancelled: true, + recoveryIsCurrent: true, + expectedOwner: owner, + currentGraphGeneration: 11, + currentEngine: staleEngine + ), + "a stop cancellation must not enter or complete graph recreation" + ) + assertFalse( + ParakeetZombieRecoveryOwnershipPolicy.canContinue( + taskIsCancelled: false, + recoveryIsCurrent: false, + expectedOwner: owner, + currentGraphGeneration: 11, + currentEngine: staleEngine + ), + "a config-change cancellation must make the old zombie generation stale" + ) + assertFalse( + ParakeetZombieRecoveryOwnershipPolicy.canContinue( + taskIsCancelled: false, + recoveryIsCurrent: true, + expectedOwner: owner, + currentGraphGeneration: 12, + currentEngine: staleEngine + ), + "a newer graph owner using the same engine must not be abandoned by an old timeout" + ) + assertFalse( + ParakeetZombieRecoveryOwnershipPolicy.canContinue( + taskIsCancelled: false, + recoveryIsCurrent: true, + expectedOwner: owner, + currentGraphGeneration: 11, + currentEngine: healthyReplacement + ), + "a stale reset completion must not mutate a healthy replacement engine" + ) + assertFalse( + ParakeetZombieRecoveryOwnershipPolicy.canContinue( + taskIsCancelled: false, + recoveryIsCurrent: true, + expectedOwner: owner, + currentGraphGeneration: 12, + currentEngine: healthyReplacement + ), + "generation and identity must both match before shared reset state changes" + ) + } + + runSuite("ParakeetAudioGraphOwnerToken preserves a newer tap after delayed cleanup") { + let retiredEngine = NSObject() + let replacementEngine = NSObject() + let delayedCleanupOwner = ParakeetAudioGraphOwnerToken(generation: 20, engine: retiredEngine) + + var replacementTapInstalled = true + if delayedCleanupOwner.matches(generation: 21, engine: replacementEngine) { + replacementTapInstalled = false + } + assertTrue( + replacementTapInstalled, + "old cleanup completion must not clear a replacement engine's installed tap" + ) + + var newerGenerationTapInstalled = true + if delayedCleanupOwner.matches(generation: 21, engine: retiredEngine) { + newerGenerationTapInstalled = false + } + assertTrue( + newerGenerationTapInstalled, + "old cleanup completion must not clear a newer generation's tap on the same engine" + ) + } + + runSuite("ParakeetTimedAudioEngineWorkOwnership moves successor work off a blocked queue") { + let retiredEngine = NSObject() + let blockedQueue = DispatchQueue(label: "test.parakeet.blocked-engine-queue") + let owner = ParakeetAudioEngineQueueOwnerToken( + generation: 30, + engine: retiredEngine, + queue: blockedQueue + ) + let ownership = ParakeetTimedAudioEngineWorkOwnership() + let blockedWorkStarted = DispatchSemaphore(value: 0) + let releaseBlockedWork = DispatchSemaphore(value: 0) + let blockedWorkFinished = DispatchSemaphore(value: 0) + let staleCompletionMutatedState = DispatchSemaphore(value: 0) + ownership.begin(owner: owner, phase: .zombieReset) + + blockedQueue.async { + blockedWorkStarted.signal() + _ = releaseBlockedWork.wait(timeout: .now() + 2) + if ownership.finish(owner: owner, phase: .zombieReset) { + staleCompletionMutatedState.signal() + } + blockedWorkFinished.signal() + } + + assertTrue( + blockedWorkStarted.wait(timeout: .now() + 1) == .success, + "the old engine helper should be suspended on its serial queue" + ) + + let claimedOwner = ownership.claimPendingWorkForSuccessor( + currentEngine: retiredEngine, + currentQueue: blockedQueue + ) + assertEqual( + claimedOwner, + ParakeetTimedAudioEngineWorkLease(owner: owner, phase: .zombieReset), + "the successor must synchronously claim pending work on the exact blocked engine and queue" + ) + + let replacementEngine = NSObject() + let replacementQueue = DispatchQueue(label: "test.parakeet.replacement-engine-queue") + let replacementOwner = ParakeetAudioEngineQueueOwnerToken( + generation: 31, + engine: replacementEngine, + queue: replacementQueue + ) + assertTrue( + replacementOwner.matches( + generation: 31, + engine: replacementEngine, + queue: replacementQueue + ), + "successor replacement should own both a fresh engine and serial queue" + ) + let successorCleanupFinished = DispatchSemaphore(value: 0) + replacementQueue.async { + successorCleanupFinished.signal() + } + assertTrue( + successorCleanupFinished.wait(timeout: .now() + 1) == .success, + "successor cleanup must run on the replacement queue while old work remains blocked" + ) + + releaseBlockedWork.signal() + assertTrue( + blockedWorkFinished.wait(timeout: .now() + 1) == .success, + "the delayed old helper should finish after the test releases it" + ) + assertTrue( + staleCompletionMutatedState.wait(timeout: .now()) == .timedOut, + "old helper completion must not reclaim ownership after successor replacement" + ) + } + + runSuite("Parakeet recovery-start cancellation replaces a blocked engine and queue") { + let blockedEngine = NSObject() + let blockedQueue = DispatchQueue(label: "test.parakeet.blocked-recovery-start") + let blockedOwner = ParakeetAudioEngineQueueOwnerToken( + generation: 35, + engine: blockedEngine, + queue: blockedQueue + ) + let ownership = ParakeetTimedAudioEngineWorkOwnership() + var startAdmission = ParakeetAudioStartAdmissionState() + let resources = ParakeetEngineQueueTestResources( + engine: blockedEngine, + queue: blockedQueue + ) + var recoveryState = ParakeetZombieRecoveryState() + let recoveryGeneration = recoveryState.begin(failureKind: "no_sample_callbacks") + assertTrue( + recoveryState.advance(to: .restart, generation: recoveryGeneration), + "the production restart stage should own the leased start operation" + ) + assertTrue( + startAdmission.begin(owner: blockedOwner), + "the blocked recovery start should own the single start-admission slot" + ) + + let blockedStartEntered = DispatchSemaphore(value: 0) + let releaseBlockedStart = DispatchSemaphore(value: 0) + let blockedStartFinished = DispatchSemaphore(value: 0) + let lateCompletionMutatedResources = DispatchSemaphore(value: 0) + ownership.begin(owner: blockedOwner, phase: .zombieRecoveryStart) + blockedQueue.async { + blockedStartEntered.signal() + _ = releaseBlockedStart.wait(timeout: .now() + 2) + if ownership.finish(owner: blockedOwner, phase: .zombieRecoveryStart) { + resources.restoreOriginalResources() + lateCompletionMutatedResources.signal() + } + blockedStartFinished.signal() + } + + assertTrue( + blockedStartEntered.wait(timeout: .now() + 1) == .success, + "recovery install/start work should be suspended on its leased engine queue" + ) + + let claimedLease = ownership.claimPendingWorkForSuccessor( + currentEngine: blockedEngine, + currentQueue: blockedQueue + ) + assertEqual( + claimedLease, + ParakeetTimedAudioEngineWorkLease( + owner: blockedOwner, + phase: .zombieRecoveryStart + ), + "cancellation should claim the exact in-flight recovery-start lease" + ) + assertTrue( + startAdmission.finish(owner: blockedOwner), + "cancellation should release only the blocked start owner's admission" + ) + + let successorEngine = NSObject() + let successorQueue = DispatchQueue(label: "test.parakeet.recovery-start-successor") + let successorOwner = ParakeetAudioEngineQueueOwnerToken( + generation: 36, + engine: successorEngine, + queue: successorQueue + ) + resources.replace(engine: successorEngine, queue: successorQueue) + assertTrue( + startAdmission.begin(owner: successorOwner), + "the replacement graph should be able to admit a successor start immediately" + ) + let terminal = recoveryState.cancelActiveAttempt() + assertTrue( + resources.matches(engine: successorEngine, queue: successorQueue), + "cancellation should synchronously replace both blocked resources before publishing terminal state" + ) + assertEqual(terminal?.stage, .restart, "cancellation should terminate the blocked restart stage") + assertEqual(terminal?.result, .cancelled, "cancellation should publish one cancelled terminal") + + let successorWorkFinished = DispatchSemaphore(value: 0) + successorQueue.async { + successorWorkFinished.signal() + } + assertTrue( + successorWorkFinished.wait(timeout: .now() + 1) == .success, + "successor cleanup should not queue behind blocked recovery-start work" + ) + + releaseBlockedStart.signal() + assertTrue( + blockedStartFinished.wait(timeout: .now() + 1) == .success, + "the delayed recovery-start callback should finish after release" + ) + assertTrue( + lateCompletionMutatedResources.wait(timeout: .now()) == .timedOut, + "late old completion must not reclaim or mutate successor resources" + ) + assertTrue( + resources.matches(engine: successorEngine, queue: successorQueue), + "late old completion must leave the successor engine and queue intact" + ) + assertFalse( + startAdmission.finish(owner: blockedOwner), + "the stale start defer must not clear the successor's admission" + ) + assertEqual( + startAdmission.owner, + successorOwner, + "the successor must remain the admitted start after stale completion" + ) + } + + runSuite("Parakeet recovery cancellation releases a pre-lease admitted start") { + let oldEngine = NSObject() + let oldQueue = DispatchQueue(label: "test.parakeet.pre-lease-recovery-start") + let oldOwner = ParakeetAudioEngineQueueOwnerToken( + generation: 37, + engine: oldEngine, + queue: oldQueue + ) + var startAdmission = ParakeetAudioStartAdmissionState() + let timedWork = ParakeetTimedAudioEngineWorkOwnership() + assertTrue( + startAdmission.begin(owner: oldOwner), + "route selection and snapshot work should hold start admission before the timed start lease begins" + ) + assertNil( + timedWork.claimPendingWorkForSuccessor(currentEngine: oldEngine, currentQueue: oldQueue), + "pre-lease cancellation should not require a timed engine-work lease" + ) + + let cancelledOwner = startAdmission.cancel() + assertEqual(cancelledOwner, oldOwner, "stop or wake should release the admitted pre-lease start") + + let successorEngine = NSObject() + let successorQueue = DispatchQueue(label: "test.parakeet.pre-lease-successor") + let successorOwner = ParakeetAudioEngineQueueOwnerToken( + generation: 38, + engine: successorEngine, + queue: successorQueue + ) + assertTrue( + startAdmission.begin(owner: successorOwner), + "a successor should start immediately after pre-lease cancellation" + ) + assertFalse( + startAdmission.finish(owner: oldOwner), + "the stale pre-lease start defer must not clear successor admission" + ) + assertEqual( + startAdmission.owner, + successorOwner, + "successor admission should survive stale pre-lease completion" + ) + } + + await runSuite("ParakeetOwnerBoundPendingState rejects a truly delayed stale restore") { + let engine = NSObject() + let staleOwner = ParakeetAudioGraphOwnerToken(generation: 40, engine: engine) + let replacementOwner = ParakeetAudioGraphOwnerToken(generation: 41, engine: engine) + let state = ParakeetPendingRestoreInterleavingHarness() + let cleanupCapturedOwner = ParakeetAsyncInterleavingGate() + let allowCleanupCompletion = ParakeetAsyncInterleavingGate() + + await state.replace("old-route", ownedBy: staleOwner) + let delayedCleanup = Task { + let capturedOwner = await state.owner() + await cleanupCapturedOwner.open() + await allowCleanupCompletion.wait() + guard let capturedOwner else { return nil as String? } + return await state.take(ownedBy: capturedOwner) + } + + await cleanupCapturedOwner.wait() + await state.replace("replacement-route", ownedBy: replacementOwner) + await allowCleanupCompletion.open() + + let staleRestore = await delayedCleanup.value + let replacementValue = await state.value(ownedBy: replacementOwner) + let replacementRestore = await state.take(ownedBy: replacementOwner) + assertNil( + staleRestore, + "cleanup delayed across an await must not consume a replacement start's restore target" + ) + assertEqual( + replacementValue, + "replacement-route", + "the replacement target must remain installed after stale cleanup resumes" + ) + assertEqual( + replacementRestore, + "replacement-route", + "only the exact newer generation+engine owner may consume its route restore" + ) + } + + await runSuite("Parakeet pending restore preserves a same-input replacement after take") { + let engine = NSObject() + let oldOwner = ParakeetAudioGraphOwnerToken(generation: 45, engine: engine) + let replacementOwner = ParakeetAudioGraphOwnerToken(generation: 46, engine: engine) + let pendingState = ParakeetPendingRestoreInterleavingHarness() + let coordinator = ParakeetSerialSystemInputWorkCoordinator( + label: "test.parakeet.system-input-interleaving" + ) + let temporaryInput = "built-in-input" + let previousInput = "bluetooth-input" + let route = ParakeetSystemInputRouteTestState( + route: temporaryInput, + recoveryMarkerIsSet: true + ) + + await pendingState.replace(previousInput, ownedBy: oldOwner) + guard let takenRestore = await pendingState.take(ownedBy: oldOwner) else { + assertTrue(false, "the old owner should take its pending restore before suspension") + return + } + + let oldRestoreEntered = ParakeetAsyncInterleavingGate() + let releaseOldRestore = DispatchSemaphore(value: 0) + let oldRestore = Task { + await coordinator.run { + Task { await oldRestoreEntered.open() } + _ = releaseOldRestore.wait(timeout: .now() + 2) + route.restoreIfStillTemporary( + temporaryInput: temporaryInput, + previousInput: takenRestore + ) + } + if !(await pendingState.hasPendingValue()) { + route.clearRecoveryMarker() + } + } + + await oldRestoreEntered.wait() + + await pendingState.replace(previousInput, ownedBy: replacementOwner) + let replacementStartScheduled = ParakeetAsyncInterleavingGate() + let replacementSelectionRan = ParakeetAsyncInterleavingGate() + let replacementStart = Task { + await replacementStartScheduled.open() + let selectedFromRoute = await coordinator.run { + Task { await replacementSelectionRan.open() } + return route.currentRoute() + } + await coordinator.run { + route.applyReplacementInput(temporaryInput) + } + return selectedFromRoute + } + + await replacementStartScheduled.wait() + assertFalse( + await replacementSelectionRan.opened(), + "replacement route selection must queue behind the already-taken restore" + ) + + releaseOldRestore.signal() + await oldRestore.value + let replacementSelection = await replacementStart.value + + assertEqual( + replacementSelection, + previousInput, + "replacement selection should run after the old restore completes" + ) + assertEqual( + route.currentRoute(), + temporaryInput, + "old restore completion must not leave the replacement on the prior route" + ) + assertTrue( + route.recoveryMarkerIsSet(), + "old restore completion must not clear the replacement owner's recovery marker" + ) + assertEqual( + await pendingState.value(ownedBy: replacementOwner), + previousInput, + "old completion must leave the replacement owner's pending restore intact" + ) + } + + await runSuite("Successful system-input fallback restores after cancel wake or graph replacement") { + enum OwnershipLoss: String, CaseIterable { + case cancel + case wake + case graphRebuild + } + + enum SuspensionPoint: String, CaseIterable { + case systemOverride + case snapshotFailure + case snapshotSuccess + } + + for ownershipLoss in OwnershipLoss.allCases { + for suspensionPoint in SuspensionPoint.allCases { + let oldEngine = NSObject() + let oldQueue = DispatchQueue(label: "test.parakeet.stale-system-input.old.\(ownershipLoss.rawValue)") + let oldOwner = ParakeetAudioEngineQueueOwnerToken( + generation: 60, + engine: oldEngine, + queue: oldQueue + ) + let successorEngine: NSObject = ownershipLoss == .graphRebuild ? NSObject() : oldEngine + let successorQueue = ownershipLoss == .graphRebuild + ? DispatchQueue(label: "test.parakeet.stale-system-input.new.graph") + : oldQueue + let successorOwner = ParakeetAudioEngineQueueOwnerToken( + generation: 61, + engine: successorEngine, + queue: successorQueue + ) + let coordinator = ParakeetSerialSystemInputWorkCoordinator( + label: "test.parakeet.stale-system-input.\(ownershipLoss.rawValue).\(suspensionPoint.rawValue)" + ) + let route = ParakeetSystemInputRouteTestState( + route: "airpods-input", + recoveryMarkerIsSet: true + ) + let overrideEntered = ParakeetAsyncInterleavingGate() + let releaseOverride = DispatchSemaphore(value: 0) + + let staleSnapshot = Task { + await coordinator.run { + route.applyReplacementInput("built-in-input") + Task { await overrideEntered.open() } + _ = releaseOverride.wait(timeout: .now() + 2) + } + guard oldOwner == successorOwner else { + await coordinator.run { + route.restoreIfStillTemporary( + temporaryInput: "built-in-input", + previousInput: "airpods-input" + ) + } + return true + } + return false + } + + await overrideEntered.wait() + releaseOverride.signal() + let didRestoreBeforeCancellation = await staleSnapshot.value + + assertTrue( + didRestoreBeforeCancellation, + "\(ownershipLoss.rawValue) should make the awaited fallback owner stale" + ) + assertEqual( + route.currentRoute(), + "airpods-input", + "\(ownershipLoss.rawValue) at \(suspensionPoint.rawValue) must restore the prior system input before stale work cancels" + ) + } + } + } + + runSuite("ParakeetZombieRecoveryState emits exactly one terminal result per attempt") { + var state = ParakeetZombieRecoveryState() + let generation = state.begin(failureKind: "no_sample_callbacks") + + assertTrue(state.advance(to: .reset, generation: generation), "active recovery should advance into reset") + assertTrue(state.advance(to: .restart, generation: generation), "active recovery should advance into restart") + let terminal = state.finish(result: .failed, generation: generation) + + assertEqual(terminal?.stage, .restart, "terminal telemetry should preserve the last actionable stage") + assertEqual(terminal?.result, .failed, "terminal telemetry should preserve the outcome") + assertEqual(terminal?.failureKind, "no_sample_callbacks", "terminal telemetry should preserve the categorical trigger") + assertEqual(state.finish(result: .failed, generation: generation), nil, "the same attempt cannot finish twice") + } + + runSuite("ParakeetZombieRecoveryState cancellation is terminal and rejects stale callbacks") { + var state = ParakeetZombieRecoveryState() + let generation = state.begin(failureKind: "silent_hfp_callbacks") + assertTrue(state.advance(to: .settle, generation: generation), "active recovery should advance into settle") + + let terminal = state.cancelActiveAttempt() + + assertEqual(terminal?.stage, .settle, "cancellation should name the stage it interrupted") + assertEqual(terminal?.result, .cancelled, "cancellation should have a categorical terminal result") + assertFalse(state.canContinue(generation: generation), "cancelled work should become stale") + assertFalse(state.advance(to: .restart, generation: generation), "late callbacks cannot revive a cancelled recovery") + } + + runSuite("ParakeetZombieRecoveryState keeps one active generation") { + var state = ParakeetZombieRecoveryState() + let first = state.begin(failureKind: "no_sample_callbacks") + let duplicate = state.begin(failureKind: "silent_hfp_callbacks") + + assertEqual(duplicate, first, "a second detector callback must not replace an unfinished recovery attempt") + assertTrue(state.canContinue(generation: first), "the original attempt should remain active") + } + runSuite("ParakeetAudioStartRecoveryPolicy.shouldRetryStartFailure — retries only normal first failures") { assertTrue( ParakeetAudioStartRecoveryPolicy.shouldRetryStartFailure(isRecoveryAttempt: false, failedAttempts: 1, retryBudget: 1), @@ -236,3 +843,140 @@ func testParakeetRecoveryState() { ) } } + +private actor ParakeetPendingRestoreInterleavingHarness { + private var state = ParakeetOwnerBoundPendingState() + + func replace(_ value: String, ownedBy owner: ParakeetAudioGraphOwnerToken) { + state.replace(value, ownedBy: owner) + } + + func owner() -> ParakeetAudioGraphOwnerToken? { + state.owner + } + + func take(ownedBy owner: ParakeetAudioGraphOwnerToken) -> String? { + state.take(ownedBy: owner) + } + + func value(ownedBy owner: ParakeetAudioGraphOwnerToken) -> String? { + state.value(ownedBy: owner) + } + + func hasPendingValue() -> Bool { + state.hasPendingValue + } +} + +private final class ParakeetEngineQueueTestResources: @unchecked Sendable { + private let lock = NSLock() + private let originalEngine: AnyObject + private let originalQueue: DispatchQueue + private var engine: AnyObject + private var queue: DispatchQueue + + init(engine: AnyObject, queue: DispatchQueue) { + originalEngine = engine + originalQueue = queue + self.engine = engine + self.queue = queue + } + + func replace(engine: AnyObject, queue: DispatchQueue) { + lock.lock() + self.engine = engine + self.queue = queue + lock.unlock() + } + + func restoreOriginalResources() { + replace(engine: originalEngine, queue: originalQueue) + } + + func matches(engine: AnyObject, queue: DispatchQueue) -> Bool { + lock.lock() + defer { lock.unlock() } + return ObjectIdentifier(self.engine) == ObjectIdentifier(engine) + && ObjectIdentifier(self.queue) == ObjectIdentifier(queue) + } +} + +private final class ParakeetSystemInputRouteTestState: @unchecked Sendable { + private let lock = NSLock() + private var route: String + private var markerIsSet: Bool + + init(route: String, recoveryMarkerIsSet: Bool) { + self.route = route + markerIsSet = recoveryMarkerIsSet + } + + func restoreIfStillTemporary(temporaryInput: String, previousInput: String) { + lock.lock() + defer { lock.unlock() } + guard route == temporaryInput else { return } + route = previousInput + } + + func applyReplacementInput(_ input: String) { + lock.lock() + route = input + lock.unlock() + } + + func currentRoute() -> String { + lock.lock() + defer { lock.unlock() } + return route + } + + func clearRecoveryMarker() { + lock.lock() + markerIsSet = false + lock.unlock() + } + + func recoveryMarkerIsSet() -> Bool { + lock.lock() + defer { lock.unlock() } + return markerIsSet + } +} + +private actor ParakeetAsyncInterleavingGate { + private var isOpen = false + private var waiters: [CheckedContinuation] = [] + + func wait() async { + guard !isOpen else { return } + await withCheckedContinuation { continuation in + waiters.append(continuation) + } + } + + func opened() -> Bool { + isOpen + } + + func open() { + guard !isOpen else { return } + isOpen = true + let pendingWaiters = waiters + waiters.removeAll() + for waiter in pendingWaiters { + waiter.resume() + } + } +} + +private func categoricalRoute( + input: String, + output: String, + shape: String +) -> ParakeetCategoricalAudioRoute { + ParakeetCategoricalAudioRoute( + inputDeviceClass: input, + outputDeviceClass: output, + routeShape: shape + ) +} diff --git a/Tests/ParakeetStartRecordingFailurePolicyTests.swift b/Tests/ParakeetStartRecordingFailurePolicyTests.swift index 9cc14e5a..c95f9859 100644 --- a/Tests/ParakeetStartRecordingFailurePolicyTests.swift +++ b/Tests/ParakeetStartRecordingFailurePolicyTests.swift @@ -748,34 +748,150 @@ func testParakeetStartRecordingFailurePolicy() { ) } - runSuite("ParakeetEngine zombie watchdog marks recording idle before graph reset") { + runSuite("ParakeetEngine zombie watchdog uses a bounded fresh-engine recovery") { let source = readParakeetEngineSource() guard let watchdogStart = source.range(of: "private func startAudioWatchdog()"), - let watchdogEnd = source.range(of: "func stopRecording()", range: watchdogStart.upperBound.. Bool"), + let recordingEnd = source.range(of: "/// Begin dictation by borrowing", range: recordingStart.upperBound.. cancelStart.lowerBound, + "normal recording starts should use the helper that preserves only their owning zombie recovery" + ) + assertTrue( + source.contains("guard !zombieRecoveryState.isActive else { return }") + && source.contains("reportZombieEngineRecoveryTerminal(terminal)"), + "duplicate detector callbacks should not replace an active attempt, and every terminal path should share one reporter" ) } + runSuite("ParakeetEngine config changes invalidate zombie ownership before cancellation can suspend") { + let source = readParakeetDeviceRecoverySource() + guard let handlerStart = source.range(of: "private func handleAudioConfigChange() async"), + let handlerEnd = source.range(of: "private func recordStableRouteChangeAnalytics", range: handlerStart.upperBound.. String { diff --git a/Tests/SentryEventPolicyTests.swift b/Tests/SentryEventPolicyTests.swift index 677b8e8c..7487055a 100644 --- a/Tests/SentryEventPolicyTests.swift +++ b/Tests/SentryEventPolicyTests.swift @@ -169,6 +169,29 @@ func testSentryEventPolicy() { assertNil(tags["transcript_text"], "transcript text should stay out of Sentry tags") } + runSuite("SentryEventPolicy keeps zombie recovery stage and result categorical") { + let tags = SentryEventPolicy.diagnosticTags( + forEngine: "parakeet", + event: "zombie_engine_recovery_failed", + context: [ + "audio_device": "Private microphone name", + "failure_kind": "no_sample_callbacks", + "input_device_class": "built_in", + "output_device_class": "bluetooth", + "result": "failed", + "route_shape": "built_in_input_to_bluetooth_output", + "sample_count": "0", + "stage": "restart", + ] + ) + + assertEqual(tags["failure_kind"], "no_sample_callbacks", "zombie trigger should be queryable") + assertEqual(tags["result"], "failed", "terminal recovery result should be queryable") + assertEqual(tags["stage"], "restart", "terminal recovery stage should be queryable") + assertNil(tags["audio_device"], "raw device labels must stay out of Sentry") + assertNil(tags["sample_count"], "exact callback counts must stay out of Sentry") + } + runSuite("SentryEventPolicy diagnosticTags keeps meeting failure triage searchable") { let tags = SentryEventPolicy.diagnosticTags( forEngine: "meeting", diff --git a/docs/posthog-product-learning-plan.md b/docs/posthog-product-learning-plan.md index 79e5dcdd..1ec43f6c 100644 --- a/docs/posthog-product-learning-plan.md +++ b/docs/posthog-product-learning-plan.md @@ -159,6 +159,7 @@ Dictation events also allow coarse route fields: `default_input_class`, | `dictation_audio_route_changed` | route fields only | | `dictation_audio_route_recovery_finished` | route fields plus `outcome` | | `dictation_audio_route_recovery_timeout` | route fields only | +| `dictation_zombie_recovery_finished` | categorical `failure_kind`, `stage`, `result`, `hfp_suspected`, and route classes/shape only | ### Meetings diff --git a/docs/privacy-first-observability.md b/docs/privacy-first-observability.md index 2a430971..4a5f8cdb 100644 --- a/docs/privacy-first-observability.md +++ b/docs/privacy-first-observability.md @@ -140,6 +140,7 @@ allowlist. - `dictation_audio_route_changed` - `dictation_audio_route_recovery_finished` - `dictation_audio_route_recovery_timeout` +- `dictation_zombie_recovery_finished` - `meeting_recording_started` - `meeting_recording_start_failed` - `meeting_detected_call_ended`