diff --git a/src/servers/sidecar/email-cron.ts b/src/servers/sidecar/email-cron.ts new file mode 100644 index 00000000..1e6dd86f --- /dev/null +++ b/src/servers/sidecar/email-cron.ts @@ -0,0 +1,87 @@ +import { getAllSyncedAccounts, getUserById, getUserIntegration } from 'officerdb'; +import * as queueRunner from './queue-runner'; + +const INTERVAL_MS = 10 * 60 * 1000; // 10 minutes + +let timer: ReturnType | null = null; + +async function tick() { + try { + const accounts = await getAllSyncedAccounts(); + if (accounts.length === 0) return; + + const allJobs = await queueRunner.listAllJobs(); + const activeEmailSyncIds = new Set( + allJobs + .filter((j) => j.type === 'email-sync' && (j.status === 'queued' || j.status === 'running')) + .map((j) => (j.meta as Record | undefined)?.emailAccountId), + ); + + for (const account of accounts) { + if (activeEmailSyncIds.has(account.id)) continue; + + const user = await getUserById(account.userId); + if (!user) continue; + + // Resolve IMAP auth + const imapAuth: Record = { user: account.email }; + if (account.authType === 'oauth') { + const integration = await getUserIntegration(account.userId, 'google'); + const config = integration?.config as Record | undefined; + const accessToken = config?.accessToken as string | undefined; + if (!accessToken) { + console.log(`[email-cron] Skipping ${account.email}: no OAuth access token`); + continue; + } + imapAuth.accessToken = accessToken; + } else { + const creds = account.credentials as Record; + imapAuth.pass = creds.password; + } + + try { + await queueRunner.enqueue({ + lane: 'email', + type: 'email-sync', + userId: user.email, + meta: { + emailAccountId: account.id, + userEmail: user.email, + account: { + id: account.id, + userId: account.userId, + email: account.email, + imapHost: account.imapHost, + imapPort: account.imapPort, + imapSecure: account.imapSecure, + provider: account.provider, + authType: account.authType, + credentials: account.credentials, + }, + imapAuth, + }, + }); + console.log(`[email-cron] Enqueued incremental sync for ${account.email}`); + } catch (err) { + console.error(`[email-cron] Failed to enqueue sync for ${account.email}:`, err instanceof Error ? err.message : err); + } + } + } catch (err) { + console.error('[email-cron] Error:', err instanceof Error ? err.message : err); + } +} + +export function initEmailCron() { + if (timer) return; + console.log(`[email-cron] Starting email sync cron (every ${INTERVAL_MS / 60_000} min)`); + timer = setInterval(tick, INTERVAL_MS); + // Run first tick after a short delay to let the queue initialize + setTimeout(tick, 30_000); +} + +export function stopEmailCron() { + if (timer) { + clearInterval(timer); + timer = null; + } +} diff --git a/src/servers/sidecar/queue-runner.ts b/src/servers/sidecar/queue-runner.ts new file mode 100644 index 00000000..d2811387 --- /dev/null +++ b/src/servers/sidecar/queue-runner.ts @@ -0,0 +1,285 @@ +import { type Job, type JobProgress, type EnqueueParams, type StepContext, PermanentError } from '../queue/types'; +import { readJob, writeJob, listAllJobs, ensureQueueDir } from '../queue/storage'; +import { getHandler } from '../queue/handler-registry'; + +// Import handlers to register them +import '../queue/handlers'; + +function formatDuration(ms: number): string { + const s = Math.floor(ms / 1000); + if (s < 60) return `${s}s`; + const m = Math.floor(s / 60); + const rem = s % 60; + if (m < 60) return rem > 0 ? `${m}m${rem}s` : `${m}m`; + const h = Math.floor(m / 60); + const remM = m % 60; + return remM > 0 ? `${h}h${remM}m` : `${h}h`; +} + +const activeLanes = new Map(); +const PROGRESS_THROTTLE_MS = 1000; + +export async function initQueue() { + await ensureQueueDir(); + await resumeInterruptedJobs(); + console.log('[sidecar:queue] initialized'); +} + +export async function enqueue(params: EnqueueParams): Promise { + const handler = getHandler(params.type); + if (!handler) throw new Error(`No handler registered for job type: ${params.type}`); + + const job: Job = { + id: crypto.randomUUID(), + lane: params.lane, + type: params.type, + userId: params.userId, + status: 'queued', + steps: handler.steps.map((s) => ({ name: s.name, status: 'pending' as const })), + currentStep: 0, + createdAt: Date.now(), + meta: params.meta, + notify: params.notify, + }; + + await writeJob(job); + console.log(`[sidecar:queue] enqueued job ${job.id} (${job.type}) in lane ${job.lane}`); + kickLane(job.lane); + return job; +} + +export async function cancelJob(id: string): Promise { + const job = await readJob(id); + if (!job) return null; + if (job.status === 'completed' || job.status === 'failed' || job.status === 'cancelled') return job; + + job.status = 'cancelled'; + job.completedAt = Date.now(); + for (const step of job.steps) { + if (step.status === 'pending' || step.status === 'running') { + step.status = 'failed'; + step.error = 'Cancelled'; + } + } + await writeJob(job); + console.log(`[sidecar:queue] cancelled job ${job.id}`); + return job; +} + +async function resumeInterruptedJobs() { + const jobs = await listAllJobs(); + const lanesToKick = new Set(); + + for (const job of jobs) { + if (job.status === 'running') { + job.status = 'queued'; + job.startedAt = undefined; + for (const step of job.steps) { + if (step.status === 'running') { + step.status = 'pending'; + step.startedAt = undefined; + } + } + await writeJob(job); + console.log(`[sidecar:queue] reset interrupted job ${job.id} back to queued`); + lanesToKick.add(job.lane); + } else if (job.status === 'queued') { + // Clear retry delay on restart — no reason to wait after a sidecar restart + if (job.retryAt) { + job.retryAt = undefined; + await writeJob(job); + console.log(`[sidecar:queue] cleared retry delay for job ${job.id}`); + } + lanesToKick.add(job.lane); + } + } + + for (const lane of lanesToKick) { + kickLane(lane); + } +} + +function kickLane(lane: string) { + if (activeLanes.get(lane)) return; + activeLanes.set(lane, true); + processNextInLane(lane); +} + +function scheduleRetry(lane: string, delayMs: number) { + setTimeout(() => kickLane(lane), delayMs); +} + +async function processNextInLane(lane: string) { + try { + const jobs = await listAllJobs(); + const now = Date.now(); + const next = jobs + .filter((j) => j.lane === lane && j.status === 'queued' && (!j.retryAt || j.retryAt <= now)) + .sort((a, b) => a.createdAt - b.createdAt)[0]; + + if (!next) { + activeLanes.set(lane, false); + return; + } + + await runJob(next); + } catch (err) { + console.error(`[sidecar:queue] lane ${lane} processing error:`, err); + } finally { + const jobs = await listAllJobs(); + const hasMore = jobs.some( + (j) => j.lane === lane && j.status === 'queued' && (!j.retryAt || j.retryAt <= Date.now()), + ); + if (hasMore) { + processNextInLane(lane); + } else { + activeLanes.set(lane, false); + } + } +} + +async function runJob(job: Job) { + const handler = getHandler(job.type); + if (!handler) { + job.status = 'failed'; + job.error = `No handler for type: ${job.type}`; + job.completedAt = Date.now(); + await writeJob(job); + return; + } + + job.status = 'running'; + job.startedAt = Date.now(); + job.retryAt = undefined; + await writeJob(job); + const isRetry = (job.retries ?? 0) > 0; + const startTime = Date.now(); + console.log( + `[sidecar:queue] ▶ ${isRetry ? 'resuming' : 'running'} job ${job.id} (${job.type})${isRetry ? ` retry ${job.retries}` : ''}`, + ); + + const sharedMeta: Record = { ...(job.meta ?? {}) }; + + for (let i = 0; i < handler.steps.length; i++) { + const fresh = await readJob(job.id); + if (!fresh || fresh.status === 'cancelled') { + console.log(`[sidecar:queue] job ${job.id} was cancelled, stopping`); + return; + } + + const handlerStep = handler.steps[i]!; + const step = fresh.steps[i]!; + + if (step.status === 'completed') continue; + + fresh.currentStep = i; + step.status = 'running'; + step.startedAt = Date.now(); + await writeJob(fresh); + + let lastProgressWrite = 0; + let pendingProgress: JobProgress | null = null; + + const updateProgress = async (progress: JobProgress) => { + step.progress = progress; + const now = Date.now(); + if (now - lastProgressWrite >= PROGRESS_THROTTLE_MS) { + lastProgressWrite = now; + pendingProgress = null; + await writeJob(fresh); + } else { + pendingProgress = progress; + } + }; + + const ctx: StepContext = { job: fresh, step, updateProgress, meta: sharedMeta }; + + try { + await handlerStep.run(ctx); + + if (pendingProgress) { + step.progress = pendingProgress; + } + step.status = 'completed'; + step.completedAt = Date.now(); + await writeJob(fresh); + } catch (err) { + const errorMessage = err instanceof Error ? err.message : String(err); + console.error(`[sidecar:queue] step "${step.name}" failed: ${errorMessage}`); + step.status = 'failed'; + step.error = errorMessage; + step.completedAt = Date.now(); + + const isPermanent = err instanceof PermanentError; + const retries = (fresh.retries ?? 0) + 1; + if (!isPermanent && handler.retry && retries <= handler.retry.maxRetries) { + step.status = 'pending'; + step.error = undefined; + step.startedAt = undefined; + step.completedAt = undefined; + step.progress = undefined; + fresh.status = 'queued'; + fresh.error = undefined; + fresh.completedAt = undefined; + fresh.startedAt = undefined; + fresh.retries = retries; + fresh.retryAt = Date.now() + handler.retry.delayMs; + await writeJob(fresh); + console.log( + `[sidecar:queue] job ${fresh.id} will retry (${retries}/${handler.retry.maxRetries}) in ${handler.retry.delayMs / 1000}s`, + ); + scheduleRetry(fresh.lane, handler.retry.delayMs); + return; + } + + fresh.status = 'failed'; + fresh.error = `Step "${step.name}" failed: ${errorMessage}`; + fresh.completedAt = Date.now(); + fresh.meta = { ...fresh.meta, ...sharedMeta }; + await writeJob(fresh); + console.error(`[sidecar:queue] ✗ job ${fresh.id} failed at step "${step.name}" in ${formatDuration(Date.now() - startTime)}:`, errorMessage); + await notifyFailure(fresh); + return; + } + } + + const final = await readJob(job.id); + if (final && final.status === 'running') { + final.status = 'completed'; + final.completedAt = Date.now(); + final.meta = { ...final.meta, ...sharedMeta }; + await writeJob(final); + console.log(`[sidecar:queue] ✓ job ${final.id} completed in ${formatDuration(Date.now() - startTime)}`); + await notifyCompletion(final); + } +} + +async function notifyCompletion(job: Job) { + try { + const { sendMail } = await import('emailer'); + await sendMail({ + template: 'JobCompleted', + subject: `Job completed: ${job.type}`, + to: job.userId, + data: { job }, + }); + } catch { + // SMTP might not be configured — non-fatal + } +} + +async function notifyFailure(job: Job) { + try { + const { sendMail } = await import('emailer'); + await sendMail({ + template: 'JobFailed', + subject: `Job failed: ${job.type}`, + to: job.userId, + data: { job }, + }); + } catch { + // SMTP might not be configured — non-fatal + } +} + +export { readJob, listAllJobs };