From 9a6914d845039ca4f8190f47fc03b159a424ba72 Mon Sep 17 00:00:00 2001 From: Tianyi Cui <53024+tianyicui@users.noreply.github.com> Date: Tue, 21 Jul 2026 01:53:08 +0800 Subject: [PATCH] feat(session): add out-of-band log appends --- packages/core/session/src/index.ts | 91 ++++++- packages/core/session/src/types.ts | 14 ++ .../core/session/tests/out-of-band.spec.ts | 227 ++++++++++++++++++ 3 files changed, 328 insertions(+), 4 deletions(-) create mode 100644 packages/core/session/tests/out-of-band.spec.ts diff --git a/packages/core/session/src/index.ts b/packages/core/session/src/index.ts index 9b2eb74a37..6350694ebe 100644 --- a/packages/core/session/src/index.ts +++ b/packages/core/session/src/index.ts @@ -13,7 +13,7 @@ import { scopeOf, scopeTarget } from '@deepseek-ai/dsh-scope' import type { Scoped } from '@deepseek-ai/dsh-scope' import type { Message } from '@deepseek-ai/dsh-llm' 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 { SurfaceManager } from './surface.ts' import type { SessionSurface } from './surface.ts' @@ -210,6 +210,7 @@ interface SessionEntry { announced: boolean announcing: boolean appending: boolean + outOfBand: boolean detachRequested: boolean detach(): void } @@ -384,7 +385,7 @@ export class Session { } finally { if (entry !== undefined) { 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, announcing: false, appending: false, + outOfBand: false, detachRequested: false, detach: () => { this.detachEntered(entry) }, } @@ -680,7 +682,7 @@ export class SessionStore extends Service { // A lifecycle listener may own the advanced detach capability. Keep the // entry and its publication hooks live until synchronous creation or append // publication unwinds, then publish the paired disposal edge. - if (entry.announcing || entry.appending) { + if (entry.announcing || entry.appending || entry.outOfBand) { entry.detachRequested = true return } @@ -734,7 +736,7 @@ export class SessionStore extends Service { } } finally { 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 } + /** + * 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( + session: Session, + type: T, + data: SessionEventMap[T], + trigger: TurnTrigger, + ): Promise> { + 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 ( + eventType: K, + eventData: SessionEventMap[K], + ) => SessionEvent + 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 | 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. */ private liveEntryFor(session: Session): SessionEntry { const entry = attachments.get(session) diff --git a/packages/core/session/src/types.ts b/packages/core/session/src/types.ts index ea0ad55a51..8f24b7ac08 100644 --- a/packages/core/session/src/types.ts +++ b/packages/core/session/src/types.ts @@ -263,9 +263,23 @@ export interface SessionEventMap { '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. */ export type SessionEventType = keyof SessionEventMap +/** Plugin-declared non-surface event types accepted by `SessionStore.appendOutOfBand()`. */ +export type OutOfBandSessionEventType = Exclude< + Extract, + SurfaceEventType +> + /** * The subset of {@link SessionEventType} values whose events produce LLM * messages and are eligible to appear on the ordered surface. Only these diff --git a/packages/core/session/tests/out-of-band.spec.ts b/packages/core/session/tests/out-of-band.spec.ts new file mode 100644 index 0000000000..2409e92f7d --- /dev/null +++ b/packages/core/session/tests/out-of-band.spec.ts @@ -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((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' } }) + }) +})