refactor(api): require the session projection registry
This commit is contained in:
parent
8645053ca0
commit
d25ace0f22
6 changed files with 32 additions and 49 deletions
|
|
@ -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<Record<SessionId, SessionProjectionBaseline>> {
|
||||
const registry = this.ctx.get('sessionProjections')
|
||||
const blocks = Object.create(null) as Record<SessionId, SessionProjectionBaseline>
|
||||
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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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')
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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 })
|
||||
})
|
||||
})
|
||||
|
|
|
|||
|
|
@ -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<Context> {
|
|||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
await ctx.plugin(AgentRegistry)
|
||||
await ctx.plugin(SessionProjectionRegistry)
|
||||
return ctx
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue