email crons
This commit is contained in:
@@ -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<typeof setInterval> | 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<string, unknown> | 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<string, unknown> = { user: account.email };
|
||||||
|
if (account.authType === 'oauth') {
|
||||||
|
const integration = await getUserIntegration(account.userId, 'google');
|
||||||
|
const config = integration?.config as Record<string, unknown> | 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<string, unknown>;
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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<string, boolean>();
|
||||||
|
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<Job> {
|
||||||
|
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<Job | null> {
|
||||||
|
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<string>();
|
||||||
|
|
||||||
|
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<string, unknown> = { ...(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 };
|
||||||
Reference in New Issue
Block a user