chore(warehouse-sources): remove the legacy cdc conversion - #110777
Conversation
Capture converted what the retired legacy lane left behind before every read. Every source that captures has converted, and the sources still on the legacy mode are marked broken, so Repair CDC is their only way back and it resets every table. - Delete cdc/legacy_conversion.py and capture's call to it. - Stop reading cdc_ingest_mode: the consumer's wait for conversion, the snapshot gate and the CDC config field go. - Drop the cdc_legacy_converted_at exemption from the buffer expiry check and the two cdc_deferred_runs guards that only the conversion's reset served. - Keep writing cdc_ingest_mode on setup and repair, and keep it across API writes, so a rollback to the previous release does not read a source without it as legacy and empty its buffer. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
😎 Merged successfully - details. |
🤖 CI report
|
| File | Comment lines | Added lines |
|---|---|---|
products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/adapter.py |
3 | 3 |
products/warehouse_sources/backend/ad_hoc_sync.py |
2 | 3 |
products/warehouse_sources/backend/presentation/views/external_data_source/helpers.py |
2 | 2 |
products/warehouse_sources/backend/temporal/data_imports/cdc/activities.py |
2 | 16 |
This check does not block merging. It updates on every push and clears when the share drops.
|
[Medium risk] Removes legacy CDC data ingestion pathway and conversion logic. Satisfy the repository’s code-comment requirement before merging; no blocking CDC behavior failure was established. Reviews (1) · Last reviewed commit: "chore(warehouse-sources): remove the leg..." |
|
Navigate logical layers of code changes, visualize relationships, and explore their blast radius. 🧰 Additional context used📚 Code guidelines (8)📝 WalkthroughWalkthroughThe CDC runtime no longer converts legacy state or gates buffered processing on Priority: ➖ Normal Merge Risk: 🟡 Moderate · up to Resolve the marker race and prevent unrepaired legacy tables from consuming buffered rows before merging. Otherwise, a broken reason can be lost or legacy rows delivered twice. Security Architecture ReviewSecurity architecture risk: 🟡 Moderate · up to The retirement changes how stopped sources recover and how old change data is handled. A failed pause or an already-running sync can bypass the intended full repair, potentially undermining data integrity. Existing ownership controls remain, but the production rollout preconditions were not independently verified. Retained concerns
Security review detailsSecurity Blast Radius
Trust Boundaries and Controls
Resilience and Maintainability Implications
Hardening Proposals
🚥 Pre-merge checks | ✅ 1✅ Passed checks (1 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Comment |
|
Risk: No findings This delta tightens the self-managed lag-sweep guard ( Sentinel reviewed |
…g is critical The lag sweep rewrote every table's broken marker to the self-managed lag reason. That marker allows Resume CDC and clears itself once the lag drops, so it could lift a stop that only Repair CDC may lift, including the one that now keeps legacy sources off Resume CDC. The sweep leaves a source alone when a table already holds a marker with another reason. Also drops test cases and fixtures that the conversion removal left without a purpose, and a gate in the ad hoc sync that is always true for a streaming table. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
There was a problem hiding this comment.
Note
Quiet mode is enabled, so only the most important comments were posted inline. Other review comments are grouped below.
🟡 Other comments (1)
products/warehouse_sources/backend/temporal/data_imports/cdc/activities.py-1699-1699 (1)
1699-1699: 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick winFilter inactive schemas from
broken_for_another_reason.
broken_for_another_reasonscans disabled and deleted schemas, butmark_cdc_brokenupdates only enabled, non-deleted CDC schemas. Disabling a schema does not clear itscdc_brokenmarker. Therefore, a marker on an inactive schema can suppress marking an active schema at critical lag.🐛 Suggested fix
- ExternalDataSchema.objects.filter(team_id=source.team_id, source=source, sync_type_config__has_key="cdc_broken") + ExternalDataSchema.objects.filter( + team_id=source.team_id, + source=source, + sync_type=ExternalDataSchema.SyncType.CDC, + should_sync=True, + sync_type_config__has_key="cdc_broken", + ) + .exclude(deleted=True) .exclude(sync_type_config__cdc_broken__reason=reason) .exists()
ℹ️ Review info
⚙️ Run configuration
Configuration used: Repository: PostHog/posthog/.coderabbit.yaml
Review profile: QUIET
Plan: Enterprise
Run ID: cab1e941-bfd1-4205-a4aa-d1d78eb45a5d
📒 Files selected for processing (10)
docs/internal/cdc-buffered-ingress-runbook.mdproducts/warehouse_sources/backend/ad_hoc_sync.pyproducts/warehouse_sources/backend/temporal/data_imports/cdc/activities.pyproducts/warehouse_sources/backend/temporal/data_imports/cdc/broken.pyproducts/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_cleanup_orphan_slots.pyproducts/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_extract_activity.pyproducts/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_source_manager.pyproducts/warehouse_sources/backend/temporal/data_imports/pipelines/test_pipeline_sync.pyproducts/warehouse_sources/backend/tests/api/test_cdc_table_mode_patch.pyproducts/warehouse_sources/backend/tests/test_stalled_schedules.py
💤 Files with no reviewable changes (2)
- products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_source_manager.py
- products/warehouse_sources/backend/temporal/data_imports/pipelines/test_pipeline_sync.py
Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 8 remain after this review.
…lback constraint The comments on the three places that keep the key described how an earlier release behaved. They now state the standing reason: a worker on a release that reads the key treats a source without it as legacy. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A new stamphog review started for this PR — the fresh verdict replaces this approval.
…sweep check mark_cdc_broken marks only CDC tables that sync, and a table whose sync is off keeps the marker it had. Counting that marker let it suppress the self-managed lag marker on the tables that still sync. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to GitHub limitations.
🟠 Major · Restore the legacy conversion before buffered consumption. · source.py:1894-1900
products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py:1894-1900
🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy liftRestore the legacy conversion before buffered consumption.
source_for_pipelineroutes every seeded streaming CDC schema here, regardless ofcdc_ingest_mode. The PR also removes the conversion that purged the legacy buffer and marked the source as buffered. A persisted legacy source can therefore read buffer rows that the legacy lane already delivered, which can load CDC changes twice.Restore the conversion gate and perform the purge, schedule handoff, and buffered-state update before this branch can read a legacy source.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Repository: PostHog/posthog/.coderabbit.yaml
Review profile: QUIET
Plan: Enterprise
Run ID: d3348da7-81ba-492a-ba36-7696b83d3b54
📒 Files selected for processing (2)
products/warehouse_sources/backend/presentation/views/external_data_source/helpers.pyproducts/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/adapter.py
Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 10 remain after this review.
A new stamphog review started for this PR — the fresh verdict replaces this approval.
There was a problem hiding this comment.
Approved.
This is a removal of dead legacy CDC conversion code by an author on the owning team with STRONG familiarity. No dangling references to the removed symbols remain, and the earlier CodeRabbit inline concern about inactive schemas is already handled in the current helper. The one real risk is a double-load of a still-legacy source if the out-of-PR "mark broken" precondition was not actually done, which I could not verify from the diff.
- Author wrote 98% of the modified lines and has 30 merged PRs in these paths (familiarity STRONG).
- 👍 on the PR from greptile-apps[bot].
- Safety depends on the one-off operation, done outside this PR, that marked every still-legacy source as broken. The diff cannot confirm it.
- The CodeRabbit comment on an older commit says a persisted legacy source could re-read buffer rows the legacy lane already delivered. The author's argument for why this cannot happen rests on production checks, not on code in the diff.
Gate mechanics and policy version
| Gate | Result | |
|---|---|---|
| prerequisites | ✓ | all clear |
| deny-list | ✓ | no deny categories matched |
| size | ✓ | 361L, 11F substantive, 774L/24F incl. docs/generated/snapshots — within ceiling |
| tier | ✓ | T1-agent / T1d-complex (774L, 24F, two-areas, chore) |
| stamphog 2.3.1 | .stamphog/policy.yml @ 9f18b5b · reviewed head 9f18b5b |
|
/trunk merge |
There was a problem hiding this comment.
Actionable comments posted: 1
Caution
Some comments are outside the diff and can’t be posted inline due to GitHub limitations.
🟠 Major · Skip buffered consumption for legacy_lane_retired schemas. · source.py:1893-1899
products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py:1893-1899
🗄️ Data Integrity & Integration | 🟠 Major | ⚡ Quick winSkip buffered consumption for
legacy_lane_retiredschemas.A persisted
cdc_brokenmarker can remain while a table schedule or run is active._buffered_cdc_source()can then consume rows already delivered by the legacy lane and merge them again. Check the marker before buffered consumption. This preserves mode-less buffered routing because it does not inspect or parsejob_inputs.repair_cdcclears the marker and restores operation.Suggested fix
+ cdc_broken = (schema.sync_type_config or {}).get("cdc_broken") or {} + if cdc_broken.get("reason") == "legacy_lane_retired": + return no_op_tick() + # Defense in depth for the v3-forcing invariant: a run that resolved its pipeline version
ℹ️ Review info
⚙️ Run configuration
Configuration used: Repository: PostHog/posthog/.coderabbit.yaml
Review profile: QUIET
Plan: Enterprise
Run ID: 86a01e3d-2416-4a32-983a-042aa43083ae
📒 Files selected for processing (3)
products/warehouse_sources/backend/temporal/data_imports/cdc/broken.pyproducts/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_cleanup_orphan_slots.pyproducts/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py
💤 Files with no reviewable changes (1)
- products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py
Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 10 remain after this review.
Merge-queue failure: both failures are unrelated flakesThe queue removed this PR after
Evidence for the first one: the identical failure ( The second is a float-equality assertion against a wall clock: No code change on this branch: the root causes sit in State of the PR right now:
What's needed from a human: re-enter the merge queue (check the Trunk box above or comment 🦉 via talyn.dev |
|
Correction to the comment above: one evidence link was wrong (comment editing is not available to me, hence this follow-up). The two PRs that hit the identical
Everything else in that comment stands. 🦉 via talyn.dev |
|
/trunk merge |
Problem
cdc/legacy_conversion.pyto convert each source's leftover state on its next capture run.Changes
Note
Precondition, done on 2026-10-02 in both regions: the legacy sources that were stopped without a broken marker are now marked broken, by a one-off operation outside this PR. A read-only check afterwards showed a broken marker on every source that still reads as legacy. A legacy source left unmarked would, after Resume CDC, capture into a buffer that its paused table schedules never read.
cdc/legacy_conversion.pyand the call capture made to it before every read.cdc_ingest_modeanymore. The table sync no longer waits for a conversion, the re-snapshot gate no longer checks the mode, andCDCConfig.ingest_modeis gone.cdc_deferred_runsguards go. They held a table back for the reset that the conversion staged, so without the conversion they would keep such a table out of the buffer for good.cdc_ingest_mode: buffered, and the API keeps the key on a PATCH. A rollback to the previous release would read a source without the key as legacy and empty its unconsumed buffer. A later PR can drop the writes._close_stranded_capture_jobsgoes with the module. It closed job rows that a legacy capture run left Running, up to 14 days old. The last legacy capture run was on 2026-09-30, and the sources that captured since have run it.CDCJobInputsUnreadableErrorfor unreadablejob_inputs. Capture still does, through the CDC config.How did you test this code?
test_patch_cdc_table_mode_adding_target_triggers_resnapshotkeeps its stalecdc_deferred_runscase. A reset still removes that key, and the table now stays in the buffer.test_critical_lag_self_managed_marks_broken_without_drop_or_pausegains two cases. A table that already holds another marker keeps it through a critical-lag sweep, and a marker on a table whose sync is off does not stop the sweep from marking the tables that sync. Each fails without its part of the guard.👉 Stay up-to-date with PostHog coding conventions for a smoother review.
Release status
Automatic notifications
Docs update
Rewrote "Leftover legacy state" in
docs/internal/cdc-buffered-ingress-runbook.md: legacy sources come back through Repair CDC only, and the two leftover keys do nothing.🤖 Agent context
Autonomy: Human-driven (agent-assisted)
Agent: Claude Code, Claude Opus 5.5
legacy cdc conversion,cdc_ingest_mode): no open PR./writing-tests,/writing-code-comments,/reviewing-with-coderabbit,/writing-pr-descriptions.cdc_ingest_modewrites: theclear_deferred_runsflag now only decides whether a dead key is removed.cdc_ingest_modedescribed an earlier release. They now state the rollback constraint.mark_cdc_brokennever marks. It now looks at the same tables.🤖 Generated with Claude Code