1
/** Host timer over stored tasks; Session activation is a delivery operation. */2
import { createUserMessage } from '@deepseek-ai/dsh-llm'3
import type { ContextFormed } from '@deepseek-ai/dsh-llm'4
declare module '@deepseek-ai/dsh-llm' {5
interface MessageSourceMap {6
'schedule': { kind: 'schedule' } & ContextFormed7
}8
}9
import type { Context } from '@deepseek-ai/cordis'10
import type {} from '@deepseek-ai/dsh-api-session-controller'11
import { isRecurringScheduleRecord, renderReminderFraming, renderRecurringReminderBatchFraming, resolveRecurringOccurrence } from './domain.ts'12
import type { DeliveryRetentionBounds, RecurringScheduleRecord } from './types.ts'13
import type { ScheduleTask } from './storage.ts'14
import { appendDelivery } from './delivery-history.ts'16
/** Largest delay Node timers represent without clamping. */17
export const MAX_TIMER_DELAY_MS = 2_147_483_64719
/** Owns at most one timer across recomputations; delivery and management share the serialized operation. */20
export class ScheduleRuntime {21
private timer: ReturnType<typeof setTimeout> | undefined22
private running: Promise<void> | undefined23
private stopping = false24
private requested = false26
/**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
) {}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) return46
this.requested = true47
this.clearTimer()48
if (this.running !== undefined) return49
let run: Promise<void>50
try {51
run = this.ctx.agents.withoutInitiator(async () => {52
while (this.requested && !this.stopping) {53
this.requested = false54
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 = false60
this.ctx.logger.warn(`schedule: dispatch stopped: ${String(error)}`)61
return62
}63
this.running = run64
void run.catch((error: unknown) => {65
this.ctx.logger.warn(`schedule: dispatch stopped: ${String(error)}`)66
}).finally(() => {67
this.running = undefined68
if (this.requested && !this.stopping) this.requestDrive()69
})70
}72
/** Stop the timer and drain an accepted delivery before storage closes. */73
async dispose(): Promise<void> {74
this.stopping = true75
this.clearTimer()76
// requestDrive reports execution failures; teardown only waits for quiescence.77
await this.running?.catch(() => undefined)78
}80
private clearTimer(): void {81
if (this.timer !== undefined) clearTimeout(this.timer)82
this.timer = undefined83
}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) return93
if (handled.has(task.record.id)) continue94
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 = group99
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.error103
// oxlint-disable-next-line typescript/no-unnecessary-condition -- Disposal can run while Session restoration is awaited.104
if (this.stopping) return105
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) continue109
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) return155
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
}