Skip to content

Phase 22: real lakehouse archival writer for payment_archival_jobs - #68

Merged
munisp merged 8 commits into
mainfrom
phase22/lakehouse-writer
Oct 1, 2026
Merged

munisp merged 8 commits into
mainfrom
phase22/lakehouse-writer

Conversation

@munisp

@munisp munisp commented Oct 1, 2026

Copy link
Copy Markdown
Owner

What

The payment-archival cron (runPaymentArchivalCron in server/_core/index.ts) honestly enqueues payment_archival_jobs rows with status='pending' (NULL storageUri, 0 bytes) when PAYMENT_ARCHIVE_SINK_BUCKET is set — but nothing ever fulfilled them. This PR adds the real writer.

Files

  • server/workers/paymentArchivalWriter.ts (new) — lakehouse writer, styled on server/paymentWorker.ts:
    • Claims pending jobs idempotently: UPDATE payment_archival_jobs SET status='running' WHERE id=? AND status='pending' RETURNING — a concurrent claim updates 0 rows and is skipped (exactly-once archival).
    • Exports the committed payment_queue rows in the job's [periodStart, periodEnd) tier window (same predicate the cron counted with), bounded at 500k rows/job.
    • PUTs the payload to s3://$PAYMENT_ARCHIVE_SINK_BUCKET/{tier}/{YYYY-MM-DD}/{jobId}.csv.
    • Marks the job completed with the real storageUri, actual bytesWritten (measured from the uploaded payload) and true transfersArchived count — only after the object exists. Any sink error → status='failed' + errorMessage, storageUri stays NULL. No fabricated URIs.
    • Lifecycle mirrors the payment worker: startPaymentArchivalWriter() / stopPaymentArchivalWriter() / status getter; disabled with a single log line when PAYMENT_ARCHIVE_SINK_BUCKET is unset.
  • server/_core/index.ts — starts/stops the writer beside startPaymentWorker(), with SIGTERM/SIGINT graceful stop. Also corrects the tier comments (CSV, not Parquet). Note: the earlier commit on this branch had left index.ts truncated at 762 lines (dropping startPaymentWorker, the archival cron and balance-drift cron); this PR restores the full file from latest main plus the writer wiring.
  • server/workers/paymentArchivalWriter.test.ts (new) — vitest, DB and S3 mocked: disabled-without-bucket, happy path (real URI/bytes/count asserted against the actual payload), sink-error → failed + errorMessage + no storageUri, and the 0-row-claim concurrency skip.

CSV vs Parquet decision

CSV. package.json has no parquet writer (no parquetjs/duckdb/arrow), so Parquet would require a new runtime dependency; the object key extension (.csv) reflects the real format. Format upgrade to Parquet can be a follow-up behind the same job contract.

Sink client

Uses the already-declared @aws-sdk/client-s3 dependency pointed at the RustFS endpoint (PAYMENT_ARCHIVE_SINK_* env overrides, falling back to RUSTFS_ENDPOINT/ACCESS_KEY/SECRET_KEY/REGION). The document-vault rustfsSvcClient.ts was deliberately not reused: it is scoped to the single vault bucket (RUSTFS_BUCKET) and force-prefixes keys with a caller namespace, so it cannot write to the operator-configured sink bucket at the required key pattern. No new runtime dependencies.

Validation

  • Schema cross-checked against drizzle/schema.ts (payment_archival_job_status enum includes pending/running/completed/failed; errorMessage, completedAt, storageUri, bytesWritten bigint, createdAt defaultNow all present).
  • Unit tests cover the honest-failure and idempotent-claim paths. tsc/vitest not run locally (no sandbox network for pnpm install); CI should run pnpm check + pnpm test.

munisp added 8 commits October 1, 2026 11:18
Polls pending payment_archival_jobs rows, claims them idempotently
(UPDATE ... WHERE status='pending' RETURNING), exports committed
payment_queue rows in the job's tier window as CSV, PUTs to the
S3-compatible sink (RustFS via @aws-sdk/client-s3, an existing
dependency) at s3://$PAYMENT_ARCHIVE_SINK_BUCKET/{tier}/{date}/{jobId}.csv,
and marks the job completed with the real storageUri/bytesWritten/
transfersArchived — or failed with errorMessage. No fabricated URIs.
…startup (part 3)

The part-1 commit left server/_core/index.ts truncated at 762 lines
(dropping startPaymentWorker, runPaymentArchivalCron, balance drift cron,
and all later startup wiring). This restores the full main content and adds
startPaymentArchivalWriter()/stopPaymentArchivalWriter() beside the payment
worker lifecycle (SIGTERM/SIGINT), plus corrects the tier comments (CSV, not
Parquet).
@munisp
munisp merged commit a6e43ab into main Oct 1, 2026
1 check failed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant