import { webPushCountBucket, WEB_PUSH_LIMITS, WEB_PUSH_PROTOCOLS, webPushFailure, webPushSuccess, type WebPushObserver, } from "../../contracts/web-push.ts"; import type { PushAssociationFenceStore } from "./push-association-fence-store.ts"; import { createPushEventAdapter, type PushEventFacade, } from "./inbound/push-event-adapter.ts"; import { createNotificationClickAdapter, type NotificationClickEventFacade, type WindowClientFacade, } from "./inbound/notification-click-adapter.ts"; import type { AssociationNotificationTagDigest, WebPushNotificationRegistry, } from "./notification-registry.ts"; import { createLinkedAbortController, nativeFailure, observeWebPush, withAbortableDeadline, type TimeoutScheduler, } from "./runtime-support.ts"; type FunctionalEventFacade = Readonly<{ waitUntil(task: Promise): void; }>; type ServiceWorkerRegistrationFacade = Readonly<{ showNotification( title: string, options: Readonly<{ body: string; data: unknown; requireInteraction: false; tag: string; }>, ): Promise; }>; type ServiceWorkerClientsFacade = Readonly<{ matchAll(input: Readonly<{ type: "window"; includeUncontrolled: boolean; }>): Promise; openWindow(url: string): Promise; }>; /** * Structural worker host used deliberately instead of exposing DOM worker * globals to application/test compilation. A selected worker entry adapts its * native scope to this facade; this factory has no registration side effect. */ export type ServiceWorkerEventHost = Readonly<{ origin: string; registration: ServiceWorkerRegistrationFacade; clients: ServiceWorkerClientsFacade; addEventListener(type: string, listener: (event: unknown) => void): void; removeEventListener(type: string, listener: (event: unknown) => void): void; }>; export type WebPushServiceWorkerRuntime = Readonly<{ dispose(): void; }>; export function createWebPushServiceWorkerRuntime( dependencies: Readonly<{ host: ServiceWorkerEventHost; fenceStore: PushAssociationFenceStore; registry: WebPushNotificationRegistry; now?: () => number; scheduler?: TimeoutScheduler; observer?: WebPushObserver; tagDigest?: AssociationNotificationTagDigest; }>, ): WebPushServiceWorkerRuntime { const lifecycle = new AbortController(); const push = createPushEventAdapter({ fenceStore: dependencies.fenceStore, registry: dependencies.registry, notifications: { showNotification: (title, options) => dependencies.host.registration.showNotification(title, options), }, now: dependencies.now, scheduler: dependencies.scheduler, observer: dependencies.observer, tagDigest: dependencies.tagDigest, signal: lifecycle.signal, }); const click = createNotificationClickAdapter({ fenceStore: dependencies.fenceStore, registry: dependencies.registry, clients: { async matchControlledWindowClients() { const candidates = await dependencies.host.clients.matchAll({ type: "window", includeUncontrolled: false, }); return candidates .map(windowClientFacade) .filter( (candidate): candidate is WindowClientFacade => candidate !== null, ); }, async openWindow(url) { return windowClientFacade( await dependencies.host.clients.openWindow(url), ); }, }, origin: dependencies.host.origin, now: dependencies.now, scheduler: dependencies.scheduler, observer: dependencies.observer, signal: lifecycle.signal, }); const onPush = (event: unknown) => { const facade = pushEventFacade(event); if (facade) void push.handle(facade); }; const onNotificationClick = (event: unknown) => { const facade = notificationClickEventFacade(event); if (facade) void click.handle(facade); }; const onSubscriptionChange = (event: unknown) => { const facade = functionalEventFacade(event); if (!facade) return; const taskControl = createLinkedAbortController(lifecycle.signal); // WP-06. Bounded fan-out is policy, but the operator must be able to see // that only part of the client set was notified. let observedClientCount = 0; let truncatedClients = false; const processing = withAbortableDeadline( async (signal) => { let clients: readonly unknown[]; try { clients = await dependencies.host.clients.matchAll({ type: "window", includeUncontrolled: true, }); } catch { return nativeFailure("SUBSCRIPTION_RECONCILE", true); } if (signal.aborted) { return webPushFailure("ABORTED", "SUBSCRIPTION_RECONCILE"); } observedClientCount = clients.length; truncatedClients = clients.length > WEB_PUSH_LIMITS.clientHandoffCount; try { for (const candidate of clients.slice( 0, WEB_PUSH_LIMITS.clientHandoffCount, )) { if (signal.aborted) { return webPushFailure( "ABORTED", "SUBSCRIPTION_RECONCILE", ); } const client = windowClientFacade(candidate); client?.postMessage( Object.freeze({ protocol: WEB_PUSH_PROTOCOLS.reconcileRequired, }), ); } } catch { return nativeFailure("SUBSCRIPTION_RECONCILE", true); } return webPushSuccess(undefined); }, { deadlineMs: WEB_PUSH_LIMITS.handlerDeadlineMs, operation: "SUBSCRIPTION_RECONCILE", signal: taskControl.signal, scheduler: dependencies.scheduler, }, ).finally(taskControl.dispose); const lifetime = processing.then((result) => { observeWebPush(dependencies.observer, { event: "web_push_subscription_rotated", outcome: result.ok && !truncatedClients ? "SUCCEEDED" : "DEGRADED", ...(result.ok ? truncatedClients ? { reason: "LIMIT_EXCEEDED" as const } : {} : { reason: result.error.code }), countBucket: webPushCountBucket(observedClientCount), truncated: truncatedClients, }); }); try { facade.waitUntil(lifetime); } catch { taskControl.abort(); void lifetime; } }; dependencies.host.addEventListener("push", onPush); dependencies.host.addEventListener( "notificationclick", onNotificationClick, ); dependencies.host.addEventListener( "pushsubscriptionchange", onSubscriptionChange, ); let disposed = false; return Object.freeze({ dispose() { if (disposed) return; disposed = true; lifecycle.abort(); dependencies.host.removeEventListener("push", onPush); dependencies.host.removeEventListener( "notificationclick", onNotificationClick, ); dependencies.host.removeEventListener( "pushsubscriptionchange", onSubscriptionChange, ); dependencies.fenceStore.close(); }, }); } function functionalEventFacade( value: unknown, ): FunctionalEventFacade | null { try { return isRecord(value) && typeof value.waitUntil === "function" ? (value as FunctionalEventFacade) : null; } catch { return null; } } function pushEventFacade(value: unknown): PushEventFacade | null { if ( !isRecord(value) || typeof value.waitUntil !== "function" || !( value.data === null || (isRecord(value.data) && typeof value.data.arrayBuffer === "function") ) ) { return null; } return value as PushEventFacade; } function notificationClickEventFacade( value: unknown, ): NotificationClickEventFacade | null { if ( !isRecord(value) || typeof value.waitUntil !== "function" || !isRecord(value.notification) || typeof value.notification.close !== "function" || !Object.hasOwn(value.notification, "data") ) { return null; } return value as NotificationClickEventFacade; } function windowClientFacade(value: unknown): WindowClientFacade | null { if ( !isRecord(value) || typeof value.url !== "string" || typeof value.focus !== "function" || typeof value.postMessage !== "function" ) { return null; } return value as WindowClientFacade; } function isRecord(value: unknown): value is Record { return Boolean(value) && typeof value === "object" && !Array.isArray(value); }