From d25ace0f22bc602bc1a4955226e56552a79902f1 Mon Sep 17 00:00:00 2001 From: _Kerman Date: Fri, 28 Aug 2026 12:50:43 +0800 Subject: [PATCH] refactor(api): require the session projection registry --- .../api/session-controller/src/control.ts | 31 ++++++++----------- packages/api/session-controller/src/list.ts | 24 +++++++------- .../tests/control-jobs.host.spec.ts | 6 ++-- .../tests/control-queue.host.spec.ts | 2 ++ .../tests/session-projections.host.spec.ts | 16 ---------- .../tests/session-search.host.spec.ts | 2 ++ 6 files changed, 32 insertions(+), 49 deletions(-) diff --git a/packages/api/session-controller/src/control.ts b/packages/api/session-controller/src/control.ts index 624c9a9493..81b1092d9f 100644 --- a/packages/api/session-controller/src/control.ts +++ b/packages/api/session-controller/src/control.ts @@ -22,15 +22,13 @@ export class SessionControlController { /** @param ctx - Host context carrying live Agent, projection, and jobs services. */ constructor(private readonly ctx: Context) { ctx.on('session/event', (session, event) => { this.onSessionEvent(session, event) }) - ctx.inject(['sessionProjections'], (projectionCtx) => { - projectionCtx.sessionProjections.onChanged((session, key, value, seq) => { - this.broadcast({ - type: 'projection', - sessionId: session.id, - key, - value: value as JsonValue, - seq, - }) + ctx.sessionProjections.onChanged((session, key, value, seq) => { + this.broadcast({ + type: 'projection', + sessionId: session.id, + key, + value: value as JsonValue, + seq, }) }) ctx.inject(['jobs'], (jobsCtx) => { @@ -83,17 +81,14 @@ export class SessionControlController { private projectionBaseline( sessions: readonly Session[], ): Readonly> { - const registry = this.ctx.get('sessionProjections') const blocks = Object.create(null) as Record for (const session of sessions) { - const snapshot = registry?.snapshot(session) - blocks[session.id] = snapshot === undefined - ? { asOfSeq: session.seq - 1, values: {} } - : { - asOfSeq: snapshot.asOfSeq, - // Every projection definition validates its value before snapshot publication. - values: snapshot.values as SessionProjectionValues, - } + const snapshot = this.ctx.sessionProjections.snapshot(session) + blocks[session.id] = { + asOfSeq: snapshot.asOfSeq, + // Every projection definition validates its value before snapshot publication. + values: snapshot.values as SessionProjectionValues, + } } return blocks } diff --git a/packages/api/session-controller/src/list.ts b/packages/api/session-controller/src/list.ts index c9a53a6bd0..739a8baf31 100644 --- a/packages/api/session-controller/src/list.ts +++ b/packages/api/session-controller/src/list.ts @@ -87,25 +87,23 @@ export class ApiSessionList { private readonly ctx: Context, private readonly coldBlankProbeMaxBytes: number, ) { - ctx.inject(['sessionProjections'], (projectionCtx) => { - projectionCtx.sessionProjections.register<'sessionListMetadata', SessionListMetadata>({ - key: 'sessionListMetadata', - stateSchema: sessionListMetadataSchema, - init: () => ({ blank: true, lastPromptAt: null }), - apply: applySessionListMetadata, - wire: { viewSchema: sessionListMetadataSchema, view: state => state }, - stateVersion: 1, - }) + ctx.sessionProjections.register<'sessionListMetadata', SessionListMetadata>({ + key: 'sessionListMetadata', + stateSchema: sessionListMetadataSchema, + init: () => ({ blank: true, lastPromptAt: null }), + apply: applySessionListMetadata, + wire: { viewSchema: sessionListMetadataSchema, view: state => state }, + stateVersion: 1, }) - ctx.inject(['sessionProjections', 'attachments'], (projectionCtx) => { - projectionCtx.sessionProjections.register<'imageLimits', null>({ + ctx.inject(['attachments'], (attachmentCtx) => { + ctx.sessionProjections.register<'imageLimits', null>({ key: 'imageLimits', stateSchema: z.null(), init: () => null, apply: state => state, wire: { viewSchema: imageLimitsSchema, - view: () => projectionCtx.attachments.imageLimits, + view: () => attachmentCtx.attachments.imageLimits, }, stateVersion: 1, }) @@ -332,7 +330,7 @@ export class ApiSessionList { try { const block = session === undefined ? this.ctx.get('sessionProjectionCache')?.cachedSnapshot(header) - : this.ctx.get('sessionProjections')?.cachedSnapshot(session) + : this.ctx.sessionProjections.cachedSnapshot(session) return block !== undefined && Object.keys(block.values).length > 0 ? { asOfSeq: block.asOfSeq, diff --git a/packages/api/session-controller/tests/control-jobs.host.spec.ts b/packages/api/session-controller/tests/control-jobs.host.spec.ts index f83fa1b3c8..f827573d03 100644 --- a/packages/api/session-controller/tests/control-jobs.host.spec.ts +++ b/packages/api/session-controller/tests/control-jobs.host.spec.ts @@ -5,6 +5,7 @@ import type { JobOutcome } from '@deepseek-ai/dsh-jobs' import LocalJobRegistry from '@deepseek-ai/dsh-jobs-local' import SessionStore, { SessionId } from '@deepseek-ai/dsh-session' import type { Session } from '@deepseek-ai/dsh-session' +import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection' import { describe, expect, it } from 'vitest' import { SessionControlController } from '../src/control.ts' import type { SessionControlFrame } from '../src/types.ts' @@ -27,7 +28,7 @@ function producer(label = 'sleep 60') { return { spec, reads, settle: (outcome: JobOutcome) => { settle(outcome) } } } -async function harness(withRegistry: boolean): Promise<{ +async function harness(withJobs: boolean): Promise<{ ctx: Context session: Session agent: Agent @@ -36,7 +37,8 @@ async function harness(withRegistry: boolean): Promise<{ const ctx = new Context() await ctx.plugin(SessionStore) await ctx.plugin(AgentRegistry) - if (withRegistry) { + await ctx.plugin(SessionProjectionRegistry) + if (withJobs) { await ctx.plugin(LocalJobRegistry) ctx.jobs.attachController('session-controller-test') } diff --git a/packages/api/session-controller/tests/control-queue.host.spec.ts b/packages/api/session-controller/tests/control-queue.host.spec.ts index b4645dc2f3..0a1e5420a7 100644 --- a/packages/api/session-controller/tests/control-queue.host.spec.ts +++ b/packages/api/session-controller/tests/control-queue.host.spec.ts @@ -3,6 +3,7 @@ import AgentRegistry, { Inbox } from '@deepseek-ai/dsh-agent' import type { Agent } from '@deepseek-ai/dsh-agent' import { createUserMessage } from '@deepseek-ai/dsh-llm' import SessionStore, { SessionId } from '@deepseek-ai/dsh-session' +import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection' import { describe, expect, it } from 'vitest' import { SessionControlController } from '../src/control.ts' @@ -15,6 +16,7 @@ async function harness(): Promise<{ const ctx = new Context() await ctx.plugin(SessionStore) await ctx.plugin(AgentRegistry) + await ctx.plugin(SessionProjectionRegistry) const session = ctx.sessions.create(SessionId('queue-session')) const inbox = new Inbox(session, { inserted: () => {}, discarded: () => {}, claimed: () => {} }) const agent = { id: session.id, session, inbox, status: 'running', ctx } as Agent diff --git a/packages/api/session-controller/tests/session-projections.host.spec.ts b/packages/api/session-controller/tests/session-projections.host.spec.ts index 43391af21f..8a31e9208a 100644 --- a/packages/api/session-controller/tests/session-projections.host.spec.ts +++ b/packages/api/session-controller/tests/session-projections.host.spec.ts @@ -26,7 +26,6 @@ import SessionProjectionCache, { projectionCacheDomainSpec } from '@deepseek-ai/ import Storage from '@deepseek-ai/dsh-storage' import * as StorageDomain from '@deepseek-ai/dsh-storage-domain' import * as StorageJson from '@deepseek-ai/dsh-storage-json' -import { SessionControlController } from '@deepseek-ai/dsh-api-session-controller/src/control.ts' import type { SessionControlFrame, SessionFollowFrame } from '@deepseek-ai/dsh-api-session-controller/types' import { createSessionTestRemote, type TestSessionRemote } from './test-remote.ts' @@ -579,19 +578,4 @@ describe('Session control projection frames', () => { const tail = await opening(proxy, session.id) expect(tail.projections.asOfSeq).toBe(pushes.at(-1)?.seq) }) - - it('emits no projection frames when the composition has no registry', async () => { - const { ctx, session } = await harness(false) - const control = new SessionControlController(ctx) - const abort = new AbortController() - const iterator = control.control(abort.signal)[Symbol.asyncIterator]() - const baseline = await iterator.next() - const next = iterator.next() - seedMessages(session, 2) - await new Promise(resolve => setTimeout(resolve, 0)) - abort.abort() - if (baseline.done) throw new Error('Control stream ended before its baseline') - expect(baseline.value.type).toBe('baseline') - await expect(next).resolves.toEqual({ done: true, value: undefined }) - }) }) diff --git a/packages/api/session-controller/tests/session-search.host.spec.ts b/packages/api/session-controller/tests/session-search.host.spec.ts index 49bc900566..02d22b4f46 100644 --- a/packages/api/session-controller/tests/session-search.host.spec.ts +++ b/packages/api/session-controller/tests/session-search.host.spec.ts @@ -10,6 +10,7 @@ import AgentRegistry from '@deepseek-ai/dsh-agent' import { createUserMessage } from '@deepseek-ai/dsh-llm' import SessionStore from '@deepseek-ai/dsh-session' import type { SessionHeader, SessionId } from '@deepseek-ai/dsh-session' +import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection' import { SessionQueryEngine, SessionQueryError, @@ -56,6 +57,7 @@ async function baseContext(): Promise { const ctx = new Context() await ctx.plugin(SessionStore) await ctx.plugin(AgentRegistry) + await ctx.plugin(SessionProjectionRegistry) return ctx }