385 lines
13 KiB
TypeScript
385 lines
13 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";
|
|
|
|
/**
|
|
* 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;
|
|
}>;
|
|
|
|
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>;
|
|
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 timerUnref(timer: ReturnType<typeof globalThis.setInterval>): void {
|
|
const candidate = timer as ReturnType<typeof globalThis.setInterval> & { unref?: () => void };
|
|
candidate.unref?.();
|
|
}
|
|
|
|
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 (!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 heartbeatTimer = setIntervalFn(() => {
|
|
heartbeatChain = heartbeatChain
|
|
.then(async () => {
|
|
if (heartbeatError !== null || controller.signal.aborted) return;
|
|
await deps.jobs.heartbeatLease({
|
|
jobId: job.id,
|
|
leaseToken: job.leaseToken!,
|
|
leaseSeconds,
|
|
});
|
|
})
|
|
.catch((error) => {
|
|
heartbeatError = error;
|
|
controller.abort(error);
|
|
});
|
|
}, heartbeatIntervalMs);
|
|
timerUnref(heartbeatTimer);
|
|
|
|
let report: PersonalReportRecord | null = null;
|
|
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 });
|
|
await heartbeatChain;
|
|
if (heartbeatError !== null) throw heartbeatError;
|
|
if (generated.status === "failed") {
|
|
throw new PersonalReportWorkerError(generated.failureCode, false);
|
|
}
|
|
|
|
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,
|
|
});
|
|
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);
|
|
}
|
|
}
|