diff --git a/packages/runtime-host/src/__tests__/host-change-feed.test.ts b/packages/runtime-host/src/__tests__/host-change-feed.test.ts new file mode 100644 index 0000000000..89f1b5cfce --- /dev/null +++ b/packages/runtime-host/src/__tests__/host-change-feed.test.ts @@ -0,0 +1,69 @@ +import assert from 'node:assert/strict'; +import { test } from 'node:test'; +import { HostChangeFeed } from '../server/host-change-feed.js'; + +test('routes each change kind only to subscribed connections', () => { + const feed = new HostChangeFeed(); + const configuration: unknown[] = []; + const project: unknown[] = []; + const all: unknown[] = []; + feed.attachConnection( + 'configuration', + { configuration: true }, + { send: async (frame) => void configuration.push(frame) }, + ); + feed.attachConnection( + 'project', + { projectCatalog: true }, + { send: async (frame) => void project.push(frame) }, + ); + feed.attachConnection( + 'all', + { configuration: true, projectCatalog: true, sessionCatalog: true, scheduledTask: true }, + { send: async (frame) => void all.push(frame) }, + ); + + feed.publishConfiguration(); + feed.publishProjectCatalog(); + feed.publishSessionCatalog('session-1'); + feed.publishScheduledTask(7, 'updated', 'task-1'); + + assert.deepEqual( + configuration.map((frame) => (frame as { kind: string }).kind), + ['configuration.changed'], + ); + assert.deepEqual( + project.map((frame) => (frame as { kind: string }).kind), + ['project.catalog.changed'], + ); + assert.equal(all.length, 4); +}); + +test('keeps catalog revisions independent and removes failed subscriptions', async () => { + const feed = new HostChangeFeed(); + const frames: unknown[] = []; + feed.attachConnection( + 'working', + { projectCatalog: true, sessionCatalog: true }, + { send: async (frame) => void frames.push(frame) }, + ); + feed.attachConnection( + 'failed', + { projectCatalog: true }, + { + send: async () => { + throw new Error('closed'); + }, + }, + ); + + feed.publishProjectCatalog(); + feed.publishSessionCatalog('session-1'); + await Promise.resolve(); + feed.publishProjectCatalog(); + + assert.deepEqual( + frames.map((frame) => (frame as { kind: string; revision: number }).revision), + [1, 1, 2], + ); +}); diff --git a/packages/runtime-host/src/__tests__/host-kernel.test.ts b/packages/runtime-host/src/__tests__/host-kernel.test.ts index 288c1bd4cc..bab8663267 100644 --- a/packages/runtime-host/src/__tests__/host-kernel.test.ts +++ b/packages/runtime-host/src/__tests__/host-kernel.test.ts @@ -65,8 +65,7 @@ import { } from '../server/candidate.js'; import type { RuntimeHostCompositionSource } from '../server/host-composition.js'; import { createUnavailableDomainOperationHandlers } from '../server/operation-dispatcher.js'; -import { HostConfigurationChangeService } from '../server/configuration-change-service.js'; -import { HostSessionCatalogChangeService } from '../server/session-catalog-change-service.js'; +import { HostChangeFeed } from '../server/host-change-feed.js'; import { FramedTransport, RuntimeHostTransportError } from '../transport/framed-transport.js'; import { prepareStorageRootControlDirectory, @@ -1880,8 +1879,7 @@ describe('non-serving Runtime Host kernel', () => { const owner = await tryAcquireInteractiveRootOwner(capability); assert.ok(owner); if (!owner) return; - const configurationChanges = new HostConfigurationChangeService(); - const sessionCatalogChanges = new HostSessionCatalogChangeService(); + const hostChanges = new HostChangeFeed(); let releaseFactory!: () => void; let markFactoryEntered!: () => void; const factoryEntered = new Promise((resolve) => { @@ -1898,8 +1896,7 @@ describe('non-serving Runtime Host kernel', () => { await factoryReleased; return { handlers: createUnavailableDomainOperationHandlers(), - configurationChanges, - sessionCatalogChanges, + hostChanges, beginDrain() {}, async recover() {}, async close() {}, @@ -1927,8 +1924,8 @@ describe('non-serving Runtime Host kernel', () => { }); releaseFactory(); host = await hostTask; - configurationChanges.publish(); - sessionCatalogChanges.publish('session-1'); + hostChanges.publishConfiguration(); + hostChanges.publishSessionCatalog('session-1'); assert.equal( await withTimeout(observed, 1_000, 'Client did not receive configuration change'), 1, diff --git a/packages/runtime-host/src/__tests__/project-catalog-coordinator.test.ts b/packages/runtime-host/src/__tests__/project-catalog-coordinator.test.ts index f20e60cf35..6ffd3c7fdf 100644 --- a/packages/runtime-host/src/__tests__/project-catalog-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/project-catalog-coordinator.test.ts @@ -5,10 +5,9 @@ import { join } from 'node:path'; import { test } from 'node:test'; import { createProjectCatalog, createSessionStore } from '@maka/storage'; import type { ConnectionContext } from '../server/operation-dispatcher.js'; -import { HostProjectCatalogChangeService } from '../server/project-catalog-change-service.js'; +import { HostChangeFeed } from '../server/host-change-feed.js'; import { HostProjectCatalogCoordinator } from '../server/project-catalog-coordinator.js'; import { HostProjectMembershipGate } from '../server/project-membership-gate.js'; -import { HostSessionCatalogChangeService } from '../server/session-catalog-change-service.js'; test('Host Project Catalog relink merges identities and reassigns every affected Session', async () => { const base = await mkdtemp(join(tmpdir(), 'maka-host-project-catalog-')); @@ -24,30 +23,36 @@ test('Host Project Catalog relink merges identities and reassigns every affected })(), }); const sessions = createSessionStore(storageRoot); - const projectChanges = new HostProjectCatalogChangeService(); - const sessionChanges = new HostSessionCatalogChangeService(); + const hostChanges = new HostChangeFeed(); const projectFrames: unknown[] = []; const sessionFrames: unknown[] = []; - projectChanges.attachConnection('desktop', { - send: async (frame) => { - projectFrames.push(frame); + hostChanges.attachConnection( + 'desktop', + { + projectCatalog: true, + sessionCatalog: true, }, - }); - projectChanges.attachConnection('tui', { - send: async (frame) => { - projectFrames.push(frame); + { + send: async (frame) => { + if (frame.kind === 'project.catalog.changed') projectFrames.push(frame); + else sessionFrames.push(frame); + }, }, - }); - sessionChanges.attachConnection('desktop', { - send: async (frame) => { - sessionFrames.push(frame); + ); + hostChanges.attachConnection( + 'tui', + { projectCatalog: true }, + { + send: async (frame) => { + projectFrames.push(frame); + }, }, - }); + ); const membership = new HostProjectMembershipGate(); const coordinator = new HostProjectCatalogCoordinator( catalog, - projectChanges, - sessionChanges, + { publish: () => hostChanges.publishProjectCatalog() }, + { publish: (sessionId: string) => hostChanges.publishSessionCatalog(sessionId) }, membership, () => assert.fail('ordinary project mutations must not drain the Host'), ); @@ -129,10 +134,11 @@ test('directory resolution failures cannot enter the unknown-commit drain path', const base = await mkdtemp(join(tmpdir(), 'maka-host-project-directory-failure-')); const catalog = createProjectCatalog(join(base, 'storage')); let drains = 0; + const hostChanges = new HostChangeFeed(); const coordinator = new HostProjectCatalogCoordinator( catalog, - new HostProjectCatalogChangeService(), - new HostSessionCatalogChangeService(), + { publish: () => hostChanges.publishProjectCatalog() }, + { publish: (sessionId: string) => hostChanges.publishSessionCatalog(sessionId) }, new HostProjectMembershipGate(), () => { drains += 1; diff --git a/packages/runtime-host/src/server/configuration-change-service.ts b/packages/runtime-host/src/server/configuration-change-service.ts deleted file mode 100644 index a989ef0e5c..0000000000 --- a/packages/runtime-host/src/server/configuration-change-service.ts +++ /dev/null @@ -1,39 +0,0 @@ -import type { ConfigurationChangedFrame } from '../protocol/index.js'; - -export interface ConfigurationChangeConnection { - close(): void; -} - -interface ConfigurationChangeSink { - send(frame: ConfigurationChangedFrame): Promise; -} - -export class HostConfigurationChangeService { - readonly #connections = new Map(); - #revision = 0; - - attachConnection( - connectionId: string, - sink: ConfigurationChangeSink, - ): ConfigurationChangeConnection { - this.#connections.set(connectionId, sink); - return { - close: () => { - if (this.#connections.get(connectionId) === sink) this.#connections.delete(connectionId); - }, - }; - } - - publish(): void { - this.#revision += 1; - const frame: ConfigurationChangedFrame = { - kind: 'configuration.changed', - revision: this.#revision, - }; - for (const [connectionId, sink] of this.#connections) { - void sink.send(frame).catch(() => { - if (this.#connections.get(connectionId) === sink) this.#connections.delete(connectionId); - }); - } - } -} diff --git a/packages/runtime-host/src/server/connection-session.ts b/packages/runtime-host/src/server/connection-session.ts index 8bb98e6d4a..b969082570 100644 --- a/packages/runtime-host/src/server/connection-session.ts +++ b/packages/runtime-host/src/server/connection-session.ts @@ -24,22 +24,7 @@ import type { ClientCapabilityConnection, ClientCapabilityService, } from './client-capability-service.js'; -import type { - ConfigurationChangeConnection, - HostConfigurationChangeService, -} from './configuration-change-service.js'; -import type { - HostSessionCatalogChangeService, - SessionCatalogChangeConnection, -} from './session-catalog-change-service.js'; -import type { - HostProjectCatalogChangeService, - ProjectCatalogChangeConnection, -} from './project-catalog-change-service.js'; -import type { - HostScheduledTaskChangeService, - ScheduledTaskChangeConnection, -} from './scheduled-task-change-service.js'; +import type { HostChangeFeed, HostChangeSubscription } from './host-change-feed.js'; import type { RuntimeHostConnectionAuthority } from './connection-authority.js'; import { authorizeClientCapabilityFrame, @@ -64,10 +49,7 @@ export interface RuntimeHostConnectionSessionOptions { resolveHandlers(): OperationHandlerMap; resolveContinuity(): SessionContinuityService | undefined; resolveClientCapabilities?(): ClientCapabilityService | undefined; - resolveConfigurationChanges?(): HostConfigurationChangeService | undefined; - resolveProjectCatalogChanges?(): HostProjectCatalogChangeService | undefined; - resolveSessionCatalogChanges?(): HostSessionCatalogChangeService | undefined; - resolveScheduledTaskChanges?(): HostScheduledTaskChangeService | undefined; + resolveHostChanges?(): HostChangeFeed | undefined; beginOperation(frame: RequestFrame): Promise; onTeardown(): void; } @@ -83,10 +65,7 @@ export class RuntimeHostConnectionSession { #clientCapabilityService: ClientCapabilityService | undefined; #clientCapabilities: ClientCapabilityConnection | undefined; #clientCapabilityCloseTask: Promise | undefined; - #configurationChanges: ConfigurationChangeConnection | undefined; - #projectCatalogChanges: ProjectCatalogChangeConnection | undefined; - #sessionCatalogChanges: SessionCatalogChangeConnection | undefined; - #scheduledTaskChanges: ScheduledTaskChangeConnection | undefined; + #hostChanges: HostChangeSubscription | undefined; #inputClosed = false; #closed = false; @@ -121,10 +100,7 @@ export class RuntimeHostConnectionSession { this.#inputClosed = true; this.#detachContinuity(); this.#detachClientCapabilities(); - this.#detachConfigurationChanges(); - this.#detachProjectCatalogChanges(); - this.#detachSessionCatalogChanges(); - this.#detachScheduledTaskChanges(); + this.#detachHostChanges(); const outcome = await Promise.race([ Promise.allSettled([...this.#requests.values()]).then(() => 'drained' as const), this.#options.transport.closed.then(() => 'closed' as const), @@ -314,104 +290,45 @@ export class RuntimeHostConnectionSession { void this.#clientCapabilityCloseTask.catch(() => undefined); } - #attachConfigurationChanges(): void { - if (!hasRuntimeHostOperationGrant(this.#options.connection.authority, 'runtime.policy.query')) { - return; - } - const service = this.#options.resolveConfigurationChanges?.(); - if (!service || this.#configurationChanges) return; - this.#configurationChanges = service.attachConnection(this.#options.connection.connectionId, { - send: (frame) => { - try { - return this.#writer.enqueue(frame).flushed; - } catch (error) { - return Promise.reject(error); - } - }, - }); - } - attachGlobalChanges(): void { if (this.#closed || this.#inputClosed) return; - this.#attachConfigurationChanges(); - this.#attachProjectCatalogChanges(); - this.#attachSessionCatalogChanges(); - this.#attachScheduledTaskChanges(); - } - - #detachConfigurationChanges(): void { - this.#configurationChanges?.close(); - this.#configurationChanges = undefined; - } - - #attachProjectCatalogChanges(): void { - if ( - !hasRuntimeHostOperationGrant(this.#options.connection.authority, 'project.catalog.query') - ) { - return; - } - const service = this.#options.resolveProjectCatalogChanges?.(); - if (!service || this.#projectCatalogChanges) return; - this.#projectCatalogChanges = service.attachConnection(this.#options.connection.connectionId, { - send: (frame) => { - try { - return this.#writer.enqueue(frame).flushed; - } catch (error) { - return Promise.reject(error); - } - }, - }); - } - - #detachProjectCatalogChanges(): void { - this.#projectCatalogChanges?.close(); - this.#projectCatalogChanges = undefined; - } - - #attachSessionCatalogChanges(): void { - if ( - !hasRuntimeHostOperationGrant(this.#options.connection.authority, 'session.catalog.query') - ) { - return; - } - const service = this.#options.resolveSessionCatalogChanges?.(); - if (!service || this.#sessionCatalogChanges) return; - this.#sessionCatalogChanges = service.attachConnection(this.#options.connection.connectionId, { - send: (frame) => { - try { - return this.#writer.enqueue(frame).flushed; - } catch (error) { - return Promise.reject(error); - } + const service = this.#options.resolveHostChanges?.(); + if (!service || this.#hostChanges) return; + this.#hostChanges = service.attachConnection( + this.#options.connection.connectionId, + { + configuration: hasRuntimeHostOperationGrant( + this.#options.connection.authority, + 'runtime.policy.query', + ), + projectCatalog: hasRuntimeHostOperationGrant( + this.#options.connection.authority, + 'project.catalog.query', + ), + sessionCatalog: hasRuntimeHostOperationGrant( + this.#options.connection.authority, + 'session.catalog.query', + ), + scheduledTask: hasRuntimeHostOperationGrant( + this.#options.connection.authority, + 'scheduled-task.query', + ), }, - }); - } - - #detachSessionCatalogChanges(): void { - this.#sessionCatalogChanges?.close(); - this.#sessionCatalogChanges = undefined; - } - - #attachScheduledTaskChanges(): void { - if (!hasRuntimeHostOperationGrant(this.#options.connection.authority, 'scheduled-task.query')) { - return; - } - const service = this.#options.resolveScheduledTaskChanges?.(); - if (!service || this.#scheduledTaskChanges) return; - this.#scheduledTaskChanges = service.attachConnection(this.#options.connection.connectionId, { - send: (frame) => { - try { - return this.#writer.enqueue(frame).flushed; - } catch (error) { - return Promise.reject(error); - } + { + send: (frame) => { + try { + return this.#writer.enqueue(frame).flushed; + } catch (error) { + return Promise.reject(error); + } + }, }, - }); + ); } - #detachScheduledTaskChanges(): void { - this.#scheduledTaskChanges?.close(); - this.#scheduledTaskChanges = undefined; + #detachHostChanges(): void { + this.#hostChanges?.close(); + this.#hostChanges = undefined; } #teardown(): void { @@ -420,10 +337,7 @@ export class RuntimeHostConnectionSession { this.#inputClosed = true; this.#detachContinuity(); this.#detachClientCapabilities(); - this.#detachConfigurationChanges(); - this.#detachProjectCatalogChanges(); - this.#detachSessionCatalogChanges(); - this.#detachScheduledTaskChanges(); + this.#detachHostChanges(); this.#writer.close(); this.#options.transport.abort(); this.#options.onTeardown(); diff --git a/packages/runtime-host/src/server/execution-composition.ts b/packages/runtime-host/src/server/execution-composition.ts index 9d1a58fc1f..fabbcfeddf 100644 --- a/packages/runtime-host/src/server/execution-composition.ts +++ b/packages/runtime-host/src/server/execution-composition.ts @@ -93,9 +93,7 @@ import { HostAgentGraphExecutionCoordinator } from './agent-graph-execution-coor import { HostScheduledTaskCoordinator } from './scheduled-task-coordinator.js'; import { recoverClientCapabilityOutcomes } from './client-capability-recovery.js'; import { HostConnectionEffectCoordinator } from './connection-effect-coordinator.js'; -import { HostConfigurationChangeService } from './configuration-change-service.js'; -import { HostSessionCatalogChangeService } from './session-catalog-change-service.js'; -import { HostScheduledTaskChangeService } from './scheduled-task-change-service.js'; +import { HostChangeFeed } from './host-change-feed.js'; import { HostConfigurationCoordinator } from './configuration-coordinator.js'; import { HostContextCoordinator } from './context-coordinator.js'; import { HostClientCapabilityCoordinator } from './client-capability-coordinator.js'; @@ -141,7 +139,6 @@ import { HostNetworkProxyCoordinator } from './network-proxy-coordinator.js'; import { HostOAuthExecutionAuthority } from './oauth-execution-authority.js'; import { HostOAuthCoordinator, type HostOAuthCoordinatorInput } from './oauth-coordinator.js'; import { HostPlanCoordinator } from './plan-coordinator.js'; -import { HostProjectCatalogChangeService } from './project-catalog-change-service.js'; import { HostProjectDirectoryAuthority, type PublishedProjectDirectoryRoot, @@ -428,15 +425,12 @@ export async function createExecutionRuntimeHostComposition( initiatingConnectionId: string, ) => Promise) | undefined; - const configurationChanges = new HostConfigurationChangeService(); - const sessionCatalogChanges = new HostSessionCatalogChangeService(); - const scheduledTaskChanges = new HostScheduledTaskChangeService(); - const projectCatalogChanges = new HostProjectCatalogChangeService(); + const hostChanges = new HostChangeFeed(); const projectMembership = new HostProjectMembershipGate(); const workspaceResolver = new HostWorkspaceResolver( openedProjectCatalog, projectMembership, - () => projectCatalogChanges.publish(), + () => hostChanges.publishProjectCatalog(), ); const skills = new HostSkillCatalogCoordinator( new SkillCatalogRepository({ @@ -476,8 +470,8 @@ export async function createExecutionRuntimeHostComposition( ); const projects = new HostProjectCatalogCoordinator( openedProjectCatalog, - projectCatalogChanges, - sessionCatalogChanges, + { publish: () => hostChanges.publishProjectCatalog() }, + { publish: (sessionId: string) => hostChanges.publishSessionCatalog(sessionId) }, projectMembership, context.requestDrain, new HostProjectDirectoryAuthority(options.projectDirectoryRoots), @@ -540,7 +534,7 @@ export async function createExecutionRuntimeHostComposition( sessionAdmission, context.requestDrain, createSessionTranscriptReader({ stores, canonicalPermissionOutcomes }), - (sessionId) => sessionCatalogChanges.publish(sessionId), + (sessionId) => hostChanges.publishSessionCatalog(sessionId), ); const continuityCoordinator = continuity; unsubscribeTranscriptChanges = stores.sessionStore.subscribeTranscriptChanges((sessionId) => @@ -1035,7 +1029,7 @@ export async function createExecutionRuntimeHostComposition( observeBackendInvalidation(manager.refreshIdleBackends()); }; const registerConfigurationMutation = (): void => { - configurationChanges.publish(); + hostChanges.publishConfiguration(); registerBackendInvalidation(); }; clientCapabilities = new HostClientCapabilityCoordinator({ @@ -1050,7 +1044,7 @@ export async function createExecutionRuntimeHostComposition( isProviderEnabled: isOAuthEnrollmentProviderEnabled, acquireResidency: () => context.acquireResidency('oauth'), invalidateBackends: () => { - configurationChanges.publish(); + hostChanges.publishConfiguration(); return manager.refreshIdleBackends(); }, onFatal: (error) => { @@ -1249,7 +1243,13 @@ export async function createExecutionRuntimeHostComposition( runtimePolicy: runtimePolicyStores, nativeEffects: clientCapabilities, createSession: (input) => sessionCatalog.createForHost(input), - changes: scheduledTaskChanges, + changes: { + publish: ( + revision: number, + reason: Parameters[1], + taskId: string, + ) => hostChanges.publishScheduledTask(revision, reason, taskId), + }, acquireResidency: () => context.acquireResidency('scheduled-task'), requestDrain: context.requestDrain, }); @@ -1622,10 +1622,7 @@ export async function createExecutionRuntimeHostComposition( workspaceExecution: requireWorkspaceExecution(workspaceExecution), continuity: continuityCoordinator, clientCapabilities, - configurationChanges, - projectCatalogChanges, - sessionCatalogChanges, - scheduledTaskChanges, + hostChanges, releaseConnection: (connectionId: string) => { for (const module of domainModules) module.releaseConnection?.(connectionId); }, diff --git a/packages/runtime-host/src/server/host-change-feed.ts b/packages/runtime-host/src/server/host-change-feed.ts new file mode 100644 index 0000000000..92e66f5217 --- /dev/null +++ b/packages/runtime-host/src/server/host-change-feed.ts @@ -0,0 +1,120 @@ +import type { + ConfigurationChangedFrame, + ProjectCatalogChangedFrame, + ScheduledTaskChangedFrame, + ScheduledTaskChangedReason, + SessionCatalogChangedFrame, +} from '../protocol/index.js'; + +export type HostChangeFrame = + | ConfigurationChangedFrame + | ProjectCatalogChangedFrame + | SessionCatalogChangedFrame + | ScheduledTaskChangedFrame; + +export type HostChangeKind = + | 'configuration' + | 'project_catalog' + | 'session_catalog' + | 'scheduled_task'; + +export interface HostChangeSubscription { + close(): void; +} + +export interface HostChangeSubscriptionMask { + readonly configuration?: boolean; + readonly projectCatalog?: boolean; + readonly sessionCatalog?: boolean; + readonly scheduledTask?: boolean; +} + +interface HostChangeSink { + send(frame: HostChangeFrame): Promise; +} + +interface Subscription { + readonly sink: HostChangeSink; + readonly mask: HostChangeSubscriptionMask; +} + +export class HostChangeFeed { + readonly #subscriptions = new Map(); + #configurationRevision = 0; + #projectCatalogRevision = 0; + #sessionCatalogRevision = 0; + + attachConnection( + connectionId: string, + mask: HostChangeSubscriptionMask, + sink: HostChangeSink, + ): HostChangeSubscription { + const subscription = { sink, mask } satisfies Subscription; + this.#subscriptions.set(connectionId, subscription); + return { + close: () => { + if (this.#subscriptions.get(connectionId) === subscription) { + this.#subscriptions.delete(connectionId); + } + }, + }; + } + + publishConfiguration(): void { + this.#configurationRevision += 1; + this.#publish('configuration', { + kind: 'configuration.changed', + revision: this.#configurationRevision, + }); + } + + publishProjectCatalog(): void { + this.#projectCatalogRevision += 1; + this.#publish('project_catalog', { + kind: 'project.catalog.changed', + revision: this.#projectCatalogRevision, + }); + } + + publishSessionCatalog(sessionId: string): void { + this.#sessionCatalogRevision += 1; + this.#publish('session_catalog', { + kind: 'session.catalog.changed', + revision: this.#sessionCatalogRevision, + sessionId, + }); + } + + publishScheduledTask(revision: number, reason: ScheduledTaskChangedReason, taskId: string): void { + this.#publish('scheduled_task', { + kind: 'scheduled-task.changed', + revision, + reason, + taskId, + }); + } + + #publish(kind: HostChangeKind, frame: HostChangeFrame): void { + for (const [connectionId, subscription] of this.#subscriptions) { + if (!isSubscribed(subscription.mask, kind)) continue; + void subscription.sink.send(frame).catch(() => { + if (this.#subscriptions.get(connectionId) === subscription) { + this.#subscriptions.delete(connectionId); + } + }); + } + } +} + +function isSubscribed(mask: HostChangeSubscriptionMask, kind: HostChangeKind): boolean { + switch (kind) { + case 'configuration': + return mask.configuration === true; + case 'project_catalog': + return mask.projectCatalog === true; + case 'session_catalog': + return mask.sessionCatalog === true; + case 'scheduled_task': + return mask.scheduledTask === true; + } +} diff --git a/packages/runtime-host/src/server/host-kernel.ts b/packages/runtime-host/src/server/host-kernel.ts index c0d9169b0d..eb5832bd7f 100644 --- a/packages/runtime-host/src/server/host-kernel.ts +++ b/packages/runtime-host/src/server/host-kernel.ts @@ -46,11 +46,8 @@ import { import type { RuntimeHostConnectionAuthority } from './connection-authority.js'; import type { SessionContinuityService } from './session-continuity-service.js'; import type { ClientCapabilityService } from './client-capability-service.js'; -import type { HostConfigurationChangeService } from './configuration-change-service.js'; -import type { HostProjectCatalogChangeService } from './project-catalog-change-service.js'; +import type { HostChangeFeed } from './host-change-feed.js'; import { runtimeHostLogBuffer } from '../process-diagnostics.js'; -import type { HostSessionCatalogChangeService } from './session-catalog-change-service.js'; -import type { HostScheduledTaskChangeService } from './scheduled-task-change-service.js'; import { type HostCompositionDescriptor, type RuntimeHostCompositionSource, @@ -101,10 +98,7 @@ export interface RuntimeHostComposition { readonly moduleIds?: readonly string[]; readonly continuity?: SessionContinuityService; readonly clientCapabilities?: ClientCapabilityService; - readonly configurationChanges?: HostConfigurationChangeService; - readonly projectCatalogChanges?: HostProjectCatalogChangeService; - readonly sessionCatalogChanges?: HostSessionCatalogChangeService; - readonly scheduledTaskChanges?: HostScheduledTaskChangeService; + readonly hostChanges?: HostChangeFeed; releaseConnection?(connectionId: string): void; beginDrain(): void; recover(): Promise; @@ -376,10 +370,7 @@ export class RuntimeHostKernel { resolveHandlers: () => this.#operationHandlers, resolveContinuity: () => this.#composition?.continuity, resolveClientCapabilities: () => this.#composition?.clientCapabilities, - resolveConfigurationChanges: () => this.#composition?.configurationChanges, - resolveProjectCatalogChanges: () => this.#composition?.projectCatalogChanges, - resolveSessionCatalogChanges: () => this.#composition?.sessionCatalogChanges, - resolveScheduledTaskChanges: () => this.#composition?.scheduledTaskChanges, + resolveHostChanges: () => this.#composition?.hostChanges, beginOperation: (request) => this.#beginOperation(request), onTeardown: releaseTransport, }); diff --git a/packages/runtime-host/src/server/project-catalog-change-service.ts b/packages/runtime-host/src/server/project-catalog-change-service.ts deleted file mode 100644 index 62f43e9ae3..0000000000 --- a/packages/runtime-host/src/server/project-catalog-change-service.ts +++ /dev/null @@ -1,39 +0,0 @@ -import type { ProjectCatalogChangedFrame } from '../protocol/index.js'; - -export interface ProjectCatalogChangeConnection { - close(): void; -} - -interface ProjectCatalogChangeSink { - send(frame: ProjectCatalogChangedFrame): Promise; -} - -export class HostProjectCatalogChangeService { - readonly #connections = new Map(); - #revision = 0; - - attachConnection( - connectionId: string, - sink: ProjectCatalogChangeSink, - ): ProjectCatalogChangeConnection { - this.#connections.set(connectionId, sink); - return { - close: () => { - if (this.#connections.get(connectionId) === sink) this.#connections.delete(connectionId); - }, - }; - } - - publish(): void { - this.#revision += 1; - const frame: ProjectCatalogChangedFrame = { - kind: 'project.catalog.changed', - revision: this.#revision, - }; - for (const [connectionId, sink] of this.#connections) { - void sink.send(frame).catch(() => { - if (this.#connections.get(connectionId) === sink) this.#connections.delete(connectionId); - }); - } - } -} diff --git a/packages/runtime-host/src/server/project-catalog-coordinator.ts b/packages/runtime-host/src/server/project-catalog-coordinator.ts index 6d977127b0..245099c43a 100644 --- a/packages/runtime-host/src/server/project-catalog-coordinator.ts +++ b/packages/runtime-host/src/server/project-catalog-coordinator.ts @@ -27,9 +27,7 @@ import { HostProjectDirectoryAuthority, type ResolvedProjectDirectoryRegistration, } from './project-directory-authority.js'; -import type { HostProjectCatalogChangeService } from './project-catalog-change-service.js'; import type { HostProjectMembershipGate } from './project-membership-gate.js'; -import type { HostSessionCatalogChangeService } from './session-catalog-change-service.js'; export class HostProjectCatalogCoordinator { readonly handlers: ProjectCatalogOperationHandlerMap = { @@ -39,8 +37,8 @@ export class HostProjectCatalogCoordinator { constructor( private readonly catalog: ProjectCatalog, - private readonly projectChanges: HostProjectCatalogChangeService, - private readonly sessionChanges: HostSessionCatalogChangeService, + private readonly projectChanges: { publish(): void }, + private readonly sessionChanges: { publish(sessionId: string): void }, private readonly membership: HostProjectMembershipGate, private readonly requestDrain: () => void, private readonly directories = new HostProjectDirectoryAuthority(), diff --git a/packages/runtime-host/src/server/scheduled-task-change-service.ts b/packages/runtime-host/src/server/scheduled-task-change-service.ts deleted file mode 100644 index daa7542925..0000000000 --- a/packages/runtime-host/src/server/scheduled-task-change-service.ts +++ /dev/null @@ -1,39 +0,0 @@ -import type { ScheduledTaskChangedFrame, ScheduledTaskChangedReason } from '../protocol/index.js'; - -export interface ScheduledTaskChangeConnection { - close(): void; -} - -interface ScheduledTaskChangeSink { - send(frame: ScheduledTaskChangedFrame): Promise; -} - -export class HostScheduledTaskChangeService { - readonly #connections = new Map(); - - attachConnection( - connectionId: string, - sink: ScheduledTaskChangeSink, - ): ScheduledTaskChangeConnection { - this.#connections.set(connectionId, sink); - return { - close: () => { - if (this.#connections.get(connectionId) === sink) this.#connections.delete(connectionId); - }, - }; - } - - publish(revision: number, reason: ScheduledTaskChangedReason, taskId: string): void { - const frame: ScheduledTaskChangedFrame = { - kind: 'scheduled-task.changed', - revision, - reason, - taskId, - }; - for (const [connectionId, sink] of this.#connections) { - void sink.send(frame).catch(() => { - if (this.#connections.get(connectionId) === sink) this.#connections.delete(connectionId); - }); - } - } -} diff --git a/packages/runtime-host/src/server/scheduled-task-coordinator.ts b/packages/runtime-host/src/server/scheduled-task-coordinator.ts index fe2f05f44b..c91695f337 100644 --- a/packages/runtime-host/src/server/scheduled-task-coordinator.ts +++ b/packages/runtime-host/src/server/scheduled-task-coordinator.ts @@ -38,7 +38,6 @@ import { import type { ScheduledTaskOperationHandlerMap } from './operation-dispatcher.js'; import type { RuntimeHostResidency } from './host-kernel.js'; import type { HostedExecutionAuthority } from './hosted-execution-authority.js'; -import type { HostScheduledTaskChangeService } from './scheduled-task-change-service.js'; import type { SessionCreateInput } from '../protocol/session-catalog.js'; const MAX_TIMER_DELAY_MS = 2_147_483_647; @@ -47,6 +46,7 @@ const NATIVE_PROVIDER_RETRY_MS = 5_000; type ScheduledTaskSessions = Pick; type ScheduledTaskRuntime = Pick; type ScheduledTaskRoot = Pick; +type HostScheduledTaskChangeServiceLike = HostScheduledTaskCoordinatorInput['changes']; interface ScheduledTaskNativeEffects { hasWorkspaceService(serviceId: string, version: string): boolean; @@ -66,7 +66,9 @@ export interface HostScheduledTaskCoordinatorInput { readonly runtimePolicy: RuntimePolicyStoresWriter; readonly nativeEffects: ScheduledTaskNativeEffects; readonly createSession: (input: SessionCreateInput) => Promise; - readonly changes: HostScheduledTaskChangeService; + readonly changes: { + publish(revision: number, reason: ScheduledTaskChangedReason, taskId: string): void; + }; readonly acquireResidency: () => RuntimeHostResidency; readonly requestDrain: () => void; readonly now?: () => number; @@ -99,7 +101,7 @@ export class HostScheduledTaskCoordinator implements ScheduledTaskToolAuthority readonly #runtimePolicy: RuntimePolicyStoresWriter; readonly #nativeEffects: ScheduledTaskNativeEffects; readonly #createSession: HostScheduledTaskCoordinatorInput['createSession']; - readonly #changes: HostScheduledTaskChangeService; + readonly #changes: HostScheduledTaskChangeServiceLike; readonly #acquireResidency: () => RuntimeHostResidency; readonly #requestDrain: () => void; readonly #now: () => number; diff --git a/packages/runtime-host/src/server/session-catalog-change-service.ts b/packages/runtime-host/src/server/session-catalog-change-service.ts deleted file mode 100644 index 3ba571eb8a..0000000000 --- a/packages/runtime-host/src/server/session-catalog-change-service.ts +++ /dev/null @@ -1,40 +0,0 @@ -import type { SessionCatalogChangedFrame } from '../protocol/index.js'; - -export interface SessionCatalogChangeConnection { - close(): void; -} - -interface SessionCatalogChangeSink { - send(frame: SessionCatalogChangedFrame): Promise; -} - -export class HostSessionCatalogChangeService { - readonly #connections = new Map(); - #revision = 0; - - attachConnection( - connectionId: string, - sink: SessionCatalogChangeSink, - ): SessionCatalogChangeConnection { - this.#connections.set(connectionId, sink); - return { - close: () => { - if (this.#connections.get(connectionId) === sink) this.#connections.delete(connectionId); - }, - }; - } - - publish(sessionId: string): void { - this.#revision += 1; - const frame: SessionCatalogChangedFrame = { - kind: 'session.catalog.changed', - revision: this.#revision, - sessionId, - }; - for (const [connectionId, sink] of this.#connections) { - void sink.send(frame).catch(() => { - if (this.#connections.get(connectionId) === sink) this.#connections.delete(connectionId); - }); - } - } -}