Staging wrote four ready chapters then marked the job schema-invalid 34ms later with no summary telemetry. Keep chapter persistence and assemble going if onProgress or a non-lease heartbeat error fails. Co-authored-by: Cursor <cursoragent@cursor.com>
436 lines
15 KiB
TypeScript
436 lines
15 KiB
TypeScript
import {
|
|
PersonalReportJobServiceError,
|
|
type PersonalReportJobRecord,
|
|
type PersonalReportJobService,
|
|
type RecoverExpiredPersonalReportJobsResult,
|
|
} from "./personal-report-job-service-core.ts";
|
|
import {
|
|
PersonalReportServiceError,
|
|
type PersonalReportFailureCode,
|
|
type PersonalReportRecord,
|
|
type PersonalReportService,
|
|
} from "./personal-report-service-core.ts";
|
|
import type { GeneratePersonalReportResult } from "./personal-report-generation.ts";
|
|
import type { PersonalReportSectionService } from "./personal-report-section-service-core.ts";
|
|
import type { ReportBillingPort } from "./personal-report-billing.ts";
|
|
|
|
/**
|
|
* Durable personal-report worker orchestration.
|
|
*
|
|
* This core owns only lease/state coordination. It deliberately delegates
|
|
* profile loading and the evidence/model pipeline so unit tests can exercise
|
|
* crash recovery without a database, HTTP server, or model provider.
|
|
*/
|
|
|
|
export const PERSONAL_REPORT_WORKER_PROGRESS = {
|
|
loadingContext: { phase: "loading_context", percent: 10 },
|
|
generatingReport: { phase: "generating_report", percent: 30 },
|
|
persistingReport: { phase: "persisting_report", percent: 90 },
|
|
} as const;
|
|
|
|
export type PersonalReportWorkerErrorCode =
|
|
| PersonalReportFailureCode
|
|
| "report_generation_failed"
|
|
| "lease_lost";
|
|
|
|
export class PersonalReportWorkerError extends Error {
|
|
readonly name = "PersonalReportWorkerError";
|
|
|
|
constructor(
|
|
readonly code: PersonalReportWorkerErrorCode,
|
|
readonly retryable: boolean,
|
|
message?: string,
|
|
) {
|
|
super(message ?? code);
|
|
}
|
|
}
|
|
|
|
export type PersonalReportWorkerGenerationContext = Readonly<{
|
|
report: PersonalReportRecord;
|
|
profile: unknown;
|
|
signal: AbortSignal;
|
|
sectionService?: PersonalReportSectionService;
|
|
billing?: ReportBillingPort;
|
|
onProgress?: (progress: Readonly<{ phase: string; completed: number; total: number }>) => Promise<void> | void;
|
|
}>;
|
|
|
|
export type PersonalReportWorkerJobPort = Pick<
|
|
PersonalReportJobService,
|
|
| "recoverExpiredLeases"
|
|
| "claimLease"
|
|
| "heartbeatLease"
|
|
| "updateProgress"
|
|
| "markReady"
|
|
| "completeReportReady"
|
|
| "completeReportFailed"
|
|
| "markFailed"
|
|
| "scheduleRetry"
|
|
>;
|
|
|
|
export type PersonalReportWorkerReportPort = Pick<
|
|
PersonalReportService,
|
|
"getByUserAndRequestId"
|
|
>;
|
|
|
|
export type PersonalReportWorkerDeps = Readonly<{
|
|
workerId: string;
|
|
jobs: PersonalReportWorkerJobPort;
|
|
reports: PersonalReportWorkerReportPort;
|
|
loadProfile: (userId: string) => Promise<unknown | null>;
|
|
generate: (
|
|
context: PersonalReportWorkerGenerationContext,
|
|
) => Promise<GeneratePersonalReportResult>;
|
|
sectionService?: PersonalReportSectionService;
|
|
billing?: ReportBillingPort;
|
|
leaseSeconds?: number;
|
|
heartbeatIntervalMs?: number;
|
|
recoveryLimit?: number;
|
|
retryDelayMs?: (job: PersonalReportJobRecord, errorCode: string) => number;
|
|
now?: () => Date;
|
|
setInterval?: typeof globalThis.setInterval;
|
|
clearInterval?: typeof globalThis.clearInterval;
|
|
}>;
|
|
|
|
export type PersonalReportWorkerTickResult = Readonly<{
|
|
recovery: RecoverExpiredPersonalReportJobsResult;
|
|
outcome:
|
|
| "idle"
|
|
| "ready"
|
|
| "failed"
|
|
| "retry_scheduled"
|
|
| "reconcile_deferred"
|
|
| "lease_lost";
|
|
jobId: string | null;
|
|
}>;
|
|
|
|
export type PersonalReportWorkerLoopOptions = Readonly<{
|
|
pollIntervalMs?: number;
|
|
signal?: AbortSignal;
|
|
onError?: (error: unknown) => void;
|
|
sleep?: (milliseconds: number, signal?: AbortSignal) => Promise<void>;
|
|
}>;
|
|
|
|
function positiveInteger(value: number | undefined, fallback: number, field: string): number {
|
|
const resolved = value ?? fallback;
|
|
if (!Number.isInteger(resolved) || resolved <= 0) {
|
|
throw new Error(`${field} must be a positive integer`);
|
|
}
|
|
return resolved;
|
|
}
|
|
|
|
function isLeaseLoss(error: unknown): boolean {
|
|
return error instanceof PersonalReportJobServiceError
|
|
&& (error.code === "lease_lost" || error.code === "terminal_immutable");
|
|
}
|
|
|
|
function stableErrorCode(error: unknown): PersonalReportWorkerErrorCode {
|
|
if (error instanceof PersonalReportWorkerError) return error.code;
|
|
if (error instanceof PersonalReportServiceError) {
|
|
if (error.code === "invalid_document") return "report_schema_invalid";
|
|
if (error.code === "not_found") return "report_not_found";
|
|
}
|
|
if (error instanceof PersonalReportJobServiceError) {
|
|
if (error.code === "storage_invalid" || error.code === "invalid_request") {
|
|
return "report_schema_invalid";
|
|
}
|
|
if (error.code === "not_found") return "report_not_found";
|
|
if (error.code === "lease_lost") return "lease_lost";
|
|
}
|
|
return "report_generation_failed";
|
|
}
|
|
|
|
function isRetryable(error: unknown): boolean {
|
|
if (error instanceof PersonalReportWorkerError) return error.retryable;
|
|
if (error instanceof PersonalReportServiceError) {
|
|
return error.code === "storage_failed";
|
|
}
|
|
if (error instanceof PersonalReportJobServiceError) {
|
|
return error.code === "storage_failed";
|
|
}
|
|
return true;
|
|
}
|
|
|
|
function reportFailureCode(code: PersonalReportWorkerErrorCode): PersonalReportFailureCode {
|
|
switch (code) {
|
|
case "profile_incomplete":
|
|
case "birth_time_not_usable":
|
|
case "calculation_unavailable":
|
|
case "model_unavailable":
|
|
case "report_schema_invalid":
|
|
case "report_guard_rejected":
|
|
case "report_not_found":
|
|
return code;
|
|
default:
|
|
return "calculation_unavailable";
|
|
}
|
|
}
|
|
|
|
function defaultRetryDelayMs(job: PersonalReportJobRecord): number {
|
|
const exponent = Math.max(0, Math.min(job.attemptCount - 1, 5));
|
|
return 5_000 * (2 ** exponent);
|
|
}
|
|
|
|
async function defaultSleep(milliseconds: number, signal?: AbortSignal): Promise<void> {
|
|
if (signal?.aborted) return;
|
|
await new Promise<void>((resolve) => {
|
|
const timer = globalThis.setTimeout(resolve, milliseconds);
|
|
const candidate = timer as ReturnType<typeof globalThis.setTimeout> & { unref?: () => void };
|
|
candidate.unref?.();
|
|
signal?.addEventListener("abort", () => {
|
|
globalThis.clearTimeout(timer);
|
|
resolve();
|
|
}, { once: true });
|
|
});
|
|
}
|
|
|
|
export function createPersonalReportWorker(deps: PersonalReportWorkerDeps) {
|
|
const leaseSeconds = positiveInteger(deps.leaseSeconds, 120, "leaseSeconds");
|
|
const heartbeatIntervalMs = positiveInteger(
|
|
deps.heartbeatIntervalMs,
|
|
Math.max(1_000, Math.floor((leaseSeconds * 1_000) / 3)),
|
|
"heartbeatIntervalMs",
|
|
);
|
|
const recoveryLimit = positiveInteger(deps.recoveryLimit, 100, "recoveryLimit");
|
|
const setIntervalFn = deps.setInterval ?? globalThis.setInterval.bind(globalThis);
|
|
const clearIntervalFn = deps.clearInterval ?? globalThis.clearInterval.bind(globalThis);
|
|
const now = deps.now ?? (() => new Date());
|
|
const retryDelayMs = deps.retryDelayMs ?? defaultRetryDelayMs;
|
|
|
|
|
|
async function settleFailure(
|
|
job: PersonalReportJobRecord,
|
|
report: PersonalReportRecord | null,
|
|
error: unknown,
|
|
): Promise<PersonalReportWorkerTickResult["outcome"]> {
|
|
if (!job.leaseToken || isLeaseLoss(error)) return "lease_lost";
|
|
|
|
const code = stableErrorCode(error);
|
|
const failureCode = reportFailureCode(code);
|
|
const terminal = !isRetryable(error) || job.attemptCount >= job.maxAttempts;
|
|
|
|
try {
|
|
if (terminal) {
|
|
if (deps.billing) {
|
|
await deps.billing.release({ userId: job.userId, requestId: job.requestId, reason: code });
|
|
}
|
|
|
|
if (!report || report.status !== "generating") {
|
|
await deps.jobs.markFailed({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken,
|
|
errorCode: code,
|
|
});
|
|
} else {
|
|
await deps.jobs.completeReportFailed({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken,
|
|
errorCode: code,
|
|
report,
|
|
failureCode,
|
|
});
|
|
}
|
|
return "failed";
|
|
}
|
|
|
|
const retryAt = new Date(now().getTime() + retryDelayMs(job, code));
|
|
const scheduled = await deps.jobs.scheduleRetry({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken,
|
|
errorCode: code,
|
|
nextAttemptAt: retryAt,
|
|
});
|
|
if (scheduled.kind === "scheduled") return "retry_scheduled";
|
|
|
|
if (!report || report.status !== "generating") return "failed";
|
|
await deps.jobs.completeReportFailed({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken,
|
|
errorCode: code,
|
|
report,
|
|
failureCode,
|
|
});
|
|
return "failed";
|
|
} catch (transitionError) {
|
|
if (isLeaseLoss(transitionError)) return "lease_lost";
|
|
throw transitionError;
|
|
}
|
|
}
|
|
|
|
async function processClaimed(job: PersonalReportJobRecord): Promise<PersonalReportWorkerTickResult["outcome"]> {
|
|
if (!job.leaseToken) {
|
|
throw new PersonalReportWorkerError("lease_lost", false, "claimed job has no lease token");
|
|
}
|
|
|
|
const controller = new AbortController();
|
|
let heartbeatError: unknown = null;
|
|
let heartbeatChain = Promise.resolve();
|
|
const refreshLease = async () => {
|
|
if (heartbeatError !== null || controller.signal.aborted) return;
|
|
await deps.jobs.heartbeatLease({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken!,
|
|
leaseSeconds,
|
|
});
|
|
};
|
|
const heartbeatTimer = setIntervalFn(() => {
|
|
heartbeatChain = heartbeatChain
|
|
.then(refreshLease)
|
|
.catch((error) => {
|
|
heartbeatError = error;
|
|
controller.abort(error);
|
|
});
|
|
}, heartbeatIntervalMs);
|
|
// Keep this timer ref'd. The worker runs from instrumentation, not an HTTP
|
|
// request; unref'd intervals in Next.js standalone are skipped while the
|
|
// only in-flight work is a long model await, which is exactly when the
|
|
// lease must be extended.
|
|
|
|
let report: PersonalReportRecord | null = null;
|
|
const generationStartedAt = now().getTime();
|
|
try {
|
|
report = await deps.reports.getByUserAndRequestId(job.userId, job.requestId);
|
|
if (!report) {
|
|
throw new PersonalReportWorkerError("report_not_found", false);
|
|
}
|
|
if (report.requestFingerprint !== job.requestFingerprint) {
|
|
throw new PersonalReportWorkerError("report_schema_invalid", false, "job/report fingerprint mismatch");
|
|
}
|
|
if (report.status === "ready") {
|
|
await deps.jobs.markReady({ jobId: job.id, leaseToken: job.leaseToken });
|
|
return "ready";
|
|
}
|
|
if (report.status === "failed") {
|
|
await deps.jobs.markFailed({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken,
|
|
errorCode: report.failureCode ?? "report_generation_failed",
|
|
});
|
|
return "failed";
|
|
}
|
|
|
|
await deps.jobs.updateProgress({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken,
|
|
...PERSONAL_REPORT_WORKER_PROGRESS.loadingContext,
|
|
});
|
|
const profile = await deps.loadProfile(job.userId);
|
|
if (!profile) {
|
|
throw new PersonalReportWorkerError("profile_incomplete", false);
|
|
}
|
|
|
|
// Refresh the lease synchronously immediately before the potentially
|
|
// long workflow/model call; the interval continues the heartbeat while
|
|
// that call is in flight.
|
|
await deps.jobs.heartbeatLease({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken,
|
|
leaseSeconds,
|
|
});
|
|
await deps.jobs.updateProgress({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken,
|
|
...PERSONAL_REPORT_WORKER_PROGRESS.generatingReport,
|
|
});
|
|
|
|
const generated = await deps.generate({
|
|
report,
|
|
profile,
|
|
signal: controller.signal,
|
|
sectionService: deps.sectionService,
|
|
onProgress: async (progress) => {
|
|
try {
|
|
await refreshLease();
|
|
if (heartbeatError !== null) throw heartbeatError;
|
|
const percent = progress.total > 0
|
|
? Math.min(89, 30 + Math.floor((progress.completed / progress.total) * 55))
|
|
: 30;
|
|
await deps.jobs.updateProgress({
|
|
jobId: job.id, leaseToken: job.leaseToken!, phase: progress.phase, percent,
|
|
});
|
|
} catch (error) {
|
|
if (isLeaseLoss(error) || heartbeatError !== null) {
|
|
heartbeatError = heartbeatError ?? error;
|
|
controller.abort(error);
|
|
throw error;
|
|
}
|
|
}
|
|
},
|
|
});
|
|
await heartbeatChain;
|
|
if (heartbeatError !== null) throw heartbeatError;
|
|
if (generated.status === "failed") {
|
|
const innerReason = generated.failureCode === "report_schema_invalid"
|
|
? generated.innerReason
|
|
: generated.failureCode;
|
|
throw new PersonalReportWorkerError(generated.failureCode, false, innerReason);
|
|
}
|
|
|
|
await deps.jobs.updateProgress({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken,
|
|
...PERSONAL_REPORT_WORKER_PROGRESS.persistingReport,
|
|
});
|
|
await deps.jobs.completeReportReady({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken,
|
|
report,
|
|
document: generated.document,
|
|
});
|
|
if (deps.billing) {
|
|
const settled = await deps.billing.complete({
|
|
userId: job.userId,
|
|
requestId: job.requestId,
|
|
usage: {
|
|
actualModelId: generated.usage?.actualModelId ?? "unknown",
|
|
modelConfigVersion: generated.usage?.modelConfigVersion,
|
|
inputTokens: generated.usage?.inputTokens ?? 0,
|
|
outputTokens: generated.usage?.outputTokens ?? 0,
|
|
durationMs: Math.max(0, now().getTime() - generationStartedAt),
|
|
},
|
|
});
|
|
if (!settled) throw new PersonalReportWorkerError("report_generation_failed", false, "report billing settlement failed");
|
|
}
|
|
return "ready";
|
|
} catch (error) {
|
|
if (report?.status === "ready") {
|
|
// A ready report is the source of truth. If the projection-only job
|
|
// reconciliation fails, keep the live lease untouched so expiry
|
|
// recovery can converge it to ready; never downgrade it to failed.
|
|
return isLeaseLoss(error) ? "lease_lost" : "reconcile_deferred";
|
|
}
|
|
return settleFailure(job, report, error);
|
|
} finally {
|
|
clearIntervalFn(heartbeatTimer);
|
|
controller.abort();
|
|
await heartbeatChain;
|
|
}
|
|
}
|
|
|
|
return {
|
|
async tick(): Promise<PersonalReportWorkerTickResult> {
|
|
const recovery = await deps.jobs.recoverExpiredLeases(recoveryLimit);
|
|
const job = await deps.jobs.claimLease({ workerId: deps.workerId, leaseSeconds });
|
|
if (!job) return { recovery, outcome: "idle", jobId: null };
|
|
const outcome = await processClaimed(job);
|
|
return { recovery, outcome, jobId: job.id };
|
|
},
|
|
};
|
|
}
|
|
|
|
export async function runPersonalReportWorkerLoop(
|
|
worker: Readonly<{ tick: () => Promise<PersonalReportWorkerTickResult> }>,
|
|
options: PersonalReportWorkerLoopOptions = {},
|
|
): Promise<void> {
|
|
const pollIntervalMs = positiveInteger(options.pollIntervalMs, 2_000, "pollIntervalMs");
|
|
const sleep = options.sleep ?? defaultSleep;
|
|
while (!options.signal?.aborted) {
|
|
try {
|
|
const result = await worker.tick();
|
|
if (result.outcome !== "idle") continue;
|
|
} catch (error) {
|
|
options.onError?.(error);
|
|
}
|
|
await sleep(pollIntervalMs, options.signal);
|
|
}
|
|
}
|