From 294a4038887bd235bd8c0de6ec3a94b9a5f12c68 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 20 Apr 2026 19:13:08 +0000 Subject: [PATCH 1/6] Add scheduler alarm diagnostics Co-authored-by: Kent C. Dodds --- packages/worker/src/jobs/manager-client.ts | 35 ++- .../worker/src/jobs/manager-do.node.test.ts | 240 ++++++++++++++++++ packages/worker/src/jobs/manager-do.ts | 125 +++++++-- .../src/jobs/process-due-jobs.node.test.ts | 31 +++ packages/worker/src/jobs/process-due-jobs.ts | 76 ++++-- packages/worker/src/jobs/scheduler-logging.ts | 95 +++++++ packages/worker/src/jobs/service.ts | 53 ++-- .../src/mcp/capabilities/jobs/shared.ts | 32 ++- 8 files changed, 631 insertions(+), 56 deletions(-) create mode 100644 packages/worker/src/jobs/manager-do.node.test.ts create mode 100644 packages/worker/src/jobs/scheduler-logging.ts diff --git a/packages/worker/src/jobs/manager-client.ts b/packages/worker/src/jobs/manager-client.ts index 5cb028f58e..7a5408cf62 100644 --- a/packages/worker/src/jobs/manager-client.ts +++ b/packages/worker/src/jobs/manager-client.ts @@ -1,4 +1,9 @@ import { type McpCallerContext } from '@kody-internal/shared/chat.ts' +import { + logJobSchedulerError, + logJobSchedulerEvent, + schedulerErrorFields, +} from './scheduler-logging.ts' import { type JobExecutionResult, type JobRepoCheckPolicy, @@ -33,15 +38,43 @@ export function jobManagerRpc(env: Env, userId: string): JobManagerRpc | null { export async function syncJobManagerAlarm(input: { env: Env; userId: string }) { const rpc = jobManagerRpc(input.env, input.userId) if (!rpc) { + logJobSchedulerEvent({ + event: 'sync_alarm_skipped_missing_binding', + userId: input.userId, + reason: 'missing_job_manager_binding', + }) return { ok: true as const, userId: input.userId, nextRunAt: null, } } - return rpc.syncAlarm({ + logJobSchedulerEvent({ + event: 'sync_alarm_requested', userId: input.userId, }) + try { + const result = await rpc.syncAlarm({ + userId: input.userId, + }) + logJobSchedulerEvent({ + event: 'sync_alarm_completed', + userId: result.userId, + nextRunAt: result.nextRunAt, + reason: + result.nextRunAt == null + ? 'no_runnable_job_found' + : 'alarm_state_updated', + }) + return result + } catch (error) { + logJobSchedulerError({ + event: 'sync_alarm_request_failed', + userId: input.userId, + ...schedulerErrorFields(error), + }) + throw error + } } export async function runJobNowViaManager(input: { diff --git a/packages/worker/src/jobs/manager-do.node.test.ts b/packages/worker/src/jobs/manager-do.node.test.ts new file mode 100644 index 0000000000..af249bccf0 --- /dev/null +++ b/packages/worker/src/jobs/manager-do.node.test.ts @@ -0,0 +1,240 @@ +import { expect, test, vi } from 'vitest' + +const mockModule = vi.hoisted(() => ({ + getNextRunnableJob: vi.fn(), + runDueJobsForUser: vi.fn(), + runJobNow: vi.fn(), + buildSentryOptions: vi.fn(), + logJobSchedulerEvent: vi.fn(), + logJobSchedulerError: vi.fn(), +})) + +vi.mock('@sentry/cloudflare', () => ({ + instrumentDurableObjectWithSentry: ( + _getOptions: unknown, + durableObjectClass: unknown, + ) => durableObjectClass, +})) + +vi.mock('cloudflare:workers', () => ({ + DurableObject: class { + protected readonly ctx: DurableObjectState + protected readonly env: Env + + constructor(ctx: DurableObjectState, env: Env) { + this.ctx = ctx + this.env = env + } + }, +})) + +vi.mock('./service.ts', () => ({ + getNextRunnableJob: (...args: Array) => + mockModule.getNextRunnableJob(...args), + runDueJobsForUser: (...args: Array) => + mockModule.runDueJobsForUser(...args), + runJobNow: (...args: Array) => mockModule.runJobNow(...args), +})) + +vi.mock('#worker/sentry-options.ts', () => ({ + buildSentryOptions: (...args: Array) => + mockModule.buildSentryOptions(...args), +})) + +vi.mock('./scheduler-logging.ts', async (importOriginal) => { + const actual = await importOriginal() + return { + ...actual, + logJobSchedulerEvent: (...args: Array) => + mockModule.logJobSchedulerEvent(...args), + logJobSchedulerError: (...args: Array) => + mockModule.logJobSchedulerError(...args), + } +}) + +const { JobManagerBase } = await import('./manager-do.ts') + +function resetMocks() { + mockModule.getNextRunnableJob.mockReset() + mockModule.runDueJobsForUser.mockReset() + mockModule.runJobNow.mockReset() + mockModule.buildSentryOptions.mockReset() + mockModule.logJobSchedulerEvent.mockReset() + mockModule.logJobSchedulerError.mockReset() +} + +function createState({ + userId = 'user-123', + currentAlarmAt = null, +}: { + userId?: string | null + currentAlarmAt?: number | null +} = {}) { + const persistedEntries = new Map() + if (userId !== undefined) { + persistedEntries.set('user-id', userId) + } + let alarmAt = currentAlarmAt + + return { + state: { + storage: { + get: vi.fn(async (key: string) => persistedEntries.get(key)), + put: vi.fn(async (key: string, value: unknown) => { + persistedEntries.set(key, value) + }), + getAlarm: vi.fn(async () => alarmAt), + setAlarm: vi.fn(async (value: Date | number) => { + alarmAt = value instanceof Date ? value.valueOf() : Number(value) + }), + deleteAlarm: vi.fn(async () => { + alarmAt = null + }), + }, + } as unknown as DurableObjectState, + persistedEntries, + getAlarmAt() { + return alarmAt + }, + } +} + +test('syncAlarm logs when it arms a new alarm for the next runnable job', async () => { + resetMocks() + const nextRunAt = '2026-04-20T18:30:00.000Z' + mockModule.getNextRunnableJob.mockResolvedValue({ + id: 'job-123', + nextRunAt, + }) + const { state, persistedEntries, getAlarmAt } = createState({ + currentAlarmAt: Date.parse('2026-04-20T18:00:00.000Z'), + }) + const manager = new JobManagerBase(state, {} as Env) + + await expect(manager.syncAlarm({ userId: 'user-123' })).resolves.toEqual({ + ok: true, + userId: 'user-123', + nextRunAt, + }) + + expect(persistedEntries.get('user-id')).toBe('user-123') + expect(getAlarmAt()).toBe(Date.parse(nextRunAt)) + expect(mockModule.logJobSchedulerEvent).toHaveBeenCalledWith({ + event: 'sync_alarm', + userId: 'user-123', + currentAlarmAt: '2026-04-20T18:00:00.000Z', + nextJobId: 'job-123', + nextRunAt, + reason: 'alarm-armed', + }) + expect(mockModule.logJobSchedulerError).not.toHaveBeenCalled() +}) + +test('syncAlarm logs when no runnable job is found and clears the alarm', async () => { + resetMocks() + mockModule.getNextRunnableJob.mockResolvedValue(null) + const { state, getAlarmAt } = createState({ + currentAlarmAt: Date.parse('2026-04-20T18:00:00.000Z'), + }) + const manager = new JobManagerBase(state, {} as Env) + + await expect(manager.syncAlarm({ userId: 'user-123' })).resolves.toEqual({ + ok: true, + userId: 'user-123', + nextRunAt: null, + }) + + expect(getAlarmAt()).toBeNull() + expect(mockModule.logJobSchedulerEvent).toHaveBeenCalledWith({ + event: 'sync_alarm', + userId: 'user-123', + currentAlarmAt: '2026-04-20T18:00:00.000Z', + nextJobId: null, + nextRunAt: null, + reason: 'no-runnable-job', + }) +}) + +test('alarm logs firing, due-job outcomes, and resyncs the next alarm', async () => { + resetMocks() + mockModule.runDueJobsForUser.mockResolvedValue({ + dueJobCount: 2, + successCount: 1, + errorCount: 1, + jobOutcomes: [ + { + jobId: 'job-success', + scheduleType: 'once', + outcome: 'success', + nextRunAt: null, + deleted: true, + }, + { + jobId: 'job-failure', + scheduleType: 'interval', + outcome: 'failure', + nextRunAt: '2026-04-20T19:00:00.000Z', + deleted: false, + error: 'boom', + }, + ], + }) + mockModule.getNextRunnableJob.mockResolvedValue({ + id: 'job-next', + nextRunAt: '2026-04-20T19:00:00.000Z', + }) + const { state } = createState() + const manager = new JobManagerBase(state, {} as Env) + + await expect( + manager.alarm({ + retryCount: 2, + isRetry: true, + }), + ).resolves.toBeUndefined() + + expect(mockModule.runDueJobsForUser).toHaveBeenCalledWith({ + env: {} as Env, + userId: 'user-123', + }) + expect(mockModule.logJobSchedulerEvent).toHaveBeenNthCalledWith(1, { + event: 'alarm_fired', + userId: 'user-123', + retryCount: 2, + isRetry: true, + }) + expect(mockModule.logJobSchedulerEvent).toHaveBeenNthCalledWith(2, { + event: 'alarm_processed_due_jobs', + userId: 'user-123', + dueJobCount: 2, + successCount: 1, + errorCount: 1, + reason: 'processed-due-jobs', + jobOutcomes: [ + { + jobId: 'job-success', + scheduleType: 'once', + outcome: 'success', + nextRunAt: null, + deleted: true, + }, + { + jobId: 'job-failure', + scheduleType: 'interval', + outcome: 'failure', + nextRunAt: '2026-04-20T19:00:00.000Z', + deleted: false, + error: 'boom', + }, + ], + }) + expect(mockModule.logJobSchedulerEvent).toHaveBeenNthCalledWith(3, { + event: 'sync_alarm', + userId: 'user-123', + currentAlarmAt: null, + nextJobId: 'job-next', + nextRunAt: '2026-04-20T19:00:00.000Z', + reason: 'alarm-armed', + }) + expect(mockModule.logJobSchedulerError).not.toHaveBeenCalled() +}) diff --git a/packages/worker/src/jobs/manager-do.ts b/packages/worker/src/jobs/manager-do.ts index 07dda828f2..90f43a45ba 100644 --- a/packages/worker/src/jobs/manager-do.ts +++ b/packages/worker/src/jobs/manager-do.ts @@ -3,48 +3,135 @@ import { type McpCallerContext } from '@kody-internal/shared/chat.ts' import { DurableObject } from 'cloudflare:workers' import { buildSentryOptions } from '#worker/sentry-options.ts' import { getNextRunnableJob, runDueJobsForUser, runJobNow } from './service.ts' +import { + logJobSchedulerError, + logJobSchedulerEvent, + schedulerErrorFields, + summarizeSchedulerJobOutcomes, +} from './scheduler-logging.ts' import { type JobRepoCheckPolicy } from './types.ts' const userIdStorageKey = 'user-id' -class JobManagerBase extends DurableObject { +export class JobManagerBase extends DurableObject { async syncAlarm(input: { userId: string }) { const userId = input.userId.trim() if (!userId) { throw new Error('Job manager requires a non-empty userId.') } - await this.ctx.storage.put(userIdStorageKey, userId) - const nextJob = await getNextRunnableJob({ - env: this.env, - userId, - }) - if (!nextJob) { - await this.ctx.storage.deleteAlarm() + try { + await this.ctx.storage.put(userIdStorageKey, userId) + const currentAlarmAt = await this.ctx.storage.getAlarm() + const nextJob = await getNextRunnableJob({ + env: this.env, + userId, + }) + if (!nextJob) { + await this.ctx.storage.deleteAlarm() + logJobSchedulerEvent({ + event: 'sync_alarm', + userId, + currentAlarmAt: + currentAlarmAt == null + ? null + : new Date(currentAlarmAt).toISOString(), + nextJobId: null, + nextRunAt: null, + reason: 'no-runnable-job', + }) + return { + ok: true as const, + userId, + nextRunAt: null, + } + } + await this.ctx.storage.setAlarm(new Date(nextJob.nextRunAt)) + logJobSchedulerEvent({ + event: 'sync_alarm', + userId, + currentAlarmAt: + currentAlarmAt == null + ? null + : new Date(currentAlarmAt).toISOString(), + nextJobId: nextJob.id, + nextRunAt: nextJob.nextRunAt, + reason: + currentAlarmAt === new Date(nextJob.nextRunAt).valueOf() + ? 'alarm-unchanged' + : 'alarm-armed', + }) return { ok: true as const, userId, - nextRunAt: null, + nextRunAt: nextJob.nextRunAt, } - } - await this.ctx.storage.setAlarm(new Date(nextJob.nextRunAt)) - return { - ok: true as const, - userId, - nextRunAt: nextJob.nextRunAt, + } catch (error) { + logJobSchedulerError({ + event: 'sync_alarm_failed', + userId, + ...schedulerErrorFields(error), + }) + throw error } } - async alarm(): Promise { + async alarm(alarmInfo?: { + retryCount?: number + isRetry?: boolean + }): Promise { const userId = await this.ctx.storage.get(userIdStorageKey) if (!userId) { await this.ctx.storage.deleteAlarm() + logJobSchedulerEvent({ + event: 'alarm_fired', + reason: 'missing-user-id', + retryCount: alarmInfo?.retryCount, + isRetry: alarmInfo?.isRetry, + }) return } - await runDueJobsForUser({ - env: this.env, + logJobSchedulerEvent({ + event: 'alarm_fired', userId, + retryCount: alarmInfo?.retryCount, + isRetry: alarmInfo?.isRetry, }) - await this.syncAlarm({ userId }) + try { + const result = await runDueJobsForUser({ + env: this.env, + userId, + }) + logJobSchedulerEvent({ + event: 'alarm_processed_due_jobs', + userId, + dueJobCount: result.dueJobCount, + successCount: result.successCount, + errorCount: result.errorCount, + reason: result.dueJobCount === 0 ? 'no-due-jobs' : 'processed-due-jobs', + ...summarizeSchedulerJobOutcomes(result.jobOutcomes), + }) + } catch (error) { + logJobSchedulerError({ + event: 'alarm_run_due_jobs_failed', + userId, + retryCount: alarmInfo?.retryCount, + isRetry: alarmInfo?.isRetry, + ...schedulerErrorFields(error), + }) + throw error + } + try { + await this.syncAlarm({ userId }) + } catch (error) { + logJobSchedulerError({ + event: 'alarm_resync_failed', + userId, + retryCount: alarmInfo?.retryCount, + isRetry: alarmInfo?.isRetry, + ...schedulerErrorFields(error), + }) + throw error + } } async runNow(input: { diff --git a/packages/worker/src/jobs/process-due-jobs.node.test.ts b/packages/worker/src/jobs/process-due-jobs.node.test.ts index 962eb9d310..43dc0504ac 100644 --- a/packages/worker/src/jobs/process-due-jobs.node.test.ts +++ b/packages/worker/src/jobs/process-due-jobs.node.test.ts @@ -55,6 +55,25 @@ test('processDueJobs records failures without aborting later jobs', async () => expect(result.deleteJobIds).toEqual([]) expect(result.saveJobs).toHaveLength(2) + expect(result.successCount).toBe(1) + expect(result.errorCount).toBe(1) + expect(result.jobOutcomes).toEqual([ + { + jobId: 'job-1', + scheduleType: 'cron', + outcome: 'failure', + nextRunAt: expect.any(String), + deleted: false, + error: 'boom', + }, + { + jobId: 'job-2', + scheduleType: 'cron', + outcome: 'success', + nextRunAt: expect.any(String), + deleted: false, + }, + ]) expect(result.saveJobs).toEqual( expect.arrayContaining([ expect.objectContaining({ @@ -107,4 +126,16 @@ test('processDueJobs deletes one-shot jobs after execution', async () => { expect(result.deleteJobIds).toEqual(['job-once']) expect(result.saveJobs).toEqual([]) + expect(result.successCount).toBe(0) + expect(result.errorCount).toBe(1) + expect(result.jobOutcomes).toEqual([ + { + jobId: 'job-once', + scheduleType: 'once', + outcome: 'failure', + nextRunAt: null, + deleted: true, + error: 'expected failure', + }, + ]) }) diff --git a/packages/worker/src/jobs/process-due-jobs.ts b/packages/worker/src/jobs/process-due-jobs.ts index 3cae9f9147..ae75fde1cc 100644 --- a/packages/worker/src/jobs/process-due-jobs.ts +++ b/packages/worker/src/jobs/process-due-jobs.ts @@ -1,3 +1,4 @@ +import { type SchedulerJobOutcomeLog } from './scheduler-logging.ts' import { computeNextRunAt, formatJobError } from './schedule.ts' import { type JobExecutionOutcome, @@ -5,9 +6,12 @@ import { type JobRunStatus, } from './types.ts' -type ProcessDueJobsResult = { +export type ProcessDueJobsResult = { deleteJobIds: Array saveJobs: Array + successCount: number + errorCount: number + jobOutcomes: Array } const maxRunHistoryEntries = 10 @@ -20,6 +24,9 @@ export async function processDueJobs(input: { const now = input.now ?? new Date() const deleteJobIds: Array = [] const saveJobs: Array = [] + const jobOutcomes: Array = [] + let successCount = 0 + let errorCount = 0 for (const job of input.jobs) { const outcome = await input.executeJob(job).catch((error) => { @@ -35,37 +42,72 @@ export async function processDueJobs(input: { durationMs: 0, } }) + const executionError = outcome.execution.ok + ? undefined + : outcome.execution.error + if (outcome.execution.ok) { + successCount += 1 + } else { + errorCount += 1 + } if (job.schedule.type === 'once') { deleteJobIds.push(job.id) + jobOutcomes.push({ + jobId: job.id, + scheduleType: job.schedule.type, + outcome: outcome.execution.ok ? 'success' : 'failure', + nextRunAt: null, + deleted: true, + ...(executionError ? { error: executionError } : {}), + }) continue } try { - saveJobs.push( - applyExecutionOutcome(job, outcome, { - updatedAt: now.toISOString(), - nextRunAt: computeNextRunAt({ - schedule: job.schedule, - timezone: job.timezone, - from: now, - }), + const updated = applyExecutionOutcome(job, outcome, { + updatedAt: now.toISOString(), + nextRunAt: computeNextRunAt({ + schedule: job.schedule, + timezone: job.timezone, + from: now, }), - ) + }) + saveJobs.push(updated) + jobOutcomes.push({ + jobId: job.id, + scheduleType: job.schedule.type, + outcome: outcome.execution.ok ? 'success' : 'failure', + nextRunAt: updated.nextRunAt, + deleted: false, + ...(executionError ? { error: executionError } : {}), + }) } catch (error) { - saveJobs.push( - applyExecutionOutcome(job, outcome, { - updatedAt: now.toISOString(), - enabled: false, - lastRunError: `Failed to reschedule job: ${formatJobError(error)}`, - }), - ) + const rescheduleError = formatJobError(error) + const updated = applyExecutionOutcome(job, outcome, { + updatedAt: now.toISOString(), + enabled: false, + lastRunError: `Failed to reschedule job: ${rescheduleError}`, + }) + saveJobs.push(updated) + jobOutcomes.push({ + jobId: job.id, + scheduleType: job.schedule.type, + outcome: outcome.execution.ok ? 'success' : 'failure', + nextRunAt: updated.nextRunAt, + deleted: false, + ...(executionError ? { error: executionError } : {}), + rescheduleError, + }) } } return { deleteJobIds, saveJobs, + successCount, + errorCount, + jobOutcomes, } } diff --git a/packages/worker/src/jobs/scheduler-logging.ts b/packages/worker/src/jobs/scheduler-logging.ts new file mode 100644 index 0000000000..1eb2ab327d --- /dev/null +++ b/packages/worker/src/jobs/scheduler-logging.ts @@ -0,0 +1,95 @@ +import { formatJobError } from './schedule.ts' +import { type JobSchedule } from './types.ts' + +const maxLoggedJobOutcomes = 10 +type SchedulerLogLevel = 'error' | 'info' + +export type SchedulerJobOutcomeLog = { + jobId: string + scheduleType: JobSchedule['type'] + outcome: 'success' | 'failure' + nextRunAt: string | null + deleted: boolean + error?: string + rescheduleError?: string +} + +type JobSchedulerLogPayload = { + event: string + userId?: string + jobId?: string | null + scheduleType?: JobSchedule['type'] + nextJobId?: string | null + nextRunAt?: string | null + currentAlarmAt?: string | null + reason?: string + dueJobCount?: number + successCount?: number + errorCount?: number + jobOutcomes?: Array + truncatedJobOutcomeCount?: number + retryCount?: number + isRetry?: boolean + errorName?: string + errorMessage?: string + timestamp?: string +} + +export function schedulerErrorFields(error: unknown): { + errorName: string + errorMessage: string +} { + if (error instanceof Error) { + return { + errorName: error.name, + errorMessage: error.message, + } + } + + return { + errorName: 'Unknown', + errorMessage: formatJobError(error), + } +} + +export function summarizeSchedulerJobOutcomes( + jobOutcomes: Array, +): Pick { + if (jobOutcomes.length <= maxLoggedJobOutcomes) { + return { jobOutcomes } + } + + return { + jobOutcomes: jobOutcomes.slice(0, maxLoggedJobOutcomes), + truncatedJobOutcomeCount: jobOutcomes.length - maxLoggedJobOutcomes, + } +} + +export function logJobSchedulerEvent(input: JobSchedulerLogPayload): void { + writeSchedulerLog('info', input) +} + +export function logJobSchedulerError(input: JobSchedulerLogPayload): void { + writeSchedulerLog('error', input) +} + +function writeSchedulerLog( + level: SchedulerLogLevel, + input: JobSchedulerLogPayload, +): void { + try { + console[level]( + 'job-scheduler', + JSON.stringify({ + timestamp: input.timestamp ?? new Date().toISOString(), + ...input, + }), + ) + } catch (error) { + console.warn('job-scheduler-log-failed', { + event: input.event, + level, + errorMessage: formatJobError(error), + }) + } +} diff --git a/packages/worker/src/jobs/service.ts b/packages/worker/src/jobs/service.ts index b395b732b8..d6a47642f3 100644 --- a/packages/worker/src/jobs/service.ts +++ b/packages/worker/src/jobs/service.ts @@ -1,8 +1,5 @@ import { type McpCallerContext } from '@kody-internal/shared/chat.ts' -import { - createMcpCallerContext, - parseMcpCallerContext, -} from '#mcp/context.ts' +import { createMcpCallerContext, parseMcpCallerContext } from '#mcp/context.ts' import { buildJobEmbedText } from '#mcp/jobs-embed.ts' import { deleteJobVector, upsertJobVector } from '#mcp/jobs-vectorize.ts' import { type ExecuteResult } from '@cloudflare/codemode' @@ -54,6 +51,10 @@ import { import { buildKodyModuleBundle } from '#worker/package-runtime/module-graph.ts' import { runBundledModuleWithRegistry } from '#mcp/run-codemode-registry.ts' import { getEntitySourceById } from '#worker/repo/entity-sources.ts' +import { + logJobSchedulerEvent, + type SchedulerJobOutcomeLog, +} from './scheduler-logging.ts' function requirePersistableJobCallerContext( callerContext: McpCallerContext, @@ -179,12 +180,10 @@ async function executeBundledJobModule(input: { source: Awaited> modulePath: string bypassLogs: Array - packageContext?: - | { - packageId: string - kodyId: string - } - | null + packageContext?: { + packageId: string + kodyId: string + } | null }): Promise { try { const sourceFiles = await loadRepoSourceFilesFromSession({ @@ -230,7 +229,9 @@ async function executeBundledJobModule(input: { storageId: input.job.storageId, writable: true, }, - ...(input.packageContext ? { packageContext: input.packageContext } : {}), + ...(input.packageContext + ? { packageContext: input.packageContext } + : {}), }, ) return { @@ -740,7 +741,10 @@ async function runRepoBackedJob(input: { } let bypassLogs: Array = [] try { - const source = await getEntitySourceById(input.env.APP_DB, input.job.sourceId) + const source = await getEntitySourceById( + input.env.APP_DB, + input.job.sourceId, + ) let session = await sessionClient.openSession(openSessionInput) if (repoSessionNeedsRefresh(session)) { const stalePublishedCommit = session.published_commit @@ -977,17 +981,29 @@ export async function runDueJobsForUser(input: { userId: string now?: Date }) { + const now = input.now ?? new Date() const dueRows = await listDueJobRows( input.env.APP_DB, input.userId, - (input.now ?? new Date()).toISOString(), + now.toISOString(), ) if (dueRows.length === 0) { - return 0 + logJobSchedulerEvent({ + event: 'run_due_jobs.empty', + userId: input.userId, + dueJobCount: 0, + reason: 'no_due_jobs_found', + }) + return { + dueJobCount: 0, + successCount: 0, + errorCount: 0, + jobOutcomes: [] satisfies Array, + } } const result = await processDueJobs({ jobs: dueRows.map((row) => row.record), - now: input.now, + now, executeJob: async (job) => { const row = dueRows.find((candidate) => candidate.record.id === job.id) const callerContext = row?.callerContext ?? null @@ -1011,7 +1027,12 @@ export async function runDueJobsForUser(input: { await deleteJobRow(input.env.APP_DB, input.userId, jobId) await deleteJobVector(input.env, jobId) } - return dueRows.length + return { + dueJobCount: dueRows.length, + successCount: result.successCount, + errorCount: result.errorCount, + jobOutcomes: result.jobOutcomes, + } } export async function getNextRunnableJob(input: { env: Env; userId: string }) { diff --git a/packages/worker/src/mcp/capabilities/jobs/shared.ts b/packages/worker/src/mcp/capabilities/jobs/shared.ts index 875d8db321..2456d97ad2 100644 --- a/packages/worker/src/mcp/capabilities/jobs/shared.ts +++ b/packages/worker/src/mcp/capabilities/jobs/shared.ts @@ -1,6 +1,11 @@ import { z } from 'zod' import { requireMcpUser } from '#mcp/capabilities/meta/require-user.ts' import { type CapabilityContext } from '#mcp/capabilities/types.ts' +import { + logJobSchedulerError, + logJobSchedulerEvent, + schedulerErrorFields, +} from '#worker/jobs/scheduler-logging.ts' import { type JobCreateInput, type JobExecutionResult, @@ -55,7 +60,9 @@ export const scheduledJobInputBaseSchema = { params: z .record(z.string(), z.unknown()) .optional() - .describe('Optional JSON params passed to the job entrypoint when it runs.'), + .describe( + 'Optional JSON params passed to the job entrypoint when it runs.', + ), timezone: z .string() .min(1) @@ -294,10 +301,29 @@ export async function createScheduledJobFromArgs(input: { callerContext: input.callerContext, body: resolveJobCreateBody(input.args, input.defaultName), }) - await syncJobManagerAlarm({ - env: input.env, + logJobSchedulerEvent({ + event: 'job-created', userId: user.userId, + jobId: created.id, + scheduleType: created.schedule.type, + nextRunAt: created.nextRunAt, }) + try { + await syncJobManagerAlarm({ + env: input.env, + userId: user.userId, + }) + } catch (error) { + logJobSchedulerError({ + event: 'job-manager-sync-after-create-failed', + userId: user.userId, + jobId: created.id, + scheduleType: created.schedule.type, + nextRunAt: created.nextRunAt, + ...schedulerErrorFields(error), + }) + throw error + } return buildJobScheduleOutput(created) } From 7c5acebace3d97a004e624f2fab31f42e627c1b4 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 20 Apr 2026 22:51:15 +0000 Subject: [PATCH 2/6] Address scheduler review feedback Co-authored-by: Kent C. Dodds --- .../src/jobs/process-due-jobs.node.test.ts | 54 +++++++++ packages/worker/src/jobs/process-due-jobs.ts | 32 ++++-- .../src/jobs/scheduler-logging.node.test.ts | 103 ++++++++++++++++++ packages/worker/src/jobs/scheduler-logging.ts | 46 +++++++- 4 files changed, 226 insertions(+), 9 deletions(-) create mode 100644 packages/worker/src/jobs/scheduler-logging.node.test.ts diff --git a/packages/worker/src/jobs/process-due-jobs.node.test.ts b/packages/worker/src/jobs/process-due-jobs.node.test.ts index 43dc0504ac..63b26aa3ca 100644 --- a/packages/worker/src/jobs/process-due-jobs.node.test.ts +++ b/packages/worker/src/jobs/process-due-jobs.node.test.ts @@ -139,3 +139,57 @@ test('processDueJobs deletes one-shot jobs after execution', async () => { }, ]) }) + +test('processDueJobs treats reschedule failures as failed outcomes', async () => { + const cronJob = createCronJob({ + id: 'job-reschedule-failure', + schedule: { + type: 'cron', + expression: '* *', + }, + nextRunAt: '2026-04-12T07:00:00.000Z', + }) + + const result = await processDueJobs({ + jobs: [cronJob], + now: new Date('2026-04-12T07:00:00.000Z'), + async executeJob() { + return { + execution: { + ok: true, + logs: ['ok'], + result: { ok: true }, + }, + startedAt: '2026-04-12T07:00:00.000Z', + finishedAt: '2026-04-12T07:00:00.000Z', + durationMs: 0, + } + }, + }) + + expect(result.saveJobs).toHaveLength(1) + expect(result.successCount).toBe(0) + expect(result.errorCount).toBe(1) + expect(result.jobOutcomes).toEqual([ + { + jobId: 'job-reschedule-failure', + scheduleType: 'cron', + outcome: 'failure', + nextRunAt: '2026-04-12T07:00:00.000Z', + deleted: false, + error: + 'Cron expressions must use standard 5-field syntax: minute hour day-of-month month day-of-week.', + rescheduleError: + 'Cron expressions must use standard 5-field syntax: minute hour day-of-month month day-of-week.', + }, + ]) + expect(result.saveJobs[0]).toEqual( + expect.objectContaining({ + id: 'job-reschedule-failure', + enabled: false, + lastRunStatus: 'error', + lastRunError: + 'Failed to reschedule job: Cron expressions must use standard 5-field syntax: minute hour day-of-month month day-of-week.', + }), + ) +}) diff --git a/packages/worker/src/jobs/process-due-jobs.ts b/packages/worker/src/jobs/process-due-jobs.ts index ae75fde1cc..45fe140832 100644 --- a/packages/worker/src/jobs/process-due-jobs.ts +++ b/packages/worker/src/jobs/process-due-jobs.ts @@ -84,19 +84,37 @@ export async function processDueJobs(input: { }) } catch (error) { const rescheduleError = formatJobError(error) - const updated = applyExecutionOutcome(job, outcome, { - updatedAt: now.toISOString(), - enabled: false, - lastRunError: `Failed to reschedule job: ${rescheduleError}`, - }) + if (outcome.execution.ok) { + successCount -= 1 + errorCount += 1 + } + const failedRescheduleError = `Failed to reschedule job: ${rescheduleError}` + const updated = applyExecutionOutcome( + job, + outcome.execution.ok + ? { + ...outcome, + execution: { + ok: false as const, + error: failedRescheduleError, + logs: outcome.execution.logs, + }, + } + : outcome, + { + updatedAt: now.toISOString(), + enabled: false, + lastRunError: failedRescheduleError, + }, + ) saveJobs.push(updated) jobOutcomes.push({ jobId: job.id, scheduleType: job.schedule.type, - outcome: outcome.execution.ok ? 'success' : 'failure', + outcome: 'failure', nextRunAt: updated.nextRunAt, deleted: false, - ...(executionError ? { error: executionError } : {}), + error: executionError ?? rescheduleError, rescheduleError, }) } diff --git a/packages/worker/src/jobs/scheduler-logging.node.test.ts b/packages/worker/src/jobs/scheduler-logging.node.test.ts new file mode 100644 index 0000000000..b595834282 --- /dev/null +++ b/packages/worker/src/jobs/scheduler-logging.node.test.ts @@ -0,0 +1,103 @@ +import { expect, test } from 'vitest' +import { logJobSchedulerError, summarizeSchedulerJobOutcomes } from './scheduler-logging.ts' + +test('summarizeSchedulerJobOutcomes keeps full fields while limiting count', () => { + const longMessage = 'x'.repeat(1300) + + const result = summarizeSchedulerJobOutcomes([ + { + jobId: 'job-123', + scheduleType: 'once', + outcome: 'failure', + nextRunAt: null, + deleted: true, + error: longMessage, + rescheduleError: longMessage, + }, + ]) + + expect(result).toEqual({ + jobOutcomes: [ + { + jobId: 'job-123', + scheduleType: 'once', + outcome: 'failure', + nextRunAt: null, + deleted: true, + error: longMessage, + rescheduleError: longMessage, + }, + ], + }) +}) + +test('logJobSchedulerError truncates top-level errorMessage before logging', () => { + const originalError = console.error + let tagArg: unknown + let jsonArg: unknown + console.error = ((tag: unknown, json?: unknown) => { + tagArg = tag + jsonArg = json + }) as typeof console.error + + try { + logJobSchedulerError({ + event: 'sync_alarm_failed', + userId: 'user-123', + errorName: 'Error', + errorMessage: 'y'.repeat(1205), + }) + } finally { + console.error = originalError + } + + expect(tagArg).toBe('job-scheduler') + expect(typeof jsonArg).toBe('string') + const payload = JSON.parse(jsonArg as string) as Record + expect(payload.errorMessage).toBe( + `${'y'.repeat(1000)}...[truncated 205 chars]`, + ) +}) + +test('logJobSchedulerError truncates per-job error fields before logging', () => { + const originalError = console.error + let jsonArg: unknown + console.error = ((_tag: unknown, json?: unknown) => { + jsonArg = json + }) as typeof console.error + + try { + logJobSchedulerError({ + event: 'alarm_processed_due_jobs', + userId: 'user-123', + jobOutcomes: [ + { + jobId: 'job-123', + scheduleType: 'once', + outcome: 'failure', + nextRunAt: null, + deleted: true, + error: 'z'.repeat(1100), + rescheduleError: 'w'.repeat(1010), + }, + ], + }) + } finally { + console.error = originalError + } + + const payload = JSON.parse(jsonArg as string) as { + jobOutcomes: Array> + } + expect(payload.jobOutcomes).toEqual([ + { + jobId: 'job-123', + scheduleType: 'once', + outcome: 'failure', + nextRunAt: null, + deleted: true, + error: `${'z'.repeat(1000)}...[truncated 100 chars]`, + rescheduleError: `${'w'.repeat(1000)}...[truncated 10 chars]`, + }, + ]) +}) diff --git a/packages/worker/src/jobs/scheduler-logging.ts b/packages/worker/src/jobs/scheduler-logging.ts index 1eb2ab327d..8eecd102b6 100644 --- a/packages/worker/src/jobs/scheduler-logging.ts +++ b/packages/worker/src/jobs/scheduler-logging.ts @@ -2,6 +2,7 @@ import { formatJobError } from './schedule.ts' import { type JobSchedule } from './types.ts' const maxLoggedJobOutcomes = 10 +const maxLoggedStringLength = 1_000 type SchedulerLogLevel = 'error' | 'info' export type SchedulerJobOutcomeLog = { @@ -78,11 +79,12 @@ function writeSchedulerLog( input: JobSchedulerLogPayload, ): void { try { + const payload = sanitizeSchedulerLogPayload(input) console[level]( 'job-scheduler', JSON.stringify({ - timestamp: input.timestamp ?? new Date().toISOString(), - ...input, + timestamp: payload.timestamp ?? new Date().toISOString(), + ...payload, }), ) } catch (error) { @@ -93,3 +95,43 @@ function writeSchedulerLog( }) } } + +function truncateLoggedString(value: string): string { + if (value.length <= maxLoggedStringLength) { + return value + } + + return `${value.slice(0, maxLoggedStringLength)}...[truncated ${value.length - maxLoggedStringLength} chars]` +} + +function sanitizeSchedulerJobOutcome( + jobOutcome: SchedulerJobOutcomeLog, +): SchedulerJobOutcomeLog { + return { + ...jobOutcome, + ...(jobOutcome.error + ? { error: truncateLoggedString(jobOutcome.error) } + : {}), + ...(jobOutcome.rescheduleError + ? { + rescheduleError: truncateLoggedString(jobOutcome.rescheduleError), + } + : {}), + } +} + +function sanitizeSchedulerLogPayload( + input: JobSchedulerLogPayload, +): JobSchedulerLogPayload { + return { + ...input, + ...(input.errorMessage + ? { errorMessage: truncateLoggedString(input.errorMessage) } + : {}), + ...(input.jobOutcomes + ? { + jobOutcomes: input.jobOutcomes.map(sanitizeSchedulerJobOutcome), + } + : {}), + } +} From 8060d221bbad6bb2513e9e760a5590d1bdc8e9b6 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 21 Apr 2026 01:42:48 +0000 Subject: [PATCH 3/6] Fix scheduler log timestamp fallback Co-authored-by: Kent C. Dodds --- .../src/jobs/scheduler-logging.node.test.ts | 29 ++++++++++++++++++- packages/worker/src/jobs/scheduler-logging.ts | 2 +- 2 files changed, 29 insertions(+), 2 deletions(-) diff --git a/packages/worker/src/jobs/scheduler-logging.node.test.ts b/packages/worker/src/jobs/scheduler-logging.node.test.ts index b595834282..771d4cbdac 100644 --- a/packages/worker/src/jobs/scheduler-logging.node.test.ts +++ b/packages/worker/src/jobs/scheduler-logging.node.test.ts @@ -1,5 +1,8 @@ import { expect, test } from 'vitest' -import { logJobSchedulerError, summarizeSchedulerJobOutcomes } from './scheduler-logging.ts' +import { + logJobSchedulerError, + summarizeSchedulerJobOutcomes, +} from './scheduler-logging.ts' test('summarizeSchedulerJobOutcomes keeps full fields while limiting count', () => { const longMessage = 'x'.repeat(1300) @@ -101,3 +104,27 @@ test('logJobSchedulerError truncates per-job error fields before logging', () => }, ]) }) + +test('logJobSchedulerError always emits a timestamp when input timestamp is undefined', () => { + const originalError = console.error + let jsonArg: unknown + console.error = ((_tag: unknown, json?: unknown) => { + jsonArg = json + }) as typeof console.error + + try { + logJobSchedulerError({ + event: 'sync_alarm_failed', + userId: 'user-123', + errorName: 'Error', + errorMessage: 'boom', + timestamp: undefined, + }) + } finally { + console.error = originalError + } + + const payload = JSON.parse(jsonArg as string) as Record + expect(typeof payload.timestamp).toBe('string') + expect(payload.timestamp).not.toBe('') +}) diff --git a/packages/worker/src/jobs/scheduler-logging.ts b/packages/worker/src/jobs/scheduler-logging.ts index 8eecd102b6..78843b7ad0 100644 --- a/packages/worker/src/jobs/scheduler-logging.ts +++ b/packages/worker/src/jobs/scheduler-logging.ts @@ -83,8 +83,8 @@ function writeSchedulerLog( console[level]( 'job-scheduler', JSON.stringify({ - timestamp: payload.timestamp ?? new Date().toISOString(), ...payload, + timestamp: payload.timestamp ?? new Date().toISOString(), }), ) } catch (error) { From 5871a7375c69260011b747885634b8e79a064c57 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 21 Apr 2026 02:36:16 +0000 Subject: [PATCH 4/6] Normalize scheduler diagnostic logging Co-authored-by: Kent C. Dodds --- packages/worker/src/jobs/manager-client.ts | 6 +- .../worker/src/jobs/manager-do.node.test.ts | 159 +++++++++++++++++- packages/worker/src/jobs/manager-do.ts | 23 ++- .../src/jobs/scheduler-logging.node.test.ts | 100 +++++------ packages/worker/src/jobs/scheduler-logging.ts | 3 + packages/worker/src/jobs/service.ts | 11 +- 6 files changed, 223 insertions(+), 79 deletions(-) diff --git a/packages/worker/src/jobs/manager-client.ts b/packages/worker/src/jobs/manager-client.ts index 7a5408cf62..0c182465c3 100644 --- a/packages/worker/src/jobs/manager-client.ts +++ b/packages/worker/src/jobs/manager-client.ts @@ -11,7 +11,10 @@ import { } from './types.ts' type JobManagerRpc = { - syncAlarm: (payload: { userId: string }) => Promise<{ + syncAlarm: (payload: { + userId: string + source?: 'alarm' | 'rpc' | 'run_now' + }) => Promise<{ ok: true userId: string nextRunAt: string | null @@ -56,6 +59,7 @@ export async function syncJobManagerAlarm(input: { env: Env; userId: string }) { try { const result = await rpc.syncAlarm({ userId: input.userId, + source: 'rpc', }) logJobSchedulerEvent({ event: 'sync_alarm_completed', diff --git a/packages/worker/src/jobs/manager-do.node.test.ts b/packages/worker/src/jobs/manager-do.node.test.ts index af249bccf0..eb8b54758d 100644 --- a/packages/worker/src/jobs/manager-do.node.test.ts +++ b/packages/worker/src/jobs/manager-do.node.test.ts @@ -1,3 +1,4 @@ +import type * as SchedulerLoggingType from './scheduler-logging.ts' import { expect, test, vi } from 'vitest' const mockModule = vi.hoisted(() => ({ @@ -42,7 +43,7 @@ vi.mock('#worker/sentry-options.ts', () => ({ })) vi.mock('./scheduler-logging.ts', async (importOriginal) => { - const actual = await importOriginal() + const actual = await importOriginal() return { ...actual, logJobSchedulerEvent: (...args: Array) => @@ -67,7 +68,7 @@ function createState({ userId = 'user-123', currentAlarmAt = null, }: { - userId?: string | null + userId?: string currentAlarmAt?: number | null } = {}) { const persistedEntries = new Map() @@ -125,7 +126,7 @@ test('syncAlarm logs when it arms a new alarm for the next runnable job', async currentAlarmAt: '2026-04-20T18:00:00.000Z', nextJobId: 'job-123', nextRunAt, - reason: 'alarm-armed', + reason: 'alarm_armed', }) expect(mockModule.logJobSchedulerError).not.toHaveBeenCalled() }) @@ -151,7 +152,7 @@ test('syncAlarm logs when no runnable job is found and clears the alarm', async currentAlarmAt: '2026-04-20T18:00:00.000Z', nextJobId: null, nextRunAt: null, - reason: 'no-runnable-job', + reason: 'no_runnable_job', }) }) @@ -204,12 +205,12 @@ test('alarm logs firing, due-job outcomes, and resyncs the next alarm', async () isRetry: true, }) expect(mockModule.logJobSchedulerEvent).toHaveBeenNthCalledWith(2, { - event: 'alarm_processed_due_jobs', + event: 'run_due_jobs_completed', userId: 'user-123', dueJobCount: 2, successCount: 1, errorCount: 1, - reason: 'processed-due-jobs', + reason: 'processed_due_jobs', jobOutcomes: [ { jobId: 'job-success', @@ -234,7 +235,151 @@ test('alarm logs firing, due-job outcomes, and resyncs the next alarm', async () currentAlarmAt: null, nextJobId: 'job-next', nextRunAt: '2026-04-20T19:00:00.000Z', - reason: 'alarm-armed', + reason: 'alarm_armed', }) expect(mockModule.logJobSchedulerError).not.toHaveBeenCalled() }) + +test('syncAlarm logs source-tagged errors when getNextRunnableJob fails', async () => { + resetMocks() + mockModule.getNextRunnableJob.mockRejectedValue( + new Error('next job lookup failed'), + ) + const { state } = createState() + const manager = new JobManagerBase(state, {} as Env) + + await expect( + manager.syncAlarm({ userId: 'user-123', source: 'alarm' }), + ).rejects.toThrow('next job lookup failed') + + expect(mockModule.logJobSchedulerError).toHaveBeenCalledWith({ + event: 'sync_alarm_failed', + userId: 'user-123', + source: 'alarm', + errorName: 'Error', + errorMessage: 'next job lookup failed', + }) +}) + +test('syncAlarm logs source-tagged errors when deleteAlarm fails', async () => { + resetMocks() + mockModule.getNextRunnableJob.mockResolvedValue(null) + const { state } = createState() + const deleteAlarm = vi.mocked( + state.storage.deleteAlarm as unknown as ( + ...args: Array + ) => Promise, + ) + deleteAlarm.mockRejectedValueOnce(new Error('delete alarm failed')) + const manager = new JobManagerBase(state, {} as Env) + + await expect( + manager.syncAlarm({ userId: 'user-123', source: 'rpc' }), + ).rejects.toThrow('delete alarm failed') + + expect(mockModule.logJobSchedulerError).toHaveBeenCalledWith({ + event: 'sync_alarm_failed', + userId: 'user-123', + source: 'rpc', + errorName: 'Error', + errorMessage: 'delete alarm failed', + }) +}) + +test('syncAlarm logs source-tagged errors when setAlarm fails', async () => { + resetMocks() + mockModule.getNextRunnableJob.mockResolvedValue({ + id: 'job-123', + nextRunAt: '2026-04-20T18:30:00.000Z', + }) + const { state } = createState() + const setAlarm = vi.mocked( + state.storage.setAlarm as unknown as ( + ...args: Array + ) => Promise, + ) + setAlarm.mockRejectedValueOnce(new Error('set alarm failed')) + const manager = new JobManagerBase(state, {} as Env) + + await expect( + manager.syncAlarm({ userId: 'user-123', source: 'run_now' }), + ).rejects.toThrow('set alarm failed') + + expect(mockModule.logJobSchedulerError).toHaveBeenCalledWith({ + event: 'sync_alarm_failed', + userId: 'user-123', + source: 'run_now', + errorName: 'Error', + errorMessage: 'set alarm failed', + }) +}) + +test('alarm logs missing_user_id when no user id was persisted', async () => { + resetMocks() + const { state } = createState() + vi.mocked( + state.storage.get as unknown as ( + key: string, + ) => Promise, + ).mockResolvedValueOnce(undefined) + const manager = new JobManagerBase(state, {} as Env) + + await expect(manager.alarm()).resolves.toBeUndefined() + + expect(mockModule.logJobSchedulerEvent).toHaveBeenCalledWith({ + event: 'alarm_fired', + reason: 'missing_user_id', + retryCount: undefined, + isRetry: undefined, + }) +}) + +test('alarm logs run_due_jobs failure details', async () => { + resetMocks() + mockModule.runDueJobsForUser.mockRejectedValue( + new Error('run due jobs failed'), + ) + const { state } = createState() + const manager = new JobManagerBase(state, {} as Env) + + await expect(manager.alarm()).rejects.toThrow('run due jobs failed') + + expect(mockModule.logJobSchedulerError).toHaveBeenCalledWith({ + event: 'alarm_run_due_jobs_failed', + userId: 'user-123', + retryCount: undefined, + isRetry: undefined, + errorName: 'Error', + errorMessage: 'run due jobs failed', + }) +}) + +test('alarm logs resync failure details after due jobs run', async () => { + resetMocks() + mockModule.runDueJobsForUser.mockResolvedValue({ + dueJobCount: 0, + successCount: 0, + errorCount: 0, + jobOutcomes: [], + }) + const { state } = createState() + const manager = new JobManagerBase(state, {} as Env) + const syncAlarmSpy = vi + .spyOn(manager, 'syncAlarm') + .mockRejectedValueOnce(new Error('resync failed')) + + await expect(manager.alarm()).rejects.toThrow('resync failed') + + expect(syncAlarmSpy).toHaveBeenCalledWith({ + userId: 'user-123', + source: 'alarm', + }) + expect(mockModule.logJobSchedulerError).toHaveBeenCalledWith({ + event: 'alarm_resync_failed', + userId: 'user-123', + retryCount: undefined, + isRetry: undefined, + errorName: 'Error', + errorMessage: 'resync failed', + }) +}) diff --git a/packages/worker/src/jobs/manager-do.ts b/packages/worker/src/jobs/manager-do.ts index 90f43a45ba..2abdcfe915 100644 --- a/packages/worker/src/jobs/manager-do.ts +++ b/packages/worker/src/jobs/manager-do.ts @@ -14,7 +14,10 @@ import { type JobRepoCheckPolicy } from './types.ts' const userIdStorageKey = 'user-id' export class JobManagerBase extends DurableObject { - async syncAlarm(input: { userId: string }) { + async syncAlarm(input: { + userId: string + source?: 'alarm' | 'rpc' | 'run_now' + }) { const userId = input.userId.trim() if (!userId) { throw new Error('Job manager requires a non-empty userId.') @@ -37,7 +40,7 @@ export class JobManagerBase extends DurableObject { : new Date(currentAlarmAt).toISOString(), nextJobId: null, nextRunAt: null, - reason: 'no-runnable-job', + reason: 'no_runnable_job', }) return { ok: true as const, @@ -57,8 +60,8 @@ export class JobManagerBase extends DurableObject { nextRunAt: nextJob.nextRunAt, reason: currentAlarmAt === new Date(nextJob.nextRunAt).valueOf() - ? 'alarm-unchanged' - : 'alarm-armed', + ? 'alarm_unchanged' + : 'alarm_armed', }) return { ok: true as const, @@ -69,6 +72,7 @@ export class JobManagerBase extends DurableObject { logJobSchedulerError({ event: 'sync_alarm_failed', userId, + source: input.source ?? 'rpc', ...schedulerErrorFields(error), }) throw error @@ -84,7 +88,7 @@ export class JobManagerBase extends DurableObject { await this.ctx.storage.deleteAlarm() logJobSchedulerEvent({ event: 'alarm_fired', - reason: 'missing-user-id', + reason: 'missing_user_id', retryCount: alarmInfo?.retryCount, isRetry: alarmInfo?.isRetry, }) @@ -102,12 +106,13 @@ export class JobManagerBase extends DurableObject { userId, }) logJobSchedulerEvent({ - event: 'alarm_processed_due_jobs', + event: 'run_due_jobs_completed', userId, dueJobCount: result.dueJobCount, successCount: result.successCount, errorCount: result.errorCount, - reason: result.dueJobCount === 0 ? 'no-due-jobs' : 'processed-due-jobs', + reason: + result.dueJobCount === 0 ? 'no_due_jobs_found' : 'processed_due_jobs', ...summarizeSchedulerJobOutcomes(result.jobOutcomes), }) } catch (error) { @@ -121,7 +126,7 @@ export class JobManagerBase extends DurableObject { throw error } try { - await this.syncAlarm({ userId }) + await this.syncAlarm({ userId, source: 'alarm' }) } catch (error) { logJobSchedulerError({ event: 'alarm_resync_failed', @@ -154,7 +159,7 @@ export class JobManagerBase extends DurableObject { originalError = error } try { - await this.syncAlarm({ userId: input.userId }) + await this.syncAlarm({ userId: input.userId, source: 'run_now' }) } catch (syncError) { console.error('[JobManager.runNow] failed to sync job alarm', { userId: input.userId, diff --git a/packages/worker/src/jobs/scheduler-logging.node.test.ts b/packages/worker/src/jobs/scheduler-logging.node.test.ts index 771d4cbdac..ffd665b447 100644 --- a/packages/worker/src/jobs/scheduler-logging.node.test.ts +++ b/packages/worker/src/jobs/scheduler-logging.node.test.ts @@ -1,9 +1,13 @@ -import { expect, test } from 'vitest' +import { afterEach, expect, test, vi } from 'vitest' import { logJobSchedulerError, summarizeSchedulerJobOutcomes, } from './scheduler-logging.ts' +afterEach(() => { + vi.restoreAllMocks() +}) + test('summarizeSchedulerJobOutcomes keeps full fields while limiting count', () => { const longMessage = 'x'.repeat(1300) @@ -35,25 +39,17 @@ test('summarizeSchedulerJobOutcomes keeps full fields while limiting count', () }) test('logJobSchedulerError truncates top-level errorMessage before logging', () => { - const originalError = console.error - let tagArg: unknown - let jsonArg: unknown - console.error = ((tag: unknown, json?: unknown) => { - tagArg = tag - jsonArg = json - }) as typeof console.error + const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => {}) - try { - logJobSchedulerError({ - event: 'sync_alarm_failed', - userId: 'user-123', - errorName: 'Error', - errorMessage: 'y'.repeat(1205), - }) - } finally { - console.error = originalError - } + logJobSchedulerError({ + event: 'sync_alarm_failed', + userId: 'user-123', + errorName: 'Error', + errorMessage: 'y'.repeat(1205), + }) + expect(errorSpy).toHaveBeenCalledTimes(1) + const [tagArg, jsonArg] = errorSpy.mock.calls[0] ?? [] expect(tagArg).toBe('job-scheduler') expect(typeof jsonArg).toBe('string') const payload = JSON.parse(jsonArg as string) as Record @@ -63,32 +59,26 @@ test('logJobSchedulerError truncates top-level errorMessage before logging', () }) test('logJobSchedulerError truncates per-job error fields before logging', () => { - const originalError = console.error - let jsonArg: unknown - console.error = ((_tag: unknown, json?: unknown) => { - jsonArg = json - }) as typeof console.error + const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => {}) - try { - logJobSchedulerError({ - event: 'alarm_processed_due_jobs', - userId: 'user-123', - jobOutcomes: [ - { - jobId: 'job-123', - scheduleType: 'once', - outcome: 'failure', - nextRunAt: null, - deleted: true, - error: 'z'.repeat(1100), - rescheduleError: 'w'.repeat(1010), - }, - ], - }) - } finally { - console.error = originalError - } + logJobSchedulerError({ + event: 'alarm_processed_due_jobs', + userId: 'user-123', + jobOutcomes: [ + { + jobId: 'job-123', + scheduleType: 'once', + outcome: 'failure', + nextRunAt: null, + deleted: true, + error: 'z'.repeat(1100), + rescheduleError: 'w'.repeat(1010), + }, + ], + }) + expect(errorSpy).toHaveBeenCalledTimes(1) + const [, jsonArg] = errorSpy.mock.calls[0] ?? [] const payload = JSON.parse(jsonArg as string) as { jobOutcomes: Array> } @@ -106,24 +96,18 @@ test('logJobSchedulerError truncates per-job error fields before logging', () => }) test('logJobSchedulerError always emits a timestamp when input timestamp is undefined', () => { - const originalError = console.error - let jsonArg: unknown - console.error = ((_tag: unknown, json?: unknown) => { - jsonArg = json - }) as typeof console.error + const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => {}) - try { - logJobSchedulerError({ - event: 'sync_alarm_failed', - userId: 'user-123', - errorName: 'Error', - errorMessage: 'boom', - timestamp: undefined, - }) - } finally { - console.error = originalError - } + logJobSchedulerError({ + event: 'sync_alarm_failed', + userId: 'user-123', + errorName: 'Error', + errorMessage: 'boom', + timestamp: undefined, + }) + expect(errorSpy).toHaveBeenCalledTimes(1) + const [, jsonArg] = errorSpy.mock.calls[0] ?? [] const payload = JSON.parse(jsonArg as string) as Record expect(typeof payload.timestamp).toBe('string') expect(payload.timestamp).not.toBe('') diff --git a/packages/worker/src/jobs/scheduler-logging.ts b/packages/worker/src/jobs/scheduler-logging.ts index 78843b7ad0..9ef3ced300 100644 --- a/packages/worker/src/jobs/scheduler-logging.ts +++ b/packages/worker/src/jobs/scheduler-logging.ts @@ -15,11 +15,14 @@ export type SchedulerJobOutcomeLog = { rescheduleError?: string } +export type SchedulerLogSource = 'alarm' | 'direct' | 'rpc' | 'run_now' + type JobSchedulerLogPayload = { event: string userId?: string jobId?: string | null scheduleType?: JobSchedule['type'] + source?: SchedulerLogSource nextJobId?: string | null nextRunAt?: string | null currentAlarmAt?: string | null diff --git a/packages/worker/src/jobs/service.ts b/packages/worker/src/jobs/service.ts index d6a47642f3..d30961df5b 100644 --- a/packages/worker/src/jobs/service.ts +++ b/packages/worker/src/jobs/service.ts @@ -987,12 +987,15 @@ export async function runDueJobsForUser(input: { input.userId, now.toISOString(), ) + const dueRowById = new Map( + dueRows.map((row) => [row.record.id, row] as const), + ) if (dueRows.length === 0) { logJobSchedulerEvent({ - event: 'run_due_jobs.empty', + event: 'run_due_jobs_empty', userId: input.userId, dueJobCount: 0, - reason: 'no_due_jobs_found', + reason: 'no_due_jobs', }) return { dueJobCount: 0, @@ -1005,7 +1008,7 @@ export async function runDueJobsForUser(input: { jobs: dueRows.map((row) => row.record), now, executeJob: async (job) => { - const row = dueRows.find((candidate) => candidate.record.id === job.id) + const row = dueRowById.get(job.id) const callerContext = row?.callerContext ?? null return executeJobOnce({ env: input.env, @@ -1015,7 +1018,7 @@ export async function runDueJobsForUser(input: { }, }) for (const job of result.saveJobs) { - const row = dueRows.find((candidate) => candidate.record.id === job.id) + const row = dueRowById.get(job.id) await updateJobRow({ db: input.env.APP_DB, userId: input.userId, From dc2230e00d2cf907d65a2da5e9e862cbe7e75022 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 21 Apr 2026 02:56:01 +0000 Subject: [PATCH 5/6] Finalize merged scheduler diagnostics fixes Co-authored-by: Kent C. Dodds --- packages/worker/src/jobs/manager-do.ts | 10 +++++----- packages/worker/src/mcp/capabilities/jobs/shared.ts | 4 ++-- 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/packages/worker/src/jobs/manager-do.ts b/packages/worker/src/jobs/manager-do.ts index ee16dfbe21..2d062ef647 100644 --- a/packages/worker/src/jobs/manager-do.ts +++ b/packages/worker/src/jobs/manager-do.ts @@ -10,10 +10,8 @@ import { summarizeSchedulerJobOutcomes, } from './scheduler-logging.ts' import { resolveJobManagerAlarmState } from './manager-state.ts' -import { - type JobManagerDebugState, - type JobRepoCheckPolicy, -} from './types.ts' +import { type JobRepoCheckPolicy } from './types.ts' +import { type JobManagerDebugState } from './manager-client.ts' const userIdStorageKey = 'user-id' @@ -83,7 +81,9 @@ export class JobManagerBase extends DurableObject { } } - async getDebugState(input: { userId: string }): Promise { + async getDebugState(input: { + userId: string + }): Promise { const userId = input.userId.trim() if (!userId) { throw new Error('Job manager requires a non-empty userId.') diff --git a/packages/worker/src/mcp/capabilities/jobs/shared.ts b/packages/worker/src/mcp/capabilities/jobs/shared.ts index b2d6c452c5..3fe3753056 100644 --- a/packages/worker/src/mcp/capabilities/jobs/shared.ts +++ b/packages/worker/src/mcp/capabilities/jobs/shared.ts @@ -421,7 +421,7 @@ export async function createScheduledJobFromArgs(input: { body: resolveJobCreateBody(input.args, input.defaultName), }) logJobSchedulerEvent({ - event: 'job-created', + event: 'job_created', userId: user.userId, jobId: created.id, scheduleType: created.schedule.type, @@ -434,7 +434,7 @@ export async function createScheduledJobFromArgs(input: { }) } catch (error) { logJobSchedulerError({ - event: 'job-manager-sync-after-create-failed', + event: 'job_manager_sync_after_create_failed', userId: user.userId, jobId: created.id, scheduleType: created.schedule.type, From ac496f8389cfb6009c57003e27a230df4e26eb0b Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 21 Apr 2026 03:37:16 +0000 Subject: [PATCH 6/6] Remove duplicate job manager sync Co-authored-by: Kent C. Dodds --- .../jobs/job-schedule.node.test.ts | 12 +--------- .../src/mcp/capabilities/jobs/shared.ts | 23 +------------------ 2 files changed, 2 insertions(+), 33 deletions(-) diff --git a/packages/worker/src/mcp/capabilities/jobs/job-schedule.node.test.ts b/packages/worker/src/mcp/capabilities/jobs/job-schedule.node.test.ts index 8c351407ff..5b3f7779b9 100644 --- a/packages/worker/src/mcp/capabilities/jobs/job-schedule.node.test.ts +++ b/packages/worker/src/mcp/capabilities/jobs/job-schedule.node.test.ts @@ -6,7 +6,6 @@ const mockModule = vi.hoisted(() => ({ createJob: vi.fn(), getJobInspection: vi.fn(), inspectJobsForUser: vi.fn(), - syncJobManagerAlarm: vi.fn(), runJobNowViaManager: vi.fn(), })) @@ -19,8 +18,6 @@ vi.mock('#worker/jobs/service.ts', () => ({ })) vi.mock('#worker/jobs/manager-client.ts', () => ({ - syncJobManagerAlarm: (...args: Array) => - mockModule.syncJobManagerAlarm(...args), runJobNowViaManager: (...args: Array) => mockModule.runJobNowViaManager(...args), })) @@ -35,12 +32,10 @@ function resetMocks() { mockModule.createJob.mockReset() mockModule.getJobInspection.mockReset() mockModule.inspectJobsForUser.mockReset() - mockModule.syncJobManagerAlarm.mockReset() mockModule.runJobNowViaManager.mockReset() - mockModule.syncJobManagerAlarm.mockResolvedValue(undefined) } -test('job_schedule creates a one-off job and syncs the job manager alarm', async () => { +test('job_schedule creates a one-off job', async () => { resetMocks() const env = {} as Env const callerContext = createMcpCallerContext({ @@ -104,10 +99,6 @@ test('job_schedule creates a one-off job and syncs the job manager alarm', async timezone: 'America/Denver', }, }) - expect(mockModule.syncJobManagerAlarm).toHaveBeenCalledWith({ - env, - userId: 'user-123', - }) expect(result).toEqual({ job_id: 'job-123', name: 'Turn lights off', @@ -538,7 +529,6 @@ test('job_schedule requires an authenticated user', async () => { ), ).rejects.toThrow('Authenticated MCP user is required for this capability.') expect(mockModule.createJob).not.toHaveBeenCalled() - expect(mockModule.syncJobManagerAlarm).not.toHaveBeenCalled() }) test('job_list returns inspectable jobs plus alarm state', async () => { diff --git a/packages/worker/src/mcp/capabilities/jobs/shared.ts b/packages/worker/src/mcp/capabilities/jobs/shared.ts index 3fe3753056..7ef90260b0 100644 --- a/packages/worker/src/mcp/capabilities/jobs/shared.ts +++ b/packages/worker/src/mcp/capabilities/jobs/shared.ts @@ -2,11 +2,7 @@ import { z } from 'zod' import { requireMcpUser } from '#mcp/capabilities/meta/require-user.ts' import { type CapabilityContext } from '#mcp/capabilities/types.ts' import { type JobManagerDebugState } from '#worker/jobs/manager-client.ts' -import { - logJobSchedulerError, - logJobSchedulerEvent, - schedulerErrorFields, -} from '#worker/jobs/scheduler-logging.ts' +import { logJobSchedulerEvent } from '#worker/jobs/scheduler-logging.ts' import { type JobCreateInput, type JobExecutionResult, @@ -414,7 +410,6 @@ export async function createScheduledJobFromArgs(input: { // Delay job runtime imports so capability registration can load without // recursively pulling the full jobs runtime back through the registry. const { createJob } = await import('#worker/jobs/service.ts') - const { syncJobManagerAlarm } = await import('#worker/jobs/manager-client.ts') const created = await createJob({ env: input.env, callerContext: input.callerContext, @@ -427,22 +422,6 @@ export async function createScheduledJobFromArgs(input: { scheduleType: created.schedule.type, nextRunAt: created.nextRunAt, }) - try { - await syncJobManagerAlarm({ - env: input.env, - userId: user.userId, - }) - } catch (error) { - logJobSchedulerError({ - event: 'job_manager_sync_after_create_failed', - userId: user.userId, - jobId: created.id, - scheduleType: created.schedule.type, - nextRunAt: created.nextRunAt, - ...schedulerErrorFields(error), - }) - throw error - } return buildJobScheduleOutput(created) }