Pipeline orchestration: spike DBOS as a checkpointing layer inside Cloud Run Jobs¶
Status: superseded — the registry tracer was withdrawn on 2026-08-16 and the web-enrichment target is replaced by ADR-0016 Date: 2026-07-20 Deciders: Chris (solo founder) Depends on: ADR-0005 (watermark discipline, run ledger, staging+swap), ADR-0007 (reuse of the ADR-0005 framework) Tracking: issue #67
Supersession of the web-enrichment half (2026-08-20)¶
The web tracer produced the evidence that this spike required. It also showed that one Cloud Run Job is the wrong production boundary for browser work, synchronous LLM calls, and asynchronous DataForSEO work. ADR-0016 keeps the domain ledgers, raw evidence, budgets, and per-company failure boundary, but replaces DBOS execution state with explicit Postgres work state, a transactional outbox, Cloud Tasks, and focused Cloud Run runtimes.
The existing test tracer can remain only until the ADR-0016 implementation cuts over. It is no longer a target architecture. The cutover removes its DBOS workflows, queues, lifecycle code, Job, and system schema. It does not add a compatibility layer.
Withdrawal of the registry half (2026-08-16)¶
The DBOS registry tracer — orchestration/registry.py, its retirement gate,
its three Cloud Run Jobs, and its schedulers — is deleted. The web-enrichment
tracer is untouched and this ADR still governs it.
The fitness report
recommended adopt for both families and reproduced none of bullet 16's
irreducible blockers. Nothing here contradicts that finding. The tracer is
removed because nothing was going to act on it: no production cutover was ever
authorised, the follow-up measurement (#195) stayed blocked, and its only live
role had become holding open _run_delta's per-batch post_upsert_hook and
the --max-docs capped loader — a second load path that bypasses the staging
swap and the guard rails, kept for a run mode the operator runbook already
tells people not to use.
Keeping a spike alive for an adoption that may happen later is the shape this repo's rules reject (CLAUDE.md 2 and 3). Re-opening the question means rebuilding the tracer, and the report and this ADR keep the evidence needed to decide it.
Context¶
This repository has two workload families with different execution shapes:
source/entity-oriented registry ingesters and shortlist/company-oriented web
enrichment. Both distinguish domain durability from execution durability, per
CONTEXT.md and ADR-0005:
- Domain durability: an
ingestion_runsbusiness/provenance ledger, watermarks advanced only on success, immutable raw artifacts, and staging-table-plus-atomic-swap correctness for backfills. - Execution durability: idempotent reruns after a failed process, currently recovered by restarting the whole Cloud Run Job rather than resuming after its last completed stage.
- Scheduling and retries wired per ingester through a Terraform-managed
Cloud Run Job and Cloud Scheduler trigger (
infra/batch_job.tf).
The domain records are intentional and remain authoritative. The narrower gap is execution recovery: the current CVR delta materializes the scroll, writes raw output, maps, loads, and stamps the ledger in one process, so interruption restarts work already completed outside the final database transaction. DBOS Transact is a durable-execution framework that checkpoints workflow and step state in Postgres. If it holds up, it is a candidate execution substrate for every ingester and enrichment pipeline in this repo, not a replacement for their domain state.
This distinction also matches the downstream web-scraping workload currently
represented by the Prefect PoC in data-collection-prephase/SCRAPER/POC.
Those workflows fan out by company, conditionally perform DNS, SERP, browser,
LLM, directory, and extraction stages, enforce per-company cost/time budgets,
share rate-limited adapters, retain a raw evidence corpus, and pass durable
artifacts between independent pipelines. DBOS may eventually replace Prefect
as their execution/checkpointing layer, but must not absorb the evidence
ledger, raw corpus, provenance, cost accounting, or published signal contract.
It also does not make scraper orchestration the registry-ingestion model. CVR
and Regnskabsdata remain source/entity-oriented bulk ingesters; web enrichment
is shortlist/campaign-oriented and may fan out by company across paid,
rate-limited adapters. The two workload families may share DBOS as a runtime
substrate while retaining different workflow graphs, checkpoint economics,
failure outcomes, concurrency controls, and deployment evolution.
No ingester has been migrated and no infra has changed yet. This ADR records the decision to spend a bounded spike finding out whether DBOS is worth adopting, and what exactly that spike must produce before any real cutover is considered.
Decision¶
- Run a bounded spike evaluating DBOS Transact strictly as an in-process
durable-execution and checkpointing layer inside the existing Cloud Run
Job topology. The spike explicitly retains Cloud Scheduler, per-workload
Cloud Run Jobs,
ingestion_runs, raw artifacts, and staging+atomic-swap. DBOS system state records execution progress; domain state continues to record what ran, what evidence was produced, what was published, and which watermark is safe for the next run. - Implement two serial tracer bullets inside the bounded spike. The first
is a CVR compatibility tracer: one coarse optimized acquisition/raw-write/
map/load step inside a DBOS workflow, retaining the in-memory hot path and
proving Job-mode startup, stable identity,
ingestion_runsmapping, and interruption recovery without production changes. The second is a narrow web-enrichment value tracer based on01-find-companyfor approximately 5–10 fixture companies: one child workflow per company, selected durable boundaries around paid or slow SERP/browser/LLM operations, external raw artifacts, and forced termination proving completed paid work is not repeated. The slices execute serially per ADR-0003. This is not a full Prefect migration or production web-enrichment deployment. - Spike the open-source DBOS Transact library in Job mode, not DBOS Cloud or DBOS Conductor. This avoids an always-on Cloud Run Service/worker pool and keeps the experiment focused on process-restart recovery. Whether the DBOS system tables can safely share the Supabase Postgres deployment—and with what database/schema and privileges—must be proven rather than assumed. Managed or self-hosted Conductor remains a separate future option if a later architecture needs distributed recovery or operational control.
- The spike runs in the
testGCP project only and does not touch theprodCVR ingester or its liveingestion_runshistory.batch_job.tfretains its Job + Scheduler shape for the duration of the spike; only a separate test-only job/configuration may opt into the DBOS entrypoint. - The spike must produce a written fitness report, not just working code, covering at minimum: domain-ledger and watermark correctness under a forced mid-run termination plus Cloud Run task restart; which expensive stages are skipped on recovery; DBOS system-data isolation and privileges; schema and code-version compatibility; rough cost at current ingester volumes; and measured repeated-call avoidance in the web-enrichment tracer. A production migration is a separate future ADR, decided from that evidence—this ADR authorizes only the two tracer bullets.
- Use an artifact-reference-only contract across durable step boundaries, without maximizing the number of boundaries. If a large or externally sourced value must cross from one durable step to another, the producing step persists it through the workload's domain raw-storage contract and returns only a small typed reference containing identifiers, content hash, source and retrieval metadata. CVR document batches, NDJSON, HTML, markdown, screenshots, PDFs, SERP payloads, and raw LLM responses must not be stored as DBOS step outputs. However, cheap CVR acquisition may keep scrolling, raw writing, mapping, and loading together in one coarse, optimized step and reuse the in-memory documents; interruption replays that whole cheap unit. Paid, slow, or rate-limited web/LLM operations justify finer steps. Checkpoint granularity follows measured recovery cost rather than a framework-driven requirement for complete step-by-step replay.
- Use workload-semantic workflow and failure boundaries. Registry
ingestion has one workflow per source, entity, run type, and logical
schedule window. Web enrichment has a lightweight campaign workflow that
coordinates one child workflow per company and pipeline/extractor version.
Company outcomes such as unresolved, blocked, or budget-exhausted are
successful domain results; infrastructure or contract failures fail that
company without replaying unrelated companies. Dependencies such as
find-employeesconsumingfind-companyuse a persisted, versioned artifact reference so either pipeline remains independently rerunnable. The shared architecture is a set of execution primitives, not one generic workflow graph or one universal failure policy. - Isolate DBOS system state within test Supabase Postgres. The spike uses
the existing test Postgres deployment but a dedicated DBOS schema and
dedicated runtime role per application, each reached through its own
DBOS_SYSTEM_DATABASE_URL; the domain-data connection remains separate. A privileged CI/CD migration step creates or upgrades DBOS objects and grants each runtime role DML/execute access only. Application startup must not perform DBOS DDL. DBOS connections use direct or session-mode endpoints compatible withLISTEN/NOTIFY, not transaction-mode pooling. DBOS objects are outside the public/read-contract schema. - Separate interruption, transient, and business retries. Each separate
test-only Cloud Run spike Job allows one task retry so an interrupted
process restarts and recovers the same deterministic DBOS workflow ID.
DBOS step retries are bounded and enabled only for classified transient
failures such as selected timeouts, connection resets,
429, and5xxresponses; robots denial, validation failures, budget exhaustion, unsupported payloads, and deterministic4xxoutcomes are not retried. An intentional operator rerun after a permanent failure, or a new web refresh/extractor version, creates a new workflow ID rather than silently resuming anERRORworkflow with different code. Registry IDs encode source/entity/run-type/logical-window; web child IDs encode campaign, pipeline, extractor version, and CVR. Domain run records retain the DBOS ID for correlation. Existing production Job retry settings remain unchanged until a later adoption ADR. - Use SQLAlchemy plus the DBOS datasource as the canonical domain-database
boundary. The target database engine uses SQLAlchemy with the psycopg
driver. Eligible ledger transitions, artifact manifests, web signal
publication, company/campaign outcomes, and budget/cost records execute as
datasource transactions so the domain mutation and DBOS completion record
commit atomically. Loaders do not own commits. High-throughput PostgreSQL
operations retain
COPYthrough the underlying psycopg driver connection; adopting SQLAlchemy does not replace efficient set-based SQL with ORM entity writes. The optimized coarse CVR acquisition/raw-write/map/load operation is the documented exception: because acquisition and loading deliberately remain one ordinary durable step, its database portion is at-least-once and must remain convergent under replay. External systems likewise remain idempotent step effects rather than pretending to participate in a Postgres transaction. DBOS datasource tracking objects live in the dedicated DBOS schema and are created by the privileged migration path, not at runtime. - Version workflow code immutably and drain before retirement. DBOS
application_versionis the deployed image's Git SHA. It represents replay compatibility and remains distinct from the logical workflow ID and from extractor/prompt/schema versions recorded on domain output. Cloud Run task retries use the same image/version. Before CD retires an application version, it checks forPENDING,ENQUEUED, orDELAYEDworkflows and either runs that immutable old image to drain them or blocks promotion. Database migrations stay backward-compatible while an old workflow version is active. A failed workflow that needs corrected code is explicitly forked onto the new version. DBOS patching is disabled for these bounded Jobs; conditional historical workflow branches are not the default upgrade mechanism. - Keep registry ingestion and web enrichment as separate execution profiles. The adapter-owned queue/broker hierarchy applies only to web enrichment: a company queue controls active company workflows; SERP brokers preserve batch submission and shared polling; browser controls global capacity plus per-domain politeness; LLM controls are partitioned by provider/model; and directory adapters remain behind their own legal gate. Company budgets are reserved in exactly-once datasource transactions before paid submission and settled afterward with stable operation identity. CVR and other registry ingesters do not acquire campaign workflows, per-company queues, or paid-web budgets: they retain scheduled source/entity Jobs, coarse optimized acquisition, watermark discipline, and bulk database loading. No generic base workflow or universal pipeline graph spans both families; only narrowly reusable DBOS configuration, transaction, identity, and observability primitives may be shared.
- Deploy the two tracer bullets as separate DBOS applications from the start. Registry ingestion and web enrichment have independent DBOS application names, system schemas, runtime roles, workflow/queue namespaces, immutable application versions, Cloud Run Jobs, and drain lifecycles. A release or pending workflow in one cannot block or alter the other. They may share only small infrastructure-neutral helpers for typed identity, datasource/configuration setup, artifact references, and telemetry. They do not share base workflow classes, queue registration, schedules, retry policy, or deployment versions.
- Implement the web tracer in this sole-writer repository. The
prephase
SCRAPER/POCremains a behavioral reference, fixture corpus, and baseline for selected outcomes; the deployed tracer neither imports the sibling checkout nor runs the PoC package. Only the narrowfind-companyvertical slice needed to test per-company DBOS recovery is ported into the separately deployable web-enrichment application. Its production-shaped contracts replace PoC file state: immutable raw artifacts use the test GCS store, domain output uses Supabase datasource transactions, and workflow state uses the web application's dedicated DBOS schema. Full PoC feature migration or archival requires later work after parity evidence. - Operate the self-hosted spike through GCP-native observability and explicit retention. Both applications emit DBOS OpenTelemetry workflow and step spans to Cloud Trace and structured JSON to Cloud Logging, with application/version, workflow/parent ID, domain run or campaign ID, source/entity or CVR, step/attempt, artifact/provider, cost/token, and outcome correlation fields. Cloud Monitoring alerts cover failed Job executions, workflow errors, over-age pending/enqueued work, web queue age or depth, failed version drains, registry staleness, and scraper error/budget anomalies. Domain ledgers remain the authoritative business summary; DBOS state is execution detail. A maintenance Job per application deletes successful DBOS histories after 30 days and failed/cancelled histories after 90 days through supported DBOS APIs, never deleting active work. Domain evidence and raw artifacts follow their independent retention contracts. Workflow arguments/outputs contain identifiers and artifact references, not raw pages, prompts, personal records, or extracted payloads.
- Treat adoption as the default outcome unless a tracer proves DBOS does not work for that application family. The fitness report records correctness, recovery, performance, security, cost, and operations evidence, but ordinary optimization findings create follow-up work rather than a rejection. Rejection requires a demonstrated blocker that cannot be corrected without abandoning the accepted architecture: failure to recover safely in Cloud Run Job mode; domain corruption or unsafe watermark movement; unavoidable duplicate paid work or budget overspend; incompatible Supabase transaction/connection behavior; unacceptable CVR hot-path degradation that cannot be optimized; a requirement for runtime DDL or excessive privilege; or inability to version, drain, and operate workflows safely. Registry ingestion and web enrichment are decided independently, so failure of one tracer does not veto adoption for the other.
Options considered¶
Orchestration substrate to spike¶
| Option | For | Against | Verdict |
|---|---|---|---|
| DBOS Transact library in Cloud Run Job mode (chosen) | Postgres-backed execution checkpoints without changing the scheduler/runtime topology; same primitive may cover ingesters and per-company web enrichment | Job restart/recovery and step boundaries are unproven; system-data isolation and schema/version lifecycle are unknown; does not provide an always-on control plane | Spike |
| DBOS Cloud (managed) | Removes self-hosting ops burden | Second paid vendor introduced before self-hosted evidence exists; no other ADR in this repo commits to a paid orchestration SaaS | Deferred — revisit only if self-hosted spike shows real ops cost |
| Prefect (web-scraping PoC baseline) | Already demonstrates named flows/tasks and can grow into managed orchestration | The PoC's actual concurrency, budgets, artifact durability, and evidence semantics remain application-owned; adopting its server/control plane is not yet justified | Baseline for web-enrichment applicability, not part of the CVR implementation spike |
| Temporal / Dagster / Airflow | Mature durable-workflow/data-orchestration tooling | Needs a separate control plane and heavier operational footprint than the current solo-founder Cloud Run Job model; no demonstrated advantage for this bounded spike | Rejected for the spike |
| Status quo (Cloud Run Job + Scheduler + domain-owned recovery) | Already built and working for CVR and Regnskabsdata; the web PoC also owns its cache/evidence recovery | Execution recovery remains workload code: registry Jobs replay coarse work, while web enrichment must coordinate paid-stage resume, batching, and per-company isolation itself | Baseline retained in prod during the spike |
Consequences¶
- Two orchestration patterns exist side by side during the spike: the
existing Cloud Run Job pattern (live,
prod, unchanged) and the DBOS spike (testonly, with one CVR compatibility tracer and one fixture-backed web-enrichment value tracer). This is deliberate and temporary. - Cloud Scheduler and Cloud Run Jobs are not replacement targets under this ADR. A successful spike may replace internal whole-process reruns with step-level recovery, but it does not authorize an always-on DBOS runtime.
- Domain durability stays explicit.
ingestion_runs, enrichment evidence, raw artifacts, budgets/costs, and published signals remain queryable without interpreting DBOS internal tables. - Workflow checkpoints contain control state and references, not a second copy of the raw corpus. Artifact retention, access control, and deletion remain domain-storage responsibilities.
- DBOS does not force the CVR hot path through a GCS read. Coarse CVR replay is accepted because acquisition is cheap; finer checkpoints are reserved for stages where repeating work has meaningful time, money, rate-limit, or anti-bot cost.
- Registry runs and individual companies are independently addressable by stable workflow identity. Campaign summaries aggregate child outcomes but do not become the only place those outcomes can be queried.
- Sharing the physical test Postgres deployment does not mean sharing database authority: each DBOS application's migration and runtime credentials are separate, and DBOS tables are not part of the hub's public schema or consumer contract.
- Cloud Run task retries recover interrupted processes; DBOS step retries handle explicitly transient operations; domain reruns represent new work. These mechanisms must not recursively amplify one another.
- Domain mutations use datasource transactions by default, giving an atomic Postgres mutation/checkpoint boundary. The coarse CVR hot path and external effects are explicit idempotent exceptions, not accidental mixtures of transaction semantics.
- Application versions, workflow IDs, and extractor versions answer different questions—code compatibility, logical-work identity, and domain provenance respectively—and are stored separately.
- Registry ingestion and web enrichment remain separate architectures. Shared DBOS infrastructure does not imply shared workflow topology, queue policy, checkpoint granularity, budgets, or domain outcomes.
- The two DBOS applications can evolve and drain independently; sharing this repository or physical Postgres deployment does not couple their runtime lifecycle.
- The prephase PoC is not a production runtime dependency. Ported behavior is owned, tested, migrated, secured, and delivered from this sole-writer repository.
- Conductor is not required to operate the spike: execution correlation, alerting, and bounded system-history retention are explicit GCP/application responsibilities. This operational evidence informs the later managed-vs- self-hosted decision.
- This ADR does not authorize migrating Regnskabsdata, retiring
batch_job.tf, replacing Prefect across the PoC, deploying web enrichment, or changingprod. A follow-up ADR is required before any of that happens, informed by the spike's fitness report. CONTEXT.mddistinguishes execution checkpoints from domain durability so later web-scraping pipelines can reuse the runtime without coupling their evidence model to it.- If the spike fails or is inconclusive, the status quo pattern continues unchanged for the affected application family and the demonstrated blocker is recorded explicitly. Inconclusive or merely suboptimal evidence extends the spike or creates optimization work; it is not by itself a rejection.
Adoption fitness criteria¶
The spike is expected to lead to adoption. The following are verification targets and architectural fitness checks, not a weighted comparison against unrelated orchestration products.
Registry compatibility tracer¶
- DBOS and non-DBOS paths produce the same rows, ledger outcome, counts, and watermark.
- Forced termination recovers the same workflow and domain run without a
duplicate
ingestion_runsrow; replay immediately around the database commit converges. - Reconciliation, staging/swap, and role-replacement guarantees remain intact.
- At least five controlled runs measure median wall time and peak memory. The targets are no more than 5% wall-time and 10% memory regression; a miss is an optimization issue unless it proves irreducible and operationally unacceptable.
Web-enrichment value tracer¶
- Selected PoC fixtures retain behaviorally equivalent evidence and outcomes.
- Forced termination after each paid boundary repeats no completed SERP, browser, or LLM operation; retry cannot overspend a company budget.
- One failed company neither replays nor fails unrelated companies.
- DataForSEO batching/shared polling, adapter rate limits, per-domain politeness, and legal gates remain effective.
- DBOS workflow state contains no raw pages, prompts/responses, screenshots, personal records, or other payloads that belong in domain artifact storage.
Database, security, and operations¶
- Eligible domain mutations and datasource checkpoints commit exactly once; coarse CVR and external-side-effect exceptions pass explicit idempotency tests.
- Separate application schemas, roles, secrets, and privileged migration paths are proven; runtime credentials cannot perform DDL.
- Connection-pool sizing leaves documented Supabase capacity headroom.
- Application-version drain and an explicit cross-version fork are tested.
- Logs and traces correlate Cloud Run execution, DBOS workflow, and domain run/company; injected failures exercise alerts.
- Retention cleanup deletes only eligible terminal histories.
- The Job-mode design requires no always-on service, worker pool, Conductor, or second database, and its projected storage/GCP cost is documented.
Open questions¶
- Managed DBOS/Conductor economics and an always-on runtime are beyond this Job-mode spike's scope.
- The web tracer selects its exact paid/slow checkpoint boundaries from measured retry cost. This does not reopen the accepted per-company workflow, artifact-reference, adapter-queue, or budget contracts.
References¶
- DBOS workflows and recovery guarantees
- DBOS steps and selective retry
- DBOS queues, concurrency, rate limits, partitions, and deduplication
- DBOS datasource transactions
- DBOS system database, schemas, and runtime privileges
- DBOS workflow code upgrades
- DBOS logging and OpenTelemetry tracing
- DBOS workflow retention
- Cloud Run Job retry and checkpoint guidance
- Cloud Run monitoring