Cross-instance
context.schedule (SynchronizedClock phía ODP core) đảm bảo chỉ 1 instance trong cluster thực thi handler mỗi tick.
const slaTick = async () => { try { await recomputeSlaStatus(ctx, SYSTEM_ACCOUNTABILITY); await autoReopenSnoozed(ctx, SYSTEM_ACCOUNTABILITY); } catch (err: any) { logger.error({ err: err?.message }, 'helpdesk SLA/snooze cron tick failed'); }};if (typeof schedule === 'function') { schedule('* * * * *', slaTick); logger.info('helpdesk SLA/snooze schedule registered (context.schedule, every minute)');
// ── Email inbound: IMAP-pull tick (P1) ── schedule('* * * * *', async () => { try { await pollImapConnections(ctx); } catch (err: any) { logger.error({ err: err?.message }, 'helpdesk IMAP inbound poll tick failed'); } }); logger.info('helpdesk IMAP inbound poll registered (context.schedule, every minute)');} else { logger.error('context.schedule unavailable — SLA/snooze + IMAP poll NOT registered (needs @odp/api >= 0.8.0)');}context.schedule (thay vì cron.CronJob tự quản lý như bản prototype cũ) dùng SynchronizedClock phía ODP core (@odp/api >= 0.8.0) — chỉ một instance trong cluster thực sự chạy handler mỗi tick, các instance khác skip. Nếu context.schedule không tồn tại (API cũ hơn), extension log lỗi và không đăng ký lịch nào cả thay vì fallback về cron package tự quản lý — tránh double-run khi chạy nhiều instance trên API cũ.
recomputeSlaStatus — quét mọi conversation mở có SLA policyexport async function recomputeSlaStatus(ctx: AppContext, accountability: any): Promise<void> { const tickStartedAt = Date.now();
// Narrow select — this scans every open conversation with an SLA policy on // every tick, so keep the row payload minimal. const conversations = await ctx.database('hd_conversations') .whereNotNull('sla_policy_id') .whereNot('status', 'resolved') .select('id', 'sla_policy_id', 'sla_status', 'status', 'priority', 'inbox_id', 'created_at', 'waiting_since');
if (!conversations.length) return; // preload policies + working hours ...Điểm đáng chú ý: label hydrate được batch, không phải N+1. Bản thân comment trong code nói rõ đây là fix so với cách làm cũ (pluck từng conversation một):
// Batch-hydrate label_ids for policy matching (mirrors m2m-hydrator.ts:// one whereIn query + Map, instead of a per-conversation pluck).const convIds = conversations.map((c) => c.id);const labelRows = await ctx.database('hd_conversation_labels') .whereIn('hd_conversations_id', convIds) .select('hd_conversations_id', 'hd_labels_id');const labelsByConversation = new Map<number, number[]>();for (const row of labelRows) { const arr = labelsByConversation.get(row.hd_conversations_id) || []; arr.push(row.hd_labels_id); labelsByConversation.set(row.hd_conversations_id, arr);}Với mỗi conversation, tính _computeStatus (dùng computeDeadlines — port verbatim từ admin/packages/helpdesk/app/utils/sla-deadline.ts để đảm bảo client/server tính giống hệt nhau) rồi banding theo phần trăm đã trôi qua so với deadline:
// _computeStatus bandingif (maxPercent >= 100) status = 'breached';else if (maxPercent >= 90) status = 'breaching';else if (maxPercent >= 75) status = 'at_risk';else status = 'on_track';Chỉ ghi lại (convSvc.updateOne) khi status thực sự đổi, và chỉ fire automation sla_warning khi vừa vượt ngưỡng 75% (tránh spam automation mỗi phút cho cùng một conversation đã ở trạng thái cảnh báo từ trước):
const enteredWarning = result.percent >= 75 && conv.sla_status !== result.status && (result.status === 'at_risk' || result.status === 'breaching' || result.status === 'breached');if (enteredWarning) { const { runGuarded } = await import('./automation/engine.js'); await runGuarded(ctx, 'sla_warning', Number(conv.id));}runGuarded dùng chung guard automationInFlight với các hook item lifecycle (xem Hooks) — nếu conversation đang được automation xử lý từ một event khác cùng lúc, tick SLA bỏ qua conversation đó thay vì chồng lấn.
Cuối tick, có cảnh báo nếu tick chạy quá lâu:
const tickDurationMs = Date.now() - tickStartedAt;if (tickDurationMs > 30_000) { ctx.log.warn({ tickDurationMs, conversations: conversations.length }, 'helpdesk SLA recompute tick took longer than 30s');}autoReopenSnoozed — mở lại conversation hết hạn snoozeexport async function autoReopenSnoozed(ctx: AppContext, accountability: any): Promise<void> { const nowIso = new Date().toISOString(); const due = await ctx.database('hd_conversations') .where('status', 'snoozed') .whereNotNull('snooze_until') .where('snooze_until', '<=', nowIso) .pluck('id');
if (!due.length) return;
const convSvc = await createItemsService(ctx, 'hd_conversations', accountability); for (const id of due) { try { await convSvc.updateOne(id, { status: 'open' }); } catch (err: any) { // Swallow per-conversation errors — don't abort the whole scan. ctx.log.error({ err: err?.message, conv: id }, 'helpdesk snooze auto-reopen failed'); } }}Mỗi conversation update lỗi bị nuốt riêng lẻ — một conversation lỗi (ví dụ vi phạm ràng buộc nào đó) không chặn việc mở lại các conversation khác trong cùng tick.
const MAX_PER_TICK = 50; // cap messages fetched per connection per tickconst FAIL_LIMIT = 3; // poison-message: skip after N failures within TTLconst FAIL_TTL = '6h';
// In-process guard so a slow tick (mailbox chậm) doesn't overlap the next tick// within the SAME instance. Cross-instance is handled by context.schedule.let pollInFlight = false;
export async function pollImapConnections(ctx: AppContext): Promise<void> { if (pollInFlight) return; pollInFlight = true; try { const conns: HdConnection[] = await ctx.database('hd_connections') .where({ inbound_provider: 'imap' }) .whereNotIn('status', ['error', 'disconnected', 'disabled']) .where((qb: any) => qb.whereNull('needs_reauth').orWhere('needs_reauth', false)); for (const conn of conns) { try { await pollOne(ctx, conn); } catch (err: any) { await onConnError(ctx, conn, err); } } } finally { pollInFlight = false; }}Hai lớp bảo vệ overlap khác nhau, dễ nhầm:
Cross-instance
context.schedule (SynchronizedClock phía ODP core) đảm bảo chỉ 1 instance trong cluster thực thi handler mỗi tick.
In-process (pollInFlight)
Biến let module-level, chặn tick tiếp theo overlap tick đang chạy trong cùng instance nếu một mailbox chậm khiến tick trước chưa xong khi tick sau đến hạn.
Lần đầu kết nối (hoặc UIDVALIDITY đổi — mailbox bị reset phía provider), extension không quét lại toàn bộ hộp thư — nó đặt watermark ở “hiện tại” rồi chỉ nhận mail đến sau đó:
let lastUid = Number(conn.imap_last_uid ?? 0);const validityChanged = uidValidity !== Number(conn.imap_uidvalidity ?? 0);if (!lastUid || validityChanged) { lastUid = Math.max(0, uidNext - 1); await updateWatermark(ctx, conn.id, { imap_uidvalidity: uidValidity, imap_last_uid: lastUid }); // Nothing older than "now" is fetched on first sync. await touchSynced(ctx, conn.id); return;}imapflow serialize command trên một connection: chạy lệnh IMAP khác (markProcessed) hoặc query DB chậm trong khi fetch stream còn mở sẽ ném lỗi. Code chủ động tách 2 pha:
// Phase 1 — DRAIN the fetch stream into memory first. ...const buffered: Array<{ uid: number; source: Buffer }> = [];let count = 0;for await (const msg of client.fetch({ uid: `${lastUid + 1}:*` }, { uid: true, source: true })) { if (count++ >= MAX_PER_TICK) break; // continue next tick (watermark advances) if (msg.source) buffered.push({ uid: Number(msg.uid), source: msg.source as Buffer });}
// Phase 2 — process (pipeline + markProcessed) with the fetch stream closed.MAX_PER_TICK = 50 giới hạn số message xử lý mỗi connection mỗi tick — nếu hộp thư có backlog lớn hơn, phần còn lại được xử lý ở các tick sau (watermark chỉ tiến tới message cuối đã xử lý thành công).
Một message parse lỗi liên tục (malformed MIME, encoding lạ) không được phép chặn watermark tiến mãi mãi. Đếm lỗi theo messageId qua ExtensionCache isolated store, TTL 6 giờ:
async function failCount(ctx: AppContext, messageId: string): Promise<number> { if (!ctx.cache) return 0; return (await ctx.cache.get<number>(`imap:fail:${messageId}`, 0)) ?? 0;}async function bumpFail(ctx: AppContext, messageId: string): Promise<void> { if (!ctx.cache) return; const n = (await ctx.cache.get<number>(`imap:fail:${messageId}`, 0)) ?? 0; await ctx.cache.set(`imap:fail:${messageId}`, n + 1, { ttl: FAIL_TTL });}Sau FAIL_LIMIT = 3 lần thất bại trong 6h, message bị markProcessed (skip) mà không xử lý — tick tiếp theo watermark đã vượt qua nó, không retry vô hạn.
// Watermark + audit writes go through raw Knex (denormalized system fields, not a// logical edit) so they don't fire hd_connections.update events (encryption filter / hooks).async function updateWatermark(ctx: AppContext, id: number, patch: Record<string, any>): Promise<void> { await ctx.database('hd_connections').where('id', id).update(patch);}Nếu đi qua ItemsService, mỗi lần cập nhật watermark (mỗi tick, mỗi connection) sẽ chạy qua filter mã hoá secret (encryptConn) trong hooks/index.ts — không cần thiết và tốn kém cho một field không phải secret.
automationInFlight guard dùng chung với runGuardedhd_connections/hd_inboxes được tạo/sửa qua APIpollInFlight/automationInFlight reset mỗi lần restart