feat(session): add out-of-band log appends
This commit is contained in:
@@ -13,7 +13,7 @@ import { scopeOf, scopeTarget } from '@deepseek-ai/dsh-scope'
|
|||||||
import type { Scoped } from '@deepseek-ai/dsh-scope'
|
import type { Scoped } from '@deepseek-ai/dsh-scope'
|
||||||
import type { Message } from '@deepseek-ai/dsh-llm'
|
import type { Message } from '@deepseek-ai/dsh-llm'
|
||||||
import { SESSION_FORMAT_VERSION, SessionId } from './types.ts'
|
import { SESSION_FORMAT_VERSION, SessionId } from './types.ts'
|
||||||
import type { CreateSessionOptions, EpochHeader, SessionEvent, SessionEventMap, SessionEventType, SessionHeader, SurfaceIntent, SurfaceEventType } from './types.ts'
|
import type { CreateSessionOptions, EpochHeader, OutOfBandSessionEventType, SessionEvent, SessionEventMap, SessionEventType, SessionHeader, SurfaceIntent, SurfaceEventType, TurnTrigger } from './types.ts'
|
||||||
import { snapshotJsonValue } from './json.ts'
|
import { snapshotJsonValue } from './json.ts'
|
||||||
import { SurfaceManager } from './surface.ts'
|
import { SurfaceManager } from './surface.ts'
|
||||||
import type { SessionSurface } from './surface.ts'
|
import type { SessionSurface } from './surface.ts'
|
||||||
@@ -210,6 +210,7 @@ interface SessionEntry {
|
|||||||
announced: boolean
|
announced: boolean
|
||||||
announcing: boolean
|
announcing: boolean
|
||||||
appending: boolean
|
appending: boolean
|
||||||
|
outOfBand: boolean
|
||||||
detachRequested: boolean
|
detachRequested: boolean
|
||||||
detach(): void
|
detach(): void
|
||||||
}
|
}
|
||||||
@@ -384,7 +385,7 @@ export class Session {
|
|||||||
} finally {
|
} finally {
|
||||||
if (entry !== undefined) {
|
if (entry !== undefined) {
|
||||||
entry.appending = false
|
entry.appending = false
|
||||||
if (entry.detachRequested && !entry.announcing) entry.detach()
|
if (entry.detachRequested && !entry.announcing && !entry.outOfBand) entry.detach()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -668,6 +669,7 @@ export class SessionStore extends Service {
|
|||||||
announced: false,
|
announced: false,
|
||||||
announcing: false,
|
announcing: false,
|
||||||
appending: false,
|
appending: false,
|
||||||
|
outOfBand: false,
|
||||||
detachRequested: false,
|
detachRequested: false,
|
||||||
detach: () => { this.detachEntered(entry) },
|
detach: () => { this.detachEntered(entry) },
|
||||||
}
|
}
|
||||||
@@ -680,7 +682,7 @@ export class SessionStore extends Service {
|
|||||||
// A lifecycle listener may own the advanced detach capability. Keep the
|
// A lifecycle listener may own the advanced detach capability. Keep the
|
||||||
// entry and its publication hooks live until synchronous creation or append
|
// entry and its publication hooks live until synchronous creation or append
|
||||||
// publication unwinds, then publish the paired disposal edge.
|
// publication unwinds, then publish the paired disposal edge.
|
||||||
if (entry.announcing || entry.appending) {
|
if (entry.announcing || entry.appending || entry.outOfBand) {
|
||||||
entry.detachRequested = true
|
entry.detachRequested = true
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -734,7 +736,7 @@ export class SessionStore extends Service {
|
|||||||
}
|
}
|
||||||
} finally {
|
} finally {
|
||||||
entry.announcing = false
|
entry.announcing = false
|
||||||
if (entry.detachRequested && !entry.appending) entry.detach()
|
if (entry.detachRequested && !entry.appending && !entry.outOfBand) entry.detach()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -778,6 +780,87 @@ export class SessionStore extends Service {
|
|||||||
if (failure !== undefined) throw failure.reason
|
if (failure !== undefined) throw failure.reason
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Append one plugin-declared log-only event without borrowing the agent
|
||||||
|
* loop's lifecycle. An open turn receives the event directly and remains
|
||||||
|
* responsible for its ordinary checkpoint. A closed log receives one
|
||||||
|
* zero-step turn around the event, followed by an awaited flush.
|
||||||
|
*
|
||||||
|
* Once the synthetic `turn/start` commits, this method always attempts its
|
||||||
|
* matching `turn/end` and flush, including when the target append fails.
|
||||||
|
* Detachment requested by an event or flush listener is deferred until that
|
||||||
|
* sequence settles, so publication cannot switch from a live scoped session
|
||||||
|
* to an unobserved bare `Session` halfway through the update.
|
||||||
|
*
|
||||||
|
* @param session - exact live session that owns the target log.
|
||||||
|
* @param type - event type opted into {@link OutOfBandSessionEventMap} by its owner.
|
||||||
|
* @param data - typed JSON payload for the target event.
|
||||||
|
* @param trigger - plugin-owned turn trigger used only when the log is closed.
|
||||||
|
* @returns the accepted target event with its assigned sequence and timestamp.
|
||||||
|
* @throws when the session is detached, another out-of-band append is active,
|
||||||
|
* event acceptance fails, the synthetic turn cannot close, or flushing fails.
|
||||||
|
*/
|
||||||
|
async appendOutOfBand<T extends OutOfBandSessionEventType>(
|
||||||
|
session: Session,
|
||||||
|
type: T,
|
||||||
|
data: SessionEventMap[T],
|
||||||
|
trigger: TurnTrigger,
|
||||||
|
): Promise<SessionEvent<T>> {
|
||||||
|
const entry = this.liveEntryFor(session)
|
||||||
|
if (entry.outOfBand) {
|
||||||
|
throw new Error(`session "${session.id}" already has an out-of-band append in progress`)
|
||||||
|
}
|
||||||
|
entry.outOfBand = true
|
||||||
|
// `T` is excluded from SurfaceEventType by OutOfBandSessionEventType, but
|
||||||
|
// TypeScript does not reduce Session.append's conditional rest parameter
|
||||||
|
// through a generic intersection. Preserve that proven two-argument call
|
||||||
|
// shape without widening the public Session.append overload.
|
||||||
|
const appendLogOnly = session.append.bind(session) as unknown as <K extends OutOfBandSessionEventType>(
|
||||||
|
eventType: K,
|
||||||
|
eventData: SessionEventMap[K],
|
||||||
|
) => SessionEvent<K>
|
||||||
|
try {
|
||||||
|
const lastBoundary = session.events.findLast(event => event.type === 'turn/start' || event.type === 'turn/end')
|
||||||
|
if (lastBoundary?.type === 'turn/start') {
|
||||||
|
return appendLogOnly(type, data)
|
||||||
|
}
|
||||||
|
|
||||||
|
const lastStart = session.events.findLast(event => event.type === 'turn/start')
|
||||||
|
const turn = (lastStart?.data.turn ?? 0) + 1
|
||||||
|
let accepted: SessionEvent<T> | undefined
|
||||||
|
let failure: unknown
|
||||||
|
let opened = false
|
||||||
|
try {
|
||||||
|
session.append('turn/start', { turn, trigger })
|
||||||
|
opened = true
|
||||||
|
accepted = appendLogOnly(type, data)
|
||||||
|
} catch (error: unknown) {
|
||||||
|
failure = error
|
||||||
|
} finally {
|
||||||
|
if (opened) {
|
||||||
|
// The only target types admitted by OutOfBandSessionEventMap are
|
||||||
|
// log-only plugin events, so the synthetic turn remains open here.
|
||||||
|
session.append('turn/end', { turn, reason: { kind: 'completed' } })
|
||||||
|
try {
|
||||||
|
await this.flush(session)
|
||||||
|
} catch (error: unknown) {
|
||||||
|
if (failure === undefined) failure = error
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (failure !== undefined) {
|
||||||
|
// eslint-disable-next-line @typescript-eslint/only-throw-error -- preserve an arbitrary flush-listener rejection exactly
|
||||||
|
throw failure
|
||||||
|
}
|
||||||
|
/* v8 ignore next -- accepted is assigned unless an append failure was captured above. */
|
||||||
|
if (accepted === undefined) throw new Error('out-of-band append completed without an accepted event')
|
||||||
|
return accepted
|
||||||
|
} finally {
|
||||||
|
entry.outOfBand = false
|
||||||
|
if (entry.detachRequested && !entry.announcing && !entry.appending) entry.detach()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/** Return the exact live entry; detached/prepared objects reject. */
|
/** Return the exact live entry; detached/prepared objects reject. */
|
||||||
private liveEntryFor(session: Session): SessionEntry {
|
private liveEntryFor(session: Session): SessionEntry {
|
||||||
const entry = attachments.get(session)
|
const entry = attachments.get(session)
|
||||||
|
|||||||
@@ -263,9 +263,23 @@ export interface SessionEventMap {
|
|||||||
'request/header': { header: EpochHeader; reason: RequestHeaderReason }
|
'request/header': { header: EpochHeader; reason: RequestHeaderReason }
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Marker map for plugin-owned log-only events accepted by
|
||||||
|
* `SessionStore.appendOutOfBand()`. A plugin extends this map with the same key
|
||||||
|
* it adds to {@link SessionEventMap}; surface and lifecycle events stay
|
||||||
|
* ineligible unless their owner explicitly opts them into this narrow seam.
|
||||||
|
*/
|
||||||
|
export interface OutOfBandSessionEventMap {}
|
||||||
|
|
||||||
/** The appendable event-type keys of {@link SessionEventMap}, plugin-merged extensions included. */
|
/** The appendable event-type keys of {@link SessionEventMap}, plugin-merged extensions included. */
|
||||||
export type SessionEventType = keyof SessionEventMap
|
export type SessionEventType = keyof SessionEventMap
|
||||||
|
|
||||||
|
/** Plugin-declared non-surface event types accepted by `SessionStore.appendOutOfBand()`. */
|
||||||
|
export type OutOfBandSessionEventType = Exclude<
|
||||||
|
Extract<SessionEventType, keyof OutOfBandSessionEventMap>,
|
||||||
|
SurfaceEventType
|
||||||
|
>
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* The subset of {@link SessionEventType} values whose events produce LLM
|
* The subset of {@link SessionEventType} values whose events produce LLM
|
||||||
* messages and are eligible to appear on the ordered surface. Only these
|
* messages and are eligible to appear on the ordered surface. Only these
|
||||||
|
|||||||
@@ -0,0 +1,227 @@
|
|||||||
|
import { Context } from 'cordis'
|
||||||
|
import { describe, expect, it } from 'vitest'
|
||||||
|
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
|
||||||
|
|
||||||
|
declare module '@deepseek-ai/dsh-session' {
|
||||||
|
interface SessionEventMap {
|
||||||
|
'test/log-only': { value: string }
|
||||||
|
}
|
||||||
|
|
||||||
|
interface OutOfBandSessionEventMap {
|
||||||
|
'test/log-only': true
|
||||||
|
}
|
||||||
|
|
||||||
|
interface TurnTriggerMap {
|
||||||
|
'test/update': { kind: 'test/update' }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('SessionStore.appendOutOfBand', () => {
|
||||||
|
it('joins an open turn without adding a boundary or flushing it', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SessionStore)
|
||||||
|
const session = ctx.sessions.create(SessionId('open'))
|
||||||
|
let flushes = 0
|
||||||
|
ctx.on('session/flush', () => { flushes += 1 })
|
||||||
|
session.append('turn/start', {
|
||||||
|
turn: 1,
|
||||||
|
trigger: { kind: 'message', source: { kind: 'user' } },
|
||||||
|
})
|
||||||
|
|
||||||
|
const event = await ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 'inside' },
|
||||||
|
{ kind: 'test/update' },
|
||||||
|
)
|
||||||
|
|
||||||
|
expect(event).toMatchObject({ type: 'test/log-only', seq: 1, data: { value: 'inside' } })
|
||||||
|
expect(session.events.map(item => item.type)).toEqual(['turn/start', 'test/log-only'])
|
||||||
|
expect(flushes).toBe(0)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('wraps a closed log in one zero-step turn and flushes the balanced update', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SessionStore)
|
||||||
|
const session = ctx.sessions.create(SessionId('closed'))
|
||||||
|
const flushedTypes: string[][] = []
|
||||||
|
ctx.on('session/flush', (flushed) => {
|
||||||
|
flushedTypes.push(flushed.events.map(event => event.type))
|
||||||
|
})
|
||||||
|
|
||||||
|
const first = await ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 'first' },
|
||||||
|
{ kind: 'test/update' },
|
||||||
|
)
|
||||||
|
const second = await ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 'second' },
|
||||||
|
{ kind: 'test/update' },
|
||||||
|
)
|
||||||
|
|
||||||
|
expect(first.seq).toBe(1)
|
||||||
|
expect(second.seq).toBe(4)
|
||||||
|
expect(session.events).toMatchObject([
|
||||||
|
{ type: 'turn/start', seq: 0, data: { turn: 1, trigger: { kind: 'test/update' } } },
|
||||||
|
{ type: 'test/log-only', seq: 1, data: { value: 'first' } },
|
||||||
|
{ type: 'turn/end', seq: 2, data: { turn: 1, reason: { kind: 'completed' } } },
|
||||||
|
{ type: 'turn/start', seq: 3, data: { turn: 2, trigger: { kind: 'test/update' } } },
|
||||||
|
{ type: 'test/log-only', seq: 4, data: { value: 'second' } },
|
||||||
|
{ type: 'turn/end', seq: 5, data: { turn: 2, reason: { kind: 'completed' } } },
|
||||||
|
])
|
||||||
|
expect(flushedTypes).toEqual([
|
||||||
|
['turn/start', 'test/log-only', 'turn/end'],
|
||||||
|
['turn/start', 'test/log-only', 'turn/end', 'turn/start', 'test/log-only', 'turn/end'],
|
||||||
|
])
|
||||||
|
})
|
||||||
|
|
||||||
|
it('closes and flushes a zero-step turn when the target event is rejected', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SessionStore)
|
||||||
|
const session = ctx.sessions.create(SessionId('rejected'))
|
||||||
|
let flushes = 0
|
||||||
|
ctx.on('session/flush', () => { flushes += 1 })
|
||||||
|
|
||||||
|
await expect(ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 1n } as never,
|
||||||
|
{ kind: 'test/update' },
|
||||||
|
)).rejects.toThrow(/non-JSON-serializable/)
|
||||||
|
|
||||||
|
expect(session.events).toMatchObject([
|
||||||
|
{ type: 'turn/start', data: { turn: 1 } },
|
||||||
|
{ type: 'turn/end', data: { turn: 1, reason: { kind: 'completed' } } },
|
||||||
|
])
|
||||||
|
expect(flushes).toBe(1)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('does not flush when the synthetic turn cannot open', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SessionStore)
|
||||||
|
const session = ctx.sessions.create(SessionId('start-failure'))
|
||||||
|
let flushes = 0
|
||||||
|
ctx.on('session/flush', () => { flushes += 1 })
|
||||||
|
|
||||||
|
await expect(ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 'unreachable' },
|
||||||
|
{ kind: 'test/update', invalid: 1n } as never,
|
||||||
|
)).rejects.toThrow(/non-JSON-serializable/)
|
||||||
|
|
||||||
|
expect(session.events).toEqual([])
|
||||||
|
expect(flushes).toBe(0)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('preserves a target rejection when the balancing flush also rejects', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SessionStore)
|
||||||
|
const session = ctx.sessions.create(SessionId('target-and-flush-failure'))
|
||||||
|
ctx.on('session/flush', () => { throw new Error('disk failed') })
|
||||||
|
|
||||||
|
await expect(ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 1n } as never,
|
||||||
|
{ kind: 'test/update' },
|
||||||
|
)).rejects.toThrow(/non-JSON-serializable/)
|
||||||
|
|
||||||
|
expect(session.events.map(event => event.type)).toEqual([
|
||||||
|
'turn/start',
|
||||||
|
'turn/end',
|
||||||
|
])
|
||||||
|
})
|
||||||
|
|
||||||
|
it('keeps the session attached through publication and its flush', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SessionStore)
|
||||||
|
const session = ctx.sessions.prepare(SessionId('dispose'))
|
||||||
|
const detach = ctx.sessions.enter(session)
|
||||||
|
ctx.sessions.announce(session)
|
||||||
|
let liveDuringFlush = false
|
||||||
|
ctx.on('session/event', (_observed, event) => {
|
||||||
|
if (event.type === 'turn/start') detach()
|
||||||
|
})
|
||||||
|
ctx.on('session/flush', () => {
|
||||||
|
liveDuringFlush = ctx.sessions.get(session.id) === session
|
||||||
|
})
|
||||||
|
|
||||||
|
await ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 'last' },
|
||||||
|
{ kind: 'test/update' },
|
||||||
|
)
|
||||||
|
|
||||||
|
expect(session.events.map(event => event.type)).toEqual([
|
||||||
|
'turn/start',
|
||||||
|
'test/log-only',
|
||||||
|
'turn/end',
|
||||||
|
])
|
||||||
|
expect(liveDuringFlush).toBe(true)
|
||||||
|
expect(ctx.sessions.get(session.id)).toBeUndefined()
|
||||||
|
})
|
||||||
|
|
||||||
|
it('rejects detached sessions before opening a turn', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SessionStore)
|
||||||
|
const session = ctx.sessions.prepare(SessionId('detached'))
|
||||||
|
|
||||||
|
await expect(ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 'nope' },
|
||||||
|
{ kind: 'test/update' },
|
||||||
|
)).rejects.toThrow('session "detached" is not live in this store')
|
||||||
|
expect(session.events).toEqual([])
|
||||||
|
})
|
||||||
|
|
||||||
|
it('leaves a balanced log when the durability checkpoint rejects', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SessionStore)
|
||||||
|
const session = ctx.sessions.create(SessionId('flush-failure'))
|
||||||
|
ctx.on('session/flush', () => { throw new Error('disk failed') })
|
||||||
|
|
||||||
|
await expect(ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 'accepted' },
|
||||||
|
{ kind: 'test/update' },
|
||||||
|
)).rejects.toThrow('disk failed')
|
||||||
|
expect(session.events.map(event => event.type)).toEqual([
|
||||||
|
'turn/start',
|
||||||
|
'test/log-only',
|
||||||
|
'turn/end',
|
||||||
|
])
|
||||||
|
})
|
||||||
|
|
||||||
|
it('rejects overlapping updates while the first append is still settling', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SessionStore)
|
||||||
|
const session = ctx.sessions.create(SessionId('overlap'))
|
||||||
|
let release!: () => void
|
||||||
|
const checkpoint = new Promise<void>((resolve) => {
|
||||||
|
release = resolve
|
||||||
|
})
|
||||||
|
ctx.on('session/flush', () => checkpoint)
|
||||||
|
|
||||||
|
const first = ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 'first' },
|
||||||
|
{ kind: 'test/update' },
|
||||||
|
)
|
||||||
|
await expect(ctx.sessions.appendOutOfBand(
|
||||||
|
session,
|
||||||
|
'test/log-only',
|
||||||
|
{ value: 'overlap' },
|
||||||
|
{ kind: 'test/update' },
|
||||||
|
)).rejects.toThrow(/out-of-band append in progress/)
|
||||||
|
release()
|
||||||
|
await expect(first).resolves.toMatchObject({ data: { value: 'first' } })
|
||||||
|
})
|
||||||
|
})
|
||||||
Reference in New Issue
Block a user