// src/manager/public-helpers.ts import type { AdvertisementObservation, ScanOptions } from '../backend-contract/advertisement' import { contractError, type CleanupRecord } from '../backend-contract/errors' import type { BackendIdentity } from '../backend-contract/identity' import type { CharacteristicPath, NotificationValue } from '../backend-contract/gatt' import type { PublicOperationOptions, SubscriptionOptions } from '../backend-contract/operations' import { capacity, type OwnedBytes, type PeerId, type Uuid } from '../backend-contract/primitives' import type { BoundedAsyncStream, BoundedAsyncStreamIterator, StreamItem, StreamTerminalNotice } from '../backend-contract/streams' import { BleManager, Connection, DiscoveredGattDatabase } from './ble-manager' type CurrentCharacteristicPath = CharacteristicPath< Attachment, string, string, string, string, 'current' > interface Successful { readonly state: 'succeeded' readonly value: Value } interface Failed { readonly state: 'failed' readonly error: unknown } type OperationOutcome = Successful | Failed interface DeadlineHandle { cancel(): void } interface DeadlineClock { monotonicNow(): number scheduleDeadline(deadline: number, action: () => void): DeadlineHandle } /** The exact public scan request and observation predicate used by find and scanUntil. */ export interface ScanUntilOptions { readonly scan: ScanOptions readonly matches: (observation: AdvertisementObservation) => boolean } /** A discovered, generation-bound GATT database that remains owned by its returned connection. */ export interface ConnectedGattDatabase> { readonly connection: Connection readonly database: DiscoveredGattDatabase readonly snapshot: Awaited['snapshot']>> } /** Bounded collection configuration for collectNotifications. */ export interface CollectNotificationsOptions { readonly subscription: SubscriptionOptions readonly maximumValues: number } /** * Scans until a caller-owned predicate accepts an observation. The exact scan * request, including active AbortSignal and deadline behavior, is forwarded * unchanged to the public manager. */ export async function scanUntil>( manager: BleManager, options: ScanUntilOptions ): Promise> { const session = await manager.scan(options.scan) return settleWithCleanup( () => withStreamIterator(session.observations, iterator => nextMatchingObservation(iterator, options, manager)), () => session.stop(), 'helpers.scan-until.stop' ) } /** Alias for scanUntil that keeps the public vocabulary compact. */ export function find>( manager: BleManager, options: ScanUntilOptions ): Promise> { return scanUntil(manager, options) } /** Connects and discovers through the public handles, releasing a partial connection if discovery fails. */ export async function connectAndDiscover>( manager: BleManager, peerId: PeerId, options: PublicOperationOptions ): Promise> { const connection = await manager.connect(peerId, options) return settleWithCleanup( async () => { const database = await connection.discover(options) const snapshot = await database.snapshot() return Object.freeze({ connection, database, snapshot }) }, () => connection.release(), 'helpers.connect-and-discover.release', false ) } /** Resolves the first notification value and removes the subscription before returning. */ export async function firstNotification>( database: DiscoveredGattDatabase, path: CurrentCharacteristicPath, options: SubscriptionOptions ): Promise { const subscription = await database.subscribe(path, options) return settleWithCleanup( () => withStreamIterator(subscription.values, async iterator => notificationValue( await nextStreamItem(iterator, options, 'helpers.first-notification', database), options, database ) ), () => subscription.remove(), 'helpers.first-notification.remove' ) } /** Collects no more than maximumValues notification payloads, then removes the subscription. */ export async function collectNotifications>( database: DiscoveredGattDatabase, path: CurrentCharacteristicPath, options: CollectNotificationsOptions ): Promise { assertCollectionBound(options.maximumValues) const subscription = await database.subscribe(path, options.subscription) return settleWithCleanup( () => withStreamIterator(subscription.values, iterator => collectSubscriptionValues(iterator, options, database)), () => subscription.remove(), 'helpers.collect-notifications.remove' ) } /** Runs an operation with one connection lease and deterministically releases that lease on every exit. */ export async function withConnection, Value>( manager: BleManager, peerId: PeerId, options: PublicOperationOptions, operation: (connection: Connection) => Promise ): Promise { const connection = await manager.connect(peerId, options) return settleWithCleanup( () => operation(connection), () => connection.release(), 'helpers.with-connection.release' ) } export function defaultScanDelivery() { return Object.freeze({ itemCapacity: capacity(32), byteCapacity: capacity(16 * 1024), reservedControlCapacity: capacity(2), overflowPolicy: 'drop-oldest' as const }) } export function scanForServices>( manager: BleManager, serviceUuids: readonly Uuid[], options: Omit, 'scan'> & { readonly scan?: Partial['scan']> } ) { const scan = options.scan ?? {} return scanUntil(manager, { matches: options.matches, scan: { filter: { serviceUuids, manufacturerData: scan.filter?.manufacturerData ?? [], localNamePrefix: scan.filter?.localNamePrefix ?? null }, duplicatePolicy: scan.duplicatePolicy ?? 'merged', timestampPolicy: scan.timestampPolicy ?? 'source-then-receipt', delivery: scan.delivery ?? defaultScanDelivery(), deadline: scan.deadline ?? null, signal: scan.signal ?? null, sharing: scan.sharing ?? { mode: 'owner', allowSharing: false } } }) } export async function withDiscoveredConnection< Attachment extends string, Identity extends BackendIdentity, Value >( manager: BleManager, peerId: PeerId, options: PublicOperationOptions, fn: (session: ConnectedGattDatabase) => Promise ): Promise { return withConnection(manager, peerId, options, async connection => { const database = await connection.discover(options) const snapshot = await database.snapshot() return fn(Object.freeze({ connection, database, snapshot })) }) } export function throwIfCleanupFailed(cleanup: CleanupRecord, operation: string): void { if (cleanup.state !== 'release-failed') { return } throw contractError('lifecycle.invalid-state', 'cleanup', operation, { domain: 'cleanup', code: 'release-failed', safeMessage: 'Owned BLE resources did not release cleanly', metadata: Object.freeze({ failureCount: cleanup.failures.length, failures: cleanup.failures.map(failure => Object.freeze({ resourceKind: failure.resourceKind, code: failure.error.code }) ) }) }) } async function nextMatchingObservation( iterator: BoundedAsyncStreamIterator>, options: ScanUntilOptions, clock: DeadlineClock ): Promise> { while (true) { const item = await nextStreamItem(iterator, options.scan, 'helpers.scan-until', clock) if (item.kind === 'value' && options.matches(item.value)) { return item.value } if (item.kind === 'terminal') { throw streamTerminalError(item, options.scan, 'helpers.scan-until', clock) } if (item.kind === 'overflow') { throw contractError('stream.overflow', 'stream', 'helpers.scan-until.overflow') } } } async function collectSubscriptionValues( iterator: BoundedAsyncStreamIterator, options: CollectNotificationsOptions, clock: DeadlineClock ): Promise { const values: OwnedBytes[] = [] for (let index = 0; index < options.maximumValues; index += 1) { const item = await nextStreamItem(iterator, options.subscription, 'helpers.collect-notifications', clock) values.push(notificationValue(item, options.subscription, clock)) } return Object.freeze(values) } async function nextStreamItem( iterator: BoundedAsyncStreamIterator, options: PublicOperationOptions, operation: string, clock: DeadlineClock ): Promise> { if (options.signal?.aborted === true) { throw contractError('operation.aborted', 'stream', operation) } if (options.deadline !== null && options.deadline <= clock.monotonicNow()) { throw contractError('operation.timed-out', 'stream', operation) } return waitForStreamItem(iterator.next(), options, operation, clock) } function waitForStreamItem( next: Promise>>, options: PublicOperationOptions, operation: string, clock: DeadlineClock ): Promise> { return new Promise((resolve, reject) => { let settled = false let deadlineHandle: DeadlineHandle | null = null const signal = options.signal const releaseWait = (): boolean => { if (settled) { return false } settled = true if (signal !== null) { signal.removeEventListener('abort', onAbort) } if (deadlineHandle !== null) { deadlineHandle.cancel() deadlineHandle = null } return true } const onAbort = () => { if (releaseWait()) { reject(contractError('operation.aborted', 'stream', operation)) } } if (signal !== null) { signal.addEventListener('abort', onAbort, { once: true }) } if (options.deadline !== null) { deadlineHandle = clock.scheduleDeadline(options.deadline, () => { if (releaseWait()) { reject(contractError('operation.timed-out', 'stream', operation)) } }) } next .then(result => { if (signal?.aborted === true || (options.deadline !== null && options.deadline <= clock.monotonicNow())) { if (releaseWait()) { reject( contractError(signal?.aborted === true ? 'operation.aborted' : 'operation.timed-out', 'stream', operation) ) } return } if (releaseWait()) { resolve(requireStreamItem(result, operation)) } }) .catch(error => { if (releaseWait()) { reject(error) } }) }) } function requireStreamItem(result: IteratorResult>, operation: string): StreamItem { if (!result.done) { return result.value } throw contractError('stream.closed', 'stream', operation) } async function withStreamIterator( stream: BoundedAsyncStream, operation: (iterator: BoundedAsyncStreamIterator) => Promise ): Promise { const iterator = stream[Symbol.asyncIterator]() const outcome = await capture(() => operation(iterator)) const iteratorCleanup = await capture(() => closeIterator(iterator)) if (outcome.state === 'failed' && iteratorCleanup.state === 'failed') { throw new AggregateError( [outcome.error, iteratorCleanup.error], 'helpers.stream-iterator: operation and iterator cleanup failed' ) } if (outcome.state === 'failed') { throw outcome.error } if (iteratorCleanup.state === 'failed') { throw iteratorCleanup.error } return outcome.value } async function closeIterator(iterator: BoundedAsyncStreamIterator): Promise { await iterator.return() } function notificationValue( item: StreamItem, options: PublicOperationOptions, clock: DeadlineClock ): OwnedBytes { if (item.kind === 'value') { return item.value.value } if (item.kind === 'terminal') { throw streamTerminalError(item, options, 'helpers.notification', clock) } throw contractError('stream.overflow', 'stream', 'helpers.notification.overflow') } function streamTerminalError( terminal: StreamTerminalNotice, options: PublicOperationOptions | null, operation: string, clock: DeadlineClock ) { if (terminal.reason === 'overflow') { return contractError('stream.overflow', 'stream', operation) } if (terminal.reason === 'connection-lost') { return contractError('connection.lost', 'connection', operation) } if (terminal.reason === 'operation-aborted' || options?.signal?.aborted === true) { return contractError('operation.aborted', 'stream', operation) } if ( terminal.reason === 'operation-timed-out' || (options?.deadline !== null && options?.deadline !== undefined && options.deadline <= clock.monotonicNow()) ) { return contractError('operation.timed-out', 'stream', operation) } if (terminal.reason === 'owner-released') { return contractError('lifecycle.destroyed', 'stream', operation) } return contractError('stream.closed', 'stream', operation) } function assertCollectionBound(maximumValues: number): void { if (!Number.isSafeInteger(maximumValues) || maximumValues < 1) { throw contractError('argument.invalid', 'stream', 'helpers.collect-notifications.maximum-values') } } async function settleWithCleanup( operation: () => Promise, cleanup: () => Promise, cleanupOperation: string, releaseOnSuccess = true ): Promise { const outcome = await capture(operation) if (outcome.state === 'succeeded' && !releaseOnSuccess) { return outcome.value } const cleanupOutcome = await captureCleanup(cleanup, cleanupOperation) if (outcome.state === 'failed' && cleanupOutcome !== null) { throw new AggregateError([outcome.error, cleanupOutcome], `${cleanupOperation}: operation and cleanup failed`) } if (outcome.state === 'failed') { throw outcome.error } if (cleanupOutcome !== null) { throw cleanupOutcome } return outcome.value } async function capture(operation: () => Promise): Promise> { try { return { state: 'succeeded', value: await operation() } } catch (error) { return { state: 'failed', error } } } async function captureCleanup(cleanup: () => Promise, operation: string): Promise { try { const record = await cleanup() if (record.state === 'released' && record.failures.length === 0) { return null } return contractError('platform.failure', 'cleanup', operation) } catch (error) { if (error instanceof Error) { return error } return contractError('platform.failure', 'cleanup', operation) } }