Skip to content

feat(data-warehouse): add durable source cursors, move xmin onto them - #108390

Open
Gilbert09 wants to merge 6 commits into
masterfrom
tom/dwh-source-cursor
Open

Gilbert09 wants to merge 6 commits into
masterfrom
tom/dwh-source-cursor

Conversation

@Gilbert09

Copy link
Copy Markdown
Member

Problem

  • A warehouse source that tracks its own read position between syncs has nowhere generic to keep it. This blocks the Kafka source stacked on this PR, whose position is a set of partition offsets.
  • Postgres xmin worked around the gap with three SourceResponse fields, three sync_type_config keys, and a pipeline step that names each one. Every new source of this kind would add the same again.
  • On the v3 pipeline, xmin stores its ceiling as soon as extraction finishes, before the loader writes the rows. A load that fails after that point leaves the cursor past rows the table never got. This comes from reading the code and was not reproduced.

Changes

  • Sources can declare a typed cursor with the CursorSource[CursorT] mixin in sources/common/cursor.py. The source loads the stored cursor and stages the next one. The framework stores it.
  • The cursor lives under one key, sync_type_config["source_cursor"], as {kind, data}. The kind stops a stored cursor of another shape from loading.
  • On v3 the cursor rides the existing staged slot next to the incremental watermark. The loader promotes both with the final batch, so the cursor never moves ahead of the data.
  • v2 runs and v3 runs with no batches store the cursor directly, because no rows are outstanding.
  • A reset or a corrupt-Delta revive loads no cursor. This reuses the activity's existing use_stored_cursors decision, so sources do not re-implement it. A reset also deletes the key.
  • Postgres xmin now stages an XminCursor. Schemas that still hold the old xmin keys read them until their next sync writes source_cursor, so no data migration runs.
  • Removed: the three xmin SourceResponse fields, advance_xmin_state, and the xmin model accessors and writers. reset_wrapped_xmin_cursors reads and clears both storage forms.
  • The incremental field is unchanged. The pipeline derives it from a column maximum, which is a different contract from a position the source states.
  • Nothing is user-visible apart from the v3 xmin timing fix. The implementing-warehouse-sources skill gains a short section on the mixin.

How did you test this code?

  • test_cursor.py (new) catches a loader that accepts another kind's cursor, crashes on a field a newer deploy added, or ignores the legacy keys.
  • test_models.py catches a source cursor lost in staging: one staged beside the watermark promotes with it, and a displaced run's cursor is parked and still promotes. A reset clears both storage forms.
  • v3 test_pipeline.py catches the cursor being stored before the final-batch notification, or a zero-batch run never storing it.
  • test_import_data.py catches a cursor that survives a reset or a revive, and checks that a legacy-key cursor loads when neither applies.
  • test_reset_wrapped_xmin_cursors.py now runs against both storage forms.
  • Run locally: the two end-to-end xmin syncs against a real Postgres, and the Postgres, Supabase, Neon, PlanetScale, v2, v3, loader and import-activity suites.
  • mypy ran scoped to products/warehouse_sources and products/data_warehouse. A repo-wide run ran out of memory on the dev machine.
  • Not checked: a real v3 load failure after extraction.

👉 Stay up-to-date with PostHog coding conventions for a smoother review.

Release status

  • No feature flag controls this change
  • This change is behind a feature flag and is not available to users
  • This change makes a previously flagged feature available to everyone

Automatic notifications

  • Publish to changelog?

Docs update

None.

Copilot AI balanced review requested due to automatic review settings September 29, 2026 14:23
@Gilbert09 Gilbert09 self-assigned this Sep 29, 2026
@trunk-io

trunk-io Bot commented Sep 29, 2026

Copy link
Copy Markdown

Merging to master in this repository is managed by Trunk.

  • To merge this pull request, check the box to the left or comment /trunk merge below.

After your PR is submitted to the merge queue, this comment will be automatically updated with its status. If the PR fails, failure details will also be posted here

@Gilbert09 Gilbert09 added the stamphog Request AI approval (no full review) label Sep 29, 2026 — with Talyn App

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@github-actions

Copy link
Copy Markdown
Contributor

Hey @Gilbert09! 👋

It looks like your git author email on this PR isn't your @posthog.com address (owerstom@gmail.com). Since you're on the PostHog team, it's worth pointing your local git author email at your @posthog.com address. Why it matters:

  • Consistent work identity in git history — internal tooling that attributes commits to team members keys off your @posthog.com address.
  • Keeps team contributions easy to tell apart from external community ones when scanning history.

You can fix it for this repo with:

git config user.email "you@posthog.com"

Or set it globally with git config --global user.email "you@posthog.com". No need to redo this PR — just a nudge for next time. 🙂

@github-actions

github-actions Bot commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

🤖 CI report

⚠️ Trunk lane — backend Python lane

This PR is assigned to the backend Python lane. It runs backend Python tests and may merge in parallel with PRs in other lanes.

⚠️ Duplication (Python) — 3 new duplicated blocks (worst 185 tokens)

New Python code duplication introduced by this branch. Fails at 70+ tokens in app code, or 150+ tokens when both copies live in test files. Advisory while the gate proves itself: extract a shared helper instead of copying.

First copy Second copy Lines Tokens
products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_import_data.py:139 products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_import_data.py:263 49 185
products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_import_data.py:150 products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_import_data.py:357 37 152
products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v2/pipeline.py:100 products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/pipeline.py:138 32 140
✅ Duplication (TypeScript) — clean

New TypeScript code duplication introduced by this branch. Fails at 70+ tokens in app code, or 150+ tokens when both copies live in test files. Advisory while the gate proves itself: extract a shared helper instead of copying.

⚠️ Comment density — 5% of added code lines are comments (23 of 492)

This section warns when comments are more than 3% of the code lines a PR adds, and alerts above 6%. Before agent-assisted PRs, the typical share was about 2%. Only full-line comments count. Docstrings, generated files, snapshots, migrations, and workflow files are left out.

Comments that restate the code, record how the change came about, or narrate the next line add noise for the next reader. Keep the comments that explain a reason the code cannot show, and remove the rest. See .agents/skills/writing-code-comments/SKILL.md for the house rules.

Files with the most added comment lines:

File Comment lines Added lines
products/warehouse_sources/backend/models/external_data_schema.py 5 30
products/warehouse_sources/backend/temporal/data_imports/sources/common/cursor.py 5 106
products/warehouse_sources/backend/temporal/data_imports/sources/postgres/xmin_cursor.py 4 20
products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py 3 16
products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/pipeline.py 2 21
products/warehouse_sources/backend/temporal/data_imports/sources/common/typings.py 2 9
products/warehouse_sources/backend/temporal/data_imports/sources/postgres/postgres.py 1 19
products/warehouse_sources/backend/temporal/data_imports/workflow_activities/import_data_sync.py 1 13

This check does not block merging. It updates on every push and clears when the share drops.

⚠️ Backend coverage — 97.0% of changed backend lines covered — 6 uncovered

🧪 Backend test coverage

Patch coverage — changed backend lines (products + core): ███████████████████░ 97.0% (273 / 279)

File Patch Uncovered changed lines
products/warehouse_sources/backend/temporal/data_imports/sources/common/typings.py 66.7% 37, 40
products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v2/pipeline.py 75.0% 94
products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/pipeline.py 80.0% 128
products/warehouse_sources/backend/temporal/data_imports/sources/common/cursor.py 97.0% 79, 99

🤖 Agents: add a test covering the lines above, or note why under "How did you test this code?". Machine-readable gap list: the patch-coverage artifact on this run (gh run download 36591553337 -n patch-coverage), or the coverage-data block at the end of this comment.

Per-product line coverage (touched products)
Product Coverage Lines
demo ███████████░░░░░░░░░ 52.8% 1,411 / 2,673
batch_exports ████████████████░░░░ 81.2% 21,528 / 26,502
cdp ██████████████████░░ 88.2% 4,545 / 5,155
mcp_analytics ██████████████████░░ 89.2% 5,038 / 5,651
product_tours ██████████████████░░ 89.3% 1,331 / 1,491
dashboards ██████████████████░░ 89.5% 6,904 / 7,714
data_warehouse ██████████████████░░ 90.0% 14,027 / 15,589
notebooks ██████████████████░░ 90.2% 15,287 / 16,945
signals ██████████████████░░ 90.2% 57,816 / 64,069
cohorts ██████████████████░░ 90.4% 8,420 / 9,316
streamlit_apps ██████████████████░░ 90.6% 2,623 / 2,895
managed_warehouse ██████████████████░░ 91.0% 10,252 / 11,263
tasks ██████████████████░░ 91.1% 75,146 / 82,503
data_modeling ██████████████████░░ 91.4% 10,554 / 11,543
exports ██████████████████░░ 91.6% 9,680 / 10,562
engineering_analytics ██████████████████░░ 91.7% 11,032 / 12,030
business_knowledge ██████████████████░░ 92.2% 7,684 / 8,330
conversations ███████████████████░ 92.5% 28,734 / 31,062
early_access_features ███████████████████░ 92.6% 1,332 / 1,439
stamphog ███████████████████░ 92.8% 8,109 / 8,742
canvas ███████████████████░ 92.8% 6,873 / 7,405
approvals ███████████████████░ 93.0% 3,974 / 4,271
mcp_registry ███████████████████░ 93.1% 1,670 / 1,794
notifications ███████████████████░ 93.2% 1,145 / 1,229
slack_app ███████████████████░ 93.2% 13,733 / 14,735
error_tracking ███████████████████░ 93.2% 16,359 / 17,547
surveys ███████████████████░ 93.3% 6,571 / 7,040
autoresearch ███████████████████░ 93.6% 8,481 / 9,061
context_layer ███████████████████░ 93.9% 3,415 / 3,638
web_analytics ███████████████████░ 94.0% 21,815 / 23,218
alerts ███████████████████░ 94.0% 8,570 / 9,114
billing_alerts ███████████████████░ 94.1% 2,094 / 2,226
mcp_store ███████████████████░ 94.4% 8,940 / 9,472
ai_observability ███████████████████░ 94.5% 24,473 / 25,896
workflows ███████████████████░ 94.6% 15,222 / 16,093
wizard ███████████████████░ 94.7% 6,150 / 6,496
reminders ███████████████████░ 94.8% 760 / 802
review_hog ███████████████████░ 95.0% 11,507 / 12,119
endpoints ███████████████████░ 95.1% 9,206 / 9,681
annotations ███████████████████░ 95.1% 817 / 859
customer_analytics ███████████████████░ 95.2% 25,081 / 26,352
marketing_analytics ███████████████████░ 95.3% 19,566 / 20,528
posthog_ai ███████████████████░ 95.3% 2,488 / 2,610
experiments ███████████████████░ 95.4% 32,457 / 34,020
actions ███████████████████░ 95.5% 756 / 792
logs ███████████████████░ 95.5% 15,399 / 16,130
data_catalog ███████████████████░ 95.5% 4,401 / 4,606
tracing ███████████████████░ 95.6% 3,536 / 3,699
replay_vision ███████████████████░ 95.6% 27,675 / 28,939
growth ███████████████████░ 95.7% 11,228 / 11,734
messaging ███████████████████░ 95.8% 3,798 / 3,963
skills ███████████████████░ 95.8% 6,972 / 7,274
product_analytics ███████████████████░ 96.0% 28,470 / 29,647
access_control ███████████████████░ 96.3% 7,113 / 7,386
revenue_analytics ███████████████████░ 96.4% 1,876 / 1,946
user_interviews ███████████████████░ 96.5% 2,867 / 2,971
feature_flags ███████████████████░ 96.5% 25,588 / 26,509
warehouse_sources ███████████████████░ 97.2% 459,764 / 472,790
data_quality ████████████████████ 97.6% 7,587 / 7,774
metrics ████████████████████ 98.0% 4,084 / 4,166
analytics_platform ████████████████████ 98.3% 2,784 / 2,833
pulse ████████████████████ 98.5% 2,046 / 2,078
live_debugger ████████████████████ 99.2% 626 / 631

Report-only. Patch coverage = changed backend lines covered vs origin/master. Sorted lowest first.
Known gaps: lines covered only by Temporal tests show as uncovered; core line numbers may drift if master changed the same file.

@Gilbert09
Gilbert09 added this pull request to stack #108392 September 29, 2026 14:24
@pr-assigner-resolver-posthog
pr-assigner-resolver-posthog Bot requested a review from a team September 29, 2026 14:24
stamphog[bot]

This comment was marked as outdated.

@stamphog
stamphog Bot dismissed their stale review September 29, 2026 14:32

A new stamphog review started for this PR — the fresh verdict replaces this approval.

stamphog[bot]

This comment was marked as outdated.

@trunk-io

trunk-io Bot commented Sep 29, 2026 •

Copy link
Copy Markdown

Static Badge   Static Badge   Static Badge

View Full Report ↗︎ ⋅ Docs

@stamphog
stamphog Bot dismissed their stale review September 29, 2026 14:46

A new stamphog review started for this PR — the fresh verdict replaces this approval.

@stamphog stamphog Bot removed the stamphog Request AI approval (no full review) label Sep 29, 2026

@stamphog stamphog Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not approved yet — waiting on the conditions below.

Re-add the stamphog label to request another review once you have addressed this.

Stamphog refused this pull request at the prerequisites gate, because the branch has merge conflicts. The other gates (deny-list, size, and tier) passed, so the conflicts are the only reason for the refusal.

To move forward, rebase or merge the base branch into this one, resolve the conflicts, and push the result. If you'd like a human to look at it in the meantime, ask a reviewer directly.

Gate mechanics and policy version
Gate Result
prerequisites ✗ merge conflicts present
deny-list ✓ no deny categories matched
size ✓ 425L, 12F substantive, 1 binary; 855L/22F incl. docs/generated/snapshots — within ceiling
tier ✓ T1-agent / T1d-complex (855L, 22F, two-areas, feat)
stamphog 2.3.0 .stamphog/policy.yml @ unknown · reviewed head fdf37cc

Rebase the source cursor implementation onto current master, preserve keyset resume behavior, and package the common source tests to avoid pytest module collisions.
@Gilbert09
Gilbert09 force-pushed the tom/dwh-source-cursor branch from fdf37cc to 0cd5d73 Compare September 29, 2026 14:48
@Gilbert09
Gilbert09 requested a review from a team September 29, 2026 14:50
test_source_cursor.py already holds the same tests, and common/test/__init__.py keeps its module name clear of the Cursor source's test_cursor.py.

@fuziontech fuziontech left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Automated review generated on behalf of @fuziontech.

APPROVE. I did not find a blocking correctness, security, data-loss, or outage issue in the diff.

Non-blocking follow-ups:

  • test_cursor.py and test_source_cursor.py are byte-for-byte duplicates; keep one to avoid maintaining the same suite twice.
  • The generic staged-cursor promotion replaces the current stored cursor without a second source-specific merge. If overlapping runs can promote out of order, consider a monotonic merge at promotion time to prevent an older cursor from moving the watermark backward. For xmin this appears recoverable as rereading, so it is not a release gate.

Focused tests could not run in this environment because hogli and pytest are unavailable.

@coderabbitai

coderabbitai Bot commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Review in Change Stack →

Navigate logical layers of code changes, visualize relationships, and explore their blast radius.

📝 Walkthrough

Walkthrough

The change introduces a shared source-cursor manager and adds cursor loading, staging, and persistence to warehouse imports. PostgreSQL XMIN state moves from dedicated schema fields to a serialized XminCursor, with legacy-key loading and clearing retained. Pipeline finalization commits staged cursors, and the wrapped-XMIN reset command reads and clears both cursor formats.

Priority: ➖ Normal

🚥 Pre-merge checks | ✅ 1
✅ Passed checks (1 passed)
Check name Status Explanation
Description check ✅ Passed The description is complete and standalone. It explains the problem, user-visible changes, testing performed, known test gap, release status, changelog status, and documentation status. No critical se…
✨ Finishing Touches
📝 Generate docstrings
  • Commit to this branch
  • Create a new PR
🧪 Generate unit tests (beta)
  • Commit to this branch
  • Create a new PR

Comment @coderabbitai help to get the list of available commands.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1


ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository: PostHog/posthog/.coderabbit.yaml

Review profile: QUIET

Plan: Enterprise

Run ID: 489d1e3c-19fe-404e-96a4-28f4d85786aa

📥 Commits

Reviewing files that changed from the base of the PR and between 2286c8a and 14a2fb8.

📒 Files selected for processing (21)
  • .agents/skills/implementing-warehouse-sources/SKILL.md
  • products/warehouse_sources/backend/management/commands/reset_wrapped_xmin_cursors.py
  • products/warehouse_sources/backend/models/external_data_schema.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/common/extract.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v2/pipeline.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/lanes.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/pipeline.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/test_pipeline.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/cursor.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/test/__init__.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/test/test_source_cursor.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/typings.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/postgres.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/test_postgres.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/xmin_cursor.py
  • products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_end_to_end.py
  • products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_import_data.py
  • products/warehouse_sources/backend/temporal/data_imports/workflow_activities/import_data_sync.py
  • products/warehouse_sources/backend/tests/management/test_reset_wrapped_xmin_cursors.py
  • products/warehouse_sources/backend/tests/test_models.py
💤 Files with no reviewable changes (1)
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/test_postgres.py

Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 5 remain after this review.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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/sources/postgres/xmin_cursor.py-24-28 (1)

24-28: 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win

Preserve partial legacy xmin state.

The field mapping is correct, but xmin_cursor_from_legacy rejects legacy state that the previous path accepted. The previous _capture_xmin_ceiling treated a missing xmin_num_wraparound as epoch 0. The current parser returns no cursor, and PostgresSource then performs a full scan.

Keep xmin_last_value as the required field. Default the epoch to 0 and derive a missing ceiling_xid8.

Suggested fix
 def xmin_cursor_from_legacy(sync_type_config: Mapping[str, Any]) -> XminCursor | None:
     ceiling_xid, ceiling_xid8, num_wraparound = (sync_type_config.get(key) for key in XMIN_LEGACY_KEYS)
-    if not isinstance(ceiling_xid, int) or not isinstance(ceiling_xid8, int) or not isinstance(num_wraparound, int):
+    if not isinstance(ceiling_xid, int):
         return None
+    if not isinstance(num_wraparound, int):
+        num_wraparound = 0
+    if not isinstance(ceiling_xid8, int):
+        ceiling_xid8 = (num_wraparound << 32) | ceiling_xid
     return XminCursor(ceiling_xid=ceiling_xid, ceiling_xid8=ceiling_xid8, num_wraparound=num_wraparound)

ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository: PostHog/posthog/.coderabbit.yaml

Review profile: QUIET

Plan: Enterprise

Run ID: 57fba4a1-c9ae-4312-b3a3-2d06824c745f

📥 Commits

Reviewing files that changed from the base of the PR and between 14a2fb8 and 0e6398a.

📒 Files selected for processing (21)
  • .agents/skills/implementing-warehouse-sources/SKILL.md
  • products/warehouse_sources/backend/management/commands/reset_wrapped_xmin_cursors.py
  • products/warehouse_sources/backend/models/external_data_schema.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/common/extract.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v2/pipeline.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/lanes.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/pipeline.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/test_pipeline.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/cursor.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/test/__init__.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/test/test_source_cursor.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/typings.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/postgres.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/test_postgres.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/xmin_cursor.py
  • products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_end_to_end.py
  • products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_import_data.py
  • products/warehouse_sources/backend/temporal/data_imports/workflow_activities/import_data_sync.py
  • products/warehouse_sources/backend/tests/management/test_reset_wrapped_xmin_cursors.py
  • products/warehouse_sources/backend/tests/test_models.py
💤 Files with no reviewable changes (2)
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/test/init.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/test_postgres.py

Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 5 remain after this review.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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/tests/management/test_reset_wrapped_xmin_cursors.py-90-90 (1)

90-90: 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win

Assert that reset removes the cursor keys.

_ceiling_xid(wrapped) is None also passes if source_cursor remains stored with a null ceiling. The reset contract requires the key to be deleted. Assert that source_cursor and the legacy XMIN keys are absent after the live run; apply the same check to the explicitly selected schema at Line 110.

🧹 Nitpick comments (1)
products/warehouse_sources/backend/temporal/data_imports/sources/common/cursor.py (1)

120-140: 🩺 Stability & Availability | 🔵 Trivial | 💤 Low value

Discard malformed payloads when the cursor constructor raises something other than TypeError.

_cursor_from_payload catches only TypeError. If a cursor class validates its fields in __post_init__ and raises ValueError, every run fails for that schema. The legacy fallback does not help, because the payload key is present. Catch (TypeError, ValueError) so that a bad stored cursor is discarded and the run continues.

Proposed fix
-    except TypeError:
+    except (TypeError, ValueError):

ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository: PostHog/posthog/.coderabbit.yaml

Review profile: QUIET

Plan: Enterprise

Run ID: 75a799b3-93d8-460d-83c4-37ee0b2a5230

📥 Commits

Reviewing files that changed from the base of the PR and between 0e6398a and 6ef9bfe.

📒 Files selected for processing (21)
  • .agents/skills/implementing-warehouse-sources/SKILL.md
  • products/warehouse_sources/backend/management/commands/reset_wrapped_xmin_cursors.py
  • products/warehouse_sources/backend/models/external_data_schema.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/common/extract.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v2/pipeline.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/lanes.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/pipeline.py
  • products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/test_pipeline.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/cursor.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/test/__init__.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/test/test_source_cursor.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/typings.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/postgres.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/test_postgres.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/xmin_cursor.py
  • products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_end_to_end.py
  • products/warehouse_sources/backend/temporal/data_imports/tests/e2e/test_import_data.py
  • products/warehouse_sources/backend/temporal/data_imports/workflow_activities/import_data_sync.py
  • products/warehouse_sources/backend/tests/management/test_reset_wrapped_xmin_cursors.py
  • products/warehouse_sources/backend/tests/test_models.py
💤 Files with no reviewable changes (2)
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/test/init.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/postgres/test_postgres.py

Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 5 remain after this review.

Suppress the UUIDv7 rule for the existing ExternalDataSchema UUIDT primary key; changing persisted primary-key generation is outside this cursor change.

This branch has not been deployed

No deployments
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.

3 participants