// src/core/unified-ble-core.ts import { BackendContractError, contractError } from '../backend-contract/errors' import { assertAttachedBackend, type BackendAttachment, type BackendEvent, type BleCentralBackend, type ConnectionOptions, type ConnectionLease, type ManagerConstruction, type ScanLease } from '../backend-contract/backend' import { createAttachmentBoundIdFactory } from '../backend-contract/primitives' import type { AdvertisementObservation, OwnerScanOptions, ScanOptions } from '../backend-contract/advertisement' import type { NormalizedScanQuery } from '../backend-contract/scan-query' import type { ScanPlan } from '../backend-contract/scan-planning' import type { CleanupFailure, CleanupRecord } from '../backend-contract/errors' import { isAuthorizationBlocking, type AdapterStateSnapshot, type AdapterStateWatch, type BackendIdentity } from '../backend-contract/identity' import type { PublicOperationOptions, SubscriptionOptions, WritePolicy, WriteReceipt } from '../backend-contract/operations' import type { AttachmentBoundIdFactory, AttachmentId, ByteLimit, OwnedBytes, PeerId } from '../backend-contract/primitives' import { deadline as createDeadline } from '../backend-contract/primitives' import type { BoundedAsyncStream } from '../backend-contract/streams' import type { FeatureId } from '../backend-contract/capabilities' import type { SecurityBackend } from '../backend-contract/security' import type { ConnectionLifecycleTerminalCause } from '../backend-contract/connection-lifecycle' import { AggregateStreamQuota } from './aggregate-stream-quota' import { CoreBoundedStream } from './bounded-stream' import { CoreOperationCoordinator } from './operation-coordinator' import { ResourceLedger } from './resource-ledger' import { CoreSubscription, SubscriptionRegistry } from './subscription-registry' import { CoreTraceRecorder, type CoreTraceResource } from './trace-recorder' import { CoreLifecycleObserver, createRetainedCleanup } from './core-lifecycle-observer' import { CoreConnection, CoreGattDatabase } from './core-gatt-handles' import { readCoreAdapterState } from './core-adapter-state' import { readCoreCharacteristic, writeCoreCharacteristic, writeCoreCharacteristicWhenReady, writeCoreLongCharacteristic } from './core-characteristic-operations' import { createCoreFeatureRegistry, observeMaximumWriteLength, planLongWrite } from './core-capabilities' import { createCoreConnectionControls, type CoreConnectionControls } from './core-connection-controls' import { readCoreDescriptor, writeCoreDescriptor } from './core-descriptor-operations' import { discoverCoreGattDatabase } from './core-discovery' import type { CurrentCharacteristicPath, CurrentDescriptorPath } from './current-gatt-paths' import { advertisementByteLength, advertisementPayloadByteLength, awaitWithOperationAdmission, activateScanLifetime, cleanupFailure, cloneObservation, deactivateScanLifetime, retryableCleanup, scheduleCoreDeadline, type CoreDeadlineHandle, type CoreDeadlineScheduler } from './unified-ble-core-helpers' import { forwardCoreBackendEvents } from './core-backend-event-stream' import { isConnectionLossCause, lifecycleCauseFromBackendDisconnect } from './connection-lifecycle-rules' import type { DiagnosticTraceDocument } from '../diagnostics/trace-format' export { DEFAULT_CORE_MAXIMUM_VALUE_BYTES } from './unified-ble-core-helpers' export type { CoreDeadlineHandle, CoreDeadlineScheduler } from './unified-ble-core-helpers' /** * Safety bound on draining a quarantined connection lane during teardown. * * Independent of any caller deadline by design: quarantine drain runs after the * operation that triggered it has already failed or been cancelled, so deriving * from that operation's deadline would mean never waiting at all and reporting a * spurious cleanup failure. One second bounds teardown without turning a slow * native release into an unbounded hang; expiry surfaces as a cleanup failure. */ const QUARANTINE_DRAIN_TIMEOUT_MS = 1_000 // Non-cancellable backend probes remain owned until they settle. Bound their accumulation. const MAX_PENDING_ADAPTER_WATCH_ACQUISITIONS = 64 export interface UnifiedBleCoreOptions { readonly now: () => number readonly maximumValueBytes: ByteLimit readonly maximumAggregateRetainedBytes: number readonly traceMaximumRecords: number readonly traceMaximumBytes: number /** Optional host-neutral scheduler used for active scan deadline enforcement. */ readonly timer?: CoreDeadlineScheduler } export interface CoreScanSession { readonly scanSessionId: ScanLease['scanSessionId'] readonly leaseId: ScanLease['leaseId'] readonly shareToken: ScanLease['shareToken'] readonly observations: BoundedAsyncStream> stop(): Promise } interface TrackedScan extends CoreScanSession { readonly lease: ScanLease readonly stream: CoreBoundedStream> readonly ownsPhysicalController: boolean stopInFlight: Promise | null released: boolean activeAbortListener: (() => void) | null activeAbortSignal: AbortSignal | null activeDeadline: CoreDeadlineHandle | null } interface CoreGattDiscovery> { readonly kind: 'discover' | 'rediscover' readonly promise: Promise> } interface TrackedAdapterStateWatch { readonly initial: AdapterStateSnapshot readonly values: BoundedAsyncStream> stop(): Promise stopAttempt: Promise | null released: CleanupRecord | null readonly onAbort: () => void } /** * One attached core owns the portable policy for one manager. It delegates only * radio mechanics to the negotiated backend and never imports a host runtime. */ export class UnifiedBleCore> { private coreState: 'new' | 'ready' | 'destroying' | 'destroyed' | 'failed' = 'new' private attachment: BackendAttachment | null = null private idFactory: AttachmentBoundIdFactory | null = null private readonly resourceLedger = new ResourceLedger() private readonly trace: CoreTraceRecorder private readonly lifecycleObserver: CoreLifecycleObserver private readonly aggregateQuota: AggregateStreamQuota private readonly operationCoordinator: CoreOperationCoordinator private readonly connectionControls: CoreConnectionControls private readonly featureRegistry private readonly scans = new Map>() private readonly connections = new Map>() private readonly connectionReleases = new Map>() private readonly subscriptions: SubscriptionRegistry private backendEventStream: BoundedAsyncStream> | null = null private destroyResult: Promise | null = null private resourceReleaseResult: Promise | null = null private backendDestroyResult: Promise | null = null private readonly discoveries = new Map>() private readonly adapterStateWatches = new Set>() private readonly pendingAdapterStateWatches = new Set<(error: BackendContractError) => void>() private readonly pendingScanAcquisitions = new Set<(error: BackendContractError) => void>() private readonly pendingConnectAcquisitions = new Set<(error: BackendContractError) => void>() private admissionEpoch = 1 private nextOperation = 1 private constructor( readonly construction: ManagerConstruction, readonly options: UnifiedBleCoreOptions ) { this.featureRegistry = createCoreFeatureRegistry(construction.attachedBackend.backend.features) this.trace = new CoreTraceRecorder(options.traceMaximumRecords, options.traceMaximumBytes) this.lifecycleObserver = new CoreLifecycleObserver(this.trace, options.now) this.aggregateQuota = new AggregateStreamQuota(options.maximumAggregateRetainedBytes) this.operationCoordinator = new CoreOperationCoordinator({ now: options.now, createCorrelation: () => this.requireIdFactory().operationCorrelation(`operation-${this.nextOperationValue()}`), resourceLedger: this.resourceLedger, trace: this.trace }) this.connectionControls = createCoreConnectionControls(this.backend, this.operationCoordinator, operation => { this.assertReady(operation) }) this.subscriptions = new SubscriptionRegistry({ backend: construction.attachedBackend.backend, attachmentId: construction.attachedBackend.attachment.attachment.attachmentId, idFactory: { clientId: value => this.requireIdFactory().clientId(value), managerId: value => this.requireIdFactory().managerId(value), connectionId: value => this.requireIdFactory().connectionId(value), leaseId: value => this.requireIdFactory().leaseId(value), scanShareToken: value => this.requireIdFactory().scanShareToken(value), scanSessionId: value => this.requireIdFactory().scanSessionId(value), databaseId: value => this.requireIdFactory().databaseId(value), subscriptionId: value => this.requireIdFactory().subscriptionId(value), operationCorrelation: value => this.requireIdFactory().operationCorrelation(value), backendOperationHandle: value => this.requireIdFactory().backendOperationHandle(value) }, operationCoordinator: this.operationCoordinator, aggregateQuota: this.aggregateQuota, resourceLedger: this.resourceLedger, trace: this.trace, now: options.now, maximumValueBytes: options.maximumValueBytes, isPathCurrent: path => this.isCurrentPath(path) }) } static async attach>( construction: ManagerConstruction, options: UnifiedBleCoreOptions ): Promise> { const core = new UnifiedBleCore(construction, options) core.adoptAttachedBackend() return core } get state(): 'new' | 'ready' | 'destroying' | 'destroyed' | 'failed' { return this.coreState } get identity(): Identity { return this.requireAttachment().identity } get attachmentId(): AttachmentId { return this.requireAttachment().attachment.attachmentId } get backend(): BleCentralBackend { return this.construction.attachedBackend.backend } get features() { return this.featureRegistry } /** True only where the instantiated backend registry exposes an invocable implementation. */ private supportsCapability(id: FeatureId): boolean { const descriptor = this.featureRegistry.descriptors.find(candidate => candidate.id === id) return descriptor !== undefined && (descriptor.state === 'supported' || descriptor.state === 'limited') } /** Manager-admitted security operations; backend ownership remains behind this lifecycle gate. */ securityBackend(): SecurityBackend | undefined { this.assertReady('security.backend') const backend = this.backend.security if (backend === undefined) return undefined return { state: (peerId, options) => { this.assertReady('security.state') this.assertOperationAdmission(options, 'security.state') return backend.state(peerId, options) }, watch: peerId => { this.assertReady('security.watch') return backend.watch(peerId) }, pair: (peerId, options) => { this.assertReady('security.pair') this.assertOperationAdmission(options, 'security.pair') return backend.pair(peerId, options) }, cancelPairing: (peerId, options) => { this.assertReady('security.cancel-pairing') this.assertOperationAdmission(options, 'security.cancel-pairing') return backend.cancelPairing(peerId, options) }, unpair: (peerId, options) => { this.assertReady('security.unpair') this.assertOperationAdmission(options, 'security.unpair') return backend.unpair(peerId, options) } } } traces(): readonly import('./trace-recorder').CoreTraceRecord[] { return this.trace.snapshot() } traceDocument(): DiagnosticTraceDocument { return this.trace.snapshotDocument() } monotonicNow(): number { return this.options.now() } scheduleDeadline(deadline: number, action: () => void): CoreDeadlineHandle { return scheduleCoreDeadline(deadline, action, this.options.timer, this.options.now) } localResourceCounters(): import('../backend-contract/backend').ResourceCounters { this.syncRetainedByteBuffers() return this.resourceLedger.snapshot() } async adapterState(): Promise> { this.assertReady('adapter-state') return readCoreAdapterState(this.backend) } async adapterStates(options: { readonly signal?: AbortSignal | null } = {}): Promise<{ readonly initial: AdapterStateSnapshot readonly values: BoundedAsyncStream> stop(): Promise }> { this.assertReady('adapter-states') if (abortRequested(options.signal)) { throw contractError('operation.aborted', 'adapter', 'adapter-states') } if (this.pendingAdapterStateWatches.size >= MAX_PENDING_ADAPTER_WATCH_ACQUISITIONS) { throw contractError('stream.quota', 'adapter', 'adapter-states.pending-acquisitions') } return new Promise((resolve, reject) => { let cancelled = false const cancel = (error: BackendContractError) => { cancelled = true options.signal?.removeEventListener('abort', onAbort) reject(error) } const onAbort = () => cancel(contractError('operation.aborted', 'adapter', 'adapter-states')) this.pendingAdapterStateWatches.add(cancel) options.signal?.addEventListener('abort', onAbort, { once: true }) const acquire = async () => { try { const watch = await this.backend.adapter.watchState() const session = this.trackAdapterStateWatch(watch, options) this.pendingAdapterStateWatches.delete(cancel) if (cancelled || this.coreState !== 'ready' || abortRequested(options.signal)) { if (!cancelled) cancel(contractError('operation.cancelled-by-destroy', 'adapter', 'adapter-states')) // Ownership has moved to the tracked session. Failed late cleanup // remains available to manager.destroy() for an explicit retry. const cleanup = session.stop() this.lifecycleObserver.observeCleanup(cleanup, 'adapter-states-late-stop') await cleanup } else { resolve(session) } } catch (error) { if (!cancelled) reject(error) else throw error } finally { this.pendingAdapterStateWatches.delete(cancel) options.signal?.removeEventListener('abort', onAbort) } } this.lifecycleObserver.observeBackground(acquire(), 'manager', 'adapter-states-acquisition') }) } private trackAdapterStateWatch( watch: AdapterStateWatch, options: { readonly signal?: AbortSignal | null } ): TrackedAdapterStateWatch { const session: TrackedAdapterStateWatch = { initial: watch.initial, values: watch.transitions, stopAttempt: null, released: null, onAbort: () => { this.lifecycleObserver.observeCleanup(session.stop(), 'adapter-states-abort-stop') }, stop: (): Promise => { if (session.released !== null) { return Promise.resolve(session.released) } if (session.stopAttempt !== null) { return session.stopAttempt } const attempt = closeAdapterStateStream(watch.transitions) .then(cleanup => { if (cleanup.state === 'released') { this.adapterStateWatches.delete(session) session.released = cleanup } return cleanup }) .finally(() => { options.signal?.removeEventListener('abort', session.onAbort) if (session.stopAttempt === attempt) { session.stopAttempt = null } }) session.stopAttempt = attempt return attempt } } this.adapterStateWatches.add(session) options.signal?.addEventListener('abort', session.onAbort, { once: true }) return session } async maximumWriteLength( database: CoreGattDatabase, path: CurrentCharacteristicPath, mode: WritePolicy['mode'] ): Promise> { this.assertReady('maximum-write-length') database.assertPath(path) const observation = await observeMaximumWriteLength(this.features, path, mode) database.assertPath(path) return observation } async writeLong( database: CoreGattDatabase, path: CurrentCharacteristicPath, bytes: Readonly, options: import('../backend-contract/operations').LongWritePolicy ): Promise> { this.assertReady('write-long') database.assertPath(path) return writeCoreLongCharacteristic( this.backend, this.operationCoordinator, this.options.maximumValueBytes, database, path, bytes, options, async () => { database.assertPath(path) const observation = await observeMaximumWriteLength(this.features, path, options.mode) database.assertPath(path) const maximumWriteLength = resolveLongWriteChunkSize(observation.maximumWriteLength, options.chunkSize) const plan = await planLongWrite( this.features, String(path.connectionId), String(path.connectionGeneration), options.mode, bytes.byteLength, maximumWriteLength ) database.assertPath(path) return Object.freeze({ maximumWriteLength, totalChunks: plan.totalChunks }) } ) } async scan(options: ScanOptions): Promise> { this.assertReady('scan') this.assertOperationAdmission(options, 'scan') if (options.platform !== undefined && !this.supportsCapability('scan:platform-options')) { // Backends that do not register scan:platform-options (deterministic // included) would silently scan with their defaults; the host-neutral // path fails closed for every caller, not only the public manager. throw contractError('capability.unsupported', 'scan', 'unified-core.scan.platform-options') } const admissionEpoch = this.admissionEpoch const stream = new CoreBoundedStream>( options.delivery, options.delivery.overflowPolicy ) this.aggregateQuota.register(stream) return new Promise((resolve, reject) => { let cancelled = false let deadlineHandle: CoreDeadlineHandle | null = null const cancel = (error: BackendContractError) => { if (cancelled) { return } cancelled = true deadlineHandle?.cancel() options.signal?.removeEventListener('abort', onAbort) reject(error) } const onAbort = () => cancel(contractError('operation.aborted', 'core', 'scan')) this.pendingScanAcquisitions.add(cancel) options.signal?.addEventListener('abort', onAbort, { once: true }) if (options.deadline !== null) { deadlineHandle = scheduleCoreDeadline( Number(options.deadline), () => cancel(contractError('operation.timed-out', 'core', 'scan')), this.options.timer, this.options.now ) } const acquire = async () => { const releaseStream = () => { stream.closeWithReason('owner-released') this.aggregateQuota.unregister(stream) } try { let lease: ScanLease try { lease = await this.startOrJoinScan(options) } catch (error) { releaseStream() throw error instanceof BackendContractError ? error : contractError('scan.start-failed', 'scan', 'unified-core.scan') } const closed = this.admissionClosedError(admissionEpoch, options, 'scan') if (cancelled || closed !== null) { releaseStream() if (!cancelled && closed !== null) { cancel(closed) } this.pendingScanAcquisitions.delete(cancel) await this.compensateUnadoptedLease(lease.stop.bind(lease), 'scan', 'scan-stale-admission-release') return } resolve(this.adoptScan(lease, stream, options)) } catch (error) { if (!cancelled) { reject(error) } else { throw error } } finally { deadlineHandle?.cancel() this.pendingScanAcquisitions.delete(cancel) options.signal?.removeEventListener('abort', onAbort) } } this.lifecycleObserver.observeBackground(acquire(), 'scan', 'scan-acquisition') }) } planScan(query: NormalizedScanQuery): ScanPlan | null { this.assertReady('scan.plan') return this.backend.scanner.plan?.(query) ?? null } async connect(peerId: PeerId, options: ConnectionOptions): Promise> { this.assertReady('connect') this.assertOperationAdmission(options, 'connect') const admissionEpoch = this.admissionEpoch return new Promise((resolve, reject) => { let cancelled = false let deadlineHandle: CoreDeadlineHandle | null = null const cancel = (error: BackendContractError) => { if (cancelled) { return } cancelled = true deadlineHandle?.cancel() options.signal?.removeEventListener('abort', onAbort) reject(error) } const onAbort = () => cancel(contractError('operation.aborted', 'core', 'connect')) this.pendingConnectAcquisitions.add(cancel) options.signal?.addEventListener('abort', onAbort, { once: true }) if (options.deadline !== null) { deadlineHandle = scheduleCoreDeadline( Number(options.deadline), () => cancel(contractError('operation.timed-out', 'core', 'connect')), this.options.timer, this.options.now ) } const acquire = async () => { try { let lease: ConnectionLease try { lease = await this.backend.connections.connect(peerId, this.construction.clientId, options) } catch (error) { throw error instanceof BackendContractError ? error : contractError('connection.failed', 'connection', 'unified-core.connect') } const closed = this.admissionClosedError(admissionEpoch, options, 'connect') if (cancelled || closed !== null) { if (!cancelled && closed !== null) { cancel(closed) } this.pendingConnectAcquisitions.delete(cancel) await this.compensateUnadoptedLease( lease.release.bind(lease), 'connection', 'connect-stale-admission-release' ) return } const connection = new CoreConnection(this, lease, this.connectionControls) this.connections.set(String(lease.connection.connectionId), connection) this.resourceLedger.increment('connectionLeases') resolve(connection) } catch (error) { if (!cancelled) { reject(error) } else { throw error } } finally { deadlineHandle?.cancel() this.pendingConnectAcquisitions.delete(cancel) options.signal?.removeEventListener('abort', onAbort) } } this.lifecycleObserver.observeBackground(acquire(), 'connection', 'connect-acquisition') }) } async destroy(): Promise { if (this.destroyResult !== null) { return this.destroyResult } const destruction = this.releaseResources().then(async released => { if (released.state === 'release-failed' || this.construction.ownerMode === 'borrowing') { return released } return this.destroyBackend() }) this.destroyResult = retryableCleanup(destruction, () => { this.destroyResult = null }) return this.destroyResult } async releaseResources(cause: ConnectionLifecycleTerminalCause = 'manager-destroyed'): Promise { if (this.resourceReleaseResult !== null) { return this.resourceReleaseResult } this.admissionEpoch += 1 this.coreState = 'destroying' this.operationCoordinator.destroy() for (const connection of this.connections.values()) { connection.finishLifecycle(cause, null) } const release = this.destroyOwnedResources(cause) this.resourceReleaseResult = retryableCleanup(release, () => { this.resourceReleaseResult = null }) return this.resourceReleaseResult } destroyBackend(): Promise { if (this.backendDestroyResult !== null) { return this.backendDestroyResult } const destruction = this.backend.destroy().then( result => result, error => { this.trace.record({ timestamp: this.options.now(), resource: 'manager', transition: 'backend-destroy-rejected', operation: null, cause: error instanceof BackendContractError ? error.normalized.code : 'platform.failure', queuedOperations: 0, dispatchedOperations: 0, quarantinedOperations: 0 }) return cleanupFailure('backend', contractError('platform.failure', 'cleanup', 'unified-core.backend-destroy')) } ) this.backendDestroyResult = retryableCleanup(destruction, () => { this.backendDestroyResult = null }) return this.backendDestroyResult } async discover( connection: CoreConnection, options: PublicOperationOptions ): Promise> { this.assertReady('discover') this.assertOperationAdmission(options, 'discover') const key = String(connection.resource.connectionId) const existing = this.discoveries.get(key) if (existing !== undefined) { const database = await awaitWithOperationAdmission(existing.promise, options, this.options.now, 'discover') if (database.isCurrent()) { return database } const registered = this.discoveries.get(key) if (registered !== undefined && registered !== existing) { const replacement = await awaitWithOperationAdmission(registered.promise, options, this.options.now, 'discover') replacement.assertCurrent() return replacement } throw contractError('gatt.stale-handle', 'gatt', 'core-gatt-database.current') } const promise = discoverCoreGattDatabase( this, this.backend, this.resourceLedger, connection, options, operation => this.assertReady(operation), (value, operation) => this.assertOperationAdmission(value, operation), (admissionEpoch, value, operation) => this.assertAdmissionCurrent(admissionEpoch, value, operation), this.admissionEpoch, null, null ) const discovery: CoreGattDiscovery = { kind: 'discover', promise } this.discoveries.set(key, discovery) try { return await promise } finally { if (this.discoveries.get(key) === discovery) { this.discoveries.delete(key) } } } async rediscoverGatt( connection: CoreConnection, options: PublicOperationOptions, reason: Extract< import('../backend-contract/gatt').GattDatabaseChangedEvent['reason'], 'service-changed' | 'manual-rediscovery' > ): Promise> { const key = String(connection.resource.connectionId) const connectionGeneration = connection.resource.connectionGeneration const startingDatabaseGeneration = connection.database?.path.databaseGeneration ?? null let lastAwaited: CoreGattDiscovery | undefined for (;;) { this.assertReady('rediscover') this.assertOperationAdmission(options, 'rediscover') connection.assertCurrent() if (connection.resource.connectionGeneration !== connectionGeneration) { throw contractError('connection.stale', 'connection', 'core-connection.current') } const current = connection.database if ( current !== null && current.isCurrent() && lastAwaited?.kind !== 'discover' && (lastAwaited?.kind === 'rediscover' || (startingDatabaseGeneration !== null && String(current.path.databaseGeneration) !== String(startingDatabaseGeneration))) ) { return current } const registered = this.discoveries.get(key) if (registered !== undefined && registered !== lastAwaited) { lastAwaited = registered try { await awaitWithOperationAdmission(registered.promise, options, this.options.now, 'rediscover') } catch (error) { if (isRediscoverCallerTerminal(error, options, this.options.now)) { throw error } } continue } this.operationCoordinator.cancelQueue(key, 'disconnected') const promise = discoverCoreGattDatabase( this, this.backend, this.resourceLedger, connection, options, operation => this.assertReady(operation), (value, operation) => this.assertOperationAdmission(value, operation), (admissionEpoch, value, operation) => this.assertAdmissionCurrent(admissionEpoch, value, operation), this.admissionEpoch, reason, () => { const drain = this.operationCoordinator.waitForQuarantineDrainCancellable(key) return awaitWithOperationAdmission(drain.promise, options, this.options.now, 'rediscover').finally(() => { drain.cancel() }) } ) const discovery: CoreGattDiscovery = { kind: 'rediscover', promise } this.discoveries.set(key, discovery) const releaseRegistration = () => { if (this.discoveries.get(key) === discovery) { this.discoveries.delete(key) } } promise.then(releaseRegistration, releaseRegistration) lastAwaited = discovery const database = await awaitWithOperationAdmission(promise, options, this.options.now, 'rediscover') database.assertCurrent() return database } } async read( database: CoreGattDatabase, path: CurrentCharacteristicPath, options: PublicOperationOptions ): Promise { return readCoreCharacteristic( this.backend, this.operationCoordinator, this.options.maximumValueBytes, database, path, options ) } async write( database: CoreGattDatabase, path: CurrentCharacteristicPath, bytes: Readonly, options: WritePolicy ): Promise> { return writeCoreCharacteristic( this.backend, this.operationCoordinator, this.options.maximumValueBytes, database, path, bytes, options ) } async writeWhenReady( database: CoreGattDatabase, path: CurrentCharacteristicPath, bytes: Readonly, options: WritePolicy ): Promise> { return writeCoreCharacteristicWhenReady( this.backend, this.operationCoordinator, this.options.maximumValueBytes, database, path, bytes, options ) } async readDescriptor( database: CoreGattDatabase, path: CurrentDescriptorPath, options: PublicOperationOptions ): Promise { return readCoreDescriptor( this.backend, this.operationCoordinator, this.options.maximumValueBytes, database, path, options ) } async writeDescriptor( database: CoreGattDatabase, path: CurrentDescriptorPath, bytes: Readonly, options: WritePolicy ): Promise> { return writeCoreDescriptor( this.backend, this.operationCoordinator, this.options.maximumValueBytes, database, path, bytes, options ) } async subscribe( database: CoreGattDatabase, path: CurrentCharacteristicPath, options: SubscriptionOptions ): Promise> { database.assertPath(path) return this.subscriptions.subscribe(path, options, String(path.connectionId)) } async releaseConnection( connection: CoreConnection, cause: ConnectionLifecycleTerminalCause ): Promise { const key = String(connection.resource.connectionId) const inFlight = this.connectionReleases.get(key) if (inFlight !== undefined) { return inFlight } if (connection.isReleased()) { return { state: 'released', failures: [] } } const release = this.releaseConnectionCurrent(connection, cause) this.connectionReleases.set(key, release) try { return await release } finally { if (this.connectionReleases.get(key) === release) { this.connectionReleases.delete(key) } } } private async releaseConnectionCurrent( connection: CoreConnection, cause: ConnectionLifecycleTerminalCause ): Promise { const key = String(connection.resource.connectionId) this.operationCoordinator.cancelQueue(key, 'disconnected') const quarantine: CleanupRecord = this.operationCoordinator.hasPendingDrain(key) ? await this.awaitQuarantineDrain(key) : { state: 'released', failures: [] } const admissionFailures = [...this.operationCoordinator.takeCleanupFailures(key), ...quarantine.failures] const mergeAdmissionFailures = (record: CleanupRecord): CleanupRecord => { if (admissionFailures.length === 0) return record return { state: 'release-failed', failures: [...admissionFailures, ...record.failures] } } const disconnect = cause === 'requested-disconnect' const reason = isConnectionLossCause(cause) ? 'connection-lost' : 'owner-released' const cleanup = await connection.cleanupChildren(reason) const mergeChildFailures = (record: CleanupRecord): CleanupRecord => { if (cleanup.state === 'released') return record return { state: 'release-failed', failures: [...cleanup.failures, ...record.failures] } } let backendResult: CleanupRecord try { backendResult = disconnect && connection.isCurrent() ? await connection.resource.disconnect() : await connection.lease.release() } catch (error) { this.trace.record({ timestamp: this.options.now(), resource: 'connection', transition: 'release-rejected', operation: null, cause: error instanceof BackendContractError ? error.normalized.code : 'platform.failure', queuedOperations: 0, dispatchedOperations: 0, quarantinedOperations: 0 }) return mergeAdmissionFailures( mergeChildFailures( cleanupFailure('connection', contractError('platform.failure', 'cleanup', 'unified-core.connection-release')) ) ) } if (backendResult.state === 'release-failed') { return mergeAdmissionFailures(mergeChildFailures(backendResult)) } connection.finishLifecycle(cause, null) if (admissionFailures.length > 0) { return mergeAdmissionFailures(mergeChildFailures(backendResult)) } if (cleanup.state === 'release-failed') { return mergeAdmissionFailures({ state: 'release-failed', failures: cleanup.failures }) } connection.markReleased() this.connections.delete(String(connection.resource.connectionId)) this.resourceLedger.decrement('connectionLeases') return mergeAdmissionFailures(backendResult) } async invalidateDatabase( database: CoreGattDatabase, reason: 'connection-lost' | 'owner-released', changeReason: import('../backend-contract/gatt').GattDatabaseChangedEvent['reason'] | null = null ): Promise { const alreadyPending = database.connection.isPendingDatabaseCleanup(database) if (!database.isAttached() && !alreadyPending) { return { state: 'released', failures: [] } } if (!alreadyPending) { database.markInvalid(changeReason) if (database.connection.retainPendingDatabaseCleanup(database)) { this.resourceLedger.decrement('databaseSnapshots') } } const subscriptionReason = changeReason === 'service-changed' ? 'service-changed' : reason const cleanup = await this.subscriptions.invalidateDatabase(database.path, subscriptionReason) if (cleanup.state === 'released') { database.connection.completeDatabaseCleanup(database) } return cleanup } private adoptAttachedBackend(): void { this.coreState = 'new' try { assertAttachedBackend(this.construction.attachedBackend) } catch (error) { this.coreState = 'failed' if (error instanceof BackendContractError) { throw error } throw contractError('protocol.incompatible', 'core', 'unified-core.adopt-attached-backend') } const attachment = this.construction.attachedBackend.attachment if ( attachment.attachment.attachmentId !== attachment.identity.attachment.attachmentId || attachment.attachment.attachmentId !== this.backend.identity.attachment.attachmentId ) { this.coreState = 'failed' throw contractError('protocol.violation', 'core', 'unified-core.adopt-identity') } this.attachment = attachment this.idFactory = createAttachmentBoundIdFactory({ attachmentId: attachment.attachment.attachmentId, backendInstanceId: attachment.attachment.backendInstanceId, backendGeneration: attachment.attachment.backendGeneration, adapterId: attachment.attachment.adapter.adapterId, adapterGeneration: attachment.attachment.adapter.adapterGeneration }) this.coreState = 'ready' const events = this.backend.events() this.backendEventStream = events this.lifecycleObserver.observeBackground( forwardCoreBackendEvents({ events, isReady: () => this.coreState === 'ready', applyEvent: event => this.applyBackendEvent(event), releaseAfterFailure: () => this.releaseResources('backend-failure'), trace: this.trace, now: this.options.now }), 'manager', 'backend-event-pump' ) } private async forwardScanSource( scan: TrackedScan, source: BoundedAsyncStream> ): Promise { try { for await (const item of source) { if (scan.released || this.coreState !== 'ready') { return } if (item.kind === 'value') { const observation = cloneObservation(item.value, this.options.maximumValueBytes) const outcome = this.aggregateQuota.emit( scan.stream, observation, advertisementByteLength(observation), String(observation.device.id), advertisementPayloadByteLength(observation) ) if (outcome.terminated) { await this.stopScan(scan) return } } if (item.kind === 'overflow') { scan.stream.observeSourceOverflow(item) } if (item.kind === 'terminal') { scan.stream.closeWithReason(item.reason, item.error ?? null) await this.stopScan(scan) return } } } catch (error) { const code = error instanceof BackendContractError ? error.normalized.code : 'platform.failure' this.trace.record({ timestamp: this.options.now(), resource: 'scan', transition: 'source-failed', operation: null, cause: code, queuedOperations: 0, dispatchedOperations: 0, quarantinedOperations: 0 }) scan.stream.closeWithReason('source-failed') await this.stopScan(scan) } } private applyBackendEvent(event: BackendEvent): void { if (event.kind === 'backend-restarted' || event.kind === 'backend-restarting') { if (event.attachment.adapter.adapterId === this.requireAttachment().attachment.adapter.adapterId) { this.lifecycleObserver.observeCleanup(this.releaseResources('backend-restart'), 'backend-restarted-cleanup') } return } if (event.attachmentId !== this.attachmentId) { return } if (event.kind === 'database-changed') { for (const connection of this.connections.values()) { const database = connection.database if (database !== null && database.matchesDatabasePath(event.database)) { this.operationCoordinator.cancelQueue(String(connection.resource.connectionId), 'disconnected') this.lifecycleObserver.observeCleanup( this.invalidateDatabase(database, 'connection-lost', 'service-changed'), 'database-changed-cleanup' ) } } return } if (event.kind === 'connection-lost') { const connection = this.connections.get(String(event.connection.connectionId)) if (connection !== undefined && connection.matchesConnectionPath(event.connection)) { connection.finishLifecycle('peer-link-loss', event.ingressOrdinal) this.lifecycleObserver.observeCleanup( this.releaseConnection(connection, 'peer-link-loss'), 'backend-event-connection-cleanup' ) } return } if (event.kind === 'connection-state-changed') { const connection = this.connections.get(String(event.connection.connectionId)) if (connection !== undefined && connection.matchesConnectionPath(event.connection)) { if (event.current === 'disconnected' || event.current === 'lost') { if (event.reason === null) { throw contractError( 'lifecycle.invariant-violation', 'connection', 'backend-event-terminal-transition-reason' ) } const cause = lifecycleCauseFromBackendDisconnect(event.reason) connection.finishBackendLifecycle(event.previous, event.current, cause, event.ingressOrdinal) this.lifecycleObserver.observeCleanup( this.releaseConnection(connection, cause), 'backend-event-connection-state-cleanup' ) } else { connection.applyBackendTransition(event.previous, event.current, event.ingressOrdinal) } } return } if (event.kind === 'disconnected') { const connection = this.connections.get(String(event.connection.connectionId)) if (connection !== undefined && connection.matchesConnectionPath(event.connection)) { const cause = lifecycleCauseFromBackendDisconnect(event.reason) connection.finishLifecycle(cause, event.ingressOrdinal) this.lifecycleObserver.observeCleanup( this.releaseConnection(connection, cause), 'backend-event-disconnected-cleanup' ) } return } if (event.kind === 'adapter-state') { this.lifecycleObserver.observeBackground(this.applyAdapterStateEvent(), 'manager', 'backend-adapter-state') } } private async applyAdapterStateEvent(): Promise { const state = await this.backend.adapter.currentState() if (state.availability !== 'available' || isAuthorizationBlocking(state.authorization) || state.power !== 'on') { await this.releaseResources('adapter-loss') } } private async stopScan(scan: TrackedScan): Promise { if (scan.released) { return { state: 'released', failures: [] } } if (scan.stopInFlight !== null) { return scan.stopInFlight } scan.stream.closeWithReason('owner-released') deactivateScanLifetime(scan) const stop = this.stopScanPhysical(scan) scan.stopInFlight = stop stop.then( result => { scan.stopInFlight = null if (result.state === 'released') { this.finalizeStoppedScan(scan) } }, () => { scan.stopInFlight = null } ) return stop } private async stopScanPhysical(scan: TrackedScan): Promise { try { return await scan.lease.stop() } catch (error) { const cause = error instanceof BackendContractError ? error.normalized.code : 'platform.failure' this.trace.record({ timestamp: this.options.now(), resource: 'scan', transition: 'stop-rejected', operation: null, cause, queuedOperations: 0, dispatchedOperations: 0, quarantinedOperations: 0 }) return cleanupFailure('scan', contractError('platform.failure', 'cleanup', 'unified-core.scan-stop')) } } private finalizeStoppedScan(scan: TrackedScan): void { if (scan.released) { return } scan.released = true this.aggregateQuota.unregister(scan.stream) this.scans.delete(String(scan.leaseId)) this.resourceLedger.decrement('scanConsumers') if (scan.ownsPhysicalController) { this.resourceLedger.decrement('activeScanControllers') } } private async startOrJoinScan(options: ScanOptions): Promise> { if (options.sharing.mode === 'owner') { const ownerOptions: OwnerScanOptions = { ...options, sharing: options.sharing } return this.backend.scanner.start(ownerOptions, this.construction.clientId) } return this.backend.scanner.join(options.sharing.sharedLeaseId, options.sharing.token, this.construction.clientId) } private adoptScan( lease: ScanLease, stream: CoreBoundedStream>, options: ScanOptions ): TrackedScan { const tracked: TrackedScan = { scanSessionId: lease.scanSessionId, leaseId: lease.leaseId, shareToken: lease.shareToken, observations: stream, lease, stream, ownsPhysicalController: options.sharing.mode === 'owner', stopInFlight: null, released: false, activeAbortListener: null, activeAbortSignal: null, activeDeadline: null, stop: () => this.stopScan(tracked) } this.scans.set(String(lease.leaseId), tracked) this.resourceLedger.increment('scanConsumers') if (tracked.ownsPhysicalController) { this.resourceLedger.increment('activeScanControllers') } activateScanLifetime( tracked, options, this.options.now, this.options.timer, () => this.stopScan(tracked), (cleanup, transition) => this.lifecycleObserver.observeCleanup(cleanup, transition) ) this.lifecycleObserver.observeBackground( this.forwardScanSource(tracked, lease.observations), 'scan', 'scan-source-pump' ) return tracked } private async compensateUnadoptedLease( retry: () => Promise, resourceKind: CoreTraceResource, transition: string ): Promise { const retained = createRetainedCleanup(resourceKind, transition, retry) this.lifecycleObserver.retainCleanup(retained) const cleanup = retained.retry() this.lifecycleObserver.observeCleanup(cleanup, transition) const result = await this.lifecycleObserver.captureCleanup(cleanup, resourceKind, transition) if (result.state === 'released') { this.lifecycleObserver.dropCleanup(retained) } } private async destroyOwnedResources(cause: ConnectionLifecycleTerminalCause): Promise { const failures: CleanupFailure[] = [] const eventClose = this.closeBackendEventStream() for (const cancel of this.pendingAdapterStateWatches) { cancel(contractError('operation.cancelled-by-destroy', 'adapter', 'adapter-states')) } if (this.pendingAdapterStateWatches.size > 0) { failures.push( ...cleanupFailure( 'adapter', contractError('lifecycle.invalid-state', 'cleanup', 'adapter-states.acquisition-pending') ).failures ) } const pendingScanCancels = [...this.pendingScanAcquisitions] for (const cancel of pendingScanCancels) { cancel(contractError('operation.cancelled-by-destroy', 'core', 'scan')) } if (this.pendingScanAcquisitions.size > 0) { failures.push( ...cleanupFailure('scan', contractError('lifecycle.invalid-state', 'cleanup', 'scan.acquisition-pending')) .failures ) } const pendingConnectCancels = [...this.pendingConnectAcquisitions] for (const cancel of pendingConnectCancels) { cancel(contractError('operation.cancelled-by-destroy', 'core', 'connect')) } if (this.pendingConnectAcquisitions.size > 0) { failures.push( ...cleanupFailure( 'connection', contractError('lifecycle.invalid-state', 'cleanup', 'connect.acquisition-pending') ).failures ) } const retained = await this.lifecycleObserver.retryRetainedCleanups() failures.push(...retained.failures) for (const watch of [...this.adapterStateWatches]) { const result = await this.lifecycleObserver.captureCleanup(watch.stop(), 'manager', 'destroy-adapter-states') failures.push(...result.failures) } for (const scan of [...this.scans.values()]) { const result = await this.lifecycleObserver.captureCleanup(this.stopScan(scan), 'scan', 'destroy-scan') failures.push(...result.failures) } for (const connection of [...this.connections.values()]) { const result = await this.lifecycleObserver.captureCleanup( this.releaseConnection(connection, cause), 'connection', 'destroy-connection' ) failures.push(...result.failures) } const subscriptions = await this.lifecycleObserver.captureCleanup( this.subscriptions.destroy(), 'subscription', 'destroy-subscriptions' ) failures.push(...subscriptions.failures) const events = await eventClose failures.push(...events.failures) if ( this.operationCoordinator.hasPendingDrain() && !failures.some(failure => failure.resourceKind === 'operation-quarantine') ) { const quarantine = await this.awaitQuarantineDrain() if (quarantine.failures.length > 0) { failures.push(...quarantine.failures) } } failures.push(...this.operationCoordinator.takeCleanupFailures()) const result: CleanupRecord = failures.length === 0 ? { state: 'released', failures: [] } : { state: 'release-failed', failures } this.syncRetainedByteBuffers() this.coreState = result.state === 'released' ? 'destroyed' : 'failed' return result } private async awaitQuarantineDrain(queueKey?: string): Promise { if (!this.operationCoordinator.hasPendingDrain(queueKey)) { return { state: 'released', failures: [] } } const drainDeadline = createDeadline(this.options.now() + QUARANTINE_DRAIN_TIMEOUT_MS) const drain = this.operationCoordinator.waitForQuarantineDrainCancellable(queueKey) try { await awaitWithOperationAdmission( drain.promise, { signal: null, deadline: drainDeadline }, this.options.now, 'unified-core.quarantine-drain' ) } catch (error) { if (!this.operationCoordinator.hasPendingDrain(queueKey)) { return { state: 'released', failures: [] } } const failure = error instanceof BackendContractError && error.normalized.code === 'operation.timed-out' ? contractError('operation.timed-out', 'cleanup', 'unified-core.quarantine-drain') : contractError('platform.failure', 'cleanup', 'unified-core.quarantine-drain') return cleanupFailure('operation-quarantine', failure) } finally { drain.cancel() } return { state: 'released', failures: [] } } private assertReady(operation: string): void { if (this.coreState !== 'ready') { throw contractError( this.coreState === 'destroyed' ? 'lifecycle.destroyed' : 'lifecycle.invalid-state', 'core', operation ) } } private assertOperationAdmission(options: PublicOperationOptions, operation: string): void { if (options.signal?.aborted === true) { throw contractError('operation.aborted', 'core', operation) } if (options.deadline !== null && options.deadline <= this.options.now()) { throw contractError('operation.timed-out', 'core', operation) } } private assertAdmissionCurrent(admissionEpoch: number, options: PublicOperationOptions, operation: string): void { const closed = this.admissionClosedError(admissionEpoch, options, operation) if (closed !== null) { throw closed } } private admissionClosedError( admissionEpoch: number, options: PublicOperationOptions, operation: string ): BackendContractError | null { if (options.signal?.aborted === true) { return contractError('operation.aborted', 'core', operation) } if (options.deadline !== null && options.deadline <= this.options.now()) { return contractError('operation.timed-out', 'core', operation) } if (this.coreState !== 'ready' || this.admissionEpoch !== admissionEpoch) { return contractError('operation.cancelled-by-destroy', 'core', operation) } return null } private isCurrentPath(path: CurrentCharacteristicPath): boolean { if (this.coreState !== 'ready' || path.attachmentId !== this.attachmentId || path.validity !== 'current') { return false } const connection = this.connections.get(String(path.connectionId)) return connection !== undefined && connection.isPathCurrent(path) } private requireAttachment(): BackendAttachment { if (this.attachment === null) { throw contractError('lifecycle.invalid-state', 'core', 'unified-core.attachment') } return this.attachment } private requireIdFactory(): AttachmentBoundIdFactory { if (this.idFactory === null) { throw contractError('lifecycle.invalid-state', 'core', 'unified-core.id-factory') } return this.idFactory } private nextOperationValue(): number { const value = this.nextOperation this.nextOperation += 1 return value } private syncRetainedByteBuffers(): void { this.resourceLedger.setRetainedStreamBytes(this.aggregateQuota.retainedPayloadBytes()) } private async closeBackendEventStream(): Promise { const events = this.backendEventStream if (events === null) { return { state: 'released', failures: [] } } const cleanup = await this.lifecycleObserver.captureCleanup(events.close(), 'manager', 'destroy-backend-events') if (cleanup.state === 'released') { this.backendEventStream = null } return cleanup } } function resolveLongWriteChunkSize(observedMaximum: number, requestedChunkSize: number | undefined): number { if (requestedChunkSize === undefined) return observedMaximum if (!Number.isSafeInteger(requestedChunkSize) || requestedChunkSize < 1 || requestedChunkSize > observedMaximum) { throw contractError('argument.invalid', 'gatt', 'unified-core.write-long.chunk-size') } return requestedChunkSize } function abortRequested(signal: AbortSignal | null | undefined): boolean { return signal !== null && signal !== undefined && signal.aborted } function isRediscoverCallerTerminal(error: unknown, options: PublicOperationOptions, now: () => number): boolean { if (!(error instanceof BackendContractError)) { return false } const code = error.normalized.code return ( (code === 'operation.aborted' && abortRequested(options.signal)) || (code === 'operation.timed-out' && options.deadline !== null && options.deadline <= now()) || code === 'operation.cancelled-by-destroy' ) } function closeAdapterStateStream(stream: BoundedAsyncStream): Promise { return stream.close() }