refactor(schedule): bound fixed-rate reminders
This commit is contained in:
@@ -6,16 +6,11 @@
|
||||
import type { Context } from 'cordis'
|
||||
import type { Agent } from '@deepseek-ai/dsh-agent'
|
||||
import { createUserMessage } from '@deepseek-ai/dsh-llm'
|
||||
import type {
|
||||
OneShotScheduleRecord,
|
||||
RecurringScheduleRecord,
|
||||
} from './types.ts'
|
||||
import type { EveryScheduleRecord, OneShotScheduleRecord } from './types.ts'
|
||||
import {
|
||||
foldScheduleEvents,
|
||||
MIN_RECURRING_INTERVAL_SECONDS,
|
||||
renderReminderBatchFraming,
|
||||
renderEveryReminderBatchFraming,
|
||||
renderReminderFraming,
|
||||
resolveCronOccurrence,
|
||||
resolveEveryOccurrence,
|
||||
ScheduleLogError,
|
||||
} from './domain.ts'
|
||||
@@ -26,68 +21,50 @@ import { runScheduleTransaction } from './transaction.ts'
|
||||
/** Largest delay that Node timers represent without clamping. */
|
||||
export const MAX_TIMER_DELAY_MS = 2_147_483_647
|
||||
|
||||
interface RecurringDue {
|
||||
readonly record: RecurringScheduleRecord
|
||||
interface EveryDue {
|
||||
readonly record: EveryScheduleRecord
|
||||
readonly occurrenceAt: string
|
||||
readonly nextScheduledAt?: string
|
||||
}
|
||||
|
||||
type DueDecision =
|
||||
| { readonly kind: 'one-shot'; readonly record: OneShotScheduleRecord }
|
||||
| { readonly kind: 'recurring'; readonly reminders: readonly RecurringDue[]; readonly acceptedAt: string }
|
||||
| { readonly kind: 'every'; readonly reminders: readonly EveryDue[]; readonly acceptedAt: string }
|
||||
| { readonly kind: 'wait'; readonly target?: number }
|
||||
|
||||
/** Select one unblocked one-shot, one complete recurring batch, or the next wake. */
|
||||
/** Select one due one-shot, one complete fixed-rate batch, or the next wake. */
|
||||
function dueDecision(folded: FoldedSchedules, now: number): DueDecision {
|
||||
const indexed = folded.active.map((record, index) => ({ record, index }))
|
||||
const dueOneShots = indexed
|
||||
const byTargetThenCreate = (
|
||||
left: { readonly record: { readonly scheduledAt: string }; readonly index: number },
|
||||
right: { readonly record: { readonly scheduledAt: string }; readonly index: number },
|
||||
): number => Date.parse(left.record.scheduledAt) - Date.parse(right.record.scheduledAt)
|
||||
|| left.index - right.index
|
||||
|
||||
const oneShot = indexed
|
||||
.filter((entry): entry is { record: OneShotScheduleRecord; index: number } =>
|
||||
entry.record.kind !== 'every' && entry.record.kind !== 'cron'
|
||||
&& Date.parse(entry.record.scheduledAt) <= now)
|
||||
.sort((left, right) =>
|
||||
Date.parse(left.record.scheduledAt) - Date.parse(right.record.scheduledAt)
|
||||
|| left.index - right.index)
|
||||
const oneShot = dueOneShots[0]?.record
|
||||
entry.record.kind !== 'every' && Date.parse(entry.record.scheduledAt) <= now)
|
||||
.sort(byTargetThenCreate)[0]?.record
|
||||
if (oneShot !== undefined) return { kind: 'one-shot', record: oneShot }
|
||||
|
||||
const recurring = indexed
|
||||
.filter((entry): entry is { record: RecurringScheduleRecord; index: number } =>
|
||||
(entry.record.kind === 'every' || entry.record.kind === 'cron')
|
||||
&& Date.parse(entry.record.scheduledAt) <= now)
|
||||
.sort((left, right) =>
|
||||
Date.parse(left.record.scheduledAt) - Date.parse(right.record.scheduledAt)
|
||||
|| left.index - right.index)
|
||||
const gate = folded.lastRecurringAcceptedAt === undefined
|
||||
? Number.NEGATIVE_INFINITY
|
||||
: Date.parse(folded.lastRecurringAcceptedAt) + MIN_RECURRING_INTERVAL_SECONDS * 1_000
|
||||
if (recurring.length > 0 && now >= gate) {
|
||||
const every = indexed
|
||||
.filter((entry): entry is { record: EveryScheduleRecord; index: number } =>
|
||||
entry.record.kind === 'every' && Date.parse(entry.record.scheduledAt) <= now)
|
||||
.sort(byTargetThenCreate)
|
||||
if (every.length > 0) {
|
||||
return {
|
||||
kind: 'recurring',
|
||||
kind: 'every',
|
||||
acceptedAt: new Date(now).toISOString(),
|
||||
reminders: recurring.map(({ record }) => {
|
||||
const occurrence = record.kind === 'every'
|
||||
? resolveEveryOccurrence(record, now)
|
||||
: resolveCronOccurrence(record, now)
|
||||
return {
|
||||
record,
|
||||
occurrenceAt: occurrence.occurrenceAt,
|
||||
...(occurrence.nextScheduledAt === undefined
|
||||
? {}
|
||||
: { nextScheduledAt: occurrence.nextScheduledAt }),
|
||||
}
|
||||
}),
|
||||
reminders: every.map(({ record }) => ({
|
||||
record,
|
||||
occurrenceAt: resolveEveryOccurrence(record, now).occurrenceAt,
|
||||
})),
|
||||
}
|
||||
}
|
||||
|
||||
const future = folded.active
|
||||
.filter(record => recurring.length === 0 || (record.kind !== 'every' && record.kind !== 'cron'))
|
||||
.map(record => Date.parse(record.scheduledAt))
|
||||
.filter(target => target > now)
|
||||
if (recurring.length > 0) future.push(gate)
|
||||
const target = future.reduce<number | undefined>(
|
||||
(selected, candidate) => selected === undefined || candidate < selected ? candidate : selected,
|
||||
undefined,
|
||||
)
|
||||
const target = folded.active.reduce<number | undefined>((selected, record) => {
|
||||
const candidate = Date.parse(record.scheduledAt)
|
||||
return candidate > now && (selected === undefined || candidate < selected) ? candidate : selected
|
||||
}, undefined)
|
||||
return { kind: 'wait', ...(target === undefined ? {} : { target }) }
|
||||
}
|
||||
|
||||
@@ -185,6 +162,11 @@ export class ScheduleOwner {
|
||||
&& this.ctx.agents.roots().includes(this.agent)
|
||||
}
|
||||
|
||||
/** Whether this owner may start or continue Schedule work. */
|
||||
private isRunnable(): boolean {
|
||||
return !this.stopping && this.isLive()
|
||||
}
|
||||
|
||||
/** Cancel the currently armed timer, if any. */
|
||||
private clearTimer(): void {
|
||||
if (this.timer === undefined) return
|
||||
@@ -235,20 +217,20 @@ export class ScheduleOwner {
|
||||
}
|
||||
}
|
||||
|
||||
/** Contain a current calendar-resolution failure without permanently faulting this owner. */
|
||||
/** Contain an invalid wall-clock decision without permanently faulting this owner. */
|
||||
private decide(folded: FoldedSchedules, now: number): DueDecision | undefined {
|
||||
try {
|
||||
return dueDecision(folded, now)
|
||||
} catch (error: unknown) {
|
||||
this.ctx.logger.warn(`tool-schedule: calendar decision failed for agent "${this.agent.id}": ${renderThrown(error)}`)
|
||||
this.ctx.logger.warn(`tool-schedule: fixed-rate decision failed for agent "${this.agent.id}": ${renderThrown(error)}`)
|
||||
return undefined
|
||||
}
|
||||
}
|
||||
|
||||
/** Preflight, fold, arm, or dispatch the next one-shot or recurring batch. */
|
||||
/** Preflight, fold, arm, or dispatch the next one-shot or fixed-rate batch. */
|
||||
private async driveOnce(): Promise<void> {
|
||||
this.clearTimer()
|
||||
if (this.stopping || !this.isLive()) return
|
||||
if (!this.isRunnable()) return
|
||||
try {
|
||||
await flushSchedulePersistence(this.ctx, this.agent.session)
|
||||
} catch (error: unknown) {
|
||||
@@ -257,8 +239,7 @@ export class ScheduleOwner {
|
||||
}
|
||||
return
|
||||
}
|
||||
// oxlint-disable-next-line typescript/no-unnecessary-condition -- disposal or replacement can win while persistence is awaited.
|
||||
if (this.stopping || !this.isLive()) return
|
||||
if (!this.isRunnable()) return
|
||||
|
||||
const folded = this.readFolded()
|
||||
if (folded === undefined) return
|
||||
@@ -273,7 +254,7 @@ export class ScheduleOwner {
|
||||
let maintenance: Promise<boolean>
|
||||
try {
|
||||
maintenance = this.agent.runMaintenance(() => {
|
||||
if (this.stopping || !this.isLive()) return Promise.resolve(false)
|
||||
if (!this.isRunnable()) return Promise.resolve(false)
|
||||
const claimed = this.readFolded()
|
||||
if (claimed === undefined) return Promise.resolve(false)
|
||||
const decisionNow = Date.now()
|
||||
@@ -286,7 +267,7 @@ export class ScheduleOwner {
|
||||
try {
|
||||
const text = decision.kind === 'one-shot'
|
||||
? renderReminderFraming(decision.record)
|
||||
: renderReminderBatchFraming(decision.reminders)
|
||||
: renderEveryReminderBatchFraming(decision.reminders)
|
||||
const message = createUserMessage({
|
||||
content: [{ type: 'text', text }],
|
||||
source: { kind: 'plugin', plugin: 'tool-schedule' },
|
||||
@@ -307,25 +288,12 @@ export class ScheduleOwner {
|
||||
})
|
||||
} else {
|
||||
for (const reminder of decision.reminders) {
|
||||
if (reminder.record.kind === 'every') {
|
||||
this.agent.session.append('schedule/change', {
|
||||
version: 1,
|
||||
operation: 'dispatch',
|
||||
id: reminder.record.id,
|
||||
acceptedAt: decision.acceptedAt,
|
||||
})
|
||||
} else {
|
||||
this.agent.session.append('schedule/change', {
|
||||
version: 1,
|
||||
operation: 'dispatch',
|
||||
id: reminder.record.id,
|
||||
occurrenceAt: reminder.occurrenceAt,
|
||||
acceptedAt: decision.acceptedAt,
|
||||
...(reminder.nextScheduledAt === undefined
|
||||
? {}
|
||||
: { nextScheduledAt: reminder.nextScheduledAt }),
|
||||
})
|
||||
}
|
||||
this.agent.session.append('schedule/change', {
|
||||
version: 1,
|
||||
operation: 'dispatch',
|
||||
id: reminder.record.id,
|
||||
acceptedAt: decision.acceptedAt,
|
||||
})
|
||||
}
|
||||
}
|
||||
} catch (error: unknown) {
|
||||
@@ -351,7 +319,6 @@ export class ScheduleOwner {
|
||||
}
|
||||
return
|
||||
}
|
||||
// oxlint-disable-next-line typescript/no-unnecessary-condition -- disposal can win while the barrier is awaited.
|
||||
if (!this.stopping && this.isLive()) this.requestDrive()
|
||||
if (this.isRunnable()) this.requestDrive()
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user