fix(session): preserve plugin turn invariants
This commit is contained in:
@@ -2,5 +2,5 @@
|
||||
# side as of the last confirmed-consistent state. Both languages carry equal authority;
|
||||
# after editing either side, bring the other along and re-record with:
|
||||
# pnpm run verify-translation-pairing --write .agents/notes/implemented/simplification/2026-07-28-remove-synthetic-log-only-turns.md
|
||||
2026-07-28-remove-synthetic-log-only-turns.md: 41ada81ef040cb7826777cd53911c3cb8bbee41e
|
||||
2026-07-28-remove-synthetic-log-only-turns.zh.md: 82d24cefc01b94ba584dfa7b1a3549eea36eedc9
|
||||
2026-07-28-remove-synthetic-log-only-turns.md: af4da00f4fe1d7aebff845cd55053bb5b807c979
|
||||
2026-07-28-remove-synthetic-log-only-turns.zh.md: bc9fc0e1637be79d75e18b678d1a7ca8b0543085
|
||||
|
||||
@@ -36,7 +36,7 @@ The historical [universal turn-enclosure decision](../../archived/architecture/2
|
||||
|
||||
## Verification
|
||||
|
||||
Core invariant tests accept an unknown plugin event between turns while continuing to reject built-in execution events there. Session-title service tests pin one direct fallback event under concurrent refresh, detached-session rejection, and newest-revision acceptance. JSONL and SQLite round trips preserve a title appended after `turn/end` through the persistence lifecycle drain, and fork tests retain a standalone log-only tail while rejecting boundaries inside an open turn. Generated API and type-equivalence catalogs contain no removed symbol.
|
||||
Core invariant tests accept an unknown plugin event between turns while continuing to reject built-in execution events there. Hook, compaction, plan-mode, Code Mode dispatch, and approval invariant companions replay existing logs and reject the same execution-scoped events before commit when no turn is open. Session-title service tests pin one direct fallback event under concurrent refresh, detached-session rejection, and newest-revision acceptance. JSONL and SQLite round trips preserve a title appended after `turn/end` through the persistence lifecycle drain, and fork tests retain a standalone log-only tail while rejecting boundaries inside an open turn. A keyless assembled ACP snapshot delays the model-backed title until after `turn/end` and pins one standalone provider title with no synthetic turn. Generated API and type-equivalence catalogs contain no removed symbol.
|
||||
|
||||
## Consequences
|
||||
|
||||
|
||||
@@ -36,7 +36,7 @@ Status: implemented
|
||||
|
||||
## 验证
|
||||
|
||||
核心不变量测试会接受轮次之间的未知插件事件,同时继续拒绝位于该处的内置执行事件。会话标题服务测试会在并发刷新、会话脱离拒绝和最新修订接受场景下,固定一个直接追加的回退事件。JSONL 和 SQLite 往返测试会通过持久化生命周期排空保留追加在 `turn/end` 之后的标题;fork 测试会保留独立纯日志尾部,同时拒绝位于开放轮次内的边界。生成的 API 和类型等价性目录不含任何已移除符号。
|
||||
核心不变量测试会接受轮次之间的未知插件事件,同时继续拒绝位于该处的内置执行事件。钩子、压缩(compaction)、plan-mode、Code Mode 分发和审批的不变量配套组件会回放既有日志,并在没有开放轮次时,于提交前拒绝相同的执行作用域事件。会话标题服务测试会在并发刷新、会话脱离拒绝和最新修订接受场景下,固定一个直接追加的回退事件。JSONL 和 SQLite 往返测试会通过持久化生命周期排空保留追加在 `turn/end` 之后的标题;fork 测试会保留独立纯日志尾部,同时拒绝位于开放轮次内的边界。一个无密钥、经完整组装的 ACP(Agent Client Protocol)快照会将模型生成的标题延迟到 `turn/end` 之后,并固定一个不含合成轮次的独立提供方标题。生成的 API 和类型等价性目录不含任何已移除符号。
|
||||
|
||||
## 后果
|
||||
|
||||
|
||||
53
examples/acp-agent/session-title.cordis.snapshot.yml
Normal file
53
examples/acp-agent/session-title.cordis.snapshot.yml
Normal file
@@ -0,0 +1,53 @@
|
||||
# Keyless session-title composition. Main-agent chunks derive from session.jsonl;
|
||||
# the auxiliary route consumes replay.override.json with pacing so its accepted
|
||||
# title commits only after the main turn has closed.
|
||||
- id: base
|
||||
name: '@cordisjs/plugin-include'
|
||||
config:
|
||||
path: ./cordis.yml
|
||||
patches:
|
||||
- id: llm-deepseek
|
||||
name: '@deepseek-ai/dsh-llm-deepseek'
|
||||
disabled: true
|
||||
- id: acp-agent
|
||||
name: '@deepseek-ai/dsh-acp-demo'
|
||||
config:
|
||||
provider: deepseek
|
||||
model: deepseek-v4-flash
|
||||
persistenceRoot: !!js process.env.DSH_SNAPSHOT_SESSIONS_ROOT ?? './.sessions'
|
||||
persistenceCompression: none
|
||||
workspaceContext:
|
||||
maxBytes: 65536
|
||||
persona: |
|
||||
You are a coding assistant powered by the {{model}} model. Your working directory is {{cwd}}. Your bash tool runs under a file sandbox — a `[sandbox: file access denied …]` result is policy, not a command bug.
|
||||
|
||||
Verify your work by running the code or tests. Keep answers brief and factual.
|
||||
- insert:
|
||||
- id: llm-replay-main
|
||||
name: '@deepseek-ai/dsh-llm-replay'
|
||||
config:
|
||||
overrideFile: ./.missing-main-replay-override.json
|
||||
providers:
|
||||
- id: deepseek
|
||||
name: DeepSeek
|
||||
models:
|
||||
- id: deepseek-v4-flash
|
||||
- id: llm-replay-title
|
||||
name: '@deepseek-ai/dsh-llm-replay'
|
||||
config:
|
||||
paceMs: 10
|
||||
providers:
|
||||
- id: title-replay
|
||||
name: Title replay
|
||||
models:
|
||||
- id: title-model
|
||||
- id: session-title-provider
|
||||
name: '@deepseek-ai/dsh-session-title-first-message-llm'
|
||||
config:
|
||||
targetWords: 5
|
||||
targetCjkCharacters: 10
|
||||
maxInputBytes: 4096
|
||||
maxOutputTokens: 32
|
||||
timeoutMs: 5000
|
||||
provider: title-replay
|
||||
model: title-model
|
||||
19
examples/acp-agent/session-title.cordis.yml
Normal file
19
examples/acp-agent/session-title.cordis.yml
Normal file
@@ -0,0 +1,19 @@
|
||||
# Session-title snapshot composition: the optional first-message provider uses
|
||||
# the ordinary DeepSeek route while the ACP app and every other capability stay
|
||||
# identical to the base example.
|
||||
- id: base
|
||||
name: '@cordisjs/plugin-include'
|
||||
config:
|
||||
path: ./cordis.yml
|
||||
patches:
|
||||
- insert:
|
||||
- id: session-title-provider
|
||||
name: '@deepseek-ai/dsh-session-title-first-message-llm'
|
||||
config:
|
||||
targetWords: 5
|
||||
targetCjkCharacters: 10
|
||||
maxInputBytes: 4096
|
||||
maxOutputTokens: 32
|
||||
timeoutMs: 5000
|
||||
provider: deepseek
|
||||
model: deepseek-v4-flash
|
||||
@@ -39,6 +39,7 @@ const PTY_CONFIG = fileURLToPath(new URL('../pty.cordis.yml', import.meta.url))
|
||||
const DEPTH_TWO_CONFIG = fileURLToPath(new URL('../depth-two.cordis.yml', import.meta.url))
|
||||
const SESSION_SANDBOX_ROOT_CONFIG = fileURLToPath(new URL('../session-sandbox-root.cordis.yml', import.meta.url))
|
||||
const RETRY_CONFIG = fileURLToPath(new URL('../retry.cordis.yml', import.meta.url))
|
||||
const SESSION_TITLE_CONFIG = fileURLToPath(new URL('../session-title.cordis.yml', import.meta.url))
|
||||
const LSP_CONFIG = fileURLToPath(new URL('./lsp.cordis.yml', import.meta.url))
|
||||
const WEB_CONFIG = fileURLToPath(new URL('../web.cordis.yml', import.meta.url))
|
||||
const SNAPSHOTS_DIR = join(dirname(fileURLToPath(import.meta.url)), 'snapshots')
|
||||
@@ -75,6 +76,13 @@ const SCENARIOS: Scenario[] = [
|
||||
// text-turn is the pinned-header scenario: the minimal single text turn.
|
||||
// Its prompt and tool-schema sidecars pin the composed header.
|
||||
{ name: 'text-turn', hasModelTurn: true, recorded: true, pinsHeader: true },
|
||||
{
|
||||
name: 'session-title-after-turn',
|
||||
hasModelTurn: true,
|
||||
recorded: false,
|
||||
overridden: true,
|
||||
configPath: SESSION_TITLE_CONFIG,
|
||||
},
|
||||
{ name: 'tool-call-turn', hasModelTurn: true, recorded: true },
|
||||
// Authored from the real PACKED_CHUNKS_SOURCE recording under the ordinary
|
||||
// app composition. The contract below pins decoded equality and all three
|
||||
|
||||
@@ -0,0 +1,8 @@
|
||||
{
|
||||
"steps": [
|
||||
{ "op": "initialize" },
|
||||
{ "op": "newSession" },
|
||||
{ "op": "prompt", "text": "Reply with exactly TITLE_DONE. Do not use tools." },
|
||||
{ "op": "waitForTitleAfterTurnEnd" }
|
||||
]
|
||||
}
|
||||
@@ -0,0 +1,15 @@
|
||||
[
|
||||
{
|
||||
"kind": "chunks",
|
||||
"chunks": [
|
||||
{ "type": "block-start", "index": 0, "blockType": "text" },
|
||||
{ "type": "text-delta", "index": 0, "text": "Late" },
|
||||
{ "type": "text-delta", "index": 0, "text": " durable" },
|
||||
{ "type": "text-delta", "index": 0, "text": " session" },
|
||||
{ "type": "text-delta", "index": 0, "text": " title" },
|
||||
{ "type": "block-end", "index": 0, "block": { "type": "text", "text": "Late durable session title" } },
|
||||
{ "type": "usage", "usage": { "inputTokens": 12, "outputTokens": 4 } },
|
||||
{ "type": "finish", "reason": { "kind": "stop" } }
|
||||
]
|
||||
}
|
||||
]
|
||||
@@ -0,0 +1,16 @@
|
||||
{"type":"session","version":0,"id":"session-title-after-turn","createdAt":0,"cwd":"/tmp/session-title-after-turn","delegationDepth":0}
|
||||
{"type":"turn/start","seq":0,"time":1785222848166,"data":{"turn":1,"trigger":{"kind":"message","source":{"kind":"user"}}}}
|
||||
{"type":"user/message","seq":1,"time":1785222848166,"data":{"content":[{"type":"text","text":"Reply with exactly TITLE_DONE. Do not use tools."}],"source":{"kind":"user"}},"surfaceOp":"append"}
|
||||
{"type":"session/title","seq":2,"time":1785222848166,"data":{"title":"Reply with exactly TITLE_DONE. Do","messageSeqs":[1],"source":{"kind":"fallback"}}}
|
||||
{"type":"step/start","seq":3,"time":1785222848199,"data":{"turn":1,"step":1}}
|
||||
{"type":"request/header","seq":4,"time":1785222848199,"data":{"header":{"config":{"provider":"deepseek","model":"deepseek-v4-flash"},"system":"{{system}}","tools":"{{tools}}"},"reason":"initial"}}
|
||||
{"type":"session/title-llm-request","seq":5,"time":1785222848201,"data":{"titleProvider":"session-title-first-message-llm","messageSeqs":[1],"route":{"provider":"title-replay","model":"title-model"},"system":"Create a concise title for an AI coding-assistant session from the supplied human messages.\nReturn only the title on one line, **in plain text of natural language**, with no quotes, prefix, explanation, Markdown, XML, or terminal control codes. No code is allowed.\nUse the language of the messages.\nAim for about 5 words in non-CJK languages or 10 CJK characters.","messages":[{"role":"user","content":[{"type":"text","text":"Generate the session title from this JSON array of human messages:\n[{\"seq\":1,\"text\":\"Reply with exactly TITLE_DONE. Do not use tools.\"}]"}]}],"maxTokens":32}}
|
||||
{"type":"assistant/chunk","seq":6,"time":1785222848208,"data":{"turn":1,"step":1,"chunk":{"type":"block-start","index":0,"blockType":"text"}}}
|
||||
{"type":"assistant/chunk","seq":7,"time":1785222848208,"data":{"turn":1,"step":1,"chunk":{"type":"text-delta","index":0,"text":"TITLE_DONE"}}}
|
||||
{"type":"assistant/chunk","seq":8,"time":1785222848208,"data":{"turn":1,"step":1,"chunk":{"type":"block-end","index":0,"block":{"type":"text","text":"TITLE_DONE"}}}}
|
||||
{"type":"assistant/chunk","seq":9,"time":1785222848208,"data":{"turn":1,"step":1,"chunk":{"type":"usage","usage":{"inputTokens":10,"outputTokens":2}}}}
|
||||
{"type":"assistant/chunk","seq":10,"time":1785222848208,"data":{"turn":1,"step":1,"chunk":{"type":"finish","reason":{"kind":"stop"}}}}
|
||||
{"type":"assistant/message","seq":11,"time":1785222848208,"data":{"turn":1,"step":1,"content":[{"type":"text","text":"TITLE_DONE"}],"provenance":{"provider":"deepseek","model":"deepseek-v4-flash"},"usage":{"inputTokens":10,"outputTokens":2}},"sourceEventSeqs":[6,7,8,9,10],"surfaceOp":"append"}
|
||||
{"type":"step/end","seq":12,"time":1785222848209,"data":{"turn":1,"step":1}}
|
||||
{"type":"turn/end","seq":13,"time":1785222848209,"data":{"turn":1,"reason":{"kind":"completed"}}}
|
||||
{"type":"session/title","seq":14,"time":1785222848209,"data":{"title":"Late durable session title","messageSeqs":[1],"source":{"kind":"provider","provider":"session-title-first-message-llm","model":{"provider":"title-replay","model":"title-model"}}}}
|
||||
@@ -0,0 +1,4 @@
|
||||
{"jsonrpc":"2.0","id":1,"result":{"protocolVersion":1,"agentInfo":{"name":"deepseek-harness-acp","version":"0.0.1"},"agentCapabilities":{"promptCapabilities":{"image":false,"audio":false,"embeddedContext":false}},"authMethods":[]}}
|
||||
{"jsonrpc":"2.0","id":2,"result":{"sessionId":"{{sessionId}}"}}
|
||||
{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"{{sessionId}}","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"TITLE_DONE"}}}}
|
||||
{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}
|
||||
@@ -17,6 +17,11 @@ interface CompactionTrace {
|
||||
summarized: boolean
|
||||
}
|
||||
|
||||
interface SessionTrace {
|
||||
openTurn: number | null
|
||||
compaction: CompactionTrace | undefined
|
||||
}
|
||||
|
||||
type CompactionTransition =
|
||||
| { kind: 'start'; turn: number }
|
||||
| { kind: 'summary'; turn: number }
|
||||
@@ -24,16 +29,27 @@ type CompactionTransition =
|
||||
|
||||
/** Validate one compaction event without advancing committed trace state. */
|
||||
function validateCompactionEvent(
|
||||
open: CompactionTrace | undefined,
|
||||
trace: SessionTrace,
|
||||
event: SessionEvent,
|
||||
fail: InvariantFailure,
|
||||
): CompactionTransition | undefined {
|
||||
if (event.type !== 'compact/start' && event.type !== 'compact/summary' && event.type !== 'compact/end') {
|
||||
return undefined
|
||||
}
|
||||
if (trace.openTurn === null) fail(`${event.type} appended outside any open turn`)
|
||||
const open = trace.compaction
|
||||
if (event.type === 'compact/start') {
|
||||
if (open !== undefined) fail(`compact/start for turn ${event.data.turn} while turn ${open.turn} is still compacting`)
|
||||
if (event.data.turn !== trace.openTurn) {
|
||||
fail(`compact/start names turn ${event.data.turn} but open turn is ${trace.openTurn}`)
|
||||
}
|
||||
return { kind: 'start', turn: event.data.turn }
|
||||
}
|
||||
if (event.type === 'compact/summary') {
|
||||
if (open === undefined) fail('compact/summary has no matching compact/start')
|
||||
if (open.turn !== trace.openTurn) {
|
||||
fail(`compact/summary belongs to turn ${open.turn} but open turn is ${trace.openTurn}`)
|
||||
}
|
||||
if (open.summarized) fail('compact/summary repeated within one compaction')
|
||||
const seqs = event.data.shadowedSeqs
|
||||
if (seqs.length === 0) fail('compact/summary shadowedSeqs must be non-empty')
|
||||
@@ -45,11 +61,13 @@ function validateCompactionEvent(
|
||||
}
|
||||
return { kind: 'summary', turn: open.turn }
|
||||
}
|
||||
if (event.type !== 'compact/end') return undefined
|
||||
if (open === undefined) fail('compact/end has no matching compact/start')
|
||||
if (event.data.turn !== open.turn) {
|
||||
fail(`compact/end turn ${event.data.turn} does not match compact/start turn ${open.turn}`)
|
||||
}
|
||||
if (event.data.turn !== trace.openTurn) {
|
||||
fail(`compact/end names turn ${event.data.turn} but open turn is ${trace.openTurn}`)
|
||||
}
|
||||
if (event.data.error === undefined && !open.summarized) {
|
||||
fail('successful compact/end requires one compact/summary')
|
||||
}
|
||||
@@ -69,29 +87,39 @@ function applyCompactionTransition(
|
||||
// Event owners keep precommit staging local so their vocabularies never move into a central helper.
|
||||
/* jscpd:ignore-start */
|
||||
const install: InvariantInstaller = Object.assign((ctx: Context, fail: InvariantFailure) => {
|
||||
const traces = new WeakMap<Session, CompactionTrace>()
|
||||
const traces = new WeakMap<Session, SessionTrace>()
|
||||
const staged = new WeakMap<SessionEvent, { session: Session; transition: CompactionTransition }>()
|
||||
const seed = (session: Session): void => {
|
||||
let open: CompactionTrace | undefined
|
||||
const seed = (session: Session): SessionTrace => {
|
||||
const trace: SessionTrace = { openTurn: null, compaction: undefined }
|
||||
traces.set(session, trace)
|
||||
for (const event of session.events) {
|
||||
const transition = validateCompactionEvent(open, event, fail)
|
||||
if (transition !== undefined) open = applyCompactionTransition(transition)
|
||||
if (event.type === 'turn/start') trace.openTurn = event.data.turn
|
||||
else if (event.type === 'turn/end') trace.openTurn = null
|
||||
const transition = validateCompactionEvent(trace, event, fail)
|
||||
if (transition !== undefined) trace.compaction = applyCompactionTransition(transition)
|
||||
}
|
||||
if (open !== undefined) traces.set(session, open)
|
||||
return trace
|
||||
}
|
||||
const traceFor = (session: Session): CompactionTrace | undefined => traces.get(session)
|
||||
const traceFor = (session: Session): SessionTrace => traces.get(session) ?? seed(session)
|
||||
|
||||
for (const session of ctx.sessions.list()) seed(session)
|
||||
ctx.on('session/created', (session) => { seed(session) }, { global: true })
|
||||
ctx.on('session/event', (session, event) => {
|
||||
const trace = traceFor(session)
|
||||
if (event.type === 'turn/start') {
|
||||
trace.openTurn = event.data.turn
|
||||
return
|
||||
}
|
||||
if (event.type === 'turn/end') {
|
||||
trace.openTurn = null
|
||||
return
|
||||
}
|
||||
if (event.type !== 'compact/start' && event.type !== 'compact/summary' && event.type !== 'compact/end') return
|
||||
const candidate = staged.get(event)
|
||||
/* v8 ignore next -- internal/dispatch stages every compaction event */
|
||||
if (candidate === undefined || candidate.session !== session) return fail('compaction event published without pre-commit validation')
|
||||
staged.delete(event)
|
||||
const next = applyCompactionTransition(candidate.transition)
|
||||
if (next === undefined) traces.delete(session)
|
||||
else traces.set(session, next)
|
||||
trace.compaction = applyCompactionTransition(candidate.transition)
|
||||
}, { global: true })
|
||||
ctx.on('internal/dispatch', (_mode, eventName, args) => {
|
||||
if (eventName !== 'session/event') return
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { Context } from 'cordis'
|
||||
import SessionStore from '@deepseek-ai/dsh-session'
|
||||
import SessionStore, { Session, SessionId } from '@deepseek-ai/dsh-session'
|
||||
import * as CompactInvariant from '@deepseek-ai/dsh-compact/invariant'
|
||||
import InvariantService from '@deepseek-ai/dsh-invariants'
|
||||
|
||||
@@ -22,15 +22,21 @@ const summary = (overrides: Record<string, unknown> = {}) => ({
|
||||
...overrides,
|
||||
})
|
||||
|
||||
function startTurn(session: ReturnType<Context['sessions']['create']>, turn = 1): void {
|
||||
session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
}
|
||||
|
||||
describe('compaction invariants', () => {
|
||||
it('accepts successful and failed compaction lifecycles', async () => {
|
||||
const ctx = await setup()
|
||||
const success = ctx.sessions.create()
|
||||
startTurn(success)
|
||||
success.append('compact/start', { turn: 1 })
|
||||
success.append('compact/summary', summary())
|
||||
success.append('compact/end', { turn: 1 })
|
||||
|
||||
const failed = ctx.sessions.create()
|
||||
startTurn(failed, 2)
|
||||
failed.append('compact/start', { turn: 2 })
|
||||
failed.append('compact/end', { turn: 2, error: 'provider failed' })
|
||||
})
|
||||
@@ -40,13 +46,68 @@ describe('compaction invariants', () => {
|
||||
await ctx.plugin(SessionStore)
|
||||
const session = ctx.sessions.create()
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
session.append('compact/start', { turn: 3 })
|
||||
session.append('compact/start', { turn: 1 })
|
||||
await ctx.plugin(InvariantService)
|
||||
await ctx.plugin(CompactInvariant)
|
||||
expect(() => session.append('compact/end', { turn: 3, error: 'resume failed' })).not.toThrow()
|
||||
expect(() => session.append('compact/end', { turn: 1, error: 'resume failed' })).not.toThrow()
|
||||
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
})
|
||||
|
||||
it('adopts a bare session and ignores unrelated committed events', async () => {
|
||||
const ctx = await setup()
|
||||
const session = new Session(SessionId('bare-compaction-session'))
|
||||
expect(() => {
|
||||
ctx.emit('session/event', session, {
|
||||
type: 'turn/start', seq: 0, time: 0,
|
||||
data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
|
||||
})
|
||||
ctx.emit('session/event', session, {
|
||||
type: 'step/start', seq: 1, time: 1, data: { turn: 1, step: 1 },
|
||||
})
|
||||
ctx.emit('session/event', session, {
|
||||
type: 'compact/start', seq: 2, time: 2, data: { turn: 1 },
|
||||
})
|
||||
}).not.toThrow()
|
||||
})
|
||||
|
||||
it('rejects compaction outside or for a different open turn', async () => {
|
||||
const ctx = await setup()
|
||||
const session = ctx.sessions.create()
|
||||
expect(() => session.append('compact/start', { turn: 1 })).toThrow(/outside any open turn/)
|
||||
startTurn(session)
|
||||
expect(() => session.append('compact/start', { turn: 2 })).toThrow(/but open turn is 1/)
|
||||
})
|
||||
|
||||
it('rejects an unenclosed compaction event when replaying an existing session', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
const session = ctx.sessions.create()
|
||||
startTurn(session)
|
||||
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
session.append('compact/start', { turn: 1 })
|
||||
await ctx.plugin(InvariantService)
|
||||
await expect(ctx.plugin(CompactInvariant).then(() => undefined)).rejects.toThrow(/outside any open turn/)
|
||||
})
|
||||
|
||||
it('rejects an open compaction that crosses into another turn', async () => {
|
||||
const ctx = await setup()
|
||||
const summarySession = ctx.sessions.create()
|
||||
startTurn(summarySession)
|
||||
summarySession.append('compact/start', { turn: 1 })
|
||||
summarySession.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
startTurn(summarySession, 2)
|
||||
expect(() => summarySession.append('compact/summary', summary()))
|
||||
.toThrow(/belongs to turn 1 but open turn is 2/)
|
||||
|
||||
const endSession = ctx.sessions.create()
|
||||
startTurn(endSession)
|
||||
endSession.append('compact/start', { turn: 1 })
|
||||
endSession.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
startTurn(endSession, 2)
|
||||
expect(() => endSession.append('compact/end', { turn: 1, error: 'late' }))
|
||||
.toThrow(/names turn 1 but open turn is 2/)
|
||||
})
|
||||
|
||||
it.each([
|
||||
['summary without start', (session: ReturnType<Context['sessions']['create']>) => {
|
||||
session.append('compact/summary', summary())
|
||||
@@ -85,6 +146,8 @@ describe('compaction invariants', () => {
|
||||
}, /requires one compact\/summary/],
|
||||
])('rejects %s', async (_name, action, message) => {
|
||||
const ctx = await setup()
|
||||
expect(() => { action(ctx.sessions.create()) }).toThrow(message)
|
||||
const session = ctx.sessions.create()
|
||||
startTurn(session)
|
||||
expect(() => { action(session) }).toThrow(message)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
/** Package-owned tool-pipeline invariants. @module @deepseek-ai/dsh-tools/invariant */
|
||||
|
||||
import type { Context } from 'cordis'
|
||||
import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
|
||||
import type { InvariantFailure, InvariantInstaller } from '@deepseek-ai/dsh-invariants'
|
||||
import type { ToolExecution, ToolExecutionResult } from './index.ts'
|
||||
|
||||
@@ -28,10 +29,40 @@ function validateResult(
|
||||
}
|
||||
}
|
||||
|
||||
/** Install monotonic pipeline and final-snapshot checks. */
|
||||
const install: InvariantInstaller = (ctx, fail) => {
|
||||
/** Install monotonic pipeline, final-snapshot, and code-dispatch enclosure checks. */
|
||||
const install: InvariantInstaller = Object.assign((ctx: Context, fail: InvariantFailure) => {
|
||||
const stages = new WeakMap<object, ToolStage>()
|
||||
const openTurns = new WeakMap<Session, number | null>()
|
||||
const seed = (session: Session): number | null => {
|
||||
let openTurn: number | null = null
|
||||
for (const event of session.events) {
|
||||
if (event.type === 'turn/start') openTurn = event.data.turn
|
||||
else if (event.type === 'turn/end') openTurn = null
|
||||
else if ((event.type === 'tool/code-dispatch-start' || event.type === 'tool/code-dispatch')
|
||||
&& openTurn === null) {
|
||||
fail(`${event.type} appended outside any open turn`)
|
||||
}
|
||||
}
|
||||
openTurns.set(session, openTurn)
|
||||
return openTurn
|
||||
}
|
||||
const openTurnFor = (session: Session): number | null => openTurns.get(session) ?? seed(session)
|
||||
|
||||
for (const session of ctx.sessions.list()) seed(session)
|
||||
ctx.on('session/created', (session) => { seed(session) }, { global: true })
|
||||
ctx.on('session/event', (session, event) => {
|
||||
if (event.type === 'turn/start') openTurns.set(session, event.data.turn)
|
||||
else if (event.type === 'turn/end') openTurns.set(session, null)
|
||||
}, { global: true })
|
||||
ctx.on('internal/dispatch', (_mode, eventName, args) => {
|
||||
if (eventName === 'session/event') {
|
||||
const [session, event] = args as [Session, SessionEvent]
|
||||
if ((event.type === 'tool/code-dispatch-start' || event.type === 'tool/code-dispatch')
|
||||
&& openTurnFor(session) === null) {
|
||||
fail(`${event.type} appended outside any open turn`)
|
||||
}
|
||||
return
|
||||
}
|
||||
if (eventName === 'tools/pre-execute') {
|
||||
const exec = args[0] as ToolExecution
|
||||
if (stages.has(exec)) fail('tools/pre-execute repeated for one execution')
|
||||
@@ -58,7 +89,7 @@ const install: InvariantInstaller = (ctx, fail) => {
|
||||
validateResult(exec, result, fail)
|
||||
stages.delete(exec)
|
||||
}, { global: true })
|
||||
}
|
||||
}, { inject: ['sessions'] })
|
||||
|
||||
/**
|
||||
* Register the tools invariant companion.
|
||||
|
||||
@@ -2,6 +2,7 @@ import { describe, expect, it } from 'vitest'
|
||||
import { Context } from 'cordis'
|
||||
import { scopeTarget } from '@deepseek-ai/dsh-scope'
|
||||
import { CallId } from '@deepseek-ai/dsh-llm'
|
||||
import SessionStore from '@deepseek-ai/dsh-session'
|
||||
import type { ToolExecution, ToolExecutionResult, ToolExecutionToken } from '@deepseek-ai/dsh-tools'
|
||||
import * as ToolsInvariant from '@deepseek-ai/dsh-tools/invariant'
|
||||
import InvariantService from '@deepseek-ai/dsh-invariants'
|
||||
@@ -10,6 +11,7 @@ const testToolSignal = new AbortController().signal
|
||||
|
||||
async function setup(): Promise<Context> {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
await ctx.plugin(InvariantService)
|
||||
await ctx.plugin(ToolsInvariant)
|
||||
return ctx
|
||||
@@ -85,4 +87,50 @@ describe('tool-pipeline invariants', () => {
|
||||
const anonymous = Object.freeze(execution({ name: '' }))
|
||||
expect(() => { emitResult(ctx, anonymous, outcome()) }).toThrow(/non-empty name and callId/)
|
||||
})
|
||||
|
||||
it('requires code-dispatch records to be turn-enclosed', async () => {
|
||||
const ctx = await setup()
|
||||
const session = ctx.sessions.create()
|
||||
const data = {
|
||||
parentCallId: CallId('parent'),
|
||||
subCallId: CallId('child'),
|
||||
name: 'echo',
|
||||
arguments: {},
|
||||
}
|
||||
expect(() => session.append('tool/code-dispatch-start', data)).toThrow(/outside any open turn/)
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
expect(() => session.append('tool/code-dispatch-start', data)).not.toThrow()
|
||||
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
})
|
||||
|
||||
it('replays enclosed code-dispatch records on late registration', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
const session = ctx.sessions.create()
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
session.append('tool/code-dispatch', {
|
||||
parentCallId: CallId('parent'),
|
||||
subCallId: CallId('child'),
|
||||
name: 'echo',
|
||||
arguments: {},
|
||||
isError: false,
|
||||
content: [{ type: 'text', text: 'ok' }],
|
||||
})
|
||||
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
await ctx.plugin(InvariantService)
|
||||
await expect(ctx.plugin(ToolsInvariant).then(() => undefined)).resolves.toBeUndefined()
|
||||
})
|
||||
|
||||
it('rejects an unenclosed code-dispatch record on late registration', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
ctx.sessions.create().append('tool/code-dispatch-start', {
|
||||
parentCallId: CallId('parent'),
|
||||
subCallId: CallId('child'),
|
||||
name: 'echo',
|
||||
arguments: {},
|
||||
})
|
||||
await ctx.plugin(InvariantService)
|
||||
await expect(ctx.plugin(ToolsInvariant).then(() => undefined)).rejects.toThrow(/outside any open turn/)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -17,6 +17,11 @@ interface HookTransition {
|
||||
delta: 1 | -1
|
||||
}
|
||||
|
||||
interface HookTrace {
|
||||
openTurn: number | null
|
||||
pending: Map<string, number>
|
||||
}
|
||||
|
||||
/** Correlation key shared by an invoked/result pair. */
|
||||
function hookKey(data: { turn: number; point: string; handlerId: string }): string {
|
||||
return `${data.turn}\0${data.point}\0${data.handlerId}`
|
||||
@@ -24,10 +29,15 @@ function hookKey(data: { turn: number; point: string; handlerId: string }): stri
|
||||
|
||||
/** Validate one hook event against committed pending invocations. */
|
||||
function validateHookEvent(
|
||||
pending: ReadonlyMap<string, number>,
|
||||
trace: HookTrace,
|
||||
event: SessionEvent,
|
||||
fail: InvariantFailure,
|
||||
): HookTransition | undefined {
|
||||
if (event.type !== 'hook/invoked' && event.type !== 'hook/result') return undefined
|
||||
if (trace.openTurn === null) fail(`${event.type} appended outside any open turn`)
|
||||
if (event.data.turn !== trace.openTurn) {
|
||||
fail(`${event.type} names turn ${event.data.turn} but open turn is ${trace.openTurn}`)
|
||||
}
|
||||
if (event.type === 'hook/invoked') {
|
||||
if (event.data.point.length === 0 || event.data.handlerId.length === 0) {
|
||||
fail('hook/invoked point and handlerId must be non-empty')
|
||||
@@ -38,9 +48,8 @@ function validateHookEvent(
|
||||
}
|
||||
return { key: hookKey(event.data), delta: 1 }
|
||||
}
|
||||
if (event.type !== 'hook/result') return undefined
|
||||
const key = hookKey(event.data)
|
||||
if ((pending.get(key) ?? 0) === 0) {
|
||||
if ((trace.pending.get(key) ?? 0) === 0) {
|
||||
fail(`hook/result has no matching hook/invoked for ${JSON.stringify(event.data.handlerId)}`)
|
||||
}
|
||||
if (!Number.isFinite(event.data.durationMs) || event.data.durationMs < 0) {
|
||||
@@ -60,28 +69,39 @@ function applyHookTransition(pending: Map<string, number>, transition: HookTrans
|
||||
// Event owners keep precommit staging local so their vocabularies never move into a central helper.
|
||||
/* jscpd:ignore-start */
|
||||
const install: InvariantInstaller = Object.assign((ctx: Context, fail: InvariantFailure) => {
|
||||
const traces = new WeakMap<Session, Map<string, number>>()
|
||||
const traces = new WeakMap<Session, HookTrace>()
|
||||
const staged = new WeakMap<SessionEvent, { session: Session; transition: HookTransition }>()
|
||||
const seed = (session: Session): Map<string, number> => {
|
||||
const pending = new Map<string, number>()
|
||||
traces.set(session, pending)
|
||||
const seed = (session: Session): HookTrace => {
|
||||
const trace: HookTrace = { openTurn: null, pending: new Map() }
|
||||
traces.set(session, trace)
|
||||
for (const event of session.events) {
|
||||
const transition = validateHookEvent(pending, event, fail)
|
||||
if (transition !== undefined) applyHookTransition(pending, transition)
|
||||
if (event.type === 'turn/start') trace.openTurn = event.data.turn
|
||||
else if (event.type === 'turn/end') trace.openTurn = null
|
||||
const transition = validateHookEvent(trace, event, fail)
|
||||
if (transition !== undefined) applyHookTransition(trace.pending, transition)
|
||||
}
|
||||
return pending
|
||||
return trace
|
||||
}
|
||||
const traceFor = (session: Session): Map<string, number> => traces.get(session) ?? seed(session)
|
||||
const traceFor = (session: Session): HookTrace => traces.get(session) ?? seed(session)
|
||||
|
||||
for (const session of ctx.sessions.list()) seed(session)
|
||||
ctx.on('session/created', (session) => { seed(session) }, { global: true })
|
||||
ctx.on('session/event', (session, event) => {
|
||||
const trace = traceFor(session)
|
||||
if (event.type === 'turn/start') {
|
||||
trace.openTurn = event.data.turn
|
||||
return
|
||||
}
|
||||
if (event.type === 'turn/end') {
|
||||
trace.openTurn = null
|
||||
return
|
||||
}
|
||||
if (event.type !== 'hook/invoked' && event.type !== 'hook/result') return
|
||||
const candidate = staged.get(event)
|
||||
/* v8 ignore next -- internal/dispatch stages every hook provenance event */
|
||||
if (candidate === undefined || candidate.session !== session) return fail('hook event published without pre-commit validation')
|
||||
staged.delete(event)
|
||||
applyHookTransition(traceFor(session), candidate.transition)
|
||||
applyHookTransition(trace.pending, candidate.transition)
|
||||
}, { global: true })
|
||||
ctx.on('internal/dispatch', (_mode, eventName, args) => {
|
||||
if (eventName !== 'session/event') return
|
||||
|
||||
@@ -29,12 +29,18 @@ const result = (overrides: Record<string, unknown> = {}) => ({
|
||||
...overrides,
|
||||
})
|
||||
|
||||
function startTurn(session: Session, turn = 1): void {
|
||||
session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
}
|
||||
|
||||
describe('hook-protocol invariants', () => {
|
||||
it('pairs serial and repeated handler invocations', async () => {
|
||||
const ctx = await setup()
|
||||
const session = ctx.sessions.create()
|
||||
startTurn(session)
|
||||
session.append('hook/invoked', invoked())
|
||||
session.append('hook/invoked', invoked())
|
||||
session.append('step/start', { turn: 1, step: 1 })
|
||||
session.append('hook/result', result())
|
||||
session.append('hook/result', result())
|
||||
})
|
||||
@@ -56,26 +62,52 @@ describe('hook-protocol invariants', () => {
|
||||
const session = new Session(SessionId('bare-hook-session'))
|
||||
expect(() => {
|
||||
ctx.emit('session/event', session, {
|
||||
type: 'hook/invoked', seq: 0, time: 0, data: invoked(),
|
||||
type: 'turn/start', seq: 0, time: 0,
|
||||
data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
|
||||
})
|
||||
ctx.emit('session/event', session, {
|
||||
type: 'hook/result', seq: 1, time: 1, data: result(),
|
||||
type: 'hook/invoked', seq: 1, time: 1, data: invoked(),
|
||||
})
|
||||
ctx.emit('session/event', session, {
|
||||
type: 'hook/result', seq: 2, time: 2, data: result(),
|
||||
})
|
||||
}).not.toThrow()
|
||||
})
|
||||
|
||||
it('rejects hook events outside or for a different open turn', async () => {
|
||||
const ctx = await setup()
|
||||
const session = ctx.sessions.create()
|
||||
expect(() => session.append('hook/invoked', invoked())).toThrow(/outside any open turn/)
|
||||
startTurn(session)
|
||||
expect(() => session.append('hook/invoked', invoked({ turn: 2 }))).toThrow(/but open turn is 1/)
|
||||
})
|
||||
|
||||
it('rejects an unenclosed hook event when replaying an existing session', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
const session = ctx.sessions.create()
|
||||
startTurn(session)
|
||||
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
session.append('hook/invoked', invoked())
|
||||
await ctx.plugin(InvariantService)
|
||||
await expect(ctx.plugin(HookInvariant).then(() => undefined)).rejects.toThrow(/outside any open turn/)
|
||||
})
|
||||
|
||||
it.each([
|
||||
[invoked({ point: '' }), /point and handlerId must be non-empty/],
|
||||
[invoked({ handlerId: '' }), /point and handlerId must be non-empty/],
|
||||
[invoked({ dialect: 'other' }), /unknown dialect/],
|
||||
])('rejects malformed hook invocation %#', async (data, message) => {
|
||||
const ctx = await setup()
|
||||
expect(() => ctx.sessions.create().append('hook/invoked', data as never)).toThrow(message)
|
||||
const session = ctx.sessions.create()
|
||||
startTurn(session)
|
||||
expect(() => session.append('hook/invoked', data as never)).toThrow(message)
|
||||
})
|
||||
|
||||
it('rejects unmatched and malformed results', async () => {
|
||||
const ctx = await setup()
|
||||
const session = ctx.sessions.create()
|
||||
startTurn(session)
|
||||
expect(() => session.append('hook/result', result())).toThrow(/no matching hook\/invoked/)
|
||||
session.append('hook/invoked', invoked())
|
||||
expect(() => session.append('hook/result', result({ durationMs: -1 })))
|
||||
|
||||
@@ -11,9 +11,10 @@ export const name = 'plan-mode-invariant'
|
||||
/** Service required before the companion can reserve package ownership. */
|
||||
export const inject = ['invariants']
|
||||
|
||||
/** Validate one `plan/mode` payload before it reaches the durable log. */
|
||||
function validateEvent(event: SessionEvent, fail: InvariantFailure): void {
|
||||
/** Validate one `plan/mode` event before it reaches the durable log. */
|
||||
function validateEvent(openTurn: number | null, event: SessionEvent, fail: InvariantFailure): void {
|
||||
if (event.type !== 'plan/mode') return
|
||||
if (openTurn === null) fail('plan/mode appended outside any open turn')
|
||||
const active = (event.data as { active?: unknown }).active
|
||||
if (typeof active !== 'boolean') {
|
||||
fail(`plan/mode carries invalid active state ${JSON.stringify(active)}; expected a boolean`)
|
||||
@@ -23,13 +24,30 @@ function validateEvent(event: SessionEvent, fail: InvariantFailure): void {
|
||||
/* jscpd:ignore-start -- package companions share replay and dispatch plumbing */
|
||||
/** Install validation for loaded and newly appended plan-mode state. */
|
||||
const install: InvariantInstaller = Object.assign((ctx: Context, fail: InvariantFailure) => {
|
||||
for (const session of ctx.sessions.list()) {
|
||||
for (const event of session.events) validateEvent(event, fail)
|
||||
const traces = new WeakMap<Session, number | null>()
|
||||
const seed = (session: Session): number | null => {
|
||||
let openTurn: number | null = null
|
||||
traces.set(session, openTurn)
|
||||
for (const event of session.events) {
|
||||
if (event.type === 'turn/start') openTurn = event.data.turn
|
||||
else if (event.type === 'turn/end') openTurn = null
|
||||
validateEvent(openTurn, event, fail)
|
||||
traces.set(session, openTurn)
|
||||
}
|
||||
return openTurn
|
||||
}
|
||||
const traceFor = (session: Session): number | null => traces.get(session) ?? seed(session)
|
||||
|
||||
for (const session of ctx.sessions.list()) seed(session)
|
||||
ctx.on('session/created', (session) => { seed(session) }, { global: true })
|
||||
ctx.on('session/event', (session, event) => {
|
||||
if (event.type === 'turn/start') traces.set(session, event.data.turn)
|
||||
else if (event.type === 'turn/end') traces.set(session, null)
|
||||
}, { global: true })
|
||||
ctx.on('internal/dispatch', (_mode, eventName, args) => {
|
||||
if (eventName !== 'session/event') return
|
||||
const event = (args as [Session, SessionEvent])[1]
|
||||
validateEvent(event, fail)
|
||||
const [session, event] = args as [Session, SessionEvent]
|
||||
validateEvent(traceFor(session), event, fail)
|
||||
}, { global: true })
|
||||
}, { inject: ['sessions'] })
|
||||
/* jscpd:ignore-end */
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { Context } from 'cordis'
|
||||
import SessionStore, { type Session, type SessionEvent } from '@deepseek-ai/dsh-session'
|
||||
import SessionStore, { Session, SessionId, type SessionEvent } from '@deepseek-ai/dsh-session'
|
||||
import * as PlanModeInvariant from '@deepseek-ai/dsh-plan-mode/invariant'
|
||||
import InvariantService from '@deepseek-ai/dsh-invariants'
|
||||
|
||||
@@ -16,24 +16,46 @@ function event(active: unknown): SessionEvent {
|
||||
return { type: 'plan/mode', seq: 0, time: 0, data: { active } } as SessionEvent
|
||||
}
|
||||
|
||||
function emitTurnStart(ctx: Context, session: Session): void {
|
||||
ctx.emit('session/event', session, {
|
||||
type: 'turn/start', seq: 0, time: 0,
|
||||
data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
|
||||
})
|
||||
}
|
||||
|
||||
describe('plan-mode stream invariants', () => {
|
||||
it('accepts either boolean state', async () => {
|
||||
const ctx = await setup()
|
||||
expect(() => { ctx.emit('session/event', {} as Session, event(true)) }).not.toThrow()
|
||||
expect(() => { ctx.emit('session/event', {} as Session, event(false)) }).not.toThrow()
|
||||
const session = new Session(SessionId('plan-state'))
|
||||
emitTurnStart(ctx, session)
|
||||
expect(() => { ctx.emit('session/event', session, event(true)) }).not.toThrow()
|
||||
expect(() => { ctx.emit('session/event', session, event(false)) }).not.toThrow()
|
||||
ctx.emit('session/event', session, {
|
||||
type: 'turn/end', seq: 3, time: 3,
|
||||
data: { turn: 1, reason: { kind: 'completed' } },
|
||||
})
|
||||
})
|
||||
|
||||
it.each([42, 'plan', undefined])('rejects invalid durable plan state %j', async (active) => {
|
||||
const ctx = await setup()
|
||||
expect(() => { ctx.emit('session/event', {} as Session, event(active)) })
|
||||
const session = new Session(SessionId(`invalid-${String(active)}`))
|
||||
emitTurnStart(ctx, session)
|
||||
expect(() => { ctx.emit('session/event', session, event(active)) })
|
||||
.toThrow(/expected a boolean/)
|
||||
})
|
||||
|
||||
it('rejects plan state outside any open turn', async () => {
|
||||
const ctx = await setup()
|
||||
expect(() => ctx.sessions.create().append('plan/mode', { active: true }))
|
||||
.toThrow(/outside any open turn/)
|
||||
})
|
||||
|
||||
it('ignores unrelated dispatches and session events', async () => {
|
||||
const ctx = await setup()
|
||||
const session = new Session(SessionId('unrelated'))
|
||||
expect(() => {
|
||||
ctx.emit('tools/change')
|
||||
ctx.emit('session/event', {} as Session, {
|
||||
ctx.emit('session/event', session, {
|
||||
type: 'turn/start', seq: 0, time: 0, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
|
||||
})
|
||||
}).not.toThrow()
|
||||
@@ -42,9 +64,33 @@ describe('plan-mode stream invariants', () => {
|
||||
it('rejects invalid existing state on late registration', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
ctx.sessions.create().append('plan/mode', { active: 'plan' as unknown as boolean })
|
||||
const session = ctx.sessions.create()
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
session.append('plan/mode', { active: 'plan' as unknown as boolean })
|
||||
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
await ctx.plugin(InvariantService, { enabled: true })
|
||||
|
||||
await expect(ctx.plugin(PlanModeInvariant).then(() => undefined)).rejects.toThrow(/expected a boolean/)
|
||||
})
|
||||
|
||||
it('replays enclosed existing plan state through its closing boundary', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
const session = ctx.sessions.create()
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
session.append('plan/mode', { active: true })
|
||||
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
await ctx.plugin(InvariantService, { enabled: true })
|
||||
|
||||
await expect(ctx.plugin(PlanModeInvariant).then(() => undefined)).resolves.toBeUndefined()
|
||||
})
|
||||
|
||||
it('rejects unenclosed existing plan state on late registration', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
ctx.sessions.create().append('plan/mode', { active: true })
|
||||
await ctx.plugin(InvariantService, { enabled: true })
|
||||
|
||||
await expect(ctx.plugin(PlanModeInvariant).then(() => undefined)).rejects.toThrow(/outside any open turn/)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -50,6 +50,7 @@ const WAIT_POLL_INTERVAL_MS = 10
|
||||
* `waitForTurnStart` waits for an open durable turn, optionally at or beyond a
|
||||
* specified turn number. `waitForTurnEnd` holds the subprocess open until the
|
||||
* selected session's latest complete raw-JSONL turn boundary is `turn/end`.
|
||||
* `waitForTitleAfterTurnEnd` additionally waits for a later durable title.
|
||||
* A standalone `cancel` may also wait for a cwd-relative readiness marker.
|
||||
* All wait timeouts default to 10s.
|
||||
*/
|
||||
@@ -67,6 +68,7 @@ export type InputStep =
|
||||
}
|
||||
| { op: 'waitForTurnStart'; minimumTurn?: number; timeoutMs?: number }
|
||||
| { op: 'waitForTurnEnd'; timeoutMs?: number }
|
||||
| { op: 'waitForTitleAfterTurnEnd'; timeoutMs?: number }
|
||||
| { op: 'cancel'; waitForFile?: { path: string; timeoutMs?: number } }
|
||||
|
||||
/** A scenario's `input.json`: an ordered list of input steps. */
|
||||
@@ -280,6 +282,7 @@ export async function runScenario(input: InputScript, opts: RunOptions): Promise
|
||||
(id) => { sessionId = id },
|
||||
(id, timeoutMs, minimumTurn) => waitForPersistedTurnStart(sessionsRoot, id, timeoutMs, minimumTurn),
|
||||
(id, timeoutMs) => waitForPersistedTurnEnd(sessionsRoot, id, timeoutMs),
|
||||
(id, timeoutMs) => waitForPersistedTitleAfterTurnEnd(sessionsRoot, id, timeoutMs),
|
||||
)
|
||||
// A permission exchange happens while a step's request is in flight, so
|
||||
// by the time the step settles any script bug it exposed is captured —
|
||||
@@ -353,6 +356,7 @@ async function runStep(
|
||||
setSessionId: (id: string) => void,
|
||||
waitForTurnStart: (sessionId: string, timeoutMs?: number, minimumTurn?: number) => Promise<void>,
|
||||
waitForTurnEnd: (sessionId: string, timeoutMs?: number) => Promise<void>,
|
||||
waitForTitleAfterTurnEnd: (sessionId: string, timeoutMs?: number) => Promise<void>,
|
||||
): Promise<void> {
|
||||
switch (step.op) {
|
||||
case 'initialize':
|
||||
@@ -430,6 +434,12 @@ async function runStep(
|
||||
await waitForTurnEnd(sessionId, step.timeoutMs)
|
||||
return
|
||||
}
|
||||
case 'waitForTitleAfterTurnEnd': {
|
||||
const sessionId = getSessionId()
|
||||
if (sessionId === undefined) throw new Error('snapshot-harness: waitForTitleAfterTurnEnd before newSession')
|
||||
await waitForTitleAfterTurnEnd(sessionId, step.timeoutMs)
|
||||
return
|
||||
}
|
||||
case 'waitForTurnStart': {
|
||||
const sessionId = getSessionId()
|
||||
if (sessionId === undefined) throw new Error('snapshot-harness: waitForTurnStart before newSession')
|
||||
@@ -497,6 +507,20 @@ async function waitForPersistedTurnEnd(
|
||||
}, { interval: WAIT_POLL_INTERVAL_MS, timeout: timeoutMs })
|
||||
}
|
||||
|
||||
/** Wait until a complete provider or fallback title record follows the latest closed turn. */
|
||||
async function waitForPersistedTitleAfterTurnEnd(
|
||||
root: string,
|
||||
sessionId: string,
|
||||
timeoutMs = DEFAULT_WAIT_TIMEOUT_MS,
|
||||
): Promise<void> {
|
||||
await vi.waitFor(async () => {
|
||||
const log = (await harvestSessionLogs(root)).find(candidate => candidate.id === sessionId)
|
||||
if (log === undefined || !latestTitleFollowsTurnEnd(log.content)) {
|
||||
throw new Error(`snapshot-harness: session "${sessionId}" did not persist session/title after turn/end within ${timeoutMs}ms`)
|
||||
}
|
||||
}, { interval: WAIT_POLL_INTERVAL_MS, timeout: timeoutMs })
|
||||
}
|
||||
|
||||
/** Wait for a cwd-relative marker proving an external action reached readiness. */
|
||||
async function waitForWorkspaceFile(
|
||||
cwd: string,
|
||||
@@ -518,6 +542,13 @@ function latestTurnIsClosed(content: string): boolean {
|
||||
> complete.lastIndexOf('\n{"type":"turn/start",')
|
||||
}
|
||||
|
||||
/** Return whether the last complete title record occurs after the last complete turn end. */
|
||||
function latestTitleFollowsTurnEnd(content: string): boolean {
|
||||
const complete = content.slice(0, content.lastIndexOf('\n') + 1)
|
||||
const turnEnd = complete.lastIndexOf('\n{"type":"turn/end",')
|
||||
return turnEnd >= 0 && complete.lastIndexOf('\n{"type":"session/title",') > turnEnd
|
||||
}
|
||||
|
||||
/** Return the latest open turn number, validating the persisted boundary record. */
|
||||
function latestOpenTurn(content: string): number | undefined {
|
||||
const complete = content.slice(0, content.lastIndexOf('\n') + 1)
|
||||
@@ -571,8 +602,8 @@ async function harvestSessionLogs(root: string): Promise<HarvestedLog[]> {
|
||||
// so session.<n>.jsonl maps to the same child on record and replay — replay
|
||||
// re-sorts childFiles by the same key, so the two stay consistent.
|
||||
logs.sort((a, b) => {
|
||||
const ap = a.parentSession === undefined ? 0 : 1
|
||||
const bp = b.parentSession === undefined ? 0 : 1
|
||||
const ap = Number(a.parentSession !== undefined)
|
||||
const bp = Number(b.parentSession !== undefined)
|
||||
return ap - bp || a.createdAt - b.createdAt || a.id.localeCompare(b.id)
|
||||
})
|
||||
return logs
|
||||
|
||||
@@ -536,7 +536,7 @@ describe('runScenario', () => {
|
||||
waitForText: 'thinking about it',
|
||||
}],
|
||||
},
|
||||
{ agent: AGENT, mode: 'replay', fixtureFile },
|
||||
{ agent: AGENT, mode: 'replay', fixtureFile, configPath: AGENT.configPath },
|
||||
)
|
||||
expect(result.rawStdout).toContain('thinking about it')
|
||||
})
|
||||
@@ -560,6 +560,32 @@ describe('runScenario', () => {
|
||||
expect(result.sessionLogs[0]?.content).toContain('"type":"turn/end"')
|
||||
})
|
||||
|
||||
it('waitForTitleAfterTurnEnd holds the app through a standalone durable title', { timeout: 20_000 }, async () => {
|
||||
const { fixtureFile } = await scenario({
|
||||
prompt: 'hang-until-cancel',
|
||||
persistLogsOnCancel: true,
|
||||
logs: [{
|
||||
file: 'project/main/session.jsonl',
|
||||
lines: [
|
||||
{ type: 'session', version: 0, id: '{{SID}}', createdAt: 1, delegationDepth: 0 },
|
||||
{ type: 'turn/end', seq: 1, time: 2, data: { turn: 1, reason: { kind: 'aborted' } } },
|
||||
{ type: 'session/title', seq: 2, time: 3, data: { title: 'Late title' } },
|
||||
],
|
||||
}],
|
||||
})
|
||||
const result = await runScenario(
|
||||
{
|
||||
steps: [
|
||||
...boot,
|
||||
{ op: 'promptAndCancel', text: 'hang' },
|
||||
{ op: 'waitForTitleAfterTurnEnd' },
|
||||
],
|
||||
},
|
||||
{ agent: AGENT, mode: 'replay', fixtureFile },
|
||||
)
|
||||
expect(result.sessionLogs[0]?.content).toMatch(/"turn\/end"[\s\S]*"session\/title"/)
|
||||
})
|
||||
|
||||
it('waitForTurnStart can require a later durable turn before continuing', { timeout: 20_000 }, async () => {
|
||||
const { fixtureFile } = await scenario({
|
||||
prompt: 'hang-until-cancel',
|
||||
@@ -692,6 +718,31 @@ describe('runScenario', () => {
|
||||
)).rejects.toThrow(/did not persist turn\/end within 20ms/)
|
||||
})
|
||||
|
||||
it('waitForTitleAfterTurnEnd times out when the title precedes the boundary', { timeout: 20_000 }, async () => {
|
||||
const { fixtureFile } = await scenario({
|
||||
prompt: 'hang-until-cancel',
|
||||
persistLogsOnCancel: true,
|
||||
logs: [{
|
||||
file: 'project/main/session.jsonl',
|
||||
lines: [
|
||||
{ type: 'session', version: 0, id: '{{SID}}', createdAt: 1, delegationDepth: 0 },
|
||||
{ type: 'session/title', seq: 1, time: 1, data: { title: 'Early title' } },
|
||||
{ type: 'turn/end', seq: 2, time: 2, data: { turn: 1, reason: { kind: 'aborted' } } },
|
||||
],
|
||||
}],
|
||||
})
|
||||
await expect(runScenario(
|
||||
{
|
||||
steps: [
|
||||
...boot,
|
||||
{ op: 'promptAndCancel', text: 'hang' },
|
||||
{ op: 'waitForTitleAfterTurnEnd', timeoutMs: 20 },
|
||||
],
|
||||
},
|
||||
{ agent: AGENT, mode: 'replay', fixtureFile },
|
||||
)).rejects.toThrow(/did not persist session\/title after turn\/end within 20ms/)
|
||||
})
|
||||
|
||||
it('promptExpectError swallows a model-error response as the expected outcome', { timeout: 20_000 }, async () => {
|
||||
const { fixtureFile } = await scenario({ prompt: 'error' })
|
||||
const result = await runScenario(
|
||||
@@ -797,6 +848,7 @@ describe('runScenario', () => {
|
||||
[{ op: 'promptAndCancel', text: 'x' }, /promptAndCancel before newSession/],
|
||||
[{ op: 'waitForTurnStart' }, /waitForTurnStart before newSession/],
|
||||
[{ op: 'waitForTurnEnd' }, /waitForTurnEnd before newSession/],
|
||||
[{ op: 'waitForTitleAfterTurnEnd' }, /waitForTitleAfterTurnEnd before newSession/],
|
||||
[{ op: 'cancel' }, /cancel before newSession/],
|
||||
] as [InputStep, RegExp][])('rejects %j before newSession', { timeout: 20_000 }, async (step, message) => {
|
||||
const { fixtureFile } = await scenario({})
|
||||
|
||||
@@ -18,19 +18,26 @@ type ApprovalTransition =
|
||||
| { kind: 'asked'; id: ApprovalRequestId }
|
||||
| { kind: 'decided'; id: ApprovalRequestId }
|
||||
|
||||
interface ApprovalTrace {
|
||||
openTurn: number | null
|
||||
pending: Set<ApprovalRequestId>
|
||||
}
|
||||
|
||||
/** Validate one approval event against committed unmatched questions. */
|
||||
function validateApprovalEvent(
|
||||
pending: ReadonlySet<ApprovalRequestId>,
|
||||
trace: ApprovalTrace,
|
||||
event: SessionEvent,
|
||||
fail: InvariantFailure,
|
||||
): ApprovalTransition | undefined {
|
||||
if (event.type === 'approval/asked') {
|
||||
if (trace.openTurn === null) fail('approval/asked appended outside any open turn')
|
||||
if (event.data.toolName.length === 0) fail('approval/asked toolName must be non-empty')
|
||||
if (pending.has(event.data.id)) fail(`approval/asked repeated open id ${JSON.stringify(event.data.id)}`)
|
||||
if (trace.pending.has(event.data.id)) fail(`approval/asked repeated open id ${JSON.stringify(event.data.id)}`)
|
||||
return { kind: 'asked', id: event.data.id }
|
||||
}
|
||||
if (event.type === 'approval/decided') {
|
||||
if (!pending.has(event.data.id)) fail(`approval/decided has no matching approval/asked for id ${JSON.stringify(event.data.id)}`)
|
||||
if (trace.openTurn === null) fail('approval/decided appended outside any open turn')
|
||||
if (!trace.pending.has(event.data.id)) fail(`approval/decided has no matching approval/asked for id ${JSON.stringify(event.data.id)}`)
|
||||
if (!APPROVAL_OUTCOMES.includes(event.data.outcome)) {
|
||||
fail(`approval/decided carries unknown outcome ${JSON.stringify(event.data.outcome)}`)
|
||||
}
|
||||
@@ -52,28 +59,39 @@ function applyApprovalTransition(pending: Set<ApprovalRequestId>, transition: Ap
|
||||
// Event owners keep precommit staging local so their vocabularies never move into a central helper.
|
||||
/* jscpd:ignore-start */
|
||||
const install: InvariantInstaller = Object.assign((ctx: Context, fail: InvariantFailure) => {
|
||||
const traces = new WeakMap<Session, Set<ApprovalRequestId>>()
|
||||
const traces = new WeakMap<Session, ApprovalTrace>()
|
||||
const staged = new WeakMap<SessionEvent, { session: Session; transition: ApprovalTransition }>()
|
||||
const seed = (session: Session): Set<ApprovalRequestId> => {
|
||||
const pending = new Set<ApprovalRequestId>()
|
||||
traces.set(session, pending)
|
||||
const seed = (session: Session): ApprovalTrace => {
|
||||
const trace: ApprovalTrace = { openTurn: null, pending: new Set() }
|
||||
traces.set(session, trace)
|
||||
for (const event of session.events) {
|
||||
const transition = validateApprovalEvent(pending, event, fail)
|
||||
if (transition !== undefined) applyApprovalTransition(pending, transition)
|
||||
if (event.type === 'turn/start') trace.openTurn = event.data.turn
|
||||
else if (event.type === 'turn/end') trace.openTurn = null
|
||||
const transition = validateApprovalEvent(trace, event, fail)
|
||||
if (transition !== undefined) applyApprovalTransition(trace.pending, transition)
|
||||
}
|
||||
return pending
|
||||
return trace
|
||||
}
|
||||
const traceFor = (session: Session): Set<ApprovalRequestId> => traces.get(session) ?? seed(session)
|
||||
const traceFor = (session: Session): ApprovalTrace => traces.get(session) ?? seed(session)
|
||||
|
||||
for (const session of ctx.sessions.list()) seed(session)
|
||||
ctx.on('session/created', (session) => { seed(session) }, { global: true })
|
||||
ctx.on('session/event', (session, event) => {
|
||||
const trace = traceFor(session)
|
||||
if (event.type === 'turn/start') {
|
||||
trace.openTurn = event.data.turn
|
||||
return
|
||||
}
|
||||
if (event.type === 'turn/end') {
|
||||
trace.openTurn = null
|
||||
return
|
||||
}
|
||||
if (event.type !== 'approval/asked' && event.type !== 'approval/decided') return
|
||||
const candidate = staged.get(event)
|
||||
/* v8 ignore next -- internal/dispatch stages every package-owned pair event */
|
||||
if (candidate === undefined || candidate.session !== session) return fail('approval audit event published without pre-commit validation')
|
||||
staged.delete(event)
|
||||
applyApprovalTransition(traceFor(session), candidate.transition)
|
||||
applyApprovalTransition(trace.pending, candidate.transition)
|
||||
}, { global: true })
|
||||
ctx.on('internal/dispatch', (_mode, eventName, args) => {
|
||||
if (eventName !== 'session/event') return
|
||||
|
||||
@@ -13,10 +13,15 @@ async function setup(): Promise<Context> {
|
||||
return ctx
|
||||
}
|
||||
|
||||
function startTurn(session: Session): void {
|
||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||
}
|
||||
|
||||
describe('approval invariants', () => {
|
||||
it('accepts paired audit events and closed policy values', async () => {
|
||||
const ctx = await setup()
|
||||
const session = ctx.sessions.create()
|
||||
startTurn(session)
|
||||
const id = ApprovalRequestId('ask-1')
|
||||
session.append('approval/asked', { id, toolName: 'bash' })
|
||||
session.append('approval/decided', { id, outcome: 'allowed-once' })
|
||||
@@ -47,14 +52,43 @@ describe('approval invariants', () => {
|
||||
type: 'approval/decided', seq: 1, time: 1, data: { id, outcome: 'rejected' as const },
|
||||
} as const
|
||||
expect(() => {
|
||||
ctx.emit('session/event', session, {
|
||||
type: 'turn/start', seq: 0, time: 0,
|
||||
data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
|
||||
})
|
||||
ctx.emit('session/event', session, asked)
|
||||
ctx.emit('session/event', session, decided)
|
||||
}).not.toThrow()
|
||||
})
|
||||
|
||||
it('rejects audit events outside any open turn', async () => {
|
||||
const ctx = await setup()
|
||||
const session = ctx.sessions.create()
|
||||
expect(() => session.append('approval/asked', {
|
||||
id: ApprovalRequestId('ask-1'), toolName: 'bash',
|
||||
})).toThrow(/outside any open turn/)
|
||||
expect(() => session.append('approval/decided', {
|
||||
id: ApprovalRequestId('ask-1'), outcome: 'rejected',
|
||||
})).toThrow(/outside any open turn/)
|
||||
})
|
||||
|
||||
it('rejects an unenclosed audit event when replaying an existing session', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SessionStore)
|
||||
const session = ctx.sessions.create()
|
||||
startTurn(session)
|
||||
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
|
||||
session.append('approval/asked', {
|
||||
id: ApprovalRequestId('ask-replay'), toolName: 'bash',
|
||||
})
|
||||
await ctx.plugin(InvariantService)
|
||||
await expect(ctx.plugin(ApprovalInvariant).then(() => undefined)).rejects.toThrow(/outside any open turn/)
|
||||
})
|
||||
|
||||
it('rejects malformed and unpaired audit events', async () => {
|
||||
const ctx = await setup()
|
||||
const session = ctx.sessions.create()
|
||||
startTurn(session)
|
||||
const id = ApprovalRequestId('ask-1')
|
||||
expect(() => session.append('approval/asked', { id, toolName: '' }))
|
||||
.toThrow(/toolName must be non-empty/)
|
||||
|
||||
Reference in New Issue
Block a user