diff --git a/docs/internal/cdc-buffered-ingress-runbook.md b/docs/internal/cdc-buffered-ingress-runbook.md index e99de330dde1..bae3f379ebad 100644 --- a/docs/internal/cdc-buffered-ingress-runbook.md +++ b/docs/internal/cdc-buffered-ingress-runbook.md @@ -134,10 +134,9 @@ columns reports no position on that run, since none of its files carries the sta residual buffer is merged once more and billed once; the next write lands with the statistic and the position reads normally from then on. -**A run stands down while any delivery for the schema is still in the queue**: a previous attempt -of this same job, or a batch the retired legacy lane queued before its source was converted. Either -would write alongside whatever this run reads, and on the append lane that is a second copy of the -same history. Two scheduled runs cannot overlap on their +**A run stands down while any delivery for the schema is still in the queue**: a batch that a +previous attempt of this same job left there. It would write alongside whatever this run reads, and +on the append lane that is a second copy of the same history. Two scheduled runs cannot overlap on their own: the v3 pipeline lock is held from the start of the workflow until the loader completes the job. The window is a retried activity, which runs under the lock its own workflow already holds, and a lock takeover, which hands the lock to a new job while the old one's batches are still @@ -147,29 +146,18 @@ consumed. The next scheduled run picks the buffer up once the queue has drained. ## Leftover legacy state -Capture used to deliver some tables' changes itself, on what is now the retired legacy lane. -That lane paused each such table's schedule while it streamed, held a snapshotting table's changes as deferred runs in `sync_type_config["cdc_deferred_runs"]`, marked its source `cdc_ingest_mode: legacy`, and created job rows of its own. -Capture converts that state before every read (`cdc/legacy_conversion.py`), so nothing needs doing by hand: - -- **A source still marked legacy** has every CDC table's buffer emptied, since it holds only copies of changes the legacy lane already delivered. - Each table is stamped `cdc_legacy_converted_at`, because it is current as of the conversion: without the stamp, its first sync would find no run that ever listed the buffer, take the buffer for expired, and re-snapshot the table (see "Buffer expiry — no partial recovery"). - Each syncing table's schedule is rebuilt unpaused, because it is now the table's consumer. - A table set slower than its source's fastest table is sped up to it (`cdc_legacy_table_frequency_raised`), because that is how often the legacy lane delivered its changes. - The source is then marked `cdc_ingest_mode: buffered`, last, so a failure repeats the whole conversion on the next run. - Legacy batches still in the load queue land first, because the consumer stands down while any are in flight. - Until then, a scheduled run of one of its tables, such as one a sync frequency change unpaused, no-ops the tick (`cdc_buffered_waiting_for_legacy_conversion`), because reading those copies would load them a second time. -- **A table with deferred runs** snapshots again in the buffer, through the same pending reset a resync hands to capture (`cdc_legacy_deferred_runs_handed_to_reset`). - Nothing merges deferred runs anymore, so the reset pauses the table's schedule and cancels its running sync, and waits while either can still hand over (`cdc_reset_waits_for_running_sync`). - Then it resets the table, drops the deferred runs, empties its buffer, and starts the new snapshot. - Capture gives the table no buffered snapshot of its own before that, because the old sync could hand over into it without the deferred changes. -- **A job row a legacy capture run left Running** is failed once it is 30 minutes old and has no batches in the queue. - -A rebuilt schedule is skipped where the pause is deliberate: the schema's status is `Paused`, an admin-triggered run holds it, the schema is halted and waits for Repair CDC, a pending reset holds it, or it has no sync frequency. -Every step logs (`cdc_legacy_source_converted`, `cdc_legacy_deferred_runs_handed_to_reset`, `cdc_stranded_capture_jobs_closed`); once none of them appears across the fleet, the module can go. - -A source whose slot is gone does not capture at all, so it converts only after Repair CDC, which already resets every table and marks the source buffered. - -`cdc_ingest_mode` must survive every API write. A PATCH that dropped it would make capture read the source as legacy and empty its unconsumed buffer. +Capture used to deliver some tables' changes itself, on the retired legacy lane. +The code that converted what that lane left behind is gone. +Every source that still read as legacy when it was removed is marked `cdc_broken`, so Repair CDC is its only way back. +Repair CDC resets every table, recreates the slot and resumes the table schedules, so no legacy state survives it. + +Two leftovers remain in the data, and nothing reads them: + +- `cdc_ingest_mode` in a source's `job_inputs`. + CDC setup and Repair CDC still write `buffered`, and the API still keeps the key on a PATCH. + A rollback to the release that converted legacy sources would read a source without it as legacy and empty its unconsumed buffer. +- `cdc_deferred_runs` in a table's `sync_type_config`. + A resync, a table-mode change and Repair CDC remove it. ## When a schedule stops firing diff --git a/products/warehouse_sources/backend/ad_hoc_sync.py b/products/warehouse_sources/backend/ad_hoc_sync.py index 18587dfed960..e1dedaeb70e7 100644 --- a/products/warehouse_sources/backend/ad_hoc_sync.py +++ b/products/warehouse_sources/backend/ad_hoc_sync.py @@ -31,7 +31,6 @@ from products.warehouse_sources.backend.temporal.data_imports.cdc.snapshot_lane import ( BUFFER_LANE, cancel_sync_that_could_hand_over, - resnapshot_stays_in_buffer, ) @@ -126,10 +125,9 @@ def trigger_ad_hoc_sync( # streaming, so ongoing CDC stays billable. The save must precede the workflow start so the # source reloads cdc_mode="snapshot" instead of racing on stale "streaming". if schema.is_cdc and schema.cdc_mode == "streaming": - # Decided while the table still streams. Without the marker, the next capture run would - # empty the buffer under the new snapshot, deleting changes an in-flight run wrote. - if resnapshot_stays_in_buffer(schema): - updates[CDC_SNAPSHOT_LANE_KEY] = BUFFER_LANE + # Without the marker, the next capture run would empty the buffer under the new snapshot, + # deleting changes an in-flight run wrote. + updates[CDC_SNAPSHOT_LANE_KEY] = BUFFER_LANE updates["cdc_mode"] = "snapshot" removes += ["cdc_last_log_position", "cdc_deferred_runs"] extra_model_fields["initial_sync_complete"] = False diff --git a/products/warehouse_sources/backend/presentation/views/external_data_source/helpers.py b/products/warehouse_sources/backend/presentation/views/external_data_source/helpers.py index cda9f9ebd9f1..bd810b0f3653 100644 --- a/products/warehouse_sources/backend/presentation/views/external_data_source/helpers.py +++ b/products/warehouse_sources/backend/presentation/views/external_data_source/helpers.py @@ -360,8 +360,8 @@ def get_nonsensitive_and_sensitive_field_names(fields: list[FieldType]) -> Field "cdc_lag_warning_threshold_mb", "cdc_lag_critical_threshold_mb", "cdc_consistent_point", - # Set by CDC setup, Repair CDC and capture, never by the API. Losing it on an unrelated PATCH - # would make capture convert the source again, which empties its unconsumed buffer. + # Set by CDC setup and Repair CDC, never by the API. A worker on a release that reads it treats a + # source without it as legacy and empties its buffer, so a PATCH must keep it for a rollback. "cdc_ingest_mode", } diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/activities.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/activities.py index 4185272dc664..3e2f8486592e 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/activities.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/activities.py @@ -56,6 +56,7 @@ from products.warehouse_sources.backend.temporal.data_imports.cdc.broken import ( AUTO_DROPPED_LAG_REASON, SELF_MANAGED_LAG_REASON, + broken_for_another_reason, clear_recovered_self_managed_lag, clear_slot_loss_markers, mark_cdc_broken, @@ -69,7 +70,6 @@ CDCSlotNotConfiguredError, classify_cdc_error, ) -from products.warehouse_sources.backend.temporal.data_imports.cdc.legacy_conversion import convert_legacy_cdc_state from products.warehouse_sources.backend.temporal.data_imports.cdc.load_resolution import has_engine_seq from products.warehouse_sources.backend.temporal.data_imports.cdc.naming import ( CDC_EXTRACTION_WORKFLOW_ID_PREFIX, @@ -510,25 +510,13 @@ def _setup(self) -> bool: return True def _prepare_buffer(self) -> None: - """Convert leftover legacy state and start pending snapshots in the buffer, before the WAL read.""" - assert self.source is not None and self.adapter is not None - convert_legacy_cdc_state( - self.source, - self.cdc_schemas, - ingest_mode=self.adapter.parse_cdc_config(self.source).ingest_mode, - logger=self.log, - ) - + """Start pending snapshots in the buffer, before the WAL read.""" for schema in self.cdc_schemas: if not captures_to_buffer(schema): # No lane writes this table mode, so the buffer could never deliver its changes. self._schema_log(schema).warning("cdc_table_mode_not_captured", cdc_table_mode=schema.cdc_table_mode) continue - # A table with deferred runs gets no buffered snapshot here: its old sync could still hand over - # into that buffer without its deferred changes. The reset the conversion staged restarts the - # snapshot instead, and holds the table out of capture until the old sync stops. Once that - # reset has run, which can be before this read, the table's changes belong in the buffer. - if snapshot_can_start_in_buffer(schema) and not (schema.sync_type_config or {}).get("cdc_deferred_runs"): + if snapshot_can_start_in_buffer(schema): self._start_snapshot_in_buffer(schema) self._buffered_table_names.add(schema.name) self.log.info("cdc_buffered_ingress_active", buffered=sorted(self._buffered_table_names)) @@ -1705,17 +1693,20 @@ def cleanup_orphan_slots_activity() -> None: elif cdc_config.management_mode == "self_managed": # Customer owns the slot: surface the broken state but keep the schedule running # and never drop — the lag may recover once they reduce load on the source. + # A marker with another reason stays. The lag marker allows Resume CDC and clears + # itself once the lag drops, which would lift a stop that only Repair CDC may lift. try: - mark_cdc_broken( - source, - SELF_MANAGED_LAG_REASON, - f"Change data capture replication lag exceeded {critical_threshold_mb} MB. " - f"This slot is self-managed, so PostHog did not drop it — reduce load or WAL " - f"retention on the source database, or it may invalidate the slot and " - f"require a full re-sync.", - pause=False, - lag_mb=round(lag_mb, 1), - ) + if not broken_for_another_reason(source, SELF_MANAGED_LAG_REASON): + mark_cdc_broken( + source, + SELF_MANAGED_LAG_REASON, + f"Change data capture replication lag exceeded {critical_threshold_mb} MB. " + f"This slot is self-managed, so PostHog did not drop it — reduce load or WAL " + f"retention on the source database, or it may invalidate the slot and " + f"require a full re-sync.", + pause=False, + lag_mb=round(lag_mb, 1), + ) except Exception: source_log.exception("failed_to_mark_self_managed_broken") metrics.get_sweeper_source_errors_metric().add(1) diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/broken.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/broken.py index 64c48d23cd88..982dfa9d66cc 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/broken.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/broken.py @@ -126,6 +126,25 @@ def mark_cdc_broken( log.warning("cdc_marked_broken", schemas=len(cdc_schemas), newly_broken=len(newly_broken), paused=pause) +def broken_for_another_reason(source: ExternalDataSource, reason: str) -> bool: + """Whether a table that `mark_cdc_broken` would mark already holds a marker with a different reason. + + A table whose sync is off keeps the marker it had, so it must not count. + """ + return ( + 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() + ) + + def clear_recovered_self_managed_lag(source: ExternalDataSource) -> int: """Lift the ``critical_lag_self_managed`` marker once the slot's lag is back under the warning threshold. diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/legacy_conversion.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/legacy_conversion.py deleted file mode 100644 index 09cdb89db066..000000000000 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/legacy_conversion.py +++ /dev/null @@ -1,229 +0,0 @@ -"""One-time conversion of the state the retired legacy CDC lane left behind. - -The legacy lane delivered some tables' changes from capture itself. It paused each such table's -schedule while it streamed, held a snapshotting table's changes as deferred runs in -`sync_type_config`, marked its source `cdc_ingest_mode: legacy`, and created job rows of its own. -Capture now only writes the S3 buffer, so it runs this before every read. Each step finds nothing -once the state it converts is gone, so this module can go once no source logs a conversion. - -Every step that can fail raises before capture reads the WAL, so a failed conversion retries on the -next run with nothing written to the buffer in between. -""" - -from __future__ import annotations - -import datetime as dt -from collections.abc import Sequence -from typing import Any - -from django.db import transaction - -import psycopg -from structlog.types import FilteringBoundLogger - -from posthog.settings import WAREHOUSE_SOURCES_DATABASE_URL - -from products.warehouse_sources.backend.models.external_data_job import ExternalDataJob -from products.warehouse_sources.backend.models.external_data_schema import ( - ExternalDataSchema, - update_sync_type_config_keys, -) -from products.warehouse_sources.backend.models.external_data_source import ExternalDataSource -from products.warehouse_sources.backend.temporal.data_imports.cdc.buffer import purge_buffer_prefix -from products.warehouse_sources.backend.temporal.data_imports.cdc.naming import CDC_EXTRACTION_WORKFLOW_ID_PREFIX -from products.warehouse_sources.backend.temporal.data_imports.cdc.snapshot_lane import ( - CDC_RESET_PENDING_KEY, - stage_handed_over_reset, -) -from products.warehouse_sources.backend.temporal.data_imports.cdc.source_manager import LEGACY_CONVERTED_AT_KEY -from products.warehouse_sources.backend.temporal.data_imports.cdc.types import IngestMode, decode_job_inputs -from products.warehouse_sources_queue.backend.sdk import BatchQueue - -# Shown as latest_error on the capture job rows `_close_stranded_capture_jobs` fails. -STRANDED_CAPTURE_JOB_MESSAGE = ( - "CDC run ended without finalizing this job (worker timeout or eviction). It was superseded by a " - "later run; no data was lost — change capture resumes from the last confirmed replication position." -) -# A healthy capture run queued its first batch within seconds, so a batch-less row older than this is -# abandoned, and one younger may still belong to a run that is finishing. -_STRANDED_JOB_MIN_AGE = dt.timedelta(minutes=30) -# Batches are pruned from the queue after 14 days, so "no batches" is only trustworthy inside that -# window. Older rows stay as they are, because an abandoned run looks the same as one whose batches aged out. -_STRANDED_JOB_MAX_AGE = dt.timedelta(days=14) - - -def convert_legacy_cdc_state( - source: ExternalDataSource, - schemas: Sequence[ExternalDataSchema], - *, - ingest_mode: IngestMode, - logger: FilteringBoundLogger, -) -> None: - """Convert whatever legacy state this source still carries. Raises on a failure that must retry.""" - # Staged first, so a legacy source's schedule rebuild below leaves these tables to their reset. - for schema in schemas: - if (schema.sync_type_config or {}).get("cdc_deferred_runs"): - _hand_deferred_runs_to_reset(schema, logger) - if ingest_mode != "buffered": - _convert_legacy_source(source, schemas, logger) - try: - _close_stranded_capture_jobs(source, schemas, logger) - except Exception: - # Cosmetic: a stranded row only misreports a run that already ended. - logger.warning("cdc_stranded_capture_jobs_close_failed", exc_info=True) - - -def _hand_deferred_runs_to_reset(schema: ExternalDataSchema, logger: FilteringBoundLogger) -> None: - """Stage a reset for a table the legacy lane still held deferred runs for, which capture then finishes. - - Nothing merges deferred runs anymore, so the table snapshots again in the buffer. The reset waits - until the old sync can no longer hand over, because one that hands over after it would flip the - table to streaming without those changes. It drops the deferred runs and starts the new snapshot. - Staging again while it waits is harmless: the reset merges into the one already pending. - """ - deferred_runs = len(schema.sync_type_config.get("cdc_deferred_runs") or []) - schema.sync_type_config = update_sync_type_config_keys(schema.id, schema.team_id, mutate=stage_handed_over_reset) - logger.info("cdc_legacy_deferred_runs_handed_to_reset", schema_id=str(schema.id), deferred_runs=deferred_runs) - - -def _convert_legacy_source( - source: ExternalDataSource, schemas: Sequence[ExternalDataSchema], logger: FilteringBoundLogger -) -> None: - """Move a source whose changes the legacy lane still delivered onto the buffer. - - Its buffer holds only copies of changes the legacy lane already delivered, so it is emptied - before capture writes the first file. What capture reads next continues from the slot's - position, which legacy delivery already reached, so each table is current as of now, and is - stamped so that its first sync reads the buffer instead of taking it for expired. The legacy lane - paused the schedule of every streaming table because capture delivered its changes. That schedule - is now the table's consumer, so it is rebuilt. Legacy batches still in the load queue land first, - because the consumer stands down while any are in flight. - """ - # See `_resume_schedule` for why this import is deferred. - from products.data_warehouse.backend.facade.api import cdc_min_interval # noqa: PLC0415 - - all_cdc_schema_ids = list( - ExternalDataSchema.objects.filter( - team_id=source.team_id, source_id=source.id, sync_type=ExternalDataSchema.SyncType.CDC - ) - .exclude(deleted=True) - .values_list("id", flat=True) - ) - for schema_id in all_cdc_schema_ids: - purge_buffer_prefix(source.team_id, str(schema_id), logger, strict=True) - converted_at = dt.datetime.now(dt.UTC).isoformat() - for schema_id in all_cdc_schema_ids: - update_sync_type_config_keys(schema_id, source.team_id, updates={LEGACY_CONVERTED_AT_KEY: converted_at}) - - capture_interval = cdc_min_interval(schema.sync_frequency_interval for schema in schemas) - for schema in schemas: - _keep_capture_cadence(schema, capture_interval, logger) - _resume_schedule(schema) - - # Written last, so a failure above repeats the whole conversion. A repeat is safe because the run - # fails before it reads the WAL, so the buffer has gained nothing a second purge would drop. - source.job_inputs = _mark_source_buffered(source) - logger.info("cdc_legacy_source_converted", schemas=len(all_cdc_schema_ids)) - - -def _keep_capture_cadence( - schema: ExternalDataSchema, capture_interval: dt.timedelta, logger: FilteringBoundLogger -) -> None: - """Speed a table up to capture's cadence, which is how often the legacy lane delivered its changes. - - Its own schedule loads it from now on, so a table set slower than its source's fastest table would - otherwise fall behind by the difference. - """ - if schema.sync_frequency_interval is None or schema.sync_frequency_interval <= capture_interval: - return - logger.info( - "cdc_legacy_table_frequency_raised", - schema_id=str(schema.id), - from_seconds=schema.sync_frequency_interval.total_seconds(), - to_seconds=capture_interval.total_seconds(), - ) - ExternalDataSchema.objects.filter(id=schema.id, team_id=schema.team_id).update( - sync_frequency_interval=capture_interval, updated_at=dt.datetime.now(dt.UTC) - ) - schema.sync_frequency_interval = capture_interval - - -def _mark_source_buffered(source: ExternalDataSource) -> dict[str, Any]: - """Mark the source buffered on its current row, so a concurrent edit to its settings survives. - - A queryset update, not `source.save()`, whose activity-log diff walks the source's whole job history. - """ - with transaction.atomic(): - locked = ExternalDataSource.objects.select_for_update(of=("self",)).get(id=source.id, team_id=source.team_id) - job_inputs = {**decode_job_inputs(locked.job_inputs), "cdc_ingest_mode": "buffered"} - ExternalDataSource.objects.filter(id=source.id, team_id=source.team_id).update( - job_inputs=job_inputs, updated_at=dt.datetime.now(dt.UTC) - ) - return job_inputs - - -def _resume_schedule(schema: ExternalDataSchema) -> None: - """Rebuild the table's schedule unpaused, and create it if it is missing. - - A plain unpause does nothing to a missing schedule. Skipped where the pause is deliberate: a - schema whose status is Paused, an admin-triggered run that holds the schedule until it finishes, - a broken source that waits for Repair CDC, which unpauses it, and a pending reset, which unpauses - it once it has reset the table. A capture pause is not checked, because capture is running, so - that pause is over. A schema with no sync frequency is skipped because the schedule builder cannot - turn a null interval into a cadence. - """ - # data_load.service imports temporalio at module scope; deferred to keep the Temporal client off - # this module's import path, as capture's other schedule calls do. - from products.data_warehouse.backend.facade.api import sync_external_data_job_workflow # noqa: PLC0415 - - config = schema.sync_type_config or {} - if ( - schema.status == ExternalDataSchema.Status.PAUSED - or config.get("cdc_broken") - or config.get("admin_unpause_schedule_after_run") - or config.get(CDC_RESET_PENDING_KEY) - or schema.sync_frequency_interval is None - ): - return - sync_external_data_job_workflow(schema, create=True, should_sync=True, trigger_immediately=False) - - -def _close_stranded_capture_jobs( - source: ExternalDataSource, schemas: Sequence[ExternalDataSchema], logger: FilteringBoundLogger -) -> None: - """Fail RUNNING job rows the legacy lane's capture runs left behind. - - Legacy capture created one per written table and closed it when the run finished. A run that died - mid-way left its row RUNNING. A row with batches in the queue belongs to the loader, which closes - it. A row without any has no outstanding work, so failing it cannot race a late load. Scoped to - capture's own workflow ids, so a table's scheduled sync is never touched. - """ - now = dt.datetime.now(tz=dt.UTC) - stranded = list( - ExternalDataJob.objects.filter( - team_id=source.team_id, - schema_id__in=[s.id for s in schemas], - status=ExternalDataJob.Status.RUNNING, - workflow_id__startswith=CDC_EXTRACTION_WORKFLOW_ID_PREFIX, - created_at__gt=now - _STRANDED_JOB_MAX_AGE, - created_at__lt=now - _STRANDED_JOB_MIN_AGE, - ).order_by("created_at")[:200] - ) - if not stranded: - return - - closed = 0 - conn = psycopg.Connection.connect(WAREHOUSE_SOURCES_DATABASE_URL, autocommit=True) - try: - for job in stranded: - if BatchQueue.count_batches_for_run(conn, job_id=str(job.id)) > 0: - continue - job.status = ExternalDataJob.Status.FAILED - job.latest_error = STRANDED_CAPTURE_JOB_MESSAGE - job.finished_at = now - job.save(update_fields=["status", "latest_error", "finished_at", "updated_at"]) - closed += 1 - finally: - conn.close() - if closed: - logger.info("cdc_stranded_capture_jobs_closed", count=closed) diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/snapshot_lane.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/snapshot_lane.py index 9d2a7b6294d2..fbae8c5a8c56 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/snapshot_lane.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/snapshot_lane.py @@ -22,7 +22,6 @@ update_sync_type_config_keys, ) from products.warehouse_sources.backend.temporal.data_imports.cdc.naming import CDC_EXTRACTION_WORKFLOW_ID_PREFIX -from products.warehouse_sources.backend.temporal.data_imports.cdc.types import parse_ingest_mode if TYPE_CHECKING: from structlog.types import FilteringBoundLogger @@ -43,20 +42,12 @@ def snapshot_in_buffer(schema: ExternalDataSchema) -> bool: def resnapshot_stays_in_buffer(schema: ExternalDataSchema) -> bool: """Whether a reset of this schema to snapshot keeps its changes in the buffer. - True for a table whose buffer already holds an unbroken run of its changes: a streaming table on - a buffered source with no deferred runs left, or one already snapshotting there. A capture run in - progress keeps adding to that buffer. Any other table is started by capture, which first empties - the buffer, because files from before a gap in capture must not be replayed. A source capture - has not converted yet is one such case: its buffer holds copies of changes the legacy lane - already delivered. + True for a table whose buffer already holds an unbroken run of its changes: a streaming table, or + one already snapshotting there. A capture run in progress keeps adding to that buffer. Any other + table is started by capture, which first empties the buffer, because files from before a gap in + capture must not be replayed. """ - if snapshot_in_buffer(schema): - return True - return ( - parse_ingest_mode(schema.source.job_inputs) == "buffered" - and schema.cdc_mode == "streaming" - and not (schema.sync_type_config or {}).get("cdc_deferred_runs") - ) + return snapshot_in_buffer(schema) or schema.cdc_mode == "streaming" def cancel_running_sync(schema: ExternalDataSchema) -> str | None: diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/source_manager.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/source_manager.py index eb415cc62e3a..000cee792348 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/source_manager.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/source_manager.py @@ -177,8 +177,6 @@ def served_lanes(schema: ExternalDataSchema) -> list[CDCLane]: BUFFER_LISTED_AT_KEY = "cdc_buffer_listed_at" # File name to ETag of the files at the highest position that listing saw, kept beside it. BUFFER_LISTED_TAIL_KEY = "cdc_buffer_listed_tail" -# When capture moved the table's legacy source onto the buffer, kept in the table's sync_type_config. -LEGACY_CONVERTED_AT_KEY = "cdc_legacy_converted_at" @frozen @@ -259,11 +257,6 @@ def buffer_may_have_expired_unread(schema: ExternalDataSchema, now: dt.datetime) if schema.last_synced_at is None: return False cutoff = now - BUFFER_FILE_RETENTION - # The conversion empties the buffer at a position the legacy lane had already delivered, so the - # table is current from then on, though no run has listed the buffer yet. - converted_at = (schema.sync_type_config or {}).get(LEGACY_CONVERTED_AT_KEY) - if converted_at is not None and dt.datetime.fromisoformat(converted_at) >= cutoff: - return False # Every completion moves it, so nothing has drained since the cutoff either. if schema.last_synced_at < cutoff: return True diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_buffered_dispatch.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_buffered_dispatch.py index 270a76db23b2..fe073353fbbe 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_buffered_dispatch.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_buffered_dispatch.py @@ -17,7 +17,7 @@ _SNAPSHOT_LANE = "products.warehouse_sources.backend.temporal.data_imports.cdc.snapshot_lane" -def _schema(ingest_mode: str = "buffered", **overrides) -> MagicMock: +def _schema(**overrides) -> MagicMock: schema = MagicMock() schema.name = "users" schema.is_cdc = overrides.get("is_cdc", True) @@ -32,7 +32,7 @@ def _schema(ingest_mode: str = "buffered", **overrides) -> MagicMock: schema.resolved_s3_folder_name = None schema.primary_key_columns = ["id"] schema.last_synced_at = None - schema.source.job_inputs = {"cdc_enabled": True, "cdc_ingest_mode": ingest_mode} + schema.source.job_inputs = {"cdc_enabled": True} return schema @@ -144,9 +144,8 @@ def test_a_missing_job_row_fails_the_run(self): with pytest.raises(ValueError, match="no job row"): _dispatch(_schema(), _inputs(), job_version=None) - @pytest.mark.parametrize("ingest_mode, in_flight", [("buffered", True), ("legacy", False)]) - def test_a_delivery_still_in_flight_or_an_unconverted_source_no_ops_the_tick(self, ingest_mode, in_flight): - response = _dispatch(_schema(ingest_mode, cdc_table_mode="both"), _inputs(), in_flight=in_flight) + def test_a_delivery_still_in_flight_no_ops_the_tick(self): + response = _dispatch(_schema(cdc_table_mode="both"), _inputs(), in_flight=True) assert list(response.items()) == [] assert response.lanes is None diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_cleanup_orphan_slots.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_cleanup_orphan_slots.py index 92002bc71830..69edd78573dc 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_cleanup_orphan_slots.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_cleanup_orphan_slots.py @@ -211,9 +211,24 @@ def test_critical_lag_auto_drop_disabled_does_not_drop(team): adapter.drop_resources.assert_not_called() -def test_critical_lag_self_managed_marks_broken_without_drop_or_pause(team): +@pytest.mark.parametrize( + "existing_reason, marked_table_syncs, expected_reason, expected_status", + [ + (None, True, "critical_lag_self_managed", ExternalDataSource.Status.ERROR), + ("legacy_lane_retired", True, "legacy_lane_retired", "Completed"), + ("legacy_lane_retired", False, "critical_lag_self_managed", ExternalDataSource.Status.ERROR), + ], +) +def test_critical_lag_self_managed_marks_broken_without_drop_or_pause( + team, existing_reason, marked_table_syncs, expected_reason, expected_status +): source = _create_source(team, job_inputs=_cdc_job_inputs(management="self_managed")) schema = _create_cdc_schema(team, source) + if existing_reason: + marked = schema if marked_table_syncs else _create_cdc_schema(team, source, name="sync_off") + marked.sync_type_config = {**marked.sync_type_config, "cdc_broken": {"reason": existing_reason}} + marked.should_sync = marked_table_syncs + marked.save() adapter = _mock_adapter(lag_bytes=5000 * 1024 * 1024) _, _, mock_pause = _run(adapter) @@ -222,9 +237,9 @@ def test_critical_lag_self_managed_marks_broken_without_drop_or_pause(team): adapter.drop_resources.assert_not_called() mock_pause.assert_not_called() source.refresh_from_db() - assert source.status == ExternalDataSource.Status.ERROR + assert source.status == expected_status schema.refresh_from_db() - assert schema.sync_type_config["cdc_broken"]["reason"] == "critical_lag_self_managed" + assert schema.sync_type_config["cdc_broken"]["reason"] == expected_reason def _job(team, source, schema, status, age, workflow_id=None): diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_extract_activity.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_extract_activity.py index 121e745c60ef..fb505a97dc3e 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_extract_activity.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_extract_activity.py @@ -144,7 +144,6 @@ def _stub_app_db_writes(): "_update_schema_sync_type_config", side_effect=_fake_update_schema_sync_type_config, ), - patch(f"{_ACTIVITIES}.convert_legacy_cdc_state"), patch(f"{_ACTIVITIES}.cancel_running_sync", return_value=None), ): yield @@ -1988,41 +1987,31 @@ def test_a_truncate_resets_the_table_to_a_snapshot_in_an_emptied_buffer(self, _n @parameterized.expand( [ - ("sync_still_stopping", None, None, {}, "waits"), + ("sync_still_stopping", None, None, "waits"), ( "sync_closed_with_batches_still_loading", RPCError("workflow not found", RPCStatusCode.NOT_FOUND, b""), 30.0, - {}, "waits", ), ( "sync_closed_and_nothing_queued", RPCError("workflow not found", RPCStatusCode.NOT_FOUND, b""), None, - {}, "resets", ), - ( - "deferred_runs_left_and_nothing_queued", - RPCError("workflow not found", RPCStatusCode.NOT_FOUND, b""), - None, - {"cdc_deferred_runs": [{"run_uuid": "r1"}]}, - "resets", - ), - ("temporal_unavailable", RPCError("unavailable", RPCStatusCode.UNAVAILABLE, b""), None, {}, "raises"), + ("temporal_unavailable", RPCError("unavailable", RPCStatusCode.UNAVAILABLE, b""), None, "raises"), ] ) @patch(f"{_ACTIVITIES}.purge_buffer_prefix") @patch("products.warehouse_sources.backend.temporal.data_imports.cdc.snapshot_lane.ExternalDataJob") def test_a_reset_stops_the_tables_running_sync_first( - self, _name, cancel_error, oldest_queued_batch_age, config, outcome, MockJob, mock_purge + self, _name, cancel_error, oldest_queued_batch_age, outcome, MockJob, mock_purge ): # A snapshot that started before a repeated reset missed the changes the reset drops, so it # must not reach its hand-over. source = _make_source() schema = _make_schema("users", cdc_mode="snapshot", source=source) - schema.sync_type_config.update(config) act = _make_extract_activity(source) running = MockJob.objects.filter.return_value.exclude.return_value.exclude.return_value running.order_by.return_value.first.return_value = MagicMock(workflow_id="users-snapshot") @@ -2059,33 +2048,17 @@ def test_a_reset_stops_the_tables_running_sync_first( @parameterized.expand( [ - ("sync_still_stopping", "users-snapshot", {"clear_deferred_runs": False}, False, True), - ("sync_stopped", None, {"clear_deferred_runs": False}, False, False), - ("sync_stopped_after_a_request_reset", None, {"clear_deferred_runs": True, "trigger": True}, False, False), - ( - "legacy_deferred_runs_old_sync_stopping", - "users-snapshot", - {"clear_deferred_runs": True, "trigger": True}, - True, - True, - ), - ( - "legacy_deferred_runs_old_sync_stopped", - None, - {"clear_deferred_runs": True, "trigger": True}, - True, - False, - ), + ("sync_still_stopping", "users-snapshot", {"clear_deferred_runs": False}, True), + ("sync_stopped", None, {"clear_deferred_runs": False}, False), + ("sync_stopped_after_a_request_reset", None, {"clear_deferred_runs": True, "trigger": True}, False), ] ) def test_a_pending_reset_finishes_before_the_read_once_the_sync_stopped( - self, _name, stopping_workflow_id, pending, deferred_runs, waits + self, _name, stopping_workflow_id, pending, waits ): source = _make_source() - schema = _make_schema("users", cdc_mode="snapshot" if deferred_runs else "streaming", source=source) + schema = _make_schema("users", cdc_mode="streaming", source=source) schema.sync_type_config["cdc_reset_pending"] = pending - if deferred_runs: - schema.sync_type_config["cdc_deferred_runs"] = [{"run_uuid": "r1"}] events = [_make_event(op="I", position="0/100", columns={"id": 1})] with ( @@ -2104,7 +2077,6 @@ def test_a_pending_reset_finishes_before_the_read_once_the_sync_stopped( assert trigger.called is (not waits and bool(pending.get("trigger"))) assert schema.sync_type_config.get("reset_pipeline") is (None if waits else True) assert ("cdc_reset_pending" in schema.sync_type_config) is waits - assert ("cdc_deferred_runs" in schema.sync_type_config) is (deferred_runs and waits) assert capture.buffer.write_batch.called is not waits capture.reader.confirm_position.assert_called_once_with("0/100") diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_legacy_conversion.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_legacy_conversion.py deleted file mode 100644 index bf21ffb07f5c..000000000000 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_legacy_conversion.py +++ /dev/null @@ -1,214 +0,0 @@ -import datetime as dt - -import pytest -from posthog.test.base import BaseTest -from unittest.mock import MagicMock, patch - -from django.utils import timezone - -from parameterized import parameterized - -from products.warehouse_sources.backend.models.external_data_job import ExternalDataJob -from products.warehouse_sources.backend.models.external_data_schema import ExternalDataSchema -from products.warehouse_sources.backend.models.external_data_source import ExternalDataSource -from products.warehouse_sources.backend.temporal.data_imports.cdc.legacy_conversion import ( - STRANDED_CAPTURE_JOB_MESSAGE, - convert_legacy_cdc_state, -) -from products.warehouse_sources.backend.temporal.data_imports.cdc.source_manager import buffer_may_have_expired_unread -from products.warehouse_sources.backend.temporal.data_imports.cdc.types import parse_ingest_mode - -_MODULE = "products.warehouse_sources.backend.temporal.data_imports.cdc.legacy_conversion" -_FACADE = "products.data_warehouse.backend.facade.api" - - -class TestLegacyConversion(BaseTest): - def _source(self, ingest_mode: str | None) -> ExternalDataSource: - job_inputs: dict = {"cdc_enabled": True, "cdc_slot_name": "slot"} - if ingest_mode is not None: - job_inputs["cdc_ingest_mode"] = ingest_mode - return ExternalDataSource.objects.create(team=self.team, source_type="Postgres", job_inputs=job_inputs) - - def _schema(self, source: ExternalDataSource, name: str, **overrides) -> ExternalDataSchema: - return ExternalDataSchema.objects.create( - team=self.team, - source=source, - name=name, - sync_type=ExternalDataSchema.SyncType.CDC, - sync_type_config=overrides.pop("sync_type_config", {"cdc_mode": "streaming"}), - initial_sync_complete=overrides.pop("initial_sync_complete", True), - should_sync=overrides.pop("should_sync", True), - sync_frequency_interval=overrides.pop("sync_frequency_interval", dt.timedelta(minutes=5)), - **overrides, - ) - - def _job( - self, - schema: ExternalDataSchema, - *, - workflow_id: str, - age: dt.timedelta, - status: str = ExternalDataJob.Status.RUNNING, - ) -> ExternalDataJob: - job = ExternalDataJob.objects.create( - team=self.team, - pipeline=schema.source, - schema=schema, - status=status, - rows_synced=0, - workflow_id=workflow_id, - schema_snapshot={"sync_type_config": schema.sync_type_config}, - ) - ExternalDataJob.objects.filter(id=job.id).update(created_at=timezone.now() - age) - return job - - def _convert(self, source, schemas, *, queued_batches: int = 0, sync_workflow=None): - purge = MagicMock() - sync_workflow = sync_workflow or MagicMock() - with ( - patch(f"{_MODULE}.purge_buffer_prefix", purge), - patch(f"{_FACADE}.sync_external_data_job_workflow", sync_workflow), - patch(f"{_MODULE}.psycopg"), - patch(f"{_MODULE}.BatchQueue.count_batches_for_run", return_value=queued_batches), - ): - convert_legacy_cdc_state( - source, schemas, ingest_mode=parse_ingest_mode(source.job_inputs), logger=MagicMock() - ) - return purge, sync_workflow - - def test_a_legacy_source_is_emptied_resumed_and_marked_buffered_once(self): - source = self._source(ingest_mode=None) - streaming = self._schema(source, "users") - sync_off = self._schema(source, "orders", should_sync=False) - edited_since_capture_loaded_it = ExternalDataSource.objects.get(id=source.id) - edited_since_capture_loaded_it.job_inputs = { - **edited_since_capture_loaded_it.job_inputs, - "cdc_slot_name": "new", - } - edited_since_capture_loaded_it.save(update_fields=["job_inputs"]) - - purge, sync_workflow = self._convert(source, [streaming]) - - assert {call.args[1] for call in purge.call_args_list} == {str(streaming.id), str(sync_off.id)} - assert all(call.kwargs["strict"] for call in purge.call_args_list) - assert [call.args[0].id for call in sync_workflow.call_args_list] == [streaming.id] - assert sync_workflow.call_args.kwargs == {"create": True, "should_sync": True, "trigger_immediately": False} - source.refresh_from_db() - assert parse_ingest_mode(source.job_inputs) == "buffered" - assert source.job_inputs["cdc_slot_name"] == "new" - - purge, sync_workflow = self._convert(source, [streaming]) - - purge.assert_not_called() - sync_workflow.assert_not_called() - - def test_a_converted_table_reads_its_buffer_instead_of_re_snapshotting(self): - source = self._source(ingest_mode="legacy") - schema = self._schema(source, "users", last_synced_at=timezone.now() - dt.timedelta(minutes=10)) - self._job( - schema, - workflow_id=f"cdc-extraction-{source.id}-run", - age=dt.timedelta(minutes=10), - status=ExternalDataJob.Status.COMPLETED, - ) - - self._convert(source, [schema]) - - schema.refresh_from_db() - assert buffer_may_have_expired_unread(schema, timezone.now()) is False - - @parameterized.expand([("legacy", dt.timedelta(minutes=5)), ("buffered", dt.timedelta(hours=6))]) - def test_a_converted_table_keeps_the_cadence_the_legacy_lane_gave_it(self, ingest_mode, expected_interval): - source = self._source(ingest_mode=ingest_mode) - fast = self._schema(source, "users", sync_frequency_interval=dt.timedelta(minutes=5)) - slow = self._schema(source, "events", sync_frequency_interval=dt.timedelta(hours=6)) - - _, sync_workflow = self._convert(source, [fast, slow]) - - slow.refresh_from_db() - assert slow.sync_frequency_interval == expected_interval - rebuilt = {call.args[0].id: call.args[0].sync_frequency_interval for call in sync_workflow.call_args_list} - assert rebuilt.get(slow.id, expected_interval) == expected_interval - - def test_a_failed_schedule_rebuild_leaves_the_source_to_convert_again(self): - source = self._source(ingest_mode="legacy") - schema = self._schema(source, "users") - - with pytest.raises(RuntimeError): - self._convert(source, [schema], sync_workflow=MagicMock(side_effect=RuntimeError("temporal down"))) - - source.refresh_from_db() - assert parse_ingest_mode(source.job_inputs) == "legacy" - - @parameterized.expand([("buffered_source", "buffered"), ("legacy_source", "legacy")]) - def test_a_table_with_deferred_runs_hands_its_snapshot_restart_to_capture(self, _name, ingest_mode): - source = self._source(ingest_mode=ingest_mode) - schema = self._schema( - source, - "users", - sync_type_config={"cdc_mode": "snapshot", "cdc_deferred_runs": [{"run_uuid": "r1"}]}, - initial_sync_complete=False, - ) - - _, sync_workflow = self._convert(source, [schema]) - - sync_workflow.assert_not_called() - schema.refresh_from_db() - pending = schema.sync_type_config["cdc_reset_pending"] - assert pending["clear_deferred_runs"] is True - assert pending["trigger"] is True - - @parameterized.expand( - [ - ("billing_paused", {"status": ExternalDataSchema.Status.PAUSED}, False), - ( - "admin_run_in_flight", - {"sync_type_config": {"cdc_mode": "streaming", "admin_unpause_schedule_after_run": True}}, - False, - ), - ( - "broken", - {"sync_type_config": {"cdc_mode": "streaming", "cdc_broken": {"reason": "slot_missing"}}}, - False, - ), - ("no_sync_frequency", {"sync_frequency_interval": None}, False), - ( - "reset_pending", - {"sync_type_config": {"cdc_mode": "streaming", "cdc_reset_pending": {"clear_deferred_runs": False}}}, - False, - ), - # The conversion runs once, so skipping this schedule would leave it paused for good. - ( - "capture_paused_earlier", - {"sync_type_config": {"cdc_mode": "streaming", "cdc_extraction_paused": {"reason": "auth_failed"}}}, - True, - ), - ] - ) - def test_the_schedule_resumes_unless_its_pause_is_deliberate(self, _name, overrides, resumed): - source = self._source(ingest_mode="legacy") - schema = self._schema(source, "users", **overrides) - - _, sync_workflow = self._convert(source, [schema]) - - assert sync_workflow.called is resumed - - @parameterized.expand( - [ - ("abandoned_capture_run", "cdc-extraction-", dt.timedelta(hours=1), 0, ExternalDataJob.Status.FAILED), - ("loader_still_owns_it", "cdc-extraction-", dt.timedelta(hours=1), 3, ExternalDataJob.Status.RUNNING), - ("possibly_still_finishing", "cdc-extraction-", dt.timedelta(minutes=5), 0, ExternalDataJob.Status.RUNNING), - ("a_scheduled_sync", "schema-", dt.timedelta(hours=1), 0, ExternalDataJob.Status.RUNNING), - ] - ) - def test_only_abandoned_capture_jobs_are_failed(self, _name, workflow_prefix, age, queued_batches, expected): - source = self._source(ingest_mode="buffered") - schema = self._schema(source, "users") - job = self._job(schema, workflow_id=f"{workflow_prefix}{source.id}", age=age) - - self._convert(source, [schema], queued_batches=queued_batches) - - job.refresh_from_db() - assert job.status == expected - if expected == ExternalDataJob.Status.FAILED: - assert job.latest_error == STRANDED_CAPTURE_JOB_MESSAGE diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_listing_proof.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_listing_proof.py index 75aa94c41e97..7e572560111d 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_listing_proof.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_listing_proof.py @@ -16,7 +16,6 @@ from products.warehouse_sources.backend.temporal.data_imports.cdc.source_manager import ( BUFFER_LISTED_AT_KEY, BUFFER_LISTED_TAIL_KEY, - LEGACY_CONVERTED_AT_KEY, ListingProof, buffer_may_have_expired_unread, clear_listing, @@ -119,19 +118,14 @@ def test_a_naive_timestamp_is_refused(self): class TestBufferExpiredUnread(BaseTest): - def _schema(self, *, synced_days_ago: int | None, converted_days_ago: int | None = None) -> ExternalDataSchema: + def _schema(self, *, synced_days_ago: int | None) -> ExternalDataSchema: source = ExternalDataSource.objects.create( team=self.team, source_id="s", connection_id="c", status="Running", source_type="Postgres" ) now = dt.datetime.now(tz=dt.UTC) synced = None if synced_days_ago is None else now - dt.timedelta(days=synced_days_ago) - config = ( - {} - if converted_days_ago is None - else {LEGACY_CONVERTED_AT_KEY: (now - dt.timedelta(days=converted_days_ago)).isoformat()} - ) return ExternalDataSchema.objects.create( - team=self.team, source=source, name="users", last_synced_at=synced, sync_type_config=config + team=self.team, source=source, name="users", last_synced_at=synced, sync_type_config={} ) def _completed_job(self, schema: ExternalDataSchema, *, days_ago: int, snapshot: dict) -> None: @@ -149,19 +143,17 @@ def _completed_job(self, schema: ExternalDataSchema, *, days_ago: int, snapshot: @parameterized.expand( [ - ("never_synced", None, None, None, False), - ("no_completion_since_retention", 15, None, None, True), - ("only_stand_downs_since_retention", 1, (1, {}), None, True), - ("a_run_listed_the_buffer", 1, (10, {BUFFER_LISTED_AT_KEY: "2026-01-01T00:00:00+00:00"}), None, False), - ("a_snapshot_reseeded_the_table", 1, (3, {"sync_type_config": {"cdc_mode": "snapshot"}}), None, False), + ("never_synced", None, None, False), + ("no_completion_since_retention", 15, None, True), + ("only_stand_downs_since_retention", 1, (1, {}), True), + ("a_run_listed_the_buffer", 1, (10, {BUFFER_LISTED_AT_KEY: "2026-01-01T00:00:00+00:00"}), False), + ("a_snapshot_reseeded_the_table", 1, (3, {"sync_type_config": {"cdc_mode": "snapshot"}}), False), ( "the_last_listing_is_older_than_retention", 1, (20, {BUFFER_LISTED_AT_KEY: "2026-01-01T00:00:00+00:00"}), - None, True, ), - ("a_legacy_conversion_older_than_retention", 1, (1, {}), 20, True), ] ) def test_only_a_listing_or_a_snapshot_proves_the_buffer_was_read( @@ -169,10 +161,9 @@ def test_only_a_listing_or_a_snapshot_proves_the_buffer_was_read( _name: str, synced_days_ago: int | None, job: tuple[int, dict] | None, - converted_days_ago: int | None, expired: bool, ) -> None: - schema = self._schema(synced_days_ago=synced_days_ago, converted_days_ago=converted_days_ago) + schema = self._schema(synced_days_ago=synced_days_ago) if job is not None: self._completed_job(schema, days_ago=job[0], snapshot=job[1]) diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_metrics.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_metrics.py index ef3cddb7c59e..8afa8e260f87 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_metrics.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_metrics.py @@ -129,7 +129,6 @@ def _extract_patches(source, schemas, events): patch(f"{_ACTIVITIES}.ExternalDataSource") as MockSource, patch.object(CDCExtractActivity, "_get_cdc_schemas", return_value=schemas), patch.object(CDCExtractActivity, "_update_schema_sync_type_config"), - patch(f"{_ACTIVITIES}.convert_legacy_cdc_state"), patch(f"{_ACTIVITIES}.get_cdc_adapter", return_value=adapter), patch(f"{_ACTIVITIES}.CDCBufferWriter") as MockBufferWriter, ): diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_source_manager.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_source_manager.py index ad7049fb707f..0f4e27addc03 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_source_manager.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/tests/test_source_manager.py @@ -359,15 +359,10 @@ def test_capture_writes_every_table_in_a_mode_the_buffer_serves(self, _name, ove {"cdc_mode": "snapshot", "sync_type_config": {"cdc_snapshot_lane": "buffer"}}, True, ), - # Its buffer holds copies of changes the legacy lane already delivered, until capture converts it. - ("streaming_on_a_source_not_converted_yet", {"job_inputs": {"cdc_ingest_mode": "legacy"}}, False), - # The previous release skipped these tables' capture, so their buffer has a gap. - ("streaming_with_deferred_runs_left", {"sync_type_config": {"cdc_deferred_runs": [{"run": 1}]}}, False), ] ) def test_a_resnapshot_stays_in_the_buffer_only_when_the_buffer_holds_every_change(self, _name, overrides, stays): - schema = _schema(**{"job_inputs": {"cdc_ingest_mode": "buffered"}, **overrides}) - assert resnapshot_stays_in_buffer(schema) is stays + assert resnapshot_stays_in_buffer(_schema(**overrides)) is stays class TestBufferedGating: @@ -386,20 +381,6 @@ def test_ineligible_schemas_do_not_consume_the_buffer(self, _name, overrides): def test_every_streaming_table_mode_serves_the_buffered_lane(self, table_mode): assert serves_buffered_lane(_schema(cdc_table_mode=table_mode)) is True - # A source capture has not converted yet still reads "legacy", and its consumer must run on v3 too. - @parameterized.expand( - [ - (f"{table_mode}_{ingest_mode}", table_mode, ingest_mode) - for table_mode in ("consolidated", "cdc_only", "both") - for ingest_mode in ("buffered", "legacy") - ] - ) - def test_a_streaming_schema_forces_the_buffered_consumer_on_its_scheduled_sync( - self, _name, table_mode, ingest_mode - ): - schema = _schema(job_inputs={"cdc_ingest_mode": ingest_mode}, cdc_table_mode=table_mode) - assert scheduled_sync_consumes_buffer(schema) is True - @parameterized.expand( [ ("unrecognized_table_mode", {"cdc_table_mode": "something_new"}), diff --git a/products/warehouse_sources/backend/temporal/data_imports/cdc/types.py b/products/warehouse_sources/backend/temporal/data_imports/cdc/types.py index f8f77b57d783..9fd466990ac6 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/cdc/types.py +++ b/products/warehouse_sources/backend/temporal/data_imports/cdc/types.py @@ -10,19 +10,14 @@ ManagementMode = Literal["posthog", "self_managed"] -# How a source's change events reach the loader. `legacy`: capture transforms and dispatches them -# itself. `buffered`: capture only writes the S3 buffer, and the normal scheduled sync consumes it. -IngestMode = Literal["legacy", "buffered"] - class CDCJobInputsUnreadableError(Exception): """`job_inputs` did not resolve to a mapping, so no CDC setting can be read from the source. Non-retryable: the stored value replays identically on every read. - Raised instead of reading the settings as absent, because a source whose stored mode cannot be - read would then route onto the lane it was never flipped to, and a buffer nothing consumes - looks the same as an idle one. + Raised instead of reading the settings as absent, because every setting would then fall back to + its default, the slot name included. """ @@ -49,12 +44,6 @@ def decode_job_inputs(job_inputs: Mapping[str, Any] | str | None) -> Mapping[str return decoded -def parse_ingest_mode(job_inputs: Mapping[str, Any] | str | None) -> IngestMode: - """An unrecognized value reads as legacy: it must not route a source onto a path it was never - flipped to. Raises ``CDCJobInputsUnreadableError`` when there is no value to read at all.""" - return "buffered" if decode_job_inputs(job_inputs).get("cdc_ingest_mode") == "buffered" else "legacy" - - @dataclass(frozen=True) class CDCConfig: """Base class for engine-specific CDC configs returned by ``parse_cdc_config``. @@ -71,7 +60,6 @@ class CDCConfig: lag_warning_threshold_mb: int lag_critical_threshold_mb: int auto_drop_slot: bool - ingest_mode: IngestMode class CDCPosition(Protocol): diff --git a/products/warehouse_sources/backend/temporal/data_imports/pipelines/test_pipeline_sync.py b/products/warehouse_sources/backend/temporal/data_imports/pipelines/test_pipeline_sync.py index d404f0d036ee..234e6a57cc6a 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/pipelines/test_pipeline_sync.py +++ b/products/warehouse_sources/backend/temporal/data_imports/pipelines/test_pipeline_sync.py @@ -853,7 +853,6 @@ def test_a_snapshot_the_buffer_carried_keeps_its_files_and_the_flip_clears_the_m sync_type="cdc", config={"cdc_mode": "snapshot", "cdc_snapshot_lane": "buffer"}, initial_sync_complete=False, - job_inputs={"cdc_ingest_mode": "buffered"}, ) with patch( @@ -883,7 +882,6 @@ def test_capture_cannot_mark_the_table_while_the_hand_over_purges(self) -> None: connection_id=str(uuid.uuid4()), status="Completed", source_type="Postgres", - job_inputs={"cdc_ingest_mode": "buffered"}, ) schema = ExternalDataSchema.objects.create( team_id=self.team.pk, diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/adapter.py b/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/adapter.py index b9487aa2664b..9da6747ed7cc 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/adapter.py +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/adapter.py @@ -245,8 +245,8 @@ def _recreate() -> str: _recreate, _retry_logger, is_retryable=_is_dropped_or_connect_timeout ) - # Every schema is reset to snapshot before this runs, so nothing from the dead slot is owed: - # the new slot starts on the buffer, as a new source does, and capture has nothing to convert. + # No code on this release reads the ingest mode. A worker on a release that reads it treats a + # source without it as legacy and empties its buffer, so a rollback needs the value. return {"cdc_consistent_point": consistent_point, "cdc_ingest_mode": "buffered"} def setup_resources( @@ -273,8 +273,7 @@ def setup_resources( "cdc_management_mode": management_mode, "cdc_slot_name": slot_name, "cdc_publication_name": pub_name, - # Written with the slot, before capture first runs, so capture never treats the new source - # as one the retired legacy lane still delivered for. + # A rollback needs it, as in `recreate_slot`. "cdc_ingest_mode": "buffered", } diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/config.py b/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/config.py index e2f186e3beec..fef28b5ac11b 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/config.py +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/config.py @@ -18,7 +18,6 @@ CDCConfig, ManagementMode, decode_job_inputs, - parse_ingest_mode, ) if TYPE_CHECKING: @@ -44,7 +43,6 @@ def from_dict(cls, job_inputs: Mapping[str, Any] | str | None) -> PostgresCDCCon management_mode: ManagementMode = ( "self_managed" if ji.get("cdc_management_mode") == "self_managed" else "posthog" ) - ingest_mode = parse_ingest_mode(ji) return cls( enabled=str_to_bool(ji.get("cdc_enabled", False)), slot_name=ji.get("cdc_slot_name") or "", @@ -54,7 +52,6 @@ def from_dict(cls, job_inputs: Mapping[str, Any] | str | None) -> PostgresCDCCon lag_critical_threshold_mb=int(ji.get("cdc_lag_critical_threshold_mb", DEFAULT_LAG_CRITICAL_THRESHOLD_MB)), auto_drop_slot=str_to_bool(ji.get("cdc_auto_drop_slot", True)), consistent_point=ji.get("cdc_consistent_point"), - ingest_mode=ingest_mode, ) @classmethod diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/tests/test_cdc_config.py b/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/tests/test_cdc_config.py index bbf7aebca502..4aebcf88664c 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/tests/test_cdc_config.py +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/cdc/tests/test_cdc_config.py @@ -59,8 +59,8 @@ def test_from_dict_defaults(): @pytest.mark.parametrize( "job_inputs", [ - {"cdc_slot_name": "slot_a", "cdc_ingest_mode": "buffered"}, - '{"cdc_slot_name": "slot_a", "cdc_ingest_mode": "buffered"}', + {"cdc_slot_name": "slot_a"}, + '{"cdc_slot_name": "slot_a"}', ], ) def test_from_dict_reads_whichever_shape_job_inputs_decrypted_to(job_inputs): @@ -68,7 +68,6 @@ def test_from_dict_reads_whichever_shape_job_inputs_decrypted_to(job_inputs): # every field here is read off a mapping. config = PostgresCDCConfig.from_dict(job_inputs) assert config.slot_name == "slot_a" - assert config.ingest_mode == "buffered" def test_from_dict_thresholds_coerce_stringified_ints(): diff --git a/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py b/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py index ed1d304e5c63..34a6286987ac 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py +++ b/products/warehouse_sources/backend/temporal/data_imports/sources/postgres/source.py @@ -1865,7 +1865,6 @@ def _buffered_cdc_source(self, schema: "ExternalDataSchema", inputs: SourceInput has_batches_in_flight, served_lanes, ) - from products.warehouse_sources.backend.temporal.data_imports.cdc.types import parse_ingest_mode if not served_lanes(schema): raise ValueError( @@ -1890,13 +1889,6 @@ def no_op_tick() -> SourceResponse: supports_resume=False, ) - if parse_ingest_mode(schema.source.job_inputs) != "buffered": - # Until capture converts this legacy source, its buffer holds copies of changes the legacy - # lane already delivered, which a read would load a second time. Conversion empties the - # buffer before it marks the source buffered. - inputs.logger.info("cdc_buffered_waiting_for_legacy_conversion", schema_name=schema.name) - return no_op_tick() - # Defense in depth for the v3-forcing invariant: a run that resolved its pipeline version # before its table started streaming, or a worker one deploy behind, would consume this # buffer on v2, which stamps no position on the rows it writes, so every later run would diff --git a/products/warehouse_sources/backend/tests/api/test_cdc_table_mode_patch.py b/products/warehouse_sources/backend/tests/api/test_cdc_table_mode_patch.py index e430cfcde7dd..eb8842ab9ad9 100644 --- a/products/warehouse_sources/backend/tests/api/test_cdc_table_mode_patch.py +++ b/products/warehouse_sources/backend/tests/api/test_cdc_table_mode_patch.py @@ -90,7 +90,6 @@ def _make_cdc_source_and_schema( cdc_last_log_position: str | None = "0/12345", cdc_deferred_runs: list[dict] | None = None, initial_sync_complete: bool = True, - ingest_mode: str | None = "buffered", ) -> tuple[ExternalDataSource, ExternalDataSchema]: job_inputs = { "schema": "public", @@ -99,8 +98,6 @@ def _make_cdc_source_and_schema( "cdc_slot_name": "test_slot", "cdc_publication_name": "test_pub", } - if ingest_mode is not None: - job_inputs["cdc_ingest_mode"] = ingest_mode source = ExternalDataSource.objects.create( team=team, source_type=ExternalDataSourceType.POSTGRES, @@ -175,7 +172,7 @@ def test_patch_cdc_table_mode_adding_target_triggers_resnapshot( schema.refresh_from_db() assert schema.cdc_table_mode == new_mode # The table's changes keep going to the buffer, which the new snapshot then replays. - assert (schema.sync_type_config.get("cdc_snapshot_lane") == "buffer") is (deferred_runs is None) + assert schema.sync_type_config.get("cdc_snapshot_lane") == "buffer" assert schema.sync_type_config.get("cdc_mode") == "snapshot" assert schema.sync_type_config.get("cdc_last_log_position") is None assert schema.sync_type_config.get("cdc_deferred_runs") is None @@ -187,7 +184,7 @@ def test_patch_cdc_table_mode_adding_target_triggers_resnapshot( @pytest.mark.parametrize("action", ["resync", "cdc_table_mode_switch", "re_enable"]) def test_a_reset_is_left_to_capture_while_the_tables_sync_can_still_hand_over(team, user, client: HttpClient, action): - source, schema = _make_cdc_source_and_schema(team, cdc_table_mode="consolidated", ingest_mode="buffered") + source, schema = _make_cdc_source_and_schema(team, cdc_table_mode="consolidated") if action == "re_enable": ExternalDataSchema.objects.filter(id=schema.id).update(should_sync=False) running_job = ExternalDataJob.objects.create( @@ -308,7 +305,7 @@ def test_a_hand_over_recreates_a_capture_schedule_that_is_gone(team, user, clien def test_toggling_sync_drops_the_snapshot_marker(team, user, client: HttpClient, should_sync_before, should_sync_after): # Capture skips a table while its sync is off, so its buffer has a gap. A marker left behind would # have the hand-over replay files from before the gap and bring deleted rows back. - _, schema = _make_cdc_source_and_schema(team, cdc_table_mode="consolidated", ingest_mode="buffered") + _, schema = _make_cdc_source_and_schema(team, cdc_table_mode="consolidated") ExternalDataSchema.objects.filter(id=schema.id).update( should_sync=should_sync_before, initial_sync_complete=False, @@ -478,7 +475,7 @@ def test_a_table_stays_in_the_publication_when_its_move_off_cdc_is_not_saved( def test_a_cdc_table_syncs_before_its_captured_changes_expire( team, user, client: HttpClient, sync_frequency, expected_status ): - _, schema = _make_cdc_source_and_schema(team, cdc_table_mode="consolidated", ingest_mode="buffered") + _, schema = _make_cdc_source_and_schema(team, cdc_table_mode="consolidated") client.force_login(user) with ( mock.patch(_PATCH_TARGETS["external_data_workflow_exists"], return_value=True), @@ -495,9 +492,8 @@ def test_a_cdc_table_syncs_before_its_captured_changes_expire( assert "must sync at least weekly" in str(response.json()) -@pytest.mark.parametrize("ingest_mode", [None, "buffered"]) -def test_resync_of_a_streaming_table_keeps_its_buffer_on_a_buffered_source(team, user, client: HttpClient, ingest_mode): - _, schema = _make_cdc_source_and_schema(team, cdc_table_mode="consolidated", ingest_mode=ingest_mode) +def test_resync_of_a_streaming_table_keeps_its_buffer(team, user, client: HttpClient): + _, schema = _make_cdc_source_and_schema(team, cdc_table_mode="consolidated") client.force_login(user) with ( mock.patch(_PATCH_TARGETS["is_any_external_data_schema_paused"], return_value=False), @@ -510,7 +506,7 @@ def test_resync_of_a_streaming_table_keeps_its_buffer_on_a_buffered_source(team, schema.refresh_from_db() assert schema.sync_type_config.get("cdc_mode") == "snapshot" # Unmarked, the next capture run empties the buffer and can delete changes the snapshot never saw. - assert (schema.sync_type_config.get("cdc_snapshot_lane") == "buffer") is (ingest_mode == "buffered") + assert schema.sync_type_config.get("cdc_snapshot_lane") == "buffer" @pytest.mark.parametrize( diff --git a/products/warehouse_sources/backend/tests/test_ad_hoc_sync.py b/products/warehouse_sources/backend/tests/test_ad_hoc_sync.py index ca52f59c940e..9646fd5d77f4 100644 --- a/products/warehouse_sources/backend/tests/test_ad_hoc_sync.py +++ b/products/warehouse_sources/backend/tests/test_ad_hoc_sync.py @@ -101,12 +101,7 @@ def test_failed_start_restores_staged_reset_state(schema, no_sync_to_wait_for): assert schema.initial_sync_complete is True -@pytest.mark.parametrize("ingest_mode", [None, "buffered"]) -def test_a_reset_of_a_streaming_cdc_table_keeps_its_buffer_and_concurrent_keys( - schema, ingest_mode, no_sync_to_wait_for -): - schema.source.job_inputs = {"cdc_ingest_mode": ingest_mode} if ingest_mode else {} - schema.source.save() +def test_a_reset_of_a_streaming_cdc_table_keeps_its_buffer_and_concurrent_keys(schema, no_sync_to_wait_for): schema.sync_type = ExternalDataSchema.SyncType.CDC schema.sync_type_config = {"cdc_mode": "streaming", "cdc_table_mode": "consolidated"} schema.initial_sync_complete = True @@ -123,7 +118,7 @@ def test_a_reset_of_a_streaming_cdc_table_keeps_its_buffer_and_concurrent_keys( schema.refresh_from_db() assert schema.sync_type_config["cdc_mode"] == "snapshot" # Unmarked, the next capture run empties the buffer and can delete changes the snapshot never saw. - assert (schema.sync_type_config.get("cdc_snapshot_lane") == "buffer") is (ingest_mode == "buffered") + assert schema.sync_type_config.get("cdc_snapshot_lane") == "buffer" assert schema.sync_type_config["cdc_last_run_at"] == "2026-09-24T00:00:00+00:00" diff --git a/products/warehouse_sources/backend/tests/test_stalled_schedules.py b/products/warehouse_sources/backend/tests/test_stalled_schedules.py index 168b075db131..0ac7cdc75cb2 100644 --- a/products/warehouse_sources/backend/tests/test_stalled_schedules.py +++ b/products/warehouse_sources/backend/tests/test_stalled_schedules.py @@ -148,7 +148,7 @@ def test_a_schema_with_a_running_job_is_a_wedged_run_not_a_stalled_schedule(self def test_a_streaming_cdc_schema_whose_consumer_stopped_is_reported_and_repairable(self) -> None: # Its schedule consumes the change buffer, so a stall lets the buffer age toward its expiry. - source = self._source(job_inputs={"cdc_ingest_mode": "buffered"}) + source = self._source() schema = self._schema( source, synced_ago=timedelta(days=5),