/** * Instance-local reactive store for an AviationPlayer. * * The store owns all snapshot state and subscription Sets that hooks read * from. A player gets exactly one store; no module-level playback state lives * here. */ import { AviationError, mapNativeErrorCode } from './errors'; import { devWarn, noopLogger, type AviationLogger } from './logger'; import { IDLE_POSITION } from './snapshots'; import type { MediaItem as MediaItemHybrid } from './specs/MediaItem.nitro'; import type { PlaybackState, MediaPosition, ItemSource, AdBreakInfo, AdInfo, PlaybackMetricEvent, CastState, CastDevice, DomainEvent, AudioRouteChangeEvent, } from './specs/types.nitro'; type Listener = () => void; type ListenerSet = Set; export interface AviationCacheEvent { uri: string; hit: boolean; bytes: number; } type CacheEventListener = (event: AviationCacheEvent) => void; type MetricListener = (event: PlaybackMetricEvent) => void; export type DomainEventListener = (event: DomainEvent) => void; type RouteChangeListener = (event: AudioRouteChangeEvent) => void; export interface CastStateSnapshot { state: CastState; device: CastDevice | undefined; isConnected: boolean; } export interface AviationStore { readonly isReady: boolean; readonly playbackState: PlaybackState; readonly mediaPosition: MediaPosition; readonly currentItem: MediaItemHybrid | undefined; readonly currentItemSource: ItemSource; readonly queueIndex: number; readonly queueRevision: number; readonly error: AviationError | undefined; readonly isPlayingAd: boolean; readonly currentAdBreak: AdBreakInfo | undefined; readonly currentAd: AdInfo | undefined; /** Last cache event, for sync reads. Subscribe for lossless delivery. */ readonly lastCacheEvent: AviationCacheEvent | undefined; /** Last playback metric, for sync reads. Subscribe for lossless delivery. */ readonly lastPlaybackMetric: PlaybackMetricEvent | undefined; /** Last interruption (began or ended), for sync reads. */ readonly lastInterruption: DomainEvent | undefined; /** Last route change, for sync reads. */ readonly lastRouteChange: AudioRouteChangeEvent | undefined; getCastStateSnapshot(): CastStateSnapshot; subscribeReady(cb: Listener): () => void; subscribeState(cb: Listener): () => void; subscribePosition(cb: Listener): () => void; subscribeCurrentItem(cb: Listener): () => void; subscribeError(cb: Listener): () => void; subscribeQueue(cb: Listener): () => void; subscribeAdState(cb: Listener): () => void; /** * Event channels deliver every event as an argument, synchronously — * unlike the snapshot subscriptions above, nothing is coalesced. */ subscribeCacheEvent(cb: CacheEventListener): () => void; subscribePlaybackMetric(cb: MetricListener): () => void; /** * Interruptions began/ended, delivered losslessly with their policy * payload. Mirrors the snapshot field `lastInterruption` for sync reads. */ subscribeInterruption(cb: DomainEventListener): () => void; /** Feed point for the native bridge's domain-event handler. */ recordInterruption(event: DomainEvent): void; /** * Route changes (device added/removed), delivered losslessly. Mirrors the * snapshot field `lastRouteChange` for sync reads. */ subscribeRouteChange(cb: RouteChangeListener): () => void; /** Feed point for the player's session wiring; not for direct consumer use. */ recordRouteChange(event: AudioRouteChangeEvent): void; /** * Every coordinator domain event — state, queue, errors, properties, * cache/preload metrics, plugin-raised ad events — delivered losslessly. * Plugins subscribe here instead of polling snapshots. */ subscribeDomainEvent(cb: DomainEventListener): () => void; /** Feed point for the native bridge; not for direct consumer use. */ dispatchDomainEvent(event: DomainEvent): void; subscribeCastState(cb: Listener): () => void; /** * Event subscription: fires when the current item finishes playing * naturally. Unlike the snapshot subscriptions above, it is not * re-notified on reset — there is no snapshot to re-read. */ subscribePlaybackEnded(cb: Listener): () => void; setReady(): void; setPlaybackState(state: PlaybackState): void; setMediaPosition(position: MediaPosition): void; setCurrentItem( item: MediaItemHybrid | undefined, source: ItemSource, queueIndex: number ): void; notifyQueueChanged(): void; setError(message: string, nativeCode: string): void; notifyQueueEnd(): void; notifyPlaybackEnded(): void; setPlaybackMetric(event: PlaybackMetricEvent): void; setAdPlaying(adBreak: AdBreakInfo): void; setAdBreakEnded(): void; setCurrentAd(ad: AdInfo | undefined): void; clearAdState(): void; setCastState(snapshot: CastStateSnapshot): void; reset(): void; resetSnapshotsForSetupFailure(): void; } /** * Read-only projection of the store for consumers that must observe but * never mutate playback state — plugin contexts, primarily. Structural * typing lets the full store stand in for it; no runtime wrapper needed. */ export interface AviationStoreReader { readonly isReady: boolean; readonly playbackState: PlaybackState; readonly mediaPosition: MediaPosition; readonly currentItem: MediaItemHybrid | undefined; readonly currentItemSource: ItemSource; readonly queueIndex: number; readonly queueRevision: number; readonly error: AviationError | undefined; readonly isPlayingAd: boolean; readonly currentAdBreak: AdBreakInfo | undefined; readonly currentAd: AdInfo | undefined; readonly lastCacheEvent: AviationCacheEvent | undefined; readonly lastPlaybackMetric: PlaybackMetricEvent | undefined; readonly lastInterruption: DomainEvent | undefined; readonly lastRouteChange: AudioRouteChangeEvent | undefined; getCastStateSnapshot(): CastStateSnapshot; subscribeReady(cb: Listener): () => void; subscribeState(cb: Listener): () => void; subscribePosition(cb: Listener): () => void; subscribeCurrentItem(cb: Listener): () => void; subscribeError(cb: Listener): () => void; subscribeQueue(cb: Listener): () => void; subscribeAdState(cb: Listener): () => void; subscribeCacheEvent(cb: CacheEventListener): () => void; subscribePlaybackMetric(cb: MetricListener): () => void; subscribeInterruption(cb: DomainEventListener): () => void; subscribeRouteChange(cb: RouteChangeListener): () => void; subscribeDomainEvent(cb: DomainEventListener): () => void; subscribeCastState(cb: Listener): () => void; subscribePlaybackEnded(cb: Listener): () => void; } const IDLE_CAST: CastStateSnapshot = { state: 'disconnected', device: undefined, isConnected: false, }; export function createAviationStore( log: AviationLogger = noopLogger ): AviationStore { let isReady = false; let playbackState: PlaybackState = 'idle'; let mediaPosition: MediaPosition = { ...IDLE_POSITION }; let currentItem: MediaItemHybrid | undefined; let currentItemSource: ItemSource = 'none'; let error: AviationError | undefined; let queueIndex = -1; let queueRevision = 0; let isPlayingAd = false; let currentAdBreak: AdBreakInfo | undefined; let currentAd: AdInfo | undefined; let castStateSnapshot: CastStateSnapshot = { ...IDLE_CAST }; let lastCacheEvent: AviationCacheEvent | undefined; let lastPlaybackMetric: PlaybackMetricEvent | undefined; let lastInterruption: DomainEvent | undefined; let lastRouteChange: AudioRouteChangeEvent | undefined; const listeners = { ready: new Set(), state: new Set(), position: new Set(), currentItem: new Set(), error: new Set(), queue: new Set(), adState: new Set(), castState: new Set(), playbackEnded: new Set(), }; // Event channels: every event delivered synchronously, nothing coalesced. const metricListeners = new Set(); const cacheEventListeners = new Set(); const domainEventListeners = new Set(); const interruptionListeners = new Set(); const routeChangeListeners = new Set(); // Snapshot-backed channels re-notify together on reset so hooks re-read the // freshly-cleared state. Derived from `listeners` so a newly added channel is // included by default; `playbackEnded` is excluded because it's an event, not // a snapshot — there is nothing to re-read. The metric/cache channels above // are event channels and never re-notified either. const snapshotListenerSets = Object.entries(listeners) .filter(([channel]) => channel !== 'playbackEnded') .map(([, set]) => set); const pendingNotificationSets = new Set(); let notificationFlushScheduled = false; function scheduleMicrotask(callback: () => void): void { if (typeof queueMicrotask === 'function') { queueMicrotask(callback); return; } Promise.resolve() .then(callback) .catch((flushError: unknown) => { setTimeout(() => { throw flushError; }, 0); }); } function flushNotifications(): void { notificationFlushScheduled = false; const sets = Array.from(pendingNotificationSets); pendingNotificationSets.clear(); let firstError: unknown; for (const set of sets) { for (const cb of Array.from(set)) { try { cb(); } catch (listenerError) { firstError ??= listenerError; } } } if (firstError !== undefined) { // A throwing subscriber must never take the JS thread down — // rethrowing here crashed the app. Surface it in dev regardless of // logLevel (an internal bug is not diagnostic chatter) and keep // serving others. devWarn('store', 'listener threw during store notification', { error: firstError, }); } } function notifySet(set: ListenerSet): void { pendingNotificationSets.add(set); if (notificationFlushScheduled) return; notificationFlushScheduled = true; scheduleMicrotask(flushNotifications); } function notifySetSync(set: ListenerSet): void { set.forEach((cb) => cb()); } function dispatchEvent( subscribers: Set<(event: T) => void>, event: T ): void { let firstError: unknown; for (const cb of Array.from(subscribers)) { try { cb(event); } catch (listenerError) { firstError ??= listenerError; } } if (firstError !== undefined) { // A throwing subscriber must never take the JS thread down — // rethrowing here crashed the app. Surface it in dev regardless of // logLevel (an internal bug is not diagnostic chatter) and keep // serving others. devWarn('store', 'listener threw during store notification', { error: firstError, }); } } function subscribe(set: ListenerSet, cb: Listener): () => void { set.add(cb); return () => { set.delete(cb); }; } function resetAdState(): void { isPlayingAd = false; currentAdBreak = undefined; currentAd = undefined; notifySet(listeners.adState); } function resetSnapshots(): void { isReady = false; playbackState = 'idle'; mediaPosition = { ...IDLE_POSITION }; currentItem = undefined; currentItemSource = 'none'; error = undefined; queueIndex = -1; queueRevision = 0; isPlayingAd = false; currentAdBreak = undefined; currentAd = undefined; castStateSnapshot = { ...IDLE_CAST }; lastCacheEvent = undefined; lastPlaybackMetric = undefined; lastInterruption = undefined; lastRouteChange = undefined; } return { get isReady() { return isReady; }, get playbackState() { return playbackState; }, get mediaPosition() { return mediaPosition; }, get currentItem() { return currentItem; }, get currentItemSource() { return currentItemSource; }, get queueIndex() { return queueIndex; }, get queueRevision() { return queueRevision; }, get error() { return error; }, get isPlayingAd() { return isPlayingAd; }, get currentAdBreak() { return currentAdBreak; }, get currentAd() { return currentAd; }, get lastCacheEvent() { return lastCacheEvent; }, get lastPlaybackMetric() { return lastPlaybackMetric; }, get lastInterruption() { return lastInterruption; }, get lastRouteChange() { return lastRouteChange; }, getCastStateSnapshot() { return castStateSnapshot; }, subscribeReady: (cb) => subscribe(listeners.ready, cb), subscribeState: (cb) => subscribe(listeners.state, cb), subscribePosition: (cb) => subscribe(listeners.position, cb), subscribeCurrentItem: (cb) => subscribe(listeners.currentItem, cb), subscribeError: (cb) => subscribe(listeners.error, cb), subscribeQueue: (cb) => subscribe(listeners.queue, cb), subscribeAdState: (cb) => subscribe(listeners.adState, cb), subscribeCacheEvent: (cb) => { cacheEventListeners.add(cb); return () => { cacheEventListeners.delete(cb); }; }, subscribePlaybackMetric: (cb) => { metricListeners.add(cb); return () => { metricListeners.delete(cb); }; }, subscribeDomainEvent: (cb) => { domainEventListeners.add(cb); return () => { domainEventListeners.delete(cb); }; }, dispatchDomainEvent: (event) => dispatchEvent(domainEventListeners, event), subscribeInterruption: (cb) => { interruptionListeners.add(cb); return () => { interruptionListeners.delete(cb); }; }, recordInterruption: (event) => { log.verbose('store', 'interruption', { began: event.began }); lastInterruption = event; dispatchEvent(interruptionListeners, event); }, subscribeRouteChange: (cb) => { routeChangeListeners.add(cb); return () => { routeChangeListeners.delete(cb); }; }, recordRouteChange: (event) => { log.verbose('store', 'route change', { reason: event.reason }); lastRouteChange = event; dispatchEvent(routeChangeListeners, event); }, subscribeCastState: (cb) => subscribe(listeners.castState, cb), subscribePlaybackEnded: (cb) => subscribe(listeners.playbackEnded, cb), setReady(): void { log.verbose('store', 'player ready'); isReady = true; notifySet(listeners.ready); }, setPlaybackState(state: PlaybackState): void { if (playbackState === state) return; log.verbose('store', 'playback state', { state }); playbackState = state; if (state !== 'error') { error = undefined; } notifySet(listeners.state); }, setMediaPosition(position: MediaPosition): void { const prev = mediaPosition; if ( prev.currentMs === position.currentMs && prev.durationMs === position.durationMs && prev.isLive === position.isLive && prev.isAtLiveEdge === position.isAtLiveEdge ) { return; } mediaPosition = position; notifySet(listeners.position); }, setCurrentItem(item, source, nextQueueIndex): void { log.verbose('store', 'current item', { uri: item?.uri, source, queueIndex: nextQueueIndex, }); currentItem = item; currentItemSource = source; queueIndex = nextQueueIndex; notifySet(listeners.currentItem); notifySet(listeners.queue); }, notifyQueueChanged(): void { log.verbose('store', 'queue changed'); queueRevision++; notifySet(listeners.queue); }, setError(message: string, nativeCode: string): void { log.error('store', 'error set', { message, nativeCode }); const code = mapNativeErrorCode(nativeCode); error = new AviationError(code, message, nativeCode); notifySet(listeners.error); notifySet(listeners.state); }, notifyQueueEnd(): void { log.verbose('store', 'queue end notified'); notifySet(listeners.queue); }, notifyPlaybackEnded(): void { log.verbose('store', 'playback ended'); notifySet(listeners.playbackEnded); }, setPlaybackMetric(event: PlaybackMetricEvent): void { log.verbose('store', 'playback metric', { type: event.type }); lastPlaybackMetric = event; dispatchEvent(metricListeners, event); if (event.type === 'cacheHit' || event.type === 'cacheMiss') { const cacheEvent: AviationCacheEvent = { uri: event.uri ?? '', hit: event.type === 'cacheHit', bytes: event.bytes ?? 0, }; lastCacheEvent = cacheEvent; dispatchEvent(cacheEventListeners, cacheEvent); } }, setAdPlaying(adBreak): void { log.verbose('store', 'ad break started'); isPlayingAd = true; currentAdBreak = adBreak; notifySet(listeners.adState); }, setAdBreakEnded(): void { log.verbose('store', 'ad break ended'); resetAdState(); }, setCurrentAd(ad): void { log.verbose('store', 'current ad updated', { hasAd: !!ad }); currentAd = ad; notifySet(listeners.adState); }, clearAdState(): void { log.verbose('store', 'ad state cleared'); resetAdState(); }, setCastState(snapshot): void { log.verbose('store', 'cast state', { state: snapshot.state, isConnected: snapshot.isConnected, }); castStateSnapshot = snapshot; notifySet(listeners.castState); }, reset(): void { log.verbose('store', 'store reset'); resetSnapshots(); for (const set of snapshotListenerSets) notifySetSync(set); }, resetSnapshotsForSetupFailure(): void { log.verbose('store', 'store reset after setup failure'); resetSnapshots(); for (const set of snapshotListenerSets) notifySet(set); }, }; }