refactor(schedule): make absolute times explicit

This commit is contained in:
Tianyi Cui
2026-08-09 16:30:11 +08:00
parent 3d6498e91b
commit b7ec8429a9
109 changed files with 1248 additions and 3219 deletions

View File

@@ -2,7 +2,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { Context } from 'cordis'
import Loader from '@cordisjs/plugin-loader'
import { createUserMessage, CallId, LlmAdapter } from '@deepseek-ai/dsh-llm'
import type { GenerateOptions, StreamChunk, UserMessage } from '@deepseek-ai/dsh-llm'
import type { GenerateOptions, StreamChunk } from '@deepseek-ai/dsh-llm'
import { Session, SessionId, type SessionEvent } from '@deepseek-ai/dsh-session'
import AgentRegistry, { agentEvents, Inbox, type Agent } from '@deepseek-ai/dsh-agent'
import { defineContentToolFixture } from '@deepseek-ai/dsh-tools'
@@ -53,11 +53,13 @@ function sessionAgent(session: Session, id = 'agent'): Agent {
}
}
function openMessageTurn(session: Session, turn: number): void {
function openMessageTurn(session: Session, turn: number, clientTimeZone?: string): void {
session.append('turn/start', { turn })
session.append('user/message', createUserMessage({
content: [{ type: 'text', text: `turn ${turn}` }],
source: { kind: 'user' },
source: clientTimeZone === undefined
? { kind: 'user' }
: { kind: 'user', rpcId: `turn-${String(turn)}`, clientTimeZone } as never,
}), { surfaceOp: 'append' })
}
@@ -79,35 +81,24 @@ async function fire(
turn: number,
step: number,
signal: AbortSignal = SIGNAL,
messages: UserMessage[] = [],
): Promise<void> {
const fallback = messages.length === 0
? createUserMessage({
content: [],
source: { kind: 'plugin', plugin: 'time-context-test-proposal' },
})
: undefined
const proposal = fallback === undefined ? messages : [fallback]
const proposed = createUserMessage({
content: [{ type: 'text', text: 'request proposal' }],
source: { kind: 'plugin', plugin: 'time-context-test' },
})
const decision = await agentEvents(ctx, agent).waterfall(
'agent/pre-step',
{ messages: proposal, turn, step, signal },
() => Promise.resolve({ kind: 'enter' as const, messages: proposal }),
{ messages: [proposed], turn, step, signal },
() => Promise.resolve({ kind: 'enter' as const, messages: [proposed] }),
)
if (decision.kind === 'enter') {
for (const message of decision.messages) {
if (message.id === fallback?.id) continue
if (message === proposed) continue
agent.session.append('user/message', message, { surfaceOp: 'append' })
}
}
}
function rpcMessage(text: string, clientTimeZone: string): UserMessage {
return createUserMessage({
content: [{ type: 'text', text }],
source: { kind: 'user', rpcId: `rpc-${text}`, clientTimeZone } as never,
})
}
function textResponse(text: string): StreamChunk[] {
return [
{ type: 'block-start', index: 0, blockType: 'text' },
@@ -161,84 +152,22 @@ function requestText(request: GenerateOptions): string {
}
describe('durable step context', () => {
it('uses the immutable Session zone and the current request message zone', async () => {
const { ctx } = await mount()
const id = SessionId('session-zone')
const session = Session.create(id, [], {
version: 0,
id,
createdAt: BASE,
timeZone: 'Asia/Shanghai',
})
session.append('turn/start', { turn: 1 })
const agent = sessionAgent(session)
await fire(ctx, agent, 1, 1, SIGNAL, [
rpcMessage('local request', 'Asia/Shanghai'),
])
expect(contextTexts(session)[0]).toContain(
'2026-07-14T08:00:00+08:00[Asia/Shanghai]',
)
expect(contextTexts(session)[0]).toContain('Session time zone: Asia/Shanghai.')
expect(contextTexts(session)[0]).toContain('Client time zone for this request: Asia/Shanghai.')
const reading = session.events.at(-1)
expect(reading).toMatchObject({
type: 'user/message',
data: {
source: { kind: 'plugin', plugin: 'time-context' },
},
})
await fire(ctx, agent, 1, 2, SIGNAL, [
rpcMessage('same zone again', 'Asia/Shanghai'),
])
expect(contextTexts(session)).toHaveLength(2)
})
it('reports sorted mixed zones from the current request chain without changing the Session zone', async () => {
const { ctx } = await mount()
const id = SessionId('mixed-zone')
const session = Session.create(id, [], {
version: 0,
id,
createdAt: BASE,
timeZone: 'Asia/Shanghai',
})
session.append('turn/start', { turn: 1 })
session.append('user/message', rpcMessage('first tab', 'Asia/Shanghai'), {
surfaceOp: 'append',
})
await fire(ctx, sessionAgent(session), 1, 1, SIGNAL, [
rpcMessage('second tab', 'America/New_York'),
])
expect(contextTexts(session)[0]).toContain('Session time zone: Asia/Shanghai.')
expect(contextTexts(session)[0]).toContain(
'Client time zone for this request: mixed ["America/New_York","Asia/Shanghai"].',
)
})
it('records turn, step, zoned time, and the preceding model-visible message baseline', async () => {
const { ctx } = await mount({ timeZone: 'Asia/Shanghai' })
const session = Session.create(SessionId('first'))
openMessageTurn(session, 1)
openMessageTurn(session, 1, 'Asia/Shanghai')
vi.setSystemTime(BASE + 90_061_000)
await fire(ctx, sessionAgent(session), 1, 1)
expect(contextTexts(session)).toEqual([
'Time sampled while preparing turn 1, step 1: 2026-07-15T09:01:01+08:00[Asia/Shanghai]\n'
+ 'Session time zone: unavailable.\n'
+ 'Client time zone for this request: missing.\n'
+ 'Browser time zone for this request: Asia/Shanghai. Interpret otherwise-unqualified dates and times in this zone.\n'
+ 'Elapsed since the preceding model-visible message: 1d 1h 1m 1s.',
])
const event = session.events.at(-1)
expect(event?.type).toBe('user/message')
if (event?.type !== 'user/message') throw new Error('missing time context')
const text = event.data.content.find(block => block.type === 'text')?.text
if (text === undefined) throw new Error('missing time-context text')
// The reading is a `snapshot`-form context: one named contribution whose
// text is exactly what the model read, so a consumer attributes it without
// re-splitting prose.
@@ -248,7 +177,9 @@ describe('durable step context', () => {
form: 'snapshot',
sections: [{
name: 'time-context',
text,
text: 'Time sampled while preparing turn 1, step 1: 2026-07-15T09:01:01+08:00[Asia/Shanghai]\n'
+ 'Browser time zone for this request: Asia/Shanghai. Interpret otherwise-unqualified dates and times in this zone.\n'
+ 'Elapsed since the preceding model-visible message: 1d 1h 1m 1s.',
}],
})
expect(event.surfaceOp).toBe('append')
@@ -281,12 +212,40 @@ describe('durable step context', () => {
expect(contextTexts(session)[1]).toBe(
'Time sampled while preparing turn 3, step 2: 2026-07-14T00:01:01+00:00[UTC]\n'
+ 'Session time zone: unavailable.\n'
+ 'Client time zone for this request: missing.\n'
+ 'Browser time zone for this request: unavailable. Ask the user to clarify otherwise-unqualified dates and times.\n'
+ 'Elapsed since the preceding step context: 1m 1s.',
)
})
it('formats in one browser zone and falls back when steering supplies mixed zones', async () => {
const { ctx } = await mount({ timeZone: 'UTC' })
const resolved = Session.create(SessionId('browser-zone-resolved'))
openMessageTurn(resolved, 1, 'America/New_York')
await fire(ctx, sessionAgent(resolved), 1, 1)
expect(contextTexts(resolved)[0]).toContain(
'2026-07-13T20:00:00-04:00[America/New_York]\n'
+ 'Browser time zone for this request: America/New_York. '
+ 'Interpret otherwise-unqualified dates and times in this zone.',
)
const mixed = Session.create(SessionId('browser-zone-mixed'))
openMessageTurn(mixed, 1, 'Asia/Shanghai')
mixed.append('user/message', createUserMessage({
content: [{ type: 'text', text: 'steering from another browser' }],
source: {
kind: 'user',
rpcId: 'mixed-steer',
clientTimeZone: 'America/New_York',
} as never,
}), { surfaceOp: 'append' })
await fire(ctx, sessionAgent(mixed), 1, 1)
expect(contextTexts(mixed)[0]).toContain(
'2026-07-14T00:00:00+00:00[UTC]\n'
+ 'Browser time zone for this request: mixed ["America/New_York","Asia/Shanghai"]. '
+ 'Ask the user to clarify otherwise-unqualified dates and times.',
)
})
it('reports an unavailable later-step baseline at the matching turn boundary', async () => {
const { ctx } = await mount()
const session = Session.create(SessionId('later-step-boundary'))
@@ -427,20 +386,6 @@ describe('configuration and lifecycle', () => {
await expect(unresolved.plugin(timeContext, {})).rejects.toThrow(/failed to resolve the system time zone/)
})
it('fails loud when a persisted Session names an invalid zone', async () => {
const { ctx } = await mount()
const id = SessionId('invalid-session-zone')
const session = Session.create(id, [], {
version: 0,
id,
createdAt: BASE,
timeZone: 'Not/A_Real_Zone',
})
openMessageTurn(session, 1)
await expect(fire(ctx, sessionAgent(session), 1, 1)).rejects.toThrow(/invalid Session time zone/)
})
it('rejects invalid refresh intervals at plugin load with one diagnostic', async () => {
const invalid = [-1, 0.5, Number.MAX_SAFE_INTEGER + 1, Number.POSITIVE_INFINITY, Number.NaN]
for (const refreshIntervalMs of invalid) {
@@ -462,26 +407,13 @@ describe('configuration and lifecycle', () => {
expect(contextTexts(session)).toHaveLength(1)
})
it('lets an already-stopped direct registration delegate without contributing', async () => {
const ctx = new Context()
await ctx.plugin(AgentRegistry)
const stop = timeContext.apply(ctx, {})
stop()
const session = Session.create(SessionId('stopped-direct-registration'))
openMessageTurn(session, 1)
await fire(ctx, sessionAgent(session), 1, 1)
expect(contextTexts(session)).toEqual([])
})
})
describe('real agent-loop request history', () => {
it.each([
['throws', 0],
['cancels', 0],
] as const)('does not persist context when a downstream pre-step listener %s', async (mode, expectedContexts) => {
['throws'],
['cancels'],
] as const)('does not commit a preparation reading when a downstream pre-step listener %s', async (mode) => {
const adapter = new ScriptedAdapter([textResponse('unused')])
const ctx = await loopHarness(adapter)
ctx.on('agent/pre-step', ({ agent: subject }, next) => {
@@ -494,170 +426,13 @@ describe('real agent-loop request history', () => {
agent.followup(createUserMessage({ content: [{ type: 'text', text: 'start' }], source: { kind: 'user' } }))
await agent.whenIdle()
expect(contextTexts(agent.session)).toHaveLength(expectedContexts)
expect(contextTexts(agent.session)).toHaveLength(0)
expect(adapter.requests).toHaveLength(0)
expect(agent.session.events.some(event => event.type === 'step/start')).toBe(false)
await ctx.fiber.dispose()
})
it('leaves steering that arrives after claim for the next step and derives fresh context', async () => {
const adapter = new ScriptedAdapter([textResponse('first'), textResponse('second')])
const ctx = await loopHarness(adapter)
const entered = Promise.withResolvers<undefined>()
const release = Promise.withResolvers<undefined>()
let blocked = true
ctx.on('system-prompt/assemble', async (_assembly, context, next) => {
if (blocked && context.agent !== undefined) {
entered.resolve(undefined)
await release.promise
}
return next()
})
const agent = ctx.agentLoop.create(SessionId('late-steering'), { provider: 'mock', model: 'mock' })
agent.followup(rpcMessage('start in Shanghai', 'Asia/Shanghai'))
await entered.promise
agent.steer(rpcMessage('switch to New York', 'America/New_York'))
blocked = false
release.resolve(undefined)
await agent.whenIdle()
expect(adapter.requests).toHaveLength(2)
expect(agent.inbox.hasPending).toBe(false)
expect(requestText(adapter.requests[0]!)).toContain('start in Shanghai')
expect(requestText(adapter.requests[0]!)).not.toContain('switch to New York')
expect(requestText(adapter.requests[0]!)).toContain('Client time zone for this request: Asia/Shanghai.')
expect(requestText(adapter.requests[1]!)).toContain('switch to New York')
expect(requestText(adapter.requests[1]!)).toContain(
'Client time zone for this request: mixed ["America/New_York","Asia/Shanghai"].',
)
expect(contextTexts(agent.session)).toHaveLength(2)
await ctx.fiber.dispose()
})
it('does not let time context create an initial step after downstream suppression', async () => {
const adapter = new ScriptedAdapter([textResponse('unused')])
const ctx = await loopHarness(adapter)
ctx.on('agent/pre-step', async (_payload, next) => {
const decision = await next()
return decision.kind === 'reject' ? decision : { kind: 'enter', messages: [] }
})
const agent = ctx.agentLoop.create(SessionId('suppressed-preparation'), {
provider: 'mock',
model: 'mock',
})
agent.followup(rpcMessage('suppress this prompt', 'Asia/Shanghai'))
await agent.whenIdle()
expect(adapter.requests).toEqual([])
expect(agent.session.events.some(event => event.type === 'step/start')).toBe(false)
expect(contextTexts(agent.session)).toEqual([])
expect(agent.inbox.hasPending).toBe(false)
await ctx.fiber.dispose()
})
it('does not revive an empty continuation after a completed step', async () => {
const adapter = new ScriptedAdapter([textResponse('done')])
const ctx = await loopHarness(adapter)
ctx.on('agent/turn-stopping', ({ agent: subject }) => {
subject.inject(createUserMessage({
content: [{ type: 'text', text: 'pending context' }],
source: { kind: 'plugin', plugin: 'test' },
}))
})
ctx.on('agent/pre-step', async ({ step }, next) => {
const decision = await next()
return step === 1 || decision.kind === 'reject'
? decision
: { kind: 'enter', messages: [] }
})
const agent = ctx.agentLoop.create(SessionId('empty-completed-continuation'), {
provider: 'mock',
model: 'mock',
})
agent.followup(rpcMessage('finish once', 'Asia/Shanghai'))
await agent.whenIdle()
expect(adapter.requests).toHaveLength(1)
expect(agent.session.events.filter(event => event.type === 'step/start')).toHaveLength(1)
expect(contextTexts(agent.session)).toHaveLength(1)
expect(agent.inbox.hasPending).toBe(false)
await ctx.fiber.dispose()
})
it('preserves post-claim steering without persisting failed-turn context', async () => {
const adapter = new ScriptedAdapter([textResponse('resumed')])
const ctx = await loopHarness(adapter)
const entered = Promise.withResolvers<undefined>()
const release = Promise.withResolvers<undefined>()
let blocked = true
ctx.on('system-prompt/assemble', async (_assembly, context, next) => {
if (blocked && context.agent !== undefined) {
entered.resolve(undefined)
await release.promise
}
return next()
})
const agent = ctx.agentLoop.create(SessionId('cancelled-assembly'), { provider: 'mock', model: 'mock' })
const steering = rpcMessage('preserve this steering', 'America/New_York')
agent.followup(rpcMessage('start', 'Asia/Shanghai'))
await entered.promise
agent.steer(steering)
agent.cancel({ kind: 'user' }, { keepInbox: true })
blocked = false
release.resolve(undefined)
await agent.whenIdle()
expect(agent.session.events.some(event => event.type === 'step/start')).toBe(false)
expect(contextTexts(agent.session)).toHaveLength(0)
expect(agent.inbox.nextStep).toEqual([steering])
expect(agent.inbox.nextStep.some(message =>
message.source.kind === 'plugin' && message.source.plugin === 'time-context')).toBe(false)
agent.followup(rpcMessage('wake', 'America/New_York'))
await agent.whenIdle()
expect(adapter.requests).toHaveLength(1)
expect(requestText(adapter.requests[0]!)).toContain('preserve this steering')
expect(requestText(adapter.requests[0]!)).toContain('Time sampled while preparing turn 2, step 1:')
await ctx.fiber.dispose()
})
it('does not contribute after its disposer wins an in-flight pre-step', async () => {
const adapter = new ScriptedAdapter([textResponse('done')])
const ctx = new Context()
await mountAgentLoopTestDependencies(ctx)
await ctx.plugin(AgentLoop, { agents: [] })
const stopTimeContext = timeContext.apply(ctx, {})
ctx.llm.registerAdapter(['mock'], adapter)
const entered = Promise.withResolvers<undefined>()
const release = Promise.withResolvers<undefined>()
ctx.on('agent/pre-step', async (_payload, next) => {
entered.resolve(undefined)
await release.promise
return next()
})
const agent = ctx.agentLoop.create(SessionId('dispose-inflight-pre-step'), {
provider: 'mock',
model: 'mock',
})
agent.followup(rpcMessage('continue without disposed context', 'Asia/Shanghai'))
await entered.promise
stopTimeContext()
release.resolve(undefined)
await agent.whenIdle()
expect(adapter.requests).toHaveLength(1)
expect(requestText(adapter.requests[0]!)).not.toContain('Time sampled while preparing')
expect(contextTexts(agent.session)).toEqual([])
expect(agent.inbox.nextStep).toEqual([])
await ctx.fiber.dispose()
})
it('does not add a reading to an empty tool continuation and leaves system headers unchanged', async () => {
it('persists one ordered context per request, accumulates readings, and leaves system headers unchanged', async () => {
const adapter = new ScriptedAdapter([toolCallResponse(), textResponse('done')])
const ctx = await loopHarness(adapter)
ctx.tools.register(defineContentToolFixture({
@@ -678,9 +453,11 @@ describe('real agent-loop request history', () => {
const contexts = agent.session.events.filter(
(event): event is SessionEvent<'user/message'> => event.type === 'user/message' && event.data.source.kind === 'plugin')
const starts = agent.session.events.filter(event => event.type === 'step/start')
expect(contexts).toHaveLength(1)
expect(contexts).toHaveLength(adapter.requests.length)
expect(starts).toHaveLength(adapter.requests.length)
expect(contexts[0]!.seq).toBeGreaterThan(starts[0]!.seq)
for (let index = 0; index < contexts.length; index += 1) {
expect(contexts[index]!.seq).toBeGreaterThan(starts[index]!.seq)
}
expect(contexts.every(event => event.data.source.kind === 'plugin'
&& event.data.source.plugin === 'time-context'
&& event.surfaceOp === 'append')).toBe(true)
@@ -691,7 +468,8 @@ describe('real agent-loop request history', () => {
expect(firstRequestText).toContain('Elapsed since the preceding model-visible message: unavailable.')
expect(firstRequestText).not.toContain('Time sampled while preparing turn 1, step 2:')
expect(secondRequestText).toContain('Time sampled while preparing turn 1, step 1:')
expect(secondRequestText).not.toContain('Time sampled while preparing turn 1, step 2:')
expect(secondRequestText).toContain('Time sampled while preparing turn 1, step 2:')
expect(secondRequestText).toContain('Elapsed since the preceding step context: 1m 1s.')
for (const request of adapter.requests) expect(request.system).not.toContain('Time sampled while preparing')
const headers = agent.session.events.filter(event => event.type === 'request/header')