Merge remote-tracking branch 'origin/master' into mergebot/pr1667

# Conflicts:
#	apps/web/tests/snapshots/queue-actions/preserved.expected.md
#	docs/cordis-catalog/services.md
#	docs/core-data-structures/session.i18n.yaml
#	packages/client/runtime/src/client/session-history/history-fold.ts
#	packages/client/ui-conversation/README.i18n.yaml
#	packages/core/session/src/index.ts
This commit is contained in:
imccyu
2026-08-05 22:56:57 +08:00
826 changed files with 18815 additions and 14777 deletions
+20 -44
View File
@@ -31,22 +31,22 @@ export { deriveEventMessage, foldSurface, isAppendSurfaceEvent, isReplacementSur
export { canonicalHeader, foldRequestHeader, headerEquals } from './request-header.ts'
/**
* Find the latest closed message-triggered turn, ignoring other triggers and
* between-turn events.
* Find the latest closed turn that entered at least one model step, ignoring
* balanced no-step turns produced by rejection, empty input, or cancellation.
* @param events - session events, or an owned suffix, to inspect.
* @returns the latest matching turn end, or `undefined`.
*/
export function findLastMessageTurnEnd(
events: readonly SessionEvent[],
): SessionEvent<'turn/end'> | undefined {
const messageTurns = new Set<number>()
const steppedTurns = new Set<number>()
let latest: SessionEvent<'turn/end'> | undefined
for (const event of events) {
if (event.type === 'turn/start') {
if (event.data.trigger.kind === 'message') messageTurns.add(event.data.turn)
if (event.type === 'step/start') {
steppedTurns.add(event.data.turn)
continue
}
if (event.type === 'turn/end' && messageTurns.delete(event.data.turn)) latest = event
if (event.type === 'turn/end' && steppedTurns.delete(event.data.turn)) latest = event
}
return latest
}
@@ -172,7 +172,6 @@ export function adoptSessionEvent<T extends SessionEvent>(event: T): T {
break
case 'assistant/message':
case 'tool/result':
case 'steering/message':
deepFreeze(event.data.message)
break
default:
@@ -208,7 +207,6 @@ function assertSessionEventEnvelope(value: Record<string, unknown>, index: numbe
throw new Error(`seed event at index ${index} has an invalid event envelope`)
}
assertCurrentLlmShape(event, index)
assertCurrentTurnEndShape(event, index)
}
/** Reject obsolete request headers and malformed messages at the seed/load boundary. */
@@ -234,7 +232,7 @@ function assertCurrentLlmShape(event: Record<string, unknown>, index: number): v
}
const type = event['type']
if (type !== 'user/message' && type !== 'assistant/message'
&& type !== 'tool/result' && type !== 'steering/message') return
&& type !== 'tool/result') return
assertMessageEventShape(event, `seed ${type} at index ${index}`)
}
@@ -262,7 +260,7 @@ function assertAdapterDefaults(
function assertMessageEventShape(event: Record<string, unknown>, subject: string): void {
const type = event['type']
if (type !== 'user/message' && type !== 'assistant/message'
&& type !== 'tool/result' && type !== 'steering/message') return
&& type !== 'tool/result') return
const data = event['data']
const record = typeof data === 'object' && data !== null
? data as Record<string, unknown>
@@ -312,22 +310,6 @@ function assertMessageEventShape(event: Record<string, unknown>, subject: string
}
}
/** Reject legacy aborted outcomes that persisted caller-owned reason detail. */
function assertCurrentTurnEndShape(event: Record<string, unknown>, index: number): void {
if (event['type'] !== 'turn/end') return
const data = event['data']
/* v8 ignore next -- this migration recognizes only the legacy object shape; format-wide payload validation is separate. */
if (typeof data !== 'object' || data === null) return
const reason = (data as Record<string, unknown>)['reason']
/* v8 ignore next -- non-object reasons cannot carry the legacy aborted detail this migration removes. */
if (typeof reason !== 'object' || reason === null || Array.isArray(reason)) return
const record = reason as Record<string, unknown>
if (record['kind'] === 'aborted'
&& (Object.keys(record).length !== 1 || !Object.hasOwn(record, 'kind'))) {
throw new Error(`seed turn/end at index ${index} uses unsupported reason-bearing aborted format`)
}
}
/** Whether an unknown value carries the current provider/model pair. */
function hasProviderModel(value: unknown): boolean {
if (typeof value !== 'object' || value === null) return false
@@ -433,9 +415,7 @@ export class Session {
* start here. Distinct from `header.seedLength`, the DURABLE fork-lineage
* boundary: a resumed session's constructor seed is its full stored log,
* while its header keeps the original fork value — this field is the
* in-process construction fact. An explicitly supplied empty seed has the
* same value as no seed (0); its `session/end-seed` event preserves the
* lifecycle distinction.
* in-process construction fact.
*
* Not persisted itself: a seeded session projects it into the log as the
* `session/end-seed` event, which is what a consumer reading STORED history
@@ -637,23 +617,18 @@ export class Session {
return this.headerFold
}
/** Cached fold of the request-context events — see {@link requestContext}. */
/** Cached fold of `request/context` events. */
private contextFold: RequestContext | undefined
/** Log position (events consumed) the context fold has reached. */
private contextFoldSeq = 0
/**
* The route metadata in force after the log's last `request/context` event —
* what the NEXT request deduplicates against — or undefined before any such
* record. Maintained incrementally like {@link requestHeader}, so a per-step
* read costs O(new events).
* @returns the folded context record, or undefined when none exists yet.
* Return the latest resolved route metadata, or `undefined` before the first
* `request/context` event. Each event is folded once.
* @returns the latest immutable route metadata.
*/
requestContext(): RequestContext | undefined {
if (this.contextFoldSeq < this.log.length) {
for (const event of this.log.slice(this.contextFoldSeq)) {
// Frozen for the same reason as the header fold: it is session state
// exposed by reference and every later dedup compares against it.
if (event.type === 'request/context') this.contextFold = deepFreeze({ ...event.data })
}
this.contextFoldSeq = this.log.length
@@ -769,7 +744,7 @@ export class SessionStore extends Service {
* {@link SessionHeader} (the store fills `version`/`id`/`createdAt`).
*
* For an agent whose session must be torn down IN ORDER with its loop (so the
* loop's final flush is captured before the store attachment ends), do NOT use this
* loop's final events are published before the store attachment ends), do NOT use this
* — fold the session lifecycle into the agent's own effect via
* {@link prepare} + {@link enter} + {@link announce} (see
* `dsh-agent-loop`'s creation transaction).
@@ -801,7 +776,7 @@ export class SessionStore extends Service {
* `ctx.effect` (the agent factory) folds the session lifecycle into that ONE
* effect so a fiber unload tears the session + agent down as a single ORDERED
* chain rather than as racing sibling effects — which would remove the publication hooks
* before the loop's closing `session/flush`, dropping the closing events.
* before the driver's closing events commit, dropping them.
*
* @param id - the session id; omitted, the store mints `session-<n>`.
* @param options - seed events and/or creation metadata for the header.
@@ -955,10 +930,11 @@ export class SessionStore extends Service {
/**
* Dispatch the awaited `session/flush` durability checkpoint for `session`,
* with the carrier captured at {@link enter}. THE flush entry point: the
* store owns the carrier, so callers (the loop's turn-end checkpoint, idle
* injection, teardown drains) must come through here rather than dispatch a
* raw `ctx.parallel('session/flush', …)` — one owner, one spelling, and the
* scoped-dispatch invariant can pin it.
* store owns the carrier, so callers (the checkpoint policy's per-request
* barrier, goal-session's idle checkpoint, teardown drains, and consumers
* that flush themselves before reading storage) must come through here
* rather than dispatch a raw `ctx.parallel('session/flush', …)` — one owner,
* one spelling, and the scoped-dispatch invariant can pin it.
* @param session - the session whose buffered events must reach durable storage.
* @returns whether at least one durability listener participated, after every
* listener has settled successfully.
+2 -3
View File
@@ -147,10 +147,9 @@ function validateEvent(
case 'session/end-seed':
// Unconstrained: an unbalanced seed legally puts it inside an open turn.
break
case 'steering/message':
case 'todo/write':
case 'request/context':
case 'request/header': {
case 'request/header':
case 'request/context': {
if (trace.openTurn === null) {
fail(`${event.type} appended outside any open turn (core execution events must be turn-enclosed)`)
}
+42 -27
View File
@@ -16,13 +16,12 @@ const SURFACE_EVENT_TYPES = new Set<string>([
'user/message',
'assistant/message',
'tool/result',
'steering/message',
])
/**
* Whether an event type can join the model-visible surface.
* @param type - event type to test.
* @returns true for one of the four message-producing event types.
* @returns true for one of the three message-producing event types.
*/
export function isSurfaceEligibleType(type: string): boolean {
return SURFACE_EVENT_TYPES.has(type)
@@ -86,21 +85,17 @@ export function deriveEventMessage(event: SessionEvent): Message | null {
// history; turn/step boundaries, chunks, usage, and errors are trace/replay
// data.
switch (event.type) {
// Ordinary prompts, injected context, and mid-turn steering project
// identically in user role: the event's model-facing content stays
// verbatim. Steering's `turn` is log-only. Do NOT re-add per-type framing
// (e.g. `<context>`/`<steering>`) here: framing is caller-owned — a
// producer bakes it into `content`, as workspace-context does with
// `<system-reminder>` — or, if reintroduced, must be driven by the event
// `meta` map and a dedicated renderer, keeping this projection a verbatim
// pass-through. See the deferred design note in
// Ordinary prompts and injected context project in user role: the event's
// model-facing content stays verbatim. Do NOT re-add per-type framing
// (e.g. `<context>`) here: framing is caller-owned — a producer bakes it
// into `content`, as workspace-context does with `<system-reminder>` — or,
// if reintroduced, must be driven by the event `meta` map and a dedicated
// renderer, keeping this projection a verbatim pass-through. See the
// deferred design note in
// ../../../../.agents/notes/implemented/simplification/2026-07-20-unwrap-injected-content-envelopes.md
case 'user/message': {
return event.data
}
case 'steering/message': {
return event.data.message
}
case 'assistant/message': {
// Skip an empty-content assistant/message: it exists only to host a
// max-tokens step's usage and must not inject a content-less assistant
@@ -293,13 +288,14 @@ function assertToolResultRewrite(
event: SessionEvent,
shadowedSeqs: readonly number[],
events: readonly SessionEvent[],
baseSeq: number,
): void {
if (event.type !== 'tool/result') return
if (shadowedSeqs.length !== 1) {
throw new Error('tool/result surface replacement must rewrite exactly one current node')
}
for (const originalSeq of shadowedSeqs) {
const original = events[originalSeq]
const original = events[originalSeq - baseSeq]
if (original?.type !== 'tool/result') {
throw new Error('tool/result surface replacement must target a current tool/result')
}
@@ -327,6 +323,7 @@ function planSurfaceEvent(
event: SessionEvent,
expectedSeq: number,
events: readonly SessionEvent[],
baseSeq: number,
): SurfacePlan | undefined {
if (event.seq !== expectedSeq) {
throw new Error(`session event seq ${event.seq} is not contiguous; expected ${expectedSeq}`)
@@ -339,7 +336,7 @@ function planSurfaceEvent(
}
const range = replacementRange(state, surfaceOp)
assertProvenance(event, range.shadowedSeqs)
assertToolResultRewrite(event, range.shadowedSeqs, events)
assertToolResultRewrite(event, range.shadowedSeqs, events, baseSeq)
return {
kind: 'replace',
seq: event.seq,
@@ -355,8 +352,9 @@ function applySurfaceEvent(
event: SessionEvent,
expectedSeq: number,
events: readonly SessionEvent[],
baseSeq: number,
): SurfaceFoldReplacement | undefined {
const plan = planSurfaceEvent(state, event, expectedSeq, events)
const plan = planSurfaceEvent(state, event, expectedSeq, events, baseSeq)
if (plan?.kind === 'append') {
state.nodes.push(plan.seq)
} else if (plan?.kind === 'replace') {
@@ -382,7 +380,7 @@ export function foldSurface(events: readonly SessionEvent[]): SurfaceFoldResult
const state = createFoldState()
const replacements: SurfaceFoldReplacement[] = []
for (const [index, event] of events.entries()) {
const replacement = applySurfaceEvent(state, event, index, events)
const replacement = applySurfaceEvent(state, event, index, events, 0)
if (replacement !== undefined) replacements.push(replacement)
}
return { nodes: [...state.nodes], replacements }
@@ -392,38 +390,55 @@ export function foldSurface(events: readonly SessionEvent[]): SurfaceFoldResult
export class SurfaceManager implements SessionSurface {
/** Shared transition state; replacement history is not retained. */
private _state = createFoldState()
/** Last processed seq; -1 folds a seeded log on first access. */
private _lastProcessedSeq = -1
/** Last processed absolute seq. */
private _lastProcessedSeq: number
constructor(private log: readonly SessionEvent[]) {}
/**
* @param log - Contiguous complete log or loaded event window.
* @param baseSeq - Absolute sequence of the window's first event.
*/
constructor(
private log: readonly SessionEvent[],
private readonly baseSeq = 0,
) {
this._lastProcessedSeq = baseSeq - 1
}
/**
* Validate the next candidate without mutating the committed surface.
* @param event - candidate event that has not entered the log yet.
*/
validateNext(event: SessionEvent): void {
if (this._lastProcessedSeq < this.log.length - 1) this._processDelta()
planSurfaceEvent(this._state, event, this.log.length, this.log)
if (this._lastProcessedSeq < this.baseSeq + this.log.length - 1) this._processDelta()
planSurfaceEvent(
this._state,
event,
this.baseSeq + this.log.length,
this.log,
this.baseSeq,
)
}
/** Monotonic count of folded positional replacements. */
get replaceGeneration(): number {
if (this._lastProcessedSeq < this.log.length - 1) this._processDelta()
if (this._lastProcessedSeq < this.baseSeq + this.log.length - 1) this._processDelta()
return this._state.replaceGeneration
}
/** Surface event sequences in model-visible order. */
get nodes(): readonly number[] {
if (this._lastProcessedSeq < this.log.length - 1) this._processDelta()
if (this._lastProcessedSeq < this.baseSeq + this.log.length - 1) this._processDelta()
return this._state.nodes
}
/** Fold events appended since the previous access. */
private _processDelta(): void {
for (let i = this._lastProcessedSeq + 1; i < this.log.length; i++) {
const tailSeq = this.baseSeq + this.log.length - 1
for (let seq = this._lastProcessedSeq + 1; seq <= tailSeq; seq++) {
const index = seq - this.baseSeq
// oxlint-disable-next-line typescript/no-non-null-assertion -- bounded by the loop condition
applySurfaceEvent(this._state, this.log[i]!, i, this.log)
this._lastProcessedSeq = i
applySurfaceEvent(this._state, this.log[index]!, seq, this.log, this.baseSeq)
this._lastProcessedSeq = seq
}
}
}
+36 -61
View File
@@ -5,7 +5,6 @@ import type {
LlmCallConfig,
LlmCallConfigAdapterDefaults,
LlmFailure,
MessageSource,
StreamChunk,
TokenUsage,
ToolResultMessage,
@@ -94,24 +93,15 @@ export interface CreateSessionOptions {
}
}
/**
* What started a turn.
* Merge-extensible sum type (same pattern as MessageSourceMap).
*/
export interface TurnTriggerMap {
message: { kind: 'message'; source: MessageSource }
/** Recovery turn reopened over the repaired current session log. */
retry: { kind: 'retry' }
/**
* An out-of-band producer explicitly enclosed injected context in a one-shot
* turn. `Agent.inject()` appends idle context directly and does not use this
* trigger; the source mirrors the producer of the enclosed `user/message`.
*/
injection: { kind: 'injection'; source: MessageSource }
}
/** Why an active agent driver was cancelled. */
export type AgentCancelCause =
| { readonly kind: 'user' }
| { readonly kind: 'parent' }
| { readonly kind: 'hook'; readonly reason: string }
| { readonly kind: 'disposed' }
/** The union over {@link TurnTriggerMap} — what started a turn; plugins extend it by merging variants into the map. */
export type TurnTrigger = TurnTriggerMap[keyof TurnTriggerMap]
/** Durable cancellation cause, including imports whose original coarse record carried no cause. */
export type TurnEndCancelCause = AgentCancelCause | { readonly kind: 'legacy' }
/**
* Why a turn ended. Merge-extensible sum type.
@@ -119,20 +109,15 @@ export type TurnTrigger = TurnTriggerMap[keyof TurnTriggerMap]
export interface TurnEndReasonMap {
completed: { kind: 'completed' }
/** A cancellation request interrupted the live turn. */
aborted: { kind: 'aborted' }
aborted: { kind: 'aborted'; reason: TurnEndCancelCause }
blocked: { kind: 'blocked' }
/**
* The turn failed: a step threw or the model reported a failure. `step` is the
* step number the failure occurred on (the operational error's location — the
* single durable record of an in-turn failure; live diagnostics also fire via
* `agent/error`). Final model-request failures retain their normalized facts
* as one `failure`; other thrown values retain their rendered message and a
* real `HarnessError` code when present.
* The turn failed. `error` is always a structured failure: the `LlmError`
* facts verbatim, or `{ message: errorChain(error), code: 'UNKNOWN' }`
* flattened from any other error.
*/
error: { kind: 'error'; step: number } & (
| { failure: LlmFailure; message?: never; code?: never }
| { message: string; code?: string; failure?: never }
)
disposed: { kind: 'disposed' }
error: { kind: 'error'; error: LlmFailure }
/** At least one step reached its output-token ceiling, even if a plugin continued the turn. */
'max-tokens': { kind: 'max-tokens' }
/**
@@ -178,17 +163,13 @@ export interface EpochHeader {
tools?: ToolSchema[]
}
/**
* Registration-bound context metadata of one resolved model route. Adapter
* metadata about a route rather than a request input, which is why it lives
* outside {@link EpochHeader}.
*/
/** Registration-bound metadata for one resolved model route. */
export interface RequestContext {
/** Registered provider route the metadata was resolved through. */
/** Registered provider route the metadata belongs to. */
provider: string
/** Provider-owned model id the metadata belongs to. */
model: string
/** Maximum combined request and response context in tokens; absent when the adapter advertises none. */
/** Maximum combined request and response context in tokens, when advertised. */
contextWindow?: number
}
@@ -208,14 +189,19 @@ export type RequestHeaderReason = 'initial' | 'resume' | 'change'
*/
export interface SessionEventMap {
/**
* Opens turn `turn`. `trigger` records what started the model loop.
* Opens turn `turn` before the loop claims queued input or runs pre-step.
* Rejection, empty input, cancellation, or failure may close it with no
* step; otherwise the following identified `user/message` event or batch
* records the messages entering the step.
*/
'turn/start': { turn: number; trigger: TurnTrigger }
'turn/start': { turn: number }
/**
* Closes turn `turn` with the {@link TurnEndReason} that ended it. The loop
* awaits `session/flush` after an ordinary turn ends before claiming the next
* queued item. Success commits the turn; rejection is reported live and does
* not prevent later work.
* Closes turn `turn` with the {@link TurnEndReason} that ended it. A turn
* with no entered step has no `step/start` or `step/end`. The loop does not await a
* flush at turn boundaries: `dsh-session-checkpoint-policy` owns the
* per-request durability checkpoint, and consumers that read storage after
* `whenIdle()` flush themselves. Success commits the turn; rejection is
* reported live and does not prevent later work.
*/
'turn/end': { turn: number; reason: TurnEndReason }
/** Opens step `step` of turn `turn` — one model call plus the tool executions it requested. */
@@ -226,9 +212,8 @@ export interface SessionEventMap {
* A user-role message on the model-visible surface: a direct human prompt
* (the queued message claimed for this turn), a synthetic `agent.inject()`
* context (file-change notices, subdir AGENTS.md, skill content, cron
* notifications, …), or an admitted goal continuation round. All three
* project their `content` verbatim; `source` tells them apart. An idle
* injection may append this event between turns without running the model.
* notifications, …), or an entered goal continuation round. All three
* project their `content` verbatim; `source` tells them apart.
*/
'user/message': UserMessage
/** Raw stream chunk — token-level replay fidelity. */
@@ -264,8 +249,6 @@ export interface SessionEventMap {
error?: { name: string; code: string }
meta?: JsonValue
}
/** Steering content injected between steps of a running turn. */
'steering/message': { turn: number; message: UserMessage }
/** Whole-list snapshot; latest write wins on replay. Log-only UI state; never derived history. */
'todo/write': { todos: TodoItem[] }
/**
@@ -274,21 +257,14 @@ export interface SessionEventMap {
*/
'request/header': { header: EpochHeader; reason: RequestHeaderReason }
/**
* Registration-bound context metadata for the route a request resolved to,
* appended inside its step beside `request/header` and only when the route
* or capacity differs from the last record. It is log-only and deliberately
* NOT part of {@link EpochHeader}: capacity is adapter metadata about a
* route, not an input the request was built from, so it must not participate
* in request reconstruction or header equality. `contextWindow` is absent
* when the route's adapter advertises no capacity.
* Route metadata for the next request, logged only when the route or capacity
* changes. It does not participate in request reconstruction or header equality.
*/
'request/context': RequestContext
/**
* Marks the end of a constructor seed. Events before it have smaller seq
* values and came from the seed (resume, fork, or replay); this lifecycle
* produced none of them. An explicitly supplied empty seed puts the marker
* at seq 0, distinguishing an empty resumed session from a fresh session.
* This log-only event is the durable projection of
* produced none of them. This log-only event is the durable projection of
* {@link Session.firstLiveSeq}. Its payload is empty — position and `time`
* carry the meaning.
*
@@ -322,7 +298,6 @@ export type SurfaceEventType =
| 'user/message'
| 'assistant/message'
| 'tool/result'
| 'steering/message'
/**
* A {@link SessionEvent} that is **on** the ordered surface — its
@@ -339,7 +314,7 @@ export type SurfaceEvent = SessionEvent<SurfaceEventType> & { surfaceOp: Surface
* How a session event entered the ordered surface. Only valid on
* {@link SurfaceEventType} events.
*
* - `'append'`: added to the tail — normal path for user/assistant/tool/steering
* - `'append'`: added to the tail — normal path for user/assistant/tool
* messages.
* - `{ op: 'replace', start, end }`: replaces surface nodes from `start`
* (inclusive) through `end` (inclusive) with this node. Both must exist as
@@ -374,7 +349,7 @@ export interface SurfaceIntent {
*
* The {@link sourceEventSeqs} and {@link surfaceOp} fields are conditional:
* they only exist on {@link SurfaceEventType} variants (`user/message`,
* `assistant/message`, `tool/result`, `steering/message`).
* `assistant/message`, `tool/result`).
* Non-surface events (boundary markers, chunks, usage, errors) never carry
* surface metadata — the compiler enforces this at `Session.append()`
* call sites.