返回源码地图

packages/schedule/schedule/src/runtime.ts

main snapshot · da00f7f5358f · 正文引用章节 13 / 13 / 17;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/** Host timer over stored tasks; Session activation is a delivery operation. */
2import { createUserMessage } from '@deepseek-ai/dsh-llm'
3import type { ContextFormed } from '@deepseek-ai/dsh-llm'
4declare module '@deepseek-ai/dsh-llm' {
5 interface MessageSourceMap {
6 'schedule': { kind: 'schedule' } & ContextFormed
7 }
8}
9import type { Context } from '@deepseek-ai/cordis'
10import type {} from '@deepseek-ai/dsh-api-session-controller'
11import { isRecurringScheduleRecord, renderReminderFraming, renderRecurringReminderBatchFraming, resolveRecurringOccurrence } from './domain.ts'
12import type { DeliveryRetentionBounds, RecurringScheduleRecord } from './types.ts'
13import type { ScheduleTask } from './storage.ts'
14import { appendDelivery } from './delivery-history.ts'
15
16/** Largest delay Node timers represent without clamping. */
17export const MAX_TIMER_DELAY_MS = 2_147_483_647
18
19/** Owns at most one timer across recomputations; delivery and management share the serialized operation. */
20export class ScheduleRuntime {
21 private timer: ReturnType<typeof setTimeout> | undefined
22 private running: Promise<void> | undefined
23 private stopping = false
24 private requested = false
25
26 /**
27 * @param ctx - Host services used to resume and enqueue.
28 * @param tasks - Current durable tasks.
29 * @param transact - Serialize delivery against management writes.
30 * @param commit - Persist task status, target, receipt, and history together after durable inbox delivery.
31 */
32 constructor(
33 private readonly ctx: Context,
34 private readonly tasks: () => readonly ScheduleTask[],
35 private readonly transact: (work: () => Promise<void>) => Promise<void>,
36 private readonly commit: (task: ScheduleTask) => Promise<void>,
37 private readonly retention: DeliveryRetentionBounds,
38 ) {}
39
40 /**
41 * Recompute the nearest obligation after startup or a durable change.
42 * Dispatch failures are logged; refused admission does not retry automatically.
43 */
44 requestDrive(): void {
45 if (this.stopping) return
46 this.requested = true
47 this.clearTimer()
48 if (this.running !== undefined) return
49 let run: Promise<void>
50 try {
51 run = this.ctx.agents.withoutInitiator(async () => {
52 while (this.requested && !this.stopping) {
53 this.requested = false
54 await this.transact(async () => { await this.drive() })
55 }
56 })
57 } catch (error: unknown) {
58 // Ancestor unloading closes initiator admission before runtime cleanup runs.
59 this.requested = false
60 this.ctx.logger.warn(`schedule: dispatch stopped: ${String(error)}`)
61 return
62 }
63 this.running = run
64 void run.catch((error: unknown) => {
65 this.ctx.logger.warn(`schedule: dispatch stopped: ${String(error)}`)
66 }).finally(() => {
67 this.running = undefined
68 if (this.requested && !this.stopping) this.requestDrive()
69 })
70 }
71
72 /** Stop the timer and drain an accepted delivery before storage closes. */
73 async dispose(): Promise<void> {
74 this.stopping = true
75 this.clearTimer()
76 // requestDrive reports execution failures; teardown only waits for quiescence.
77 await this.running?.catch(() => undefined)
78 }
79
80 private clearTimer(): void {
81 if (this.timer !== undefined) clearTimeout(this.timer)
82 this.timer = undefined
83 }
84
85 private async drive(): Promise<void> {
86 this.clearTimer()
87 const failed = new Set<string>()
88 const handled = new Set<string>()
89 const scanNow = Date.now()
90 const due = this.tasks().filter(task => task.status === 'active' && Date.parse(task.record.scheduledAt) <= scanNow)
91 for (const task of due) {
92 if (this.stopping) return
93 if (handled.has(task.record.id)) continue
94 const group = isRecurringScheduleRecord(task.record)
95 ? due.filter(candidate => candidate.sessionId === task.sessionId && isRecurringScheduleRecord(candidate.record))
96 : [task]
97 for (const member of group) handled.add(member.record.id)
98 let admitted = group
99 const committed = new Set<ScheduleTask['record']['id']>()
100 try {
101 const resolved = await this.ctx.sessionController.resolveAgent(task.sessionId)
102 if ('error' in resolved) throw resolved.error
103 // oxlint-disable-next-line typescript/no-unnecessary-condition -- Disposal can run while Session restoration is awaited.
104 if (this.stopping) return
105 const now = Date.now()
106 // Session restoration can span a wall-clock rollback; future members keep their timer obligation.
107 admitted = group.filter(member => Date.parse(member.record.scheduledAt) <= now)
108 if (admitted.length === 0) continue
109 const recurring = admitted.filter((member): member is ScheduleTask & { record: RecurringScheduleRecord } =>
110 isRecurringScheduleRecord(member.record))
111 const occurrences = recurring.map(member => ({
112 task: member, occurrence: resolveRecurringOccurrence(member.record, now),
113 }))
114 const text = isRecurringScheduleRecord(task.record)
115 ? renderRecurringReminderBatchFraming(occurrences.map(({ task: member, occurrence }) => ({
116 record: member.record, occurrenceAt: occurrence.occurrenceAt,
117 })))
118 : renderReminderFraming(task.record)
119 const message = createUserMessage({
120 content: [{ type: 'text', text }], source: { kind: 'schedule' },
121 })
122 // followup synchronously appends the inbox splice before flush observes the Session.
123 resolved.agent.followup(message)
124 const flushed = await this.ctx.sessions.flush(resolved.agent.session)
125 if (!flushed) throw new Error('Session persistence did not acknowledge the reminder')
126 const deliveredAt = new Date(Date.now()).toISOString()
127 if (!isRecurringScheduleRecord(task.record)) {
128 await this.commit({
129 ...task, status: 'inactive',
130 ...appendDelivery(task, { scheduledAt: task.record.scheduledAt, deliveredAt, messageId: message.id }, this.retention),
131 })
132 committed.add(task.record.id)
133 }
134 for (const { task: member, occurrence } of occurrences) {
135 await this.commit({
136 ...member,
137 record: { ...member.record, scheduledAt: occurrence.nextScheduledAt ?? occurrence.occurrenceAt },
138 status: occurrence.nextScheduledAt === undefined ? 'inactive' : 'active',
139 ...appendDelivery(member, { scheduledAt: occurrence.occurrenceAt, deliveredAt, messageId: message.id }, this.retention),
140 })
141 committed.add(member.record.id)
142 }
143 } catch (error: unknown) {
144 // Successful commits and targets made future by clock rollback keep their timer obligation.
145 const pending = admitted.filter(member => !committed.has(member.record.id))
146 const failedAt = Date.now()
147 for (const member of pending) {
148 if (Date.parse(member.record.scheduledAt) <= failedAt) failed.add(member.record.id)
149 }
150 const ids = pending.map(member => member.record.id)
151 this.ctx.logger.warn(`schedule: reminders ${JSON.stringify(ids)} were not acknowledged: ${String(error)}`)
152 }
153 }
154 if (this.stopping) return
155 const next = this.tasks().filter(task => task.status === 'active' && !failed.has(task.record.id))
156 .reduce<number | undefined>((at, task) => {
157 const target = Date.parse(task.record.scheduledAt)
158 return at === undefined ? target : Math.min(at, target)
159 }, undefined)
160 if (next !== undefined) {
161 this.timer = setTimeout(() => { this.timer = undefined; this.requestDrive() },
162 Math.max(0, Math.min(next - Date.now(), MAX_TIMER_DELAY_MS)))
163 this.timer.unref()
164 }
165 }
166}