Skip to content
File

Blob: src/worker/do/repo/scheduler.ts

typescript128 lines
1import type { RepoStateSchema } from "./repoState";
2import { asTypedStorage } from "./repoState";
3import { getConfig } from "./repoConfig";
4import { createLogger } from "@/worker/common";
5import { activeLeaseOrUndefined } from "./catalog/activity";
6import { COMPACTION_REARM_DELAY_MS } from "./catalog/shared";
7 
8/**
9 * Plan the next alarm time purely from existing DO state and repo config.
10 * Priority: compaction wake/retry, then idle cleanup.
11 */
12export async function planNextAlarm(
13 state: DurableObjectState,
14 env: Env,
15 now = Date.now()
16): Promise<{
17 when: number;
18 reason: "compaction" | "idle";
19} | null> {
20 const log = createLogger(env.LOG_LEVEL, { service: "Scheduler", doId: state.id.toString() });
21 const store = asTypedStorage<RepoStateSchema>(state.storage);
22 const cfg = getConfig(env);
23 
24 // 1) Re-arm compaction via alarms when a compaction request or lease is active.
25 try {
26 const [compactionWantedAt, receiveLease, compactLease] = await Promise.all([
27 store.get("compactionWantedAt"),
28 store.get("receiveLease"),
29 store.get("compactLease"),
30 ]);
31 
32 const activeReceiveLease = activeLeaseOrUndefined(receiveLease, now);
33 if (activeReceiveLease) {
34 return { when: activeReceiveLease.expiresAt, reason: "compaction" };
35 }
36 
37 const activeCompactLease = activeLeaseOrUndefined(compactLease, now);
38 if (activeCompactLease) {
39 return { when: activeCompactLease.expiresAt, reason: "compaction" };
40 }
41 
42 if (typeof compactionWantedAt === "number") {
43 return { when: now + COMPACTION_REARM_DELAY_MS, reason: "compaction" };
44 }
45 } catch (e) {
46 log.warn("sched:read-compaction-state-failed", { error: String(e) });
47 }
48 
49 // 2) Idle cleanup planning
50 try {
51 const lastAccess = await store.get("lastAccessMs");
52 // Guard against past deadlines (e.g., host slept). If an idle
53 // deadline is already in the past, push it forward by its interval so we
54 // do not immediately re-schedule and tight-loop alarms.
55 const nextIdleAt = (lastAccess ?? now) + cfg.idleMs;
56 const when = nextIdleAt <= now ? now + cfg.idleMs : nextIdleAt;
57 return { when, reason: "idle" };
58 } catch (e) {
59 log.error("sched:plan-idle-failed", { error: String(e) });
60 return null;
61 }
62}
63 
64/**
65 * Set the DO alarm only if this would fire sooner than the existing one.
66 */
67export async function scheduleAlarmIfSooner(
68 state: DurableObjectState,
69 env: Env,
70 when: number,
71 now = Date.now()
72): Promise<{ scheduled: boolean; prev: number | null; next: number }> {
73 const log = createLogger(env.LOG_LEVEL, { service: "Scheduler", doId: state.id.toString() });
74 let prev: number | null = null;
75 try {
76 prev = (await state.storage.getAlarm()) as number | null;
77 } catch (e) {
78 log.warn("sched:get-alarm-failed", { error: String(e) });
79 prev = null;
80 }
81 
82 // Avoid redundant reset to the same timestamp (even if in the past)
83 if (prev !== null && prev === when) {
84 return { scheduled: false, prev, next: prev };
85 }
86 
87 if (!prev || prev < now || prev > when) {
88 try {
89 await state.storage.setAlarm(when);
90 log.debug("sched:set-alarm", { when });
91 return { scheduled: true, prev: prev ?? null, next: when };
92 } catch (e) {
93 log.error("sched:set-alarm-failed", { error: String(e), when });
94 return { scheduled: false, prev: prev ?? null, next: prev ?? when };
95 }
96 }
97 return { scheduled: false, prev: prev ?? null, next: prev };
98}
99 
100/**
101 * Compute and schedule in one step. No-ops if nothing to schedule.
102 */
103export async function ensureScheduled(
104 state: DurableObjectState,
105 env: Env,
106 now = Date.now()
107): Promise<{
108 scheduled: boolean;
109 when?: number;
110 reason?: "compaction" | "idle";
111}> {
112 const log = createLogger(env.LOG_LEVEL, { service: "Scheduler", doId: state.id.toString() });
113 try {
114 const plan = await planNextAlarm(state, env, now);
115 if (!plan) return { scheduled: false };
116 // Clamp to a near-future time to avoid repeatedly scheduling past alarms
117 const targetWhen = Math.max(plan.when, now + 5);
118 const res = await scheduleAlarmIfSooner(state, env, targetWhen, now);
119 if (res.scheduled) {
120 log.debug("sched:alarm-set", { when: res.next, reason: plan.reason });
121 }
122 return { scheduled: res.scheduled, when: res.next, reason: plan.reason };
123 } catch (e) {
124 log.error("sched:ensure-failed", { error: String(e) });
125 return { scheduled: false };
126 }
127}