fix(ai-moderation): matikan hot requeue loop saat upload attachment in-flight
Batch race guard balikin {ok:true, rows:[]} tanpa sinyal saat semua target
masih upload-pending -> processor klasifikasi semua incomplete -> fanout ke
individual queue -> di situ requeue + reschedule 250ms -> balik ke batch:
hot loop ~300ms sepanjang upload (10 siklus/3 dtk di log prod 08:13).
Fix: worker batch kini return uploadPendingIds eksplisit; classifier pure
baru (partitionBatchOutcome) partisi completed/upload_pending/incomplete/
parse_failed/api_failed; target upload-pending DEFERRED dengan poll backoff
linear (AI_ANALYSIS_UPLOAD_POLL_MS 1500 base, cap AI_ANALYSIS_MAX_UPLOAD_POLL_MS
8000), tidak pernah masuk fanout; tail shouldScheduleNext tak menimpa defer.
Test: tests/batchOutcomeClassifier.test.ts (8 kasus, pure tanpa DB/Piscina).
This commit is contained in:
@@ -0,0 +1,108 @@
|
||||
/**
|
||||
* batchOutcomeClassifier.ts
|
||||
*
|
||||
* Pure partitioner of the batch worker response (2026-08-25).
|
||||
*
|
||||
* Bug history: the batch race guard returned `{ok:true, rows:[]}` when every
|
||||
* target's attachment upload was still in-flight. The processor classified all
|
||||
* of them as "incomplete" and fanned out to the individual queue, where the
|
||||
* guard there requeued + rescheduled at the 250ms debounce — a hot ~300ms loop
|
||||
* for the entire upload duration (~10 cycles in 3s in prod logs). Root fix:
|
||||
* the worker now reports `uploadPendingIds` explicitly and this pure function
|
||||
* partitions the outcome so upload-pending targets NEVER enter the fanout.
|
||||
*/
|
||||
|
||||
export interface BatchRowLike {
|
||||
id?: string;
|
||||
ai_status?: string | null;
|
||||
ai_moderation_flags?: string | null;
|
||||
}
|
||||
|
||||
export interface BatchWorkerResponseLike {
|
||||
ok?: boolean;
|
||||
rows?: BatchRowLike[];
|
||||
/** Explicit race-guard signal from the worker (2026-08-25). */
|
||||
uploadPendingIds?: string[];
|
||||
error?: string;
|
||||
}
|
||||
|
||||
/** One target's per-message disposition after a batch attempt. */
|
||||
export type BatchTargetKind =
|
||||
| "completed"
|
||||
| "upload_pending"
|
||||
| "incomplete"
|
||||
| "parse_failed"
|
||||
| "api_failed";
|
||||
|
||||
function flagsOf(row: { ai_moderation_flags?: string | null }): string[] {
|
||||
if (!row.ai_moderation_flags) return [];
|
||||
try {
|
||||
const parsed = JSON.parse(row.ai_moderation_flags) as unknown;
|
||||
return Array.isArray(parsed) ? (parsed as string[]) : [];
|
||||
} catch {
|
||||
return [] as string[];
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Partition the input message ids into per-message dispositions for one batch
|
||||
* worker response. Pure: no DB/Piscina/logger — unit-testable directly.
|
||||
*
|
||||
* Priority per id: explicit uploadPendingIds → completed row → flag-based
|
||||
* failure kinds → unexplained missing (treated like incomplete).
|
||||
*/
|
||||
export function partitionBatchOutcome(
|
||||
messages: ReadonlyArray<{ id: string }>,
|
||||
response: BatchWorkerResponseLike,
|
||||
): Map<string, BatchTargetKind> {
|
||||
const pendingSet = new Set(response.uploadPendingIds ?? []);
|
||||
const rowsById = new Map(
|
||||
(response.rows ?? [])
|
||||
.filter((r): r is BatchRowLike & { id: string } => Boolean(r?.id))
|
||||
.map((r) => [r.id, r]),
|
||||
);
|
||||
|
||||
const out = new Map<string, BatchTargetKind>();
|
||||
for (const msg of messages) {
|
||||
if (pendingSet.has(msg.id)) {
|
||||
out.set(msg.id, "upload_pending");
|
||||
continue;
|
||||
}
|
||||
const row = rowsById.get(msg.id);
|
||||
if (!row) {
|
||||
// Unexplained drop: LLM silently omitted it. Same retryable bucket as
|
||||
// analysis_incomplete — never a silent success.
|
||||
out.set(msg.id, "incomplete");
|
||||
continue;
|
||||
}
|
||||
if (row.ai_status !== "error") {
|
||||
out.set(msg.id, "completed");
|
||||
continue;
|
||||
}
|
||||
const flags = flagsOf(row);
|
||||
if (flags.includes("analysis_incomplete")) {
|
||||
out.set(msg.id, "incomplete");
|
||||
} else if (flags.includes("analysis_parse_failed")) {
|
||||
out.set(msg.id, "parse_failed");
|
||||
} else if (flags.includes("analysis_api_failed")) {
|
||||
out.set(msg.id, "api_failed");
|
||||
} else {
|
||||
out.set(msg.id, "incomplete");
|
||||
}
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/**
|
||||
* Linear backoff ramp for consecutive upload-pending polls:
|
||||
* poll N (1-based) waits min(base × N, cap). Keeps latency low for fast
|
||||
* uploads while bounding total polling cost for long uploads.
|
||||
*/
|
||||
export function computeUploadPollDelayMs(
|
||||
consecutivePolls: number,
|
||||
baseMs: number,
|
||||
capMs: number,
|
||||
): number {
|
||||
const n = Math.max(1, Math.floor(consecutivePolls));
|
||||
return Math.min(Math.round(baseMs * n), Math.round(capMs));
|
||||
}
|
||||
Reference in New Issue
Block a user