File
Blob: src/worker/do/repo/scheduler.ts
| 1 | import type { RepoStateSchema } from "./repoState"; |
| 2 | import { asTypedStorage } from "./repoState"; |
| 3 | import { getConfig } from "./repoConfig"; |
| 4 | import { createLogger } from "@/worker/common"; |
| 5 | import { activeLeaseOrUndefined } from "./catalog/activity"; |
| 6 | import { 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 | */ |
| 12 | export 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 | */ |
| 67 | export 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 | */ |
| 103 | export 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 | } |