diff --git a/apps/journal/app/jobs/backfill-user-keypairs.ts b/apps/journal/app/jobs/backfill-user-keypairs.ts index cd267cf..1f79e4b 100644 --- a/apps/journal/app/jobs/backfill-user-keypairs.ts +++ b/apps/journal/app/jobs/backfill-user-keypairs.ts @@ -1,4 +1,4 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { ensureUserKeypair, listUsersWithoutKeypair } from "../lib/federation-keys.server.ts"; import { logger } from "../lib/logger.server.ts"; @@ -10,7 +10,7 @@ import { logger } from "../lib/logger.server.ts"; * only touches users whose public_key IS NULL, so re-runs are no-ops. * New users get keys at registration and never appear in this workload. */ -export const backfillUserKeypairsJob: JobDefinition = { +export const backfillUserKeypairsJob = defineJournalJob({ name: "backfill-user-keypairs", retryLimit: 3, expireInSeconds: 300, @@ -23,4 +23,4 @@ export const backfillUserKeypairsJob: JobDefinition = { logger.info({ candidates: ids.length, generated }, "backfill-user-keypairs"); return { candidates: ids.length, generated }; }, -}; +}); diff --git a/apps/journal/app/jobs/consumed-jti-sweep.ts b/apps/journal/app/jobs/consumed-jti-sweep.ts index 1b8a7c0..5be848c 100644 --- a/apps/journal/app/jobs/consumed-jti-sweep.ts +++ b/apps/journal/app/jobs/consumed-jti-sweep.ts @@ -1,4 +1,4 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { lt } from "drizzle-orm"; import { consumedJwtJti } from "@trails-cool/db/schema/journal"; import { getDb } from "../lib/db.ts"; @@ -13,7 +13,7 @@ import { logger } from "../lib/logger.server.ts"; * * See planner-audit #2 Phase B. */ -export const consumedJtiSweepJob: JobDefinition = { +export const consumedJtiSweepJob = defineJournalJob({ name: "consumed-jti-sweep", cron: "45 3 * * *", // daily at 03:45 UTC (offset from notifications-purge) retryLimit: 1, @@ -28,4 +28,4 @@ export const consumedJtiSweepJob: JobDefinition = { logger.info({ purged }, "consumed-jti-sweep"); return { purged }; }, -}; +}); diff --git a/apps/journal/app/jobs/deliver-activity.ts b/apps/journal/app/jobs/deliver-activity.ts index c40ab44..b75631c 100644 --- a/apps/journal/app/jobs/deliver-activity.ts +++ b/apps/journal/app/jobs/deliver-activity.ts @@ -1,4 +1,4 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { and, eq } from "drizzle-orm"; import { activities, users } from "@trails-cool/db/schema/journal"; import { getDb } from "../lib/db.ts"; @@ -39,12 +39,12 @@ async function paceHost(host: string): Promise { * failed and pg-boss retries; exhausting the budget is the permanent * failure, logged by the final catch. */ -export const deliverActivityJob: JobDefinition = { +export const deliverActivityJob = defineJournalJob({ name: "deliver-activity", expireInSeconds: 60, async handler(jobs) { for (const job of jobs) { - const p = job.data as DeliveryPayload; + const p = job.data; try { await deliverOne(p); } catch (err) { @@ -56,7 +56,7 @@ export const deliverActivityJob: JobDefinition = { } } }, -}; +}); async function deliverOne(p: DeliveryPayload): Promise { const federation = getFederation(); diff --git a/apps/journal/app/jobs/demo-bot-generate.ts b/apps/journal/app/jobs/demo-bot-generate.ts index a25020a..7257a3d 100644 --- a/apps/journal/app/jobs/demo-bot-generate.ts +++ b/apps/journal/app/jobs/demo-bot-generate.ts @@ -1,4 +1,4 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { DEMO_BACKFILL_TARGET, DEMO_DAILY_CAP, @@ -21,7 +21,7 @@ import { logger } from "../lib/logger.server.ts"; * 3. Otherwise apply the decide-to-walk gate (local hour + p=0.09) and * daily cap; on pass, insert one route+activity via `generateOneWalk`. */ -export const demoBotGenerateJob: JobDefinition = { +export const demoBotGenerateJob = defineJournalJob({ name: "demo-bot-generate", cron: "0,30 * * * *", retryLimit: 1, @@ -57,4 +57,4 @@ export const demoBotGenerateJob: JobDefinition = { await refreshDemoBotGauges(); return { mode: "single", routeId: id }; }, -}; +}); diff --git a/apps/journal/app/jobs/demo-bot-prune.ts b/apps/journal/app/jobs/demo-bot-prune.ts index 6ac2855..b065bf8 100644 --- a/apps/journal/app/jobs/demo-bot-prune.ts +++ b/apps/journal/app/jobs/demo-bot-prune.ts @@ -1,4 +1,4 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { demoRetentionDays, isDemoBotEnabled, @@ -11,7 +11,7 @@ import { logger } from "../lib/logger.server.ts"; * Daily prune. Deletes synthetic rows older than * `DEMO_BOT_RETENTION_DAYS` (default 14). Never touches real users. */ -export const demoBotPruneJob: JobDefinition = { +export const demoBotPruneJob = defineJournalJob({ name: "demo-bot-prune", cron: "15 3 * * *", retryLimit: 1, @@ -24,4 +24,4 @@ export const demoBotPruneJob: JobDefinition = { await refreshDemoBotGauges(); return { days, ...counts }; }, -}; +}); diff --git a/apps/journal/app/jobs/federation-kv-sweep.ts b/apps/journal/app/jobs/federation-kv-sweep.ts index b21a7bb..6bec691 100644 --- a/apps/journal/app/jobs/federation-kv-sweep.ts +++ b/apps/journal/app/jobs/federation-kv-sweep.ts @@ -1,4 +1,4 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { PostgresKvStore } from "../lib/federation-kv.server.ts"; import { logger } from "../lib/logger.server.ts"; @@ -7,7 +7,7 @@ import { logger } from "../lib/logger.server.ts"; * nonces and caches carry TTLs; reads already filter expired rows, this * keeps the table from growing unbounded). */ -export const federationKvSweepJob: JobDefinition = { +export const federationKvSweepJob = defineJournalJob({ name: "federation-kv-sweep", cron: "15 4 * * *", // daily at 04:15 UTC (offset from the other sweeps) retryLimit: 1, @@ -17,4 +17,4 @@ export const federationKvSweepJob: JobDefinition = { logger.info({ purged }, "federation-kv-sweep"); return { purged }; }, -}; +}); diff --git a/apps/journal/app/jobs/garmin-import-activity.ts b/apps/journal/app/jobs/garmin-import-activity.ts index edff614..b396150 100644 --- a/apps/journal/app/jobs/garmin-import-activity.ts +++ b/apps/journal/app/jobs/garmin-import-activity.ts @@ -1,9 +1,6 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { logger } from "../lib/logger.server.ts"; -import { - runGarminActivityImport, - type GarminImportData, -} from "../lib/connected-services/providers/garmin/import.server.ts"; +import { runGarminActivityImport } from "../lib/connected-services/providers/garmin/import.server.ts"; // Garmin webhook notifications enqueue here (spec: garmin-import, // "Push-notification activity import"): the webhook answers 200 @@ -11,14 +8,14 @@ import { // download, FIT→GPX, activity creation. Backfill bursts deliver many // notifications at once; the queue absorbs them and pg-boss retries // transient download failures. -export const garminImportActivityJob: JobDefinition = { +export const garminImportActivityJob = defineJournalJob({ name: "garmin-import-activity", retryLimit: 3, expireInSeconds: 300, async handler(jobs) { const batch = Array.isArray(jobs) ? jobs : [jobs]; for (const job of batch) { - const data = job.data as GarminImportData; + const data = job.data; try { await runGarminActivityImport(data); } catch (err) { @@ -30,4 +27,4 @@ export const garminImportActivityJob: JobDefinition = { } } }, -}; +}); diff --git a/apps/journal/app/jobs/import-batches-sweep.ts b/apps/journal/app/jobs/import-batches-sweep.ts index 2cb405d..5bee88e 100644 --- a/apps/journal/app/jobs/import-batches-sweep.ts +++ b/apps/journal/app/jobs/import-batches-sweep.ts @@ -1,4 +1,4 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { and, lt, inArray, count } from "drizzle-orm"; import { getDb } from "../lib/db.ts"; import { importBatches, type ImportBatchStatus } from "@trails-cool/db/schema/journal"; @@ -7,7 +7,7 @@ import { logger } from "../lib/logger.server.ts"; const STALE_MS = 10 * 60 * 1000; const STALE_STATUSES: ImportBatchStatus[] = ["pending", "running"]; -export const importBatchesSweepJob: JobDefinition = { +export const importBatchesSweepJob = defineJournalJob({ name: "import-batches-sweep", cron: "* * * * *", retryLimit: 0, @@ -39,4 +39,4 @@ export const importBatchesSweepJob: JobDefinition = { logger.info({ count: result.length }, "import-batches-sweep: marked stale batches as failed"); }, -}; +}); diff --git a/apps/journal/app/jobs/komoot-bulk-import.test.ts b/apps/journal/app/jobs/komoot-bulk-import.test.ts new file mode 100644 index 0000000..bcf3d42 --- /dev/null +++ b/apps/journal/app/jobs/komoot-bulk-import.test.ts @@ -0,0 +1,62 @@ +import { describe, it, expect, vi, beforeEach } from "vitest"; + +vi.mock("../lib/logger.server.ts", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})); + +const withFreshCredentials = vi.fn(); +vi.mock("../lib/connected-services/manager.ts", () => ({ + withFreshCredentials: (...args: unknown[]) => withFreshCredentials(...args), +})); + +const runKomootBulkImport = vi.fn(); +const markBatchFailed = vi.fn(); +vi.mock("../lib/komoot-bulk-import.server.ts", () => ({ + runKomootBulkImport: (...args: unknown[]) => runKomootBulkImport(...args), + markBatchFailed: (...args: unknown[]) => markBatchFailed(...args), +})); + +import { komootBulkImportJob } from "./komoot-bulk-import.ts"; + +type HandlerJobs = Parameters[0]; + +function jobWith(data: unknown): HandlerJobs { + return [{ id: "j1", data }] as unknown as HandlerJobs; +} + +const PAYLOAD = { batchId: "batch-1", userId: "user-1", serviceId: "svc-1" }; + +describe("komoot-bulk-import job", () => { + beforeEach(() => { + vi.clearAllMocks(); + }); + + it("resolves credentials through the manager — only the serviceId crosses the queue", async () => { + const creds = { mode: "public", komootUserId: "k1" }; + withFreshCredentials.mockImplementation(async (_serviceId, fn) => fn(creds)); + runKomootBulkImport.mockResolvedValue(undefined); + + await komootBulkImportJob.handler(jobWith(PAYLOAD)); + + expect(withFreshCredentials).toHaveBeenCalledWith("svc-1", expect.any(Function)); + expect(runKomootBulkImport).toHaveBeenCalledWith("batch-1", "user-1", creds); + expect(markBatchFailed).not.toHaveBeenCalled(); + }); + + it("marks the batch failed and rethrows when credential resolution fails", async () => { + withFreshCredentials.mockRejectedValue(new Error("needs relink")); + + await expect(komootBulkImportJob.handler(jobWith(PAYLOAD))).rejects.toThrow("needs relink"); + + expect(runKomootBulkImport).not.toHaveBeenCalled(); + expect(markBatchFailed).toHaveBeenCalledWith("batch-1", "needs relink"); + }); + + it("also marks failed when the import itself rejects (markBatchFailed is a no-op on terminal batches)", async () => { + withFreshCredentials.mockImplementation(async (_serviceId, fn) => fn({ mode: "public" })); + runKomootBulkImport.mockRejectedValue(new Error("komoot 500")); + + await expect(komootBulkImportJob.handler(jobWith(PAYLOAD))).rejects.toThrow("komoot 500"); + expect(markBatchFailed).toHaveBeenCalledWith("batch-1", "komoot 500"); + }); +}); diff --git a/apps/journal/app/jobs/komoot-bulk-import.ts b/apps/journal/app/jobs/komoot-bulk-import.ts index 87d788d..fa0c9a2 100644 --- a/apps/journal/app/jobs/komoot-bulk-import.ts +++ b/apps/journal/app/jobs/komoot-bulk-import.ts @@ -1,29 +1,36 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { logger } from "../lib/logger.server.ts"; -import { runKomootBulkImport } from "../lib/komoot-bulk-import.server.ts"; +import { withFreshCredentials } from "../lib/connected-services/manager.ts"; +import { + markBatchFailed, + runKomootBulkImport, + type KomootCreds, +} from "../lib/komoot-bulk-import.server.ts"; -type KomootCreds = - | { mode: "public"; komootUserId: string } - | { mode: "authenticated"; email: string; encryptedPassword: string; komootUserId: string }; - -interface KomootBulkImportData { - batchId: string; - userId: string; - creds: KomootCreds; -} - -export const komootBulkImportJob: JobDefinition = { +export const komootBulkImportJob = defineJournalJob({ name: "komoot-bulk-import", retryLimit: 1, expireInSeconds: 1800, async handler(jobs) { const batch = Array.isArray(jobs) ? jobs : [jobs]; for (const job of batch) { - // pg-boss serialized payload — caller (enqueueOptional) wrote it - // as KomootBulkImportData. Narrow at the boundary. - const { batchId, userId, creds } = job.data as KomootBulkImportData; + const { batchId, userId, serviceId } = job.data; logger.info({ batchId, userId }, "komoot bulk import job started"); - await runKomootBulkImport(batchId, userId, creds); + try { + // Credentials are resolved through the ConnectedServiceManager at + // execution time — the payload carries only the serviceId, so + // nothing credential-shaped sits in the job table and a relink + // between enqueue and execution is picked up here. + await withFreshCredentials(serviceId, (creds) => + runKomootBulkImport(batchId, userId, creds as KomootCreds), + ); + } catch (err) { + // runKomootBulkImport marks its own failures; this covers errors + // before it ran (service missing/not active/needs relink), where + // the batch would otherwise stay "pending" forever. + await markBatchFailed(batchId, err instanceof Error ? err.message : String(err)); + throw err; + } } }, -}; +}); diff --git a/apps/journal/app/jobs/notifications-fanout.ts b/apps/journal/app/jobs/notifications-fanout.ts index ef96f78..70caae9 100644 --- a/apps/journal/app/jobs/notifications-fanout.ts +++ b/apps/journal/app/jobs/notifications-fanout.ts @@ -1,21 +1,17 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { and, eq, isNotNull } from "drizzle-orm"; import { getDb } from "../lib/db.ts"; import { activities, follows, users } from "@trails-cool/db/schema/journal"; import { createNotification } from "../lib/notifications.server.ts"; import { logger } from "../lib/logger.server.ts"; -interface FanoutData { - activityId: string; -} - /** * Fan out an `activity_published` notification to every accepted * follower of the activity's owner. Idempotent at the DB level via the * `(recipient_user_id, type, subject_id)` unique partial index — a * retry after partial failure won't double-insert. */ -export const notificationsFanoutJob: JobDefinition = { +export const notificationsFanoutJob = defineJournalJob({ name: "notifications-fanout", retryLimit: 3, expireInSeconds: 300, @@ -24,11 +20,10 @@ export const notificationsFanoutJob: JobDefinition = { // process whichever shape we get. const batch = Array.isArray(job) ? job : [job]; for (const item of batch) { - const data = item.data as FanoutData; - await fanout(data.activityId); + await fanout(item.data.activityId); } }, -}; +}); export async function fanout(activityId: string): Promise { const db = getDb(); diff --git a/apps/journal/app/jobs/notifications-purge.ts b/apps/journal/app/jobs/notifications-purge.ts index cf87fa9..625689c 100644 --- a/apps/journal/app/jobs/notifications-purge.ts +++ b/apps/journal/app/jobs/notifications-purge.ts @@ -1,4 +1,4 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { purgeReadOlderThan } from "../lib/notifications.server.ts"; import { logger } from "../lib/logger.server.ts"; @@ -7,7 +7,7 @@ import { logger } from "../lib/logger.server.ts"; * than 90 days; unread rows are kept indefinitely so users never miss * an event. */ -export const notificationsPurgeJob: JobDefinition = { +export const notificationsPurgeJob = defineJournalJob({ name: "notifications-purge", cron: "30 3 * * *", // daily at 03:30 UTC (offset from demo-bot-prune to spread load) retryLimit: 1, @@ -18,4 +18,4 @@ export const notificationsPurgeJob: JobDefinition = { logger.info({ days, purged }, "notifications-purge"); return { days, purged }; }, -}; +}); diff --git a/apps/journal/app/jobs/payloads.test.ts b/apps/journal/app/jobs/payloads.test.ts new file mode 100644 index 0000000..00e4844 --- /dev/null +++ b/apps/journal/app/jobs/payloads.test.ts @@ -0,0 +1,47 @@ +import { describe, it, expect } from "vitest"; + +// Importing every job module pulls in their (transitively heavy) +// server deps; mock the leaf modules with side effects so this stays a +// unit test of the registry wiring. +import { vi } from "vitest"; +vi.mock("../lib/db.ts", () => ({ getDb: vi.fn() })); +vi.mock("../lib/logger.server.ts", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})); + +describe("job registry", () => { + it("every job's queue name is unique and a key of JobPayloads", async () => { + const modules = await Promise.all([ + import("./backfill-user-keypairs.ts"), + import("./consumed-jti-sweep.ts"), + import("./deliver-activity.ts"), + import("./demo-bot-generate.ts"), + import("./demo-bot-prune.ts"), + import("./federation-kv-sweep.ts"), + import("./garmin-import-activity.ts"), + import("./import-batches-sweep.ts"), + import("./komoot-bulk-import.ts"), + import("./notifications-fanout.ts"), + import("./notifications-purge.ts"), + import("./poll-remote-actor.ts"), + import("./poll-remote-outboxes.ts"), + import("./send-welcome-email.ts"), + ]); + + const names = modules.flatMap((m) => + Object.values(m as Record) + .filter( + (v): v is { name: string; handler: unknown } => + typeof v === "object" && v !== null && "handler" in v && "name" in v, + ) + .map((def) => def.name), + ); + + expect(names).toHaveLength(14); + expect(new Set(names).size).toBe(names.length); + // pg-boss v11+ queue-name constraint, mirrored from packages/jobs + for (const name of names) { + expect(name).toMatch(/^[A-Za-z0-9_.-]+$/); + } + }); +}); diff --git a/apps/journal/app/jobs/payloads.ts b/apps/journal/app/jobs/payloads.ts new file mode 100644 index 0000000..6178f5c --- /dev/null +++ b/apps/journal/app/jobs/payloads.ts @@ -0,0 +1,41 @@ +import { defineJob, type JobDefinition, type TypedJobDefinition } from "@trails-cool/jobs"; +import type { DeliveryPayload } from "../lib/federation-delivery.server.ts"; +import type { GarminImportData } from "../lib/connected-services/providers/garmin/import.server.ts"; + +/** + * Every journal job queue and its payload shape, in one place. The + * typed `enqueue` / `enqueueOptional` in boss.server.ts and the + * `defineJournalJob` helper below both key off this map, so an enqueue + * site and its handler cannot drift apart, and a queue-name typo is a + * compile error instead of an orphaned queue. + * + * `void` marks cron-only jobs that are never enqueued with data. + */ +export interface JobPayloads { + "backfill-user-keypairs": Record; + "consumed-jti-sweep": void; + "deliver-activity": DeliveryPayload; + "demo-bot-generate": void; + "demo-bot-prune": void; + "federation-kv-sweep": void; + "garmin-import-activity": GarminImportData; + "import-batches-sweep": void; + "komoot-bulk-import": { batchId: string; userId: string; serviceId: string }; + "notifications-fanout": { activityId: string }; + "notifications-purge": void; + "poll-remote-actor": { actorIri: string }; + "poll-remote-outboxes": void; + "send-welcome-email": { email: string; username: string }; +} + +export type JobName = keyof JobPayloads; + +/** + * defineJob, constrained to the journal's queue map: the name must be + * a known queue and the handler's payload type follows from it. + */ +export function defineJournalJob( + definition: TypedJobDefinition & { name: K }, +): JobDefinition { + return defineJob(definition); +} diff --git a/apps/journal/app/jobs/poll-remote-actor.ts b/apps/journal/app/jobs/poll-remote-actor.ts index 31582c4..18dc1d7 100644 --- a/apps/journal/app/jobs/poll-remote-actor.ts +++ b/apps/journal/app/jobs/poll-remote-actor.ts @@ -1,26 +1,23 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { pollRemoteActor } from "../lib/federation-ingest.server.ts"; import { logger } from "../lib/logger.server.ts"; -interface PollPayload { - actorIri?: string; -} - /** * Poll one remote trails actor's outbox (spec §7). Enqueued by the * inbox Accept(Follow) listener (first poll, 7.5) and fanned out by * the poll-remote-outboxes cron sweep (7.1). */ -export const pollRemoteActorJob: JobDefinition = { +export const pollRemoteActorJob = defineJournalJob({ name: "poll-remote-actor", retryLimit: 2, expireInSeconds: 120, async handler(jobs) { for (const job of jobs) { - const { actorIri } = (job.data ?? {}) as PollPayload; + // Defensive: jobs enqueued before the typed seam may carry no data. + const actorIri = job.data?.actorIri; if (!actorIri) continue; const result = await pollRemoteActor(actorIri); logger.info({ actorIri, result }, "poll-remote-actor"); } }, -}; +}); diff --git a/apps/journal/app/jobs/poll-remote-outboxes.ts b/apps/journal/app/jobs/poll-remote-outboxes.ts index c6cf085..1847ad2 100644 --- a/apps/journal/app/jobs/poll-remote-outboxes.ts +++ b/apps/journal/app/jobs/poll-remote-outboxes.ts @@ -1,4 +1,4 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { listActorsDuePolling } from "../lib/federation-ingest.server.ts"; import { enqueueOptional } from "../lib/boss.server.ts"; import { logger } from "../lib/logger.server.ts"; @@ -9,7 +9,7 @@ import { logger } from "../lib/logger.server.ts"; * been polled within the last hour, and fan out one poll-remote-actor * job each. Per-host pacing lives in the poll itself. */ -export const pollRemoteOutboxesJob: JobDefinition = { +export const pollRemoteOutboxesJob = defineJournalJob({ name: "poll-remote-outboxes", cron: "*/5 * * * *", retryLimit: 1, @@ -22,4 +22,4 @@ export const pollRemoteOutboxesJob: JobDefinition = { if (due.length > 0) logger.info({ due: due.length }, "poll-remote-outboxes sweep"); return { due: due.length }; }, -}; +}); diff --git a/apps/journal/app/jobs/send-welcome-email.ts b/apps/journal/app/jobs/send-welcome-email.ts index b8a426d..69551d3 100644 --- a/apps/journal/app/jobs/send-welcome-email.ts +++ b/apps/journal/app/jobs/send-welcome-email.ts @@ -1,26 +1,21 @@ -import type { JobDefinition } from "@trails-cool/jobs"; +import { defineJournalJob } from "./payloads.ts"; import { sendWelcome } from "../lib/email.server.ts"; import { logger } from "../lib/logger.server.ts"; -interface WelcomeEmailData { - email: string; - username: string; -} - /** * Queue-backed welcome email send. The old code did * `sendWelcome(...).catch(log)` inline, which silently dropped failures. * pg-boss retries on transient failure (3 attempts) and surfaces persistent * failures via the dead-letter queue / logs. */ -export const sendWelcomeEmailJob: JobDefinition = { +export const sendWelcomeEmailJob = defineJournalJob({ name: "send-welcome-email", retryLimit: 3, expireInSeconds: 120, async handler(job) { const batch = Array.isArray(job) ? job : [job]; for (const item of batch) { - const { email, username } = item.data as WelcomeEmailData; + const { email, username } = item.data; try { await sendWelcome(email, username); } catch (err) { @@ -29,4 +24,4 @@ export const sendWelcomeEmailJob: JobDefinition = { } } }, -}; +}); diff --git a/apps/journal/app/lib/boss.server.ts b/apps/journal/app/lib/boss.server.ts index a224d38..94970a4 100644 --- a/apps/journal/app/lib/boss.server.ts +++ b/apps/journal/app/lib/boss.server.ts @@ -16,6 +16,7 @@ // global registry key. import { logger } from "./logger.server.ts"; +import type { JobName, JobPayloads } from "../jobs/payloads.ts"; // Structurally typed (we only need `send`) so we don't have to pull // pg-boss into the journal app's dep graph just for the typedef. @@ -51,20 +52,33 @@ export function getBoss(): BossLike { return boss; } +/** + * Enqueue a job. The queue name must be a key of JobPayloads and the + * payload must match its declared shape. Throws if the queue is down — + * use this when the caller's correctness depends on the job existing + * (e.g. a batch row that would otherwise wait forever). + */ +export async function enqueue( + queue: K, + data: JobPayloads[K], + options?: BossSendOptions, +): Promise { + await getBoss().send(queue, data, options); +} + /** * Best-effort enqueue: log + swallow errors so a downstream queue * outage doesn't fail the user-visible request that triggered the * fan-out. Use this for "fire and forget" notifications work. */ -export async function enqueueOptional( - queue: string, - data: unknown, +export async function enqueueOptional( + queue: K, + data: JobPayloads[K], ctx: Record = {}, options?: BossSendOptions, ): Promise { try { - const boss = getBoss(); - await boss.send(queue, data, options); + await enqueue(queue, data, options); } catch (err) { logger.warn({ err, queue, ...ctx }, "boss.send failed; continuing"); } diff --git a/apps/journal/app/lib/komoot-bulk-import.server.ts b/apps/journal/app/lib/komoot-bulk-import.server.ts index 7468fd7..474be09 100644 --- a/apps/journal/app/lib/komoot-bulk-import.server.ts +++ b/apps/journal/app/lib/komoot-bulk-import.server.ts @@ -1,4 +1,4 @@ -import { eq } from "drizzle-orm"; +import { and, eq, notInArray } from "drizzle-orm"; import { getDb } from "./db.ts"; import { importBatches } from "@trails-cool/db/schema/journal"; import { fetchKomootTours, fetchKomootTourGpx } from "./komoot.server.ts"; @@ -7,10 +7,29 @@ import { createActivity } from "./activities.server.ts"; import { decrypt } from "./crypto.server.ts"; import { logger } from "./logger.server.ts"; -type KomootCreds = +export type KomootCreds = | { mode: "public"; komootUserId: string } | { mode: "authenticated"; email: string; encryptedPassword: string; komootUserId: string }; +/** + * Marks a batch failed unless it already reached a terminal state. + * Used by the job handler when credential resolution fails before + * runKomootBulkImport (which owns failure-marking for its own errors) + * ever runs — otherwise the batch would sit "pending" forever. + */ +export async function markBatchFailed(batchId: string, message: string): Promise { + const db = getDb(); + await db + .update(importBatches) + .set({ status: "failed", errorMessage: message, completedAt: new Date() }) + .where( + and( + eq(importBatches.id, batchId), + notInArray(importBatches.status, ["completed", "failed"]), + ), + ); +} + function getBasicAuthToken(creds: KomootCreds): string | undefined { if (creds.mode !== "authenticated") return undefined; const password = decrypt(creds.encryptedPassword); diff --git a/apps/journal/app/routes/api.sync.komoot.import.ts b/apps/journal/app/routes/api.sync.komoot.import.ts index 4c84a95..2394eeb 100644 --- a/apps/journal/app/routes/api.sync.komoot.import.ts +++ b/apps/journal/app/routes/api.sync.komoot.import.ts @@ -7,7 +7,7 @@ import { data } from "react-router"; import type { Route } from "./+types/api.sync.komoot.import"; import { requireSessionUser } from "~/lib/auth/session.server"; import { getService } from "~/lib/connected-services/manager"; -import { getBoss } from "~/lib/boss.server"; +import { enqueue } from "~/lib/boss.server"; import { getDb } from "~/lib/db"; import { importBatches } from "@trails-cool/db/schema/journal"; @@ -27,11 +27,12 @@ export async function action({ request }: Route.ActionArgs) { status: "pending", }); - const boss = getBoss(); - await boss.send("komoot-bulk-import", { + // Only the serviceId crosses the queue — the job resolves fresh + // credentials through the ConnectedServiceManager at execution time. + await enqueue("komoot-bulk-import", { batchId, userId: user.id, - creds: service.credentials, + serviceId: service.id, }); return data({ batchId }); diff --git a/apps/journal/server.ts b/apps/journal/server.ts index 682d93f..02ca41d 100644 --- a/apps/journal/server.ts +++ b/apps/journal/server.ts @@ -193,7 +193,7 @@ server.listen(port, async () => { await startWorker(boss, jobs); // Register the started boss so feature code can enqueue jobs against // the same instance via getBoss() / enqueueOptional(). - const { setBoss } = await import("./app/lib/boss.server.ts"); + const { setBoss, enqueue } = await import("./app/lib/boss.server.ts"); setBoss(boss); logger.info("Background job worker started"); @@ -201,7 +201,7 @@ server.listen(port, async () => { // users get keys before any federation traffic). Each run only // touches users whose public_key IS NULL, so repeats are no-ops. if (process.env.FEDERATION_ENABLED === "true") { - await boss.send("backfill-user-keypairs", {}); + await enqueue("backfill-user-keypairs", {}); logger.info("federation keypair backfill enqueued"); } }); diff --git a/packages/jobs/src/index.ts b/packages/jobs/src/index.ts index 4d47f98..8aa2f13 100644 --- a/packages/jobs/src/index.ts +++ b/packages/jobs/src/index.ts @@ -1,3 +1,4 @@ export { createBoss } from "./boss.ts"; export { startWorker } from "./worker.ts"; -export type { JobDefinition } from "./types.ts"; +export { defineJob } from "./types.ts"; +export type { JobDefinition, TypedJobDefinition } from "./types.ts"; diff --git a/packages/jobs/src/types.test.ts b/packages/jobs/src/types.test.ts new file mode 100644 index 0000000..4310424 --- /dev/null +++ b/packages/jobs/src/types.test.ts @@ -0,0 +1,28 @@ +import { describe, it, expect } from "vitest"; +import type { Job } from "pg-boss"; +import { defineJob } from "./types.ts"; + +describe("defineJob", () => { + it("returns the definition unchanged (the cast is type-level only)", async () => { + const seen: string[] = []; + const def = defineJob<{ widgetId: string }>({ + name: "widget-job", + retryLimit: 2, + async handler(jobs) { + for (const job of jobs) { + // job.data is typed { widgetId: string } — no cast needed + seen.push(job.data.widgetId); + } + return null; + }, + }); + + expect(def.name).toBe("widget-job"); + expect(def.retryLimit).toBe(2); + + // The returned JobDefinition's handler is the same function and + // receives whatever pg-boss hands it. + await def.handler([{ data: { widgetId: "w1" } } as Job]); + expect(seen).toEqual(["w1"]); + }); +}); diff --git a/packages/jobs/src/types.ts b/packages/jobs/src/types.ts index 3f75205..ae071a2 100644 --- a/packages/jobs/src/types.ts +++ b/packages/jobs/src/types.ts @@ -1,12 +1,15 @@ import type { Job } from "pg-boss"; /** - * A pg-boss job definition. The payload type is intentionally - * `unknown` at this boundary: a heterogeneous `JobDefinition[]` array - * (e.g. the one in `apps/journal/server.ts`) was previously forced to - * cast each typed job to `any` because of contravariance — a handler - * taking `Job` is not assignable to one taking - * `Job`. Handlers narrow internally instead. + * A pg-boss job definition. The payload type is `unknown` at this + * boundary: a heterogeneous `JobDefinition[]` array (e.g. the one in + * `apps/journal/server.ts`) would otherwise be impossible because of + * contravariance — a handler taking `Job` is not + * assignable to one taking `Job`. + * + * Don't write handlers against this type directly; author them with + * `defineJob`, which keeps the payload typed inside the handler and + * performs the contravariance cast once, here in the package. */ export interface JobDefinition { name: string; @@ -15,3 +18,23 @@ export interface JobDefinition { retryLimit?: number; expireInSeconds?: number; } + +/** A JobDefinition whose handler sees the payload type it was enqueued with. */ +export interface TypedJobDefinition { + name: string; + handler: (jobs: Job[]) => Promise; + cron?: string; + retryLimit?: number; + expireInSeconds?: number; +} + +/** + * Author a job with a typed payload. The handler receives + * `Job[]` — no `job.data as ...` casts at the call sites. + * This is the single place the contravariance cast happens. + */ +export function defineJob( + definition: TypedJobDefinition, +): JobDefinition { + return definition as unknown as JobDefinition; +}