feat(hk): connect delivery failures to precise diagnostic records
CI / check (push) Successful in 3m31s
CI / deploy (push) Successful in 18s
CI / publish-container (push) Successful in 1m12s

This commit is contained in:
Simo committed 2026-09-13 20:25:03 +02:00
1 parent 275a574203
commit 60e49c1ec9
17 files changed
+253 -20

No files matched your search

+21 -5
View File
@@ -2,6 +2,8 @@ import "server-only";
import { randomUUID } from "node:crypto";
import { sql } from "drizzle-orm";
import { db } from "@/lib/db";
import { logger } from "@/lib/logger";
import { isDeliveryReference } from "./delivery-diagnostics";
import {
type EffectClaim,
type EffectTopic,
@@ -60,11 +62,12 @@ export const effectRepository = {
async claim(): Promise<EffectClaim | null> {
return db.transaction(async (tx) => {
const [rows] = await tx.execute(
sql`SELECT id,topic,attempts FROM cms_outbox WHERE (status='pending' AND available_at<=UTC_TIMESTAMP(3)) OR (status='running' AND lease_until<UTC_TIMESTAMP(3)) ORDER BY available_at,id LIMIT 1 FOR UPDATE`,
sql`SELECT id,operation_id AS operationId,topic,attempts FROM cms_outbox WHERE (status='pending' AND available_at<=UTC_TIMESTAMP(3)) OR (status='running' AND lease_until<UTC_TIMESTAMP(3)) ORDER BY available_at,id LIMIT 1 FOR UPDATE`,
);
const row = (
rows as unknown as Array<{
id: string;
operationId: string;
topic: EffectTopic;
attempts: number;
}>
@@ -82,15 +85,28 @@ export const effectRepository = {
sql`UPDATE cms_outbox SET status='done',lease_token=NULL,lease_until=NULL,last_error=NULL WHERE id=${claim.id} AND status='running' AND lease_token=${claim.token}`,
);
},
async fail(claim: EffectClaim) {
async fail(claim: EffectClaim, error: unknown) {
const errorId = logger.error("Operation delivery failed", {
module: "operations",
deliveryId: claim.id,
durableOperationId: claim.operationId,
topic: claim.topic,
attempt: claim.attempts,
error,
});
const reference = isDeliveryReference(errorId)
? `Diagnostic reference: ${errorId}`
: "Delivery failed; service diagnostics unavailable.";
await db.execute(
sql`UPDATE cms_outbox SET status=${claim.attempts >= 8 ? "failed" : "pending"},available_at=DATE_ADD(UTC_TIMESTAMP(3),INTERVAL ${retryDelay(claim.attempts)} SECOND),lease_token=NULL,lease_until=NULL,last_error='Delivery failed; inspect the linked operation and service diagnostics.' WHERE id=${claim.id} AND status='running' AND lease_token=${claim.token}`,
sql`UPDATE cms_outbox SET status=${claim.attempts >= 8 ? "failed" : "pending"},available_at=DATE_ADD(UTC_TIMESTAMP(3),INTERVAL ${retryDelay(claim.attempts)} SECOND),lease_token=NULL,lease_until=NULL,last_error=${reference} WHERE id=${claim.id} AND status='running' AND lease_token=${claim.token}`,
);
},
};
export async function listEffects() {
export async function listEffects(id?: string) {
if (id !== undefined && !isDeliveryReference(id)) return [];
const filter = id === undefined ? sql`` : sql`WHERE e.id=${id}`;
const [rows] = await db.execute(
sql`SELECT e.id,e.operation_id AS operationId,e.topic,e.status,e.attempts,e.last_error AS lastError,e.created_at AS createdAt,e.available_at AS availableAt,o.actor_id AS actorId,o.kind,o.result_json AS resultJson FROM cms_outbox e JOIN cms_operations o ON o.id=e.operation_id ORDER BY e.created_at DESC,e.id DESC LIMIT 100`,
sql`SELECT e.id,e.operation_id AS operationId,e.topic,e.status,e.attempts,e.last_error AS lastError,e.created_at AS createdAt,e.available_at AS availableAt,o.actor_id AS actorId,o.kind,o.result_json AS resultJson FROM cms_outbox e JOIN cms_operations o ON o.id=e.operation_id ${filter} ORDER BY e.created_at DESC,e.id DESC LIMIT 100`,
);
return rows as unknown as Array<{
id: string;