Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- fix(db): phase-offset the model-sync interval (+45 min, period unchanged) from the cleanup 6h scheduler so they never fire in the same second, and add the missing `conversation_turn_nodes.last_seen_at` index used by the retention DELETE (#13973)
8 changes: 4 additions & 4 deletions src/lib/db/cleanup.ts
Original file line number Diff line number Diff line change
Expand Up @@ -583,10 +583,10 @@ export async function cleanupExpiredFiles(): Promise<CleanupResult> {
* separate compliance cleanup path and does not override this window.
* Deleting an old node only affects reconnect anchors: a conversation resumed
* after the window mints a new id, which is already the documented
* anchor-miss behavior of resolveConversationId. `last_seen_at` has no index
* (migration 156), so each DELETE is a table scan. Bounded batches yield
* between writes so an existing large table cannot park the event loop for
* the whole cleanup pass.
* anchor-miss behavior of resolveConversationId. `last_seen_at` is indexed
* (migration 186, #13973 — migration 156 originally missed it). Bounded
* batches yield between writes so an existing large table cannot park the
* event loop for the whole cleanup pass.
*/
export async function cleanupConversationTurnNodes(): Promise<CleanupResult> {
const retention = getRetentionSettings();
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
-- 186_conversation_turn_nodes_last_seen_index.sql
-- conversation_turn_nodes retention cleanup (cleanup.ts::cleanupConversationTurnNodes,
-- #13973) deletes rows WHERE last_seen_at < cutoff. Migration 156 indexed
-- conversation_id/parent_id/content_hash but never last_seen_at, so every
-- periodic retention pass was a full table scan on top of the batch delete's
-- own cost. Purely additive — only speeds up the existing DELETE, no schema
-- or behavior change for callers.

CREATE INDEX IF NOT EXISTS idx_turn_nodes_last_seen
ON conversation_turn_nodes(last_seen_at);
35 changes: 30 additions & 5 deletions src/shared/services/modelSyncScheduler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,15 @@ export const DEFAULT_INTERVAL_MS = 6 * 60 * 60 * 1000; // 6 hours
export const MODEL_SYNC_CYCLE_CONCURRENCY = 4;
/** First cycle after boot. Past cleanup's 30s so the two jobs do not overlap. */
export const MODEL_SYNC_STARTUP_DELAY_MS = 90_000;
/**
* Phase offset (not a period change) between this scheduler's recurring tick
* and cleanup.ts's own 6h scheduler (#13973 — both are started back-to-back
* in the same boot sequence, so with the same period they collide every 6h
* for the process lifetime). The recurring `setInterval` is armed only after
* this delay, so every periodic tick lands at boot + offset + k * interval
* while the configured interval itself stays exactly as configured.
*/
export const MODEL_SYNC_STAGGER_OFFSET_MS = 45 * 60 * 1000; // 45 minutes
const MODEL_SYNC_SETTING_KEY = "model_sync_last_run";
const MODEL_SYNC_INTERNAL_AUTH_HEADER = "x-model-sync-internal-auth";

Expand Down Expand Up @@ -121,6 +130,8 @@ const globalState = globalThis as typeof globalThis & {
};

let schedulerTimer: NodeJS.Timeout | null = null;
/** Pending one-shot that arms `schedulerTimer` after the phase offset. */
let phaseTimer: NodeJS.Timeout | null = null;
let isRunning = false;
let internalAuthToken: string | null = null;

Expand Down Expand Up @@ -303,7 +314,7 @@ export function startModelSyncScheduler(
apiBaseUrl = getModelSyncInternalBaseUrl(),
intervalMs = DEFAULT_INTERVAL_MS
): void {
if (schedulerTimer) {
if (schedulerTimer || phaseTimer) {
console.log("[ModelSync] Scheduler already running — skipping start");
return;
}
Expand All @@ -314,7 +325,10 @@ export function startModelSyncScheduler(
!isNaN(envHours) && envHours > 0 ? envHours * 60 * 60 * 1000 : intervalMs;
const trustedApiBaseUrl = resolveModelSyncInternalBaseUrl(apiBaseUrl);

console.log(`[ModelSync] Scheduler started — interval: ${effectiveIntervalMs / 3_600_000}h`);
console.log(
`[ModelSync] Scheduler started — interval: ${effectiveIntervalMs / 3_600_000}h ` +
`(phase offset +${MODEL_SYNC_STAGGER_OFFSET_MS / 60_000}m)`
);

// Serve traffic first; cleanup's first pass is +30s, so stay past that window.
const startupDelay = setTimeout(
Expand All @@ -332,15 +346,26 @@ export function startModelSyncScheduler(
// silent
});

// Then run on the regular interval
schedulerTimer = setInterval(() => runSyncCycle(trustedApiBaseUrl), effectiveIntervalMs);
schedulerTimer.unref?.();
// Then run on the regular interval, phase-shifted against cleanup.ts's
// own 6h scheduler (#13973). The period stays `effectiveIntervalMs`; only
// the moment the interval is armed moves.
phaseTimer = setTimeout(() => {
phaseTimer = null;
schedulerTimer = setInterval(() => runSyncCycle(trustedApiBaseUrl), effectiveIntervalMs);
schedulerTimer.unref?.();
}, MODEL_SYNC_STAGGER_OFFSET_MS);
phaseTimer.unref?.();
}

/**
* Stop the model sync scheduler.
*/
export function stopModelSyncScheduler(): void {
if (phaseTimer) {
clearTimeout(phaseTimer);
phaseTimer = null;
console.log("[ModelSync] Scheduler stopped");
}
if (schedulerTimer) {
clearInterval(schedulerTimer);
schedulerTimer = null;
Expand Down
57 changes: 54 additions & 3 deletions tests/unit/model-sync-scheduler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -365,10 +365,17 @@ test("modelSyncScheduler starts once, honors env interval and syncs only active
scheduler.startModelSyncScheduler("http://127.0.0.1:7777", 1000);
scheduler.startModelSyncScheduler("http://127.0.0.1:8888", 9999);

assert.equal(timers.timeouts.length, 1);
assert.equal(timers.timeouts.length, 2);
assert.equal(timers.timeouts[0].ms, scheduler.MODEL_SYNC_STARTUP_DELAY_MS);
assert.equal(timers.timeouts[0].unrefCalled, true);
// #13973: the recurring interval is phase-shifted against cleanup.ts's
// own 6h scheduler — armed only after the offset, never at boot.
assert.equal(timers.timeouts[1].ms, scheduler.MODEL_SYNC_STAGGER_OFFSET_MS);
assert.equal(timers.timeouts[1].unrefCalled, true);
assert.equal(timers.intervals.length, 0);
timers.timeouts[1].fn();
assert.equal(timers.intervals.length, 1);
// The period itself is exactly the configured interval (no stagger added).
assert.equal(timers.intervals[0].ms, 6 * 60 * 60 * 1000);
assert.equal(timers.intervals[0].unrefCalled, true);

Expand Down Expand Up @@ -396,6 +403,50 @@ test("modelSyncScheduler starts once, honors env interval and syncs only active
}
});

test("#13973: phase offset delays the first periodic tick but never changes the configured period", async () => {
const timers = installTimerStubs();
const previousHours = process.env.MODEL_SYNC_INTERVAL_HOURS;
process.env.MODEL_SYNC_INTERVAL_HOURS = "4";

try {
const scheduler = await loadScheduler("phase-offset-period");
scheduler.startModelSyncScheduler("http://127.0.0.1:7777");
// Let the fire-and-forget codex revalidation import settle while the timer
// stubs are still installed, so its setTimeout(0) is captured here instead
// of arming a real loopback poll that leaks fetches into later tests.
await flushMicrotasks();

// No recurring timer exists at boot: arming it at boot is what made it
// collide with cleanup.ts's boot-anchored 6h interval.
assert.equal(timers.intervals.length, 0);
const phase = timers.timeouts.find((t) => t.ms === scheduler.MODEL_SYNC_STAGGER_OFFSET_MS);
assert.ok(phase, "expected a one-shot phase timer of MODEL_SYNC_STAGGER_OFFSET_MS");
assert.ok(scheduler.MODEL_SYNC_STAGGER_OFFSET_MS > 0);

phase.fn();
assert.equal(timers.intervals.length, 1);
// The operator-configured MODEL_SYNC_INTERVAL_HOURS=4 must be honored exactly.
assert.equal(timers.intervals[0].ms, 4 * 60 * 60 * 1000);

// Stopping before the phase fires must cancel the pending arm too.
scheduler.stopModelSyncScheduler();
assert.equal(timers.intervals[0].cleared, true);

timers.timeouts.length = 0;
timers.intervals.length = 0;
scheduler.startModelSyncScheduler("http://127.0.0.1:7777");
await flushMicrotasks();
const pending = timers.timeouts.find((t) => t.ms === scheduler.MODEL_SYNC_STAGGER_OFFSET_MS);
scheduler.stopModelSyncScheduler();
assert.equal(pending.cleared, true);
assert.equal(timers.intervals.length, 0);
} finally {
if (previousHours === undefined) delete process.env.MODEL_SYNC_INTERVAL_HOURS;
else process.env.MODEL_SYNC_INTERVAL_HOURS = previousHours;
timers.restore();
}
});

test("modelSyncScheduler skips empty cycles and tolerates failing sync requests", async () => {
const timers = installTimerStubs();
const originalFetch = globalThis.fetch;
Expand Down Expand Up @@ -445,7 +496,7 @@ test("modelSyncScheduler skips empty cycles and tolerates failing sync requests"
test("test 12: default interval is 6h; env hours override; no-arg uses default", async () => {
const source = fs.readFileSync(
path.join(process.cwd(), "src/shared/services/modelSyncScheduler.ts"),
"utf8",
"utf8"
);
assert.match(source, /DEFAULT_INTERVAL_MS\s*=\s*6\s*\*\s*60\s*\*\s*60\s*\*\s*1000/);
assert.doesNotMatch(source, /DEFAULT_INTERVAL_MS\s*=\s*24\s*\*\s*60\s*\*\s*60\s*\*\s*1000/);
Expand All @@ -457,7 +508,7 @@ test("test 12: default interval is 6h; env hours override; no-arg uses default",
test("test 12: MODEL_SYNC_INTERVAL_HOURS still wins over default", () => {
const source = fs.readFileSync(
path.join(process.cwd(), "src/shared/services/modelSyncScheduler.ts"),
"utf8",
"utf8"
);
assert.match(source, /MODEL_SYNC_INTERVAL_HOURS/);
assert.match(source, /envHours \* 60 \* 60 \* 1000/);
Expand Down
106 changes: 106 additions & 0 deletions tests/unit/scheduler-6h-stagger-13973.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,106 @@
// Regression test for issue #13973 (remaining scope after PR #14005):
// the cleanup scheduler and the model-sync scheduler are both anchored to
// process-start with an *identical* 6h period and *zero* jitter/offset
// between them. Since both are started within milliseconds of each other
// during server bootstrap (src/instrumentation-node.ts calls
// startCleanupScheduler() in the same Promise.all-driven boot sequence that
// calls ensureCloudSyncInitialized() -> startModelSyncScheduler()), they
// would otherwise keep firing in the same second every 6 hours for the
// lifetime of the process.
//
// This test does NOT touch the DB or execute the scheduled work: it spies on
// the global timer constructors to capture the delay each scheduler
// registers. The stagger must be a PHASE offset — the model-sync interval is
// armed only after MODEL_SYNC_STAGGER_OFFSET_MS — never a change of PERIOD
// (a 6h45m period merely drifts back into collision and silently overrides
// MODEL_SYNC_INTERVAL_HOURS). It must run as its own file, not as part of a
// suite.

import { test, before, after } from "node:test";
import assert from "node:assert/strict";

process.env.OMNIROUTE_DISABLE_BACKGROUND_SERVICES = "1";

let startCleanupScheduler: () => void;
let stopCleanupScheduler: () => void;
let startModelSyncScheduler: (apiBaseUrl?: string, intervalMs?: number) => void;
let stopModelSyncScheduler: () => void;

before(async () => {
({ startCleanupScheduler, stopCleanupScheduler } = await import("../../src/lib/db/cleanup.ts"));
({ startModelSyncScheduler, stopModelSyncScheduler } =
await import("../../src/shared/services/modelSyncScheduler.ts"));
});

after(() => {
try {
stopCleanupScheduler();
} catch {}
try {
stopModelSyncScheduler();
} catch {}
});

test("cleanup and model-sync 6h schedulers keep the same period but different phase — issue #13973 remaining scope", async () => {
const { MODEL_SYNC_STAGGER_OFFSET_MS } =
await import("../../src/shared/services/modelSyncScheduler.ts");
const originalSetInterval = global.setInterval;
const originalSetTimeout = global.setTimeout;

type TimerFn = (...a: unknown[]) => void;
type TimerSetter = (fn: TimerFn, delay?: number, ...rest: unknown[]) => unknown;

const bootIntervals: number[] = [];
const bootTimeouts: { fn: TimerFn; delay: number }[] = [];

(global as unknown as { setInterval: TimerSetter }).setInterval = (
fn: TimerFn,
delay?: number,
...rest: unknown[]
) => {
bootIntervals.push(delay ?? -1);
// Capture only: arm an inert, unref'd timer so no scheduled DB/HTTP work
// runs and the process is not kept alive by the real 30s/90s first passes.
void fn;
return originalSetInterval(() => {}, delay, ...rest).unref();
};
(global as unknown as { setTimeout: TimerSetter }).setTimeout = (
fn: TimerFn,
delay?: number,
...rest: unknown[]
) => {
bootTimeouts.push({ fn, delay: delay ?? -1 });
return originalSetTimeout(() => {}, delay, ...rest).unref();
};

let modelSyncIntervals: number[] = [];
try {
startCleanupScheduler();
const cleanupIntervals = bootIntervals.length;
startModelSyncScheduler("http://127.0.0.1:0");

const SIX_HOURS_MS = 6 * 60 * 60 * 1000;
assert.equal(cleanupIntervals, 1, "cleanup arms its 6h interval at boot");
assert.deepEqual(
bootIntervals,
[SIX_HOURS_MS],
"model-sync must NOT arm its recurring interval at boot (it would share cleanup's phase), " +
`got ${JSON.stringify(bootIntervals)}`
);

assert.ok(MODEL_SYNC_STAGGER_OFFSET_MS > 0);
const phase = bootTimeouts.find((t) => t.delay === MODEL_SYNC_STAGGER_OFFSET_MS);
assert.ok(phase, "model-sync must schedule a one-shot phase-offset timer");

// Fire the phase timer (still under the capture stubs): the recurring
// interval it arms must keep the configured 6h period exactly.
bootIntervals.length = 0;
phase.fn();
modelSyncIntervals = [...bootIntervals];
} finally {
global.setInterval = originalSetInterval;
global.setTimeout = originalSetTimeout;
}

assert.deepEqual(modelSyncIntervals, [6 * 60 * 60 * 1000]);
});
Loading