fix(tui): correlate referenced prompt snapshots

This commit is contained in:
_Kerman
2026-07-27 19:06:44 +08:00
parent 9152858fd7
commit 613eaba80c
2 changed files with 39 additions and 22 deletions

View File

@@ -2910,38 +2910,50 @@ export function createTuiChat(
// Idle: the snapshot rides the prompt's admission transaction so a // Idle: the snapshot rides the prompt's admission transaction so a
// blocking hook discards both together. // blocking hook discards both together.
let cleanedUp = false let cleanedUp = false
let acceptedId: AgentMessageId | undefined
let acceptedContent: ContentBlock[] | undefined
const enqueued = new Map<AgentMessageId, ContentBlock[]>()
const discarded = new Set<AgentMessageId>()
const cleanup = (): void => { const cleanup = (): void => {
// Each trigger detaches both listeners, so a second call needs a // Every completion path detaches all three listeners. Keep this
// future third trigger; kept so adding one cannot double-release. // idempotent so later cleanup paths cannot double-release them.
/* v8 ignore next -- unreachable idempotence guard, see above */ /* v8 ignore next -- unreachable idempotence guard, see above */
if (cleanedUp) return if (cleanedUp) return
cleanedUp = true cleanedUp = true
detachEnqueue()
detachSubmit() detachSubmit()
detachDiscard() detachDiscard()
} }
// send() snapshots input before publishing it, and publishes enqueue
// before returning its id. Capture that snapshot by id so admission can
// use exact reference identity without depending on caller-owned input.
const detachEnqueue = ctx.on('agent/inbox/enqueue', (subject, message) => {
if (subject === agent) enqueued.set(message.id, message.content)
})
// Prepended so this wrapper is outermost: it observes the admission // Prepended so this wrapper is outermost: it observes the admission
// whether a downstream hook allows or blocks, and detaches either way. // whether a downstream hook allows or blocks, and detaches either way.
const detachSubmit = ctx.on('agent/prompt-submit', async (subject, submitted, _source, _signal, next) => { const detachSubmit = ctx.on('agent/prompt-submit', async (subject, submitted, _source, _signal, next) => {
if (subject !== agent || submitted !== content) return next() if (subject !== agent || submitted !== acceptedContent) return next()
cleanup() cleanup()
const decision = await next() const decision = await next()
if (decision.kind !== 'allow') return decision if (decision.kind !== 'allow') return decision
return { ...decision, additionalContexts: [...decision.additionalContexts ?? [], attachedContext] } return { ...decision, additionalContexts: [...decision.additionalContexts ?? [], attachedContext] }
}, { prepend: true }) }, { prepend: true })
// Installed BEFORE followup(): admission runs synchronously inside it on // Installed before followup(): an enqueue listener can synchronously
// the common path, and a listener registered after cleanup() already ran // cancel and discard before followup() returns its id.
// would never be released. Match on the `content` reference, not the
// returned id: an enqueue listener that synchronously cancels emits
// discard before followup() returns to assign the id, and content is the
// same reference send() carries onto the message (mirrors detachSubmit).
const detachDiscard = ctx.on('agent/inbox/discard', (subject, messages) => { const detachDiscard = ctx.on('agent/inbox/discard', (subject, messages) => {
if (subject === agent && messages.some(message => message.content === content)) cleanup() if (subject !== agent) return
for (const message of messages) discarded.add(message.id)
if (acceptedId !== undefined && discarded.has(acceptedId)) cleanup()
}) })
// followup() accepts any typed input and contains listener failures; // followup() accepts any typed input and contains listener failures;
// this guards a future synchronous throw so the wrapper cannot leak. // this guards a future synchronous throw so the wrapper cannot leak.
/* v8 ignore start -- future-proofing guard, see above */ /* v8 ignore start -- future-proofing guard, see above */
try { try {
agent.followup({ content, source: { kind: 'user' } }) acceptedId = agent.followup({ content, source: { kind: 'user' } })
acceptedContent = enqueued.get(acceptedId) ?? content
detachEnqueue()
if (discarded.has(acceptedId)) cleanup()
} catch (error: unknown) { } catch (error: unknown) {
cleanup() cleanup()
throw error throw error

View File

@@ -2039,16 +2039,21 @@ describe('pi-tui chat lifecycle and transcript', () => {
appendUser(source, 'source background') appendUser(source, 'source background')
}, },
}) })
// Real send() emits agent/inbox/discard when an enqueue listener cancels // Real send() publishes its snapshotted message, then an enqueue listener
// synchronously, before followup() returns to assign the message id. This // may synchronously cancel and discard it before followup() returns the
// stub reproduces that timing: the wrapper must match on content, since id // already-assigned id. This stub reproduces that ordering.
// is not yet observable at discard time. const foreign = { ...result.agent, id: SessionId('foreign') } as unknown as Agent
result.agent.followup = (input) => { result.agent.followup = (input) => {
result.agent.sent.push(input.content) result.agent.sent.push(input.content)
result.ctx.emit('agent/inbox/discard', result.agent, [{ const message = {
id: AgentMessageId('unassigned'), content: input.content, source: input.source, id: AgentMessageId('stub'),
}]) content: structuredClone(input.content),
return AgentMessageId('stub') source: structuredClone(input.source),
}
result.ctx.emit('agent/inbox/enqueue', foreign, message, 'queued')
result.ctx.emit('agent/inbox/enqueue', result.agent, message, 'queued')
result.ctx.emit('agent/inbox/discard', result.agent, [message])
return message.id
} }
result.terminal.send('@sync-source') result.terminal.send('@sync-source')
@@ -2058,9 +2063,9 @@ describe('pi-tui chat lifecycle and transcript', () => {
result.terminal.send('\r') result.terminal.send('\r')
await vi.waitFor(() => { expect(result.agent.sent).toHaveLength(1) }) await vi.waitFor(() => { expect(result.agent.sent).toHaveLength(1) })
// The synchronous discard released both listeners despite the id being // The synchronous discard released the listeners even though followup()
// unassigned: replaying the prompt's admission attaches no stranded // had not returned the id yet: replaying the prompt's admission attaches
// snapshot, and nothing leaks for the TUI lifetime. // no stranded snapshot, and nothing leaks for the TUI lifetime.
const replay = await agentEvents(result.ctx, result.agent).waterfall( const replay = await agentEvents(result.ctx, result.agent).waterfall(
'agent/prompt-submit', result.agent.sent[0]!, { kind: 'user' }, 'agent/prompt-submit', result.agent.sent[0]!, { kind: 'user' },
new AbortController().signal, () => Promise.resolve({ kind: 'allow' as const }), new AbortController().signal, () => Promise.resolve({ kind: 'allow' as const }),