From 600e4b9650c3d7da48ca90b79bd53e8a98763f7a Mon Sep 17 00:00:00 2001 From: kbrandt-hs Date: Thu, 8 Oct 2026 09:42:24 +0200 Subject: [PATCH 1/5] fix: fix current_table_topic_type attribution - should always look at current --- src/dfe_engine/api/v1/sources.py | 2 +- src/dfe_engine/source/models.py | 14 ++++---- tests/unit/test_api/test_sources.py | 49 +++++++++++++++++++++++++++ tests/unit/test_source/test_models.py | 32 +++++++++++++++++ 4 files changed, 90 insertions(+), 7 deletions(-) diff --git a/src/dfe_engine/api/v1/sources.py b/src/dfe_engine/api/v1/sources.py index 018fa91d..6f98bf13 100644 --- a/src/dfe_engine/api/v1/sources.py +++ b/src/dfe_engine/api/v1/sources.py @@ -2318,7 +2318,7 @@ def _to_summary(raw: dict[str, Any]) -> SourceSummaryObject: has_transform=bool(raw.get("has_transform")), has_fetcher=bool(raw.get("has_fetcher")), origin=raw.get("origin"), - current_table_topic_type=raw.get("current_table_topic_type", "own"), + current_table_topic_type=raw.get("current_table_topic_type", "main"), views=list(raw.get("views") or []), ) diff --git a/src/dfe_engine/source/models.py b/src/dfe_engine/source/models.py index 7d18102c..e07419df 100644 --- a/src/dfe_engine/source/models.py +++ b/src/dfe_engine/source/models.py @@ -1425,12 +1425,13 @@ def landing_label(self) -> str: @property def table_topic_type(self) -> FetcherTopic: - """Whether records land on this source's own table or the shared main one. + """Where the working ``current`` version's records land: own table or main. - A pinned ``meta_schema`` means the source has its own ClickHouse table. - Without one, records share the platform landing table (``main``). + A pinned ``meta_schema`` on ``current`` means that version has its own + ClickHouse table. Without one, it shares the platform landing table + (``main``) -- even when a still-deployed older version pinned a schema. """ - return "own" if self.schema_config.meta_schema else "main" + return "own" if self.version(self.current).effective_schema().meta_schema else "main" # ----------------------------------------------------------------- # Derived properties @@ -1556,8 +1557,9 @@ class SourceSummaryObject(BaseModel): ) current_table_topic_type: FetcherTopic = Field( description=( - "Where records land: ``own`` when a meta schema pins a per-source table, " - "``main`` when there is none and records share the landing table" + "Where the working ``current`` version lands: ``own`` when that version " + "pins a meta schema, ``main`` when it does not (even if a deployed older " + "version still has its own table)" ), ) views: list[str] = Field( diff --git a/tests/unit/test_api/test_sources.py b/tests/unit/test_api/test_sources.py index ba493f97..cf395d16 100644 --- a/tests/unit/test_api/test_sources.py +++ b/tests/unit/test_api/test_sources.py @@ -68,6 +68,55 @@ def test_list_reports_own_when_a_meta_schema_is_pinned( by_name = {item["name"]: item for item in resp.json()["items"]} assert by_name["aws-own"]["current_table_topic_type"] == "own" + def test_list_table_topic_type_tracks_the_working_current_version( + self, client: TestClient, admin_headers: dict, sample_source: dict + ): + # Deployed 1.0.0 owns a table; working current 2.0.0 dropped the meta schema. + client.post( + "/api/v1/sources", + json={ + **sample_source, + "source": "azure-drift", + "schema_config": { + "meta_schema": "meta/azure", + "meta_schema_version": "1.0.0", + }, + }, + headers=admin_headers, + ) + registry = _registries["source"] + registry.set_deployed_version("azure-drift", "1.0.0") + registry.save_source( + { + "source": "azure-drift", + "deployed_version": "1.0.0", + "current": "2.0.0", + "versions": { + "1.0.0": { + "date_time": "2026-01-01", + "match": sample_source["match"], + "schema": { + "meta_schema": "meta/azure", + "meta_schema_version": "1.0.0", + }, + }, + "2.0.0": { + "date_time": "2026-01-02", + "match": sample_source["match"], + "schema": {}, + }, + }, + }, + created_by="test", + description="e2e: current without meta schema", + ) + + resp = client.get("/api/v1/sources", headers=admin_headers) + by_name = {item["name"]: item for item in resp.json()["items"]} + assert by_name["azure-drift"]["current"] == "2.0.0" + assert by_name["azure-drift"]["deployed_version"] == "1.0.0" + assert by_name["azure-drift"]["current_table_topic_type"] == "main" + def test_list_pagination(self, client: TestClient, admin_headers: dict): # Create 5 sources for i in range(5): diff --git a/tests/unit/test_source/test_models.py b/tests/unit/test_source/test_models.py index 9b7506b0..e0466a99 100644 --- a/tests/unit/test_source/test_models.py +++ b/tests/unit/test_source/test_models.py @@ -412,6 +412,38 @@ def test_the_landing_source_is_the_main_table(self): s = Source.model_validate({"source": "main", "resource_type": "core"}) assert s.table_topic_type == "main" + def test_table_topic_type_follows_the_current_version_not_the_deployed_one(self): + """An undeployed current without a meta schema shares main, even if deployed owns a table.""" + s = Source.model_validate( + { + "source": "azure", + "deployed_version": "1.0.0", + "current": "2.0.0", + "versions": { + "1.0.0": { + "date_time": "2026-01-01", + "match": {"field": "azure", "value": "azure"}, + "schema": { + "meta_schema": "meta/azure", + "meta_schema_version": "1.0.0", + }, + }, + "2.0.0": { + "date_time": "2026-01-02", + "match": {"field": "azure", "value": "azure"}, + "schema": {}, + }, + }, + } + ) + assert s.version("1.0.0").effective_schema().meta_schema == "meta/azure" + assert s.version("2.0.0").effective_schema().meta_schema is None + assert s.table_topic_type == "main" + + # Flip current back to the deployed pin: own table again. + s = s.model_copy(update={"current": "1.0.0"}) + assert s.table_topic_type == "own" + def test_the_write_body_needs_exactly_one_origin(self): with pytest.raises(ValueError, match="exactly one"): SourceWriteRequest(source="x") From 2f28319f985a8f89d192e083679c13d78054667e Mon Sep 17 00:00:00 2001 From: kbrandt-hs Date: Thu, 8 Oct 2026 09:54:06 +0200 Subject: [PATCH 2/5] fix: add validation to source update - don't allow removing meta schema and don't allow changing from fetcher <-> receiver --- src/dfe_engine/api/v1/sources.py | 3 ++ src/dfe_engine/source/models.py | 25 ++++++++++ tests/unit/test_api/test_sources.py | 37 +++++++++++++++ tests/unit/test_source/test_models.py | 66 +++++++++++++++++++++++++++ 4 files changed, 131 insertions(+) diff --git a/src/dfe_engine/api/v1/sources.py b/src/dfe_engine/api/v1/sources.py index 6f98bf13..fd5bd04a 100644 --- a/src/dfe_engine/api/v1/sources.py +++ b/src/dfe_engine/api/v1/sources.py @@ -1837,6 +1837,9 @@ def update_source( ``views``, or ``transform`` change. Draft versions (``current`` not deployed) update in place. + Origin (receiver vs fetcher) cannot change, and a pinned ``meta_schema`` cannot be + cleared: once a source owns its own table it cannot move back onto ``main``. + ``header`` is optional: when omitted, no header is stored on the written version snapshot (same as create). Send ``header`` explicitly to set or change it. diff --git a/src/dfe_engine/source/models.py b/src/dfe_engine/source/models.py index e07419df..62585997 100644 --- a/src/dfe_engine/source/models.py +++ b/src/dfe_engine/source/models.py @@ -1156,11 +1156,36 @@ def draft_build_version_to_invalidate(existing: Source, updated: Source) -> str return None +def _reject_origin_or_table_paradigm_change( + previous: SourceVersion, updated: SourceVersion +) -> None: + """Refuse an update that swaps intake origin or drops back onto the shared table. + + Origin (receiver vs fetcher) is fixed for the life of a source: the apps and + routing compiled around it assume one intake. A pinned ``meta_schema`` means + the source owns its own ClickHouse table; clearing it would move records onto + ``main``, which an update is not allowed to do. + """ + if previous.origin != updated.origin: + raise ValueError( + f"Cannot change source origin from {previous.origin!r} to {updated.origin!r}: " + "a source is receiver-based or fetcher-based for its whole life" + ) + previous_owns_table = bool(previous.effective_schema().meta_schema) + updated_owns_table = bool(updated.effective_schema().meta_schema) + if previous_owns_table and not updated_owns_table: + raise ValueError( + "Cannot remove the meta schema: once a source owns its own table it " + "cannot move back to the shared main table" + ) + + def apply_source_write_update(existing: Source, write: SourceWriteRequest) -> Source: """Persist a write: bump major version only when ``current`` is deployed and pins/mappings change.""" snapshot = _build_merged_version_snapshot(existing, write) current_id = existing.current previous = existing.versions[current_id] + _reject_origin_or_table_paradigm_change(previous, snapshot) append_version = ( existing.deployed_version is not None diff --git a/tests/unit/test_api/test_sources.py b/tests/unit/test_api/test_sources.py index cf395d16..40872767 100644 --- a/tests/unit/test_api/test_sources.py +++ b/tests/unit/test_api/test_sources.py @@ -772,6 +772,43 @@ def test_update_success(self, client: TestClient, admin_headers: dict, sample_so assert body["current"] == "1.0.0" assert body["versions"]["1.0.0"]["schema"]["engine"] == "MergeTree" + def test_update_refuses_to_swap_origin( + self, client: TestClient, admin_headers: dict, sample_source: dict + ): + client.post("/api/v1/sources", json=sample_source, headers=admin_headers) + resp = client.put( + "/api/v1/sources/test-source", + json={ + "source": "test-source", + "fetcher": {"source_type": "okta", "topic": "own"}, + }, + headers=admin_headers, + ) + assert resp.status_code == 422, resp.text + assert "Cannot change source origin" in resp.json()["message"] + + def test_update_refuses_to_clear_a_meta_schema( + self, client: TestClient, admin_headers: dict, sample_source: dict + ): + client.post( + "/api/v1/sources", + json={ + **sample_source, + "schema_config": { + "meta_schema": "meta/aws_cloudtrail", + "meta_schema_version": "1.0.0", + }, + }, + headers=admin_headers, + ) + resp = client.put( + "/api/v1/sources/test-source", + json={**sample_source, "schema_config": {"engine": "MergeTree"}}, + headers=admin_headers, + ) + assert resp.status_code == 422, resp.text + assert "Cannot remove the meta schema" in resp.json()["message"] + def test_update_rejects_versions_payload( self, client: TestClient, admin_headers: dict, sample_source: dict ): diff --git a/tests/unit/test_source/test_models.py b/tests/unit/test_source/test_models.py index e0466a99..0859a96b 100644 --- a/tests/unit/test_source/test_models.py +++ b/tests/unit/test_source/test_models.py @@ -450,6 +450,72 @@ def test_the_write_body_needs_exactly_one_origin(self): body = SourceWriteRequest(source="x", fetcher=SourceFetcher(source_type="pypi")) assert body.to_version_snapshot().origin == "fetcher" + def test_an_update_cannot_swap_receiver_for_fetcher(self): + existing = Source.model_validate( + {"source": "syslog", "match": {"field": "f", "value": "v"}} + ) + write = SourceWriteRequest.model_validate( + {"fetcher": {"source_type": "okta", "topic": "own"}} + ) + with pytest.raises(ValueError, match="Cannot change source origin"): + apply_source_write_update(existing, write) + + def test_an_update_cannot_swap_fetcher_for_receiver(self): + existing = Source.model_validate( + {"source": "okta", "fetcher": {"source_type": "okta", "topic": "own"}} + ) + write = SourceWriteRequest.model_validate({"match": {"field": "f", "value": "v"}}) + with pytest.raises(ValueError, match="Cannot change source origin"): + apply_source_write_update(existing, write) + + def test_an_update_cannot_clear_a_pinned_meta_schema(self): + existing = Source.model_validate( + { + "source": "aws", + "match": {"field": "f", "value": "v"}, + "schema": {"meta_schema": "meta/aws", "meta_schema_version": "1.0.0"}, + } + ) + write = SourceWriteRequest.model_validate( + { + "match": {"field": "f", "value": "v"}, + "schema": {"engine": "MergeTree"}, + } + ) + with pytest.raises(ValueError, match="Cannot remove the meta schema"): + apply_source_write_update(existing, write) + + def test_an_update_may_still_change_which_meta_schema_is_pinned(self): + existing = Source.model_validate( + { + "source": "aws", + "match": {"field": "f", "value": "v"}, + "schema": {"meta_schema": "meta/aws", "meta_schema_version": "1.0.0"}, + } + ) + write = SourceWriteRequest.model_validate( + { + "match": {"field": "f", "value": "v"}, + "schema": {"meta_schema": "meta/aws_cloudtrail", "meta_schema_version": "2.0.0"}, + } + ) + updated = apply_source_write_update(existing, write) + assert updated.versions["1.0.0"].schema_config.meta_schema == "meta/aws_cloudtrail" + assert updated.table_topic_type == "own" + + def test_a_source_without_a_meta_schema_may_still_gain_one(self): + existing = Source.model_validate( + {"source": "azure", "match": {"field": "azure", "value": "azure"}, "schema": {}} + ) + write = SourceWriteRequest.model_validate( + { + "match": {"field": "azure", "value": "azure"}, + "schema": {"meta_schema": "meta/azure", "meta_schema_version": "1.0.0"}, + } + ) + updated = apply_source_write_update(existing, write) + assert updated.table_topic_type == "own" + # --------------------------------------------------------------------------- # SourceView From 312d5262bdf15e6b14553bf66b49fa341021dffd Mon Sep 17 00:00:00 2001 From: kbrandt-hs Date: Thu, 8 Oct 2026 10:01:33 +0200 Subject: [PATCH 3/5] fix: update criteria for what causes a source version bump --- src/dfe_engine/api/v1/sources.py | 9 +- src/dfe_engine/source/models.py | 17 +++- tests/unit/test_source/test_models.py | 124 +++++++++++++++++++++++++- 3 files changed, 143 insertions(+), 7 deletions(-) diff --git a/src/dfe_engine/api/v1/sources.py b/src/dfe_engine/api/v1/sources.py index fd5bd04a..a1ccd794 100644 --- a/src/dfe_engine/api/v1/sources.py +++ b/src/dfe_engine/api/v1/sources.py @@ -1832,10 +1832,11 @@ def update_source( """Update a source from a flat revision body. Before the first deploy, edits update the working version in place. After deploy, a new - major version is created only when ``current`` equals ``deployed_version`` and schema pins - (``meta_schema``, ``meta_schema_version``, ``derived_schema``, ``additional_fields``), - ``views``, or ``transform`` change. Draft versions (``current`` not deployed) update - in place. + major version is created only when ``current`` equals ``deployed_version`` and a + table-shaping field changes: schema pins (``meta_schema``, ``meta_schema_version``, + ``derived_schema``, ``additional_fields``), ``ttl_days``, ``engine``, common-header + type/version, ``views``, or ``transform``. Draft versions (``current`` not deployed) + update in place. Origin (receiver vs fetcher) cannot change, and a pinned ``meta_schema`` cannot be cleared: once a source owns its own table it cannot move back onto ``main``. diff --git a/src/dfe_engine/source/models.py b/src/dfe_engine/source/models.py index 62585997..4bccd144 100644 --- a/src/dfe_engine/source/models.py +++ b/src/dfe_engine/source/models.py @@ -1086,16 +1086,25 @@ def _build_merged_version_snapshot(existing: Source, write: SourceWriteRequest) def _schema_pin_for_bump(schema: SourceSchema) -> tuple[Any, ...]: - """Schema references that change composed columns / deploy DDL.""" + """Schema fields that change the table shape or retention / engine DDL.""" return ( schema.meta_schema, schema.meta_schema_version, schema.derived_schema, schema.derived_schema_version, schema.additional_fields, + schema.ttl_days, + schema.engine, ) +def _header_pin_for_bump(header: SourceHeader | None) -> tuple[str | None, str | None]: + """Common-header type and version; absent header is a distinct pin from a defaulted one.""" + if header is None: + return (None, None) + return (header.type, header.version) + + def _transform_pin_for_bump(transform: SourceTransform | None) -> dict[str, Any] | None: """The transform fields that change the table, so a deployed edit needs a version.""" if transform is None: @@ -1106,6 +1115,10 @@ def _transform_pin_for_bump(transform: SourceTransform | None) -> dict[str, Any] def source_version_bump_required(previous: SourceVersion, updated: SourceVersion) -> bool: """True when a deployed source needs a new major version id for this snapshot change. + Bumps when the composed table would change: schema pins (meta / derived / + additional), ``ttl_days``, ``engine``, common-header type/version, ``views``, + or table-affecting ``transform`` fields. + ``transport``, ``archive``, ``transform.variant`` and ``fetcher.routes`` are deliberately absent: they move records around, and a version id exists to pin the shape of the table those records land in. @@ -1114,6 +1127,8 @@ def source_version_bump_required(previous: SourceVersion, updated: SourceVersion new_schema = updated.effective_schema() if _schema_pin_for_bump(prev_schema) != _schema_pin_for_bump(new_schema): return True + if _header_pin_for_bump(previous.header) != _header_pin_for_bump(updated.header): + return True # Any view change bumps: the standard set, a field_map pin, per-view # custom_mappings, or the sigma view's taxonomy/category/service. prev_views = [v.model_dump(mode="json") for v in previous.views] diff --git a/tests/unit/test_source/test_models.py b/tests/unit/test_source/test_models.py index 0859a96b..d9e5b1d6 100644 --- a/tests/unit/test_source/test_models.py +++ b/tests/unit/test_source/test_models.py @@ -1129,17 +1129,137 @@ def test_apply_write_update_in_place_after_deploy_non_bump_fields(self): }, } ) + # Match alone does not shape the table, so the deployed version stays current. write = SourceWriteRequest.model_validate( { "match": {"field": "g", "value": "w"}, - "schema": {"ttl_days": 30, "meta_schema_version": "1.0.0"}, + "schema": {"ttl_days": 90, "meta_schema_version": "1.0.0"}, } ) updated = apply_source_write_update(existing, write) assert updated.versions.keys() == {"1.0.0"} - assert updated.versions["1.0.0"].schema_config.ttl_days == 30 + assert updated.versions["1.0.0"].schema_config.ttl_days == 90 assert updated.versions["1.0.0"].match.field == "g" + def test_apply_write_update_bumps_on_ttl_days_after_deploy(self): + existing = Source.model_validate( + { + "source": "src-a", + "deployed_version": "1.0.0", + "current": "1.0.0", + "versions": { + "1.0.0": { + "date_time": "2026-01-01", + "match": {"field": "f", "value": "v"}, + "schema": {"ttl_days": 90}, + } + }, + } + ) + write = SourceWriteRequest.model_validate( + { + "match": {"field": "f", "value": "v"}, + "schema": {"ttl_days": 30}, + } + ) + updated = apply_source_write_update(existing, write) + assert updated.current == "2.0.0" + assert updated.versions["1.0.0"].schema_config.ttl_days == 90 + assert updated.versions["2.0.0"].schema_config.ttl_days == 30 + + def test_apply_write_update_bumps_on_engine_after_deploy(self): + existing = Source.model_validate( + { + "source": "src-a", + "deployed_version": "1.0.0", + "current": "1.0.0", + "versions": { + "1.0.0": { + "date_time": "2026-01-01", + "match": {"field": "f", "value": "v"}, + "schema": {"engine": "MergeTree"}, + } + }, + } + ) + write = SourceWriteRequest.model_validate( + { + "match": {"field": "f", "value": "v"}, + "schema": {"engine": "ReplacingMergeTree"}, + } + ) + updated = apply_source_write_update(existing, write) + assert updated.current == "2.0.0" + assert updated.versions["2.0.0"].schema_config.engine == "ReplacingMergeTree" + + def test_apply_write_update_bumps_on_header_after_deploy(self): + existing = Source.model_validate( + { + "source": "src-a", + "deployed_version": "1.0.0", + "current": "1.0.0", + "versions": { + "1.0.0": { + "date_time": "2026-01-01", + "header": {"type": "timeseries", "version": "1.0.0"}, + "match": {"field": "f", "value": "v"}, + "schema": {}, + } + }, + } + ) + write = SourceWriteRequest.model_validate( + { + "header": {"type": "timeseries", "version": "1.0.1"}, + "match": {"field": "f", "value": "v"}, + "schema": {}, + } + ) + updated = apply_source_write_update(existing, write) + assert updated.current == "2.0.0" + assert updated.versions["2.0.0"].header is not None + assert updated.versions["2.0.0"].header.version == "1.0.1" + + def test_source_version_bump_required_ttl_engine_and_header(self): + base = { + "date_time": "2026-01-01", + "header": {"type": "timeseries", "version": "1.0.0"}, + "match": {"field": "f", "value": "v"}, + "schema": {"ttl_days": 90, "engine": "MergeTree"}, + } + prev = SourceVersion.model_validate(base) + assert ( + source_version_bump_required( + prev, + SourceVersion.model_validate( + {**base, "schema": {"ttl_days": 30, "engine": "MergeTree"}} + ), + ) + is True + ) + assert ( + source_version_bump_required( + prev, + SourceVersion.model_validate( + {**base, "schema": {"ttl_days": 90, "engine": "ReplacingMergeTree"}} + ), + ) + is True + ) + assert ( + source_version_bump_required( + prev, + SourceVersion.model_validate( + { + **base, + "header": {"type": "minimal", "version": "1.0.0"}, + } + ), + ) + is True + ) + assert source_version_bump_required(prev, SourceVersion.model_validate(base)) is False + def test_apply_write_update_bumps_on_meta_schema_path_after_deploy(self): existing = Source.model_validate( { From 73f7c638c34f4880ea391c8f9b3a837cb9af4cc6 Mon Sep 17 00:00:00 2001 From: kbrandt-hs Date: Thu, 8 Oct 2026 11:37:42 +0200 Subject: [PATCH 4/5] fix: add current_table_topic_type to source version detail --- openapi-spec/openapi.json | 14 ++++- src/dfe_engine/api/v1/sources.py | 1 + src/dfe_engine/source/models.py | 7 +++ tests/unit/test_api/test_sources.py | 74 +++++++++++++++++++++++++++ tests/unit/test_source/test_models.py | 2 + 5 files changed, 96 insertions(+), 2 deletions(-) diff --git a/openapi-spec/openapi.json b/openapi-spec/openapi.json index 0730acdb..56b4e395 100644 --- a/openapi-spec/openapi.json +++ b/openapi-spec/openapi.json @@ -5287,7 +5287,7 @@ "Sources" ], "summary": "Update Source", - "description": "Update a source from a flat revision body.\n\nBefore the first deploy, edits update the working version in place. After deploy, a new\nmajor version is created only when ``current`` equals ``deployed_version`` and schema pins\n(``meta_schema``, ``meta_schema_version``, ``derived_schema``, ``additional_fields``),\n``views``, or ``transform`` change. Draft versions (``current`` not deployed) update\nin place.\n\n``header`` is optional: when omitted, no header is stored on the written version snapshot\n(same as create). Send ``header`` explicitly to set or change it.\n\nA schema pin naming nothing is accepted here and refused by ``/sources/import``.\nSee ``create_source`` and #504; the asymmetry is known and undecided.", + "description": "Update a source from a flat revision body.\n\nBefore the first deploy, edits update the working version in place. After deploy, a new\nmajor version is created only when ``current`` equals ``deployed_version`` and a\ntable-shaping field changes: schema pins (``meta_schema``, ``meta_schema_version``,\n``derived_schema``, ``additional_fields``), ``ttl_days``, ``engine``, common-header\ntype/version, ``views``, or ``transform``. Draft versions (``current`` not deployed)\nupdate in place.\n\nOrigin (receiver vs fetcher) cannot change, and a pinned ``meta_schema`` cannot be\ncleared: once a source owns its own table it cannot move back onto ``main``.\n\n``header`` is optional: when omitted, no header is stored on the written version snapshot\n(same as create). Send ``header`` explicitly to set or change it.\n\nA schema pin naming nothing is accepted here and refused by ``/sources/import``.\nSee ``create_source`` and #504; the asymmetry is known and undecided.", "operationId": "update_source_api_v1_sources__name__put", "parameters": [ { @@ -37790,7 +37790,7 @@ "main" ], "title": "Current Table Topic Type", - "description": "Where records land: ``own`` when a meta schema pins a per-source table, ``main`` when there is none and records share the landing table" + "description": "Where the working ``current`` version lands: ``own`` when that version pins a meta schema, ``main`` when it does not (even if a deployed older version still has its own table)" }, "views": { "items": { @@ -38102,6 +38102,15 @@ "title": "Previous Deployed Versions", "description": "Version ids with a successful deploy in source-deploys history, excluding the live deployed_version" }, + "current_table_topic_type": { + "type": "string", + "enum": [ + "own", + "main" + ], + "title": "Current Table Topic Type", + "description": "Where the working ``current`` version lands: ``own`` when that version pins a meta schema, ``main`` when it does not (even if a deployed older version still has its own table)" + }, "version": { "$ref": "#/components/schemas/SourceVersionDetail", "description": "Configuration snapshot for ``selected`` plus pipeline artifacts" @@ -38113,6 +38122,7 @@ "current", "selected", "versions", + "current_table_topic_type", "version" ], "title": "SourceVersionGetDetailResponse", diff --git a/src/dfe_engine/api/v1/sources.py b/src/dfe_engine/api/v1/sources.py index a1ccd794..adfa6eea 100644 --- a/src/dfe_engine/api/v1/sources.py +++ b/src/dfe_engine/api/v1/sources.py @@ -1004,6 +1004,7 @@ def get_source_version( selected=version, versions=version_ids, previous_deployed_versions=prev_deployed, + current_table_topic_type=source.table_topic_type, version=_version_detail_from_snapshot(source.source, version, snap, store), ) diff --git a/src/dfe_engine/source/models.py b/src/dfe_engine/source/models.py index 4bccd144..46f4e4ad 100644 --- a/src/dfe_engine/source/models.py +++ b/src/dfe_engine/source/models.py @@ -1556,6 +1556,13 @@ class SourceVersionGetResponse(BaseModel): "excluding the live deployed_version" ), ) + current_table_topic_type: FetcherTopic = Field( + description=( + "Where the working ``current`` version lands: ``own`` when that version " + "pins a meta schema, ``main`` when it does not (even if a deployed older " + "version still has its own table)" + ), + ) version: SourceVersion = Field( ..., description="Immutable configuration snapshot for ``selected``" ) diff --git a/tests/unit/test_api/test_sources.py b/tests/unit/test_api/test_sources.py index 40872767..e14bd70d 100644 --- a/tests/unit/test_api/test_sources.py +++ b/tests/unit/test_api/test_sources.py @@ -652,11 +652,85 @@ def test_get_version_after_create( assert body["current"] == "1.0.0" assert body["versions"] == ["1.0.0"] assert body["previous_deployed_versions"] == [] + assert body["current_table_topic_type"] == "main" assert body["version"]["schema"]["engine"] == "MergeTree" assert body["version"]["origin"] == "receiver" assert body["version"]["source_build"] is None assert body["version"]["source_deployment"] is None + def test_get_version_reports_own_when_current_pins_meta_schema( + self, client: TestClient, admin_headers: dict, sample_source: dict + ): + client.post( + "/api/v1/sources", + json={ + **sample_source, + "schema_config": { + "meta_schema": "meta/aws_cloudtrail", + "meta_schema_version": "1.0.0", + }, + }, + headers=admin_headers, + ) + resp = client.get( + "/api/v1/sources/test-source/versions/1.0.0", + headers=admin_headers, + ) + assert resp.status_code == 200 + assert resp.json()["current_table_topic_type"] == "own" + + def test_get_version_table_topic_type_tracks_working_current( + self, client: TestClient, admin_headers: dict, sample_source: dict + ): + # Selected 1.0.0 owns a table; working current 2.0.0 does not. + client.post( + "/api/v1/sources", + json={ + **sample_source, + "schema_config": { + "meta_schema": "meta/azure", + "meta_schema_version": "1.0.0", + }, + }, + headers=admin_headers, + ) + registry = _registries["source"] + registry.set_deployed_version("test-source", "1.0.0") + registry.save_source( + { + "source": "test-source", + "deployed_version": "1.0.0", + "current": "2.0.0", + "versions": { + "1.0.0": { + "date_time": "2026-01-01", + "match": sample_source["match"], + "schema": { + "meta_schema": "meta/azure", + "meta_schema_version": "1.0.0", + }, + }, + "2.0.0": { + "date_time": "2026-01-02", + "match": sample_source["match"], + "schema": {}, + }, + }, + }, + created_by="test", + description="current without meta schema", + ) + + resp = client.get( + "/api/v1/sources/test-source/versions/1.0.0", + headers=admin_headers, + ) + assert resp.status_code == 200 + body = resp.json() + assert body["selected"] == "1.0.0" + assert body["current"] == "2.0.0" + assert body["current_table_topic_type"] == "main" + def test_get_version_after_update_preserves_history( self, client: TestClient, admin_headers: dict, sample_source: dict ): diff --git a/tests/unit/test_source/test_models.py b/tests/unit/test_source/test_models.py index d9e5b1d6..2828da24 100644 --- a/tests/unit/test_source/test_models.py +++ b/tests/unit/test_source/test_models.py @@ -1840,6 +1840,7 @@ def test_round_trip_fields(self): deployed_version="1.0.0", selected="1.0.0", versions=["1.0.0", "2.0.0"], + current_table_topic_type="main", version=SourceVersion( date_time="2026-06-12", header=SourceHeader(type="timeseries", version="1.0.0"), @@ -1849,6 +1850,7 @@ def test_round_trip_fields(self): ), ) assert resp.source == "my-source" + assert resp.current_table_topic_type == "main" assert resp.version.schema_config.engine == "MergeTree" assert resp.version.match is not None assert resp.version.match.value == "x" From f3fa14d96bb22d33cfb03fac1df644dba68083e6 Mon Sep 17 00:00:00 2001 From: kbrandt-hs Date: Thu, 8 Oct 2026 12:32:53 +0200 Subject: [PATCH 5/5] fix: fix ci --- src/dfe_engine/source/registry.py | 4 ++++ tests/unit/test_query/test_client.py | 5 ++++- tests/unit/test_schema/test_ddl_writer.py | 5 ++++- 3 files changed, 12 insertions(+), 2 deletions(-) diff --git a/src/dfe_engine/source/registry.py b/src/dfe_engine/source/registry.py index 73188022..48b86549 100644 --- a/src/dfe_engine/source/registry.py +++ b/src/dfe_engine/source/registry.py @@ -504,6 +504,10 @@ def update_source_from_write( raise SourceValidationError(f"Invalid source definition: {e}") from e existing = self.get_source(source_name) + # Refuse engine-owned sources before merge/origin checks so PUT main + # stays a 409 conflict, not a 422 about origin None -> receiver. + if self.is_core(source_name): + raise SourceCoreResourceError(action="modify", source=source_name) try: updated = apply_source_write_update(existing, write) except ValueError as e: diff --git a/tests/unit/test_query/test_client.py b/tests/unit/test_query/test_client.py index 09ead6bf..bd6eabbc 100644 --- a/tests/unit/test_query/test_client.py +++ b/tests/unit/test_query/test_client.py @@ -130,7 +130,10 @@ def test_metadata_from_headers(self): class TestQueryClientDirectMode: - def test_direct_mode_requires_view_executor(self): + def test_direct_mode_requires_view_executor(self, monkeypatch): client = QueryClient(direct=True) + # Hermetic: a local .env that can build a ViewExecutor against a down + # ClickHouse would otherwise raise ServiceUnavailable instead. + monkeypatch.setattr(client, "_get_view_executor", lambda: None) with pytest.raises(RuntimeError, match="ViewExecutor not available"): client.query("analytics/user_activity") diff --git a/tests/unit/test_schema/test_ddl_writer.py b/tests/unit/test_schema/test_ddl_writer.py index 64e150bc..f4a9f3c5 100644 --- a/tests/unit/test_schema/test_ddl_writer.py +++ b/tests/unit/test_schema/test_ddl_writer.py @@ -258,7 +258,10 @@ def test_all_values_are_sql_strings(self, writer): class TestWriteAll: """Test file writing.""" - def test_creates_files(self, writer, tmp_path): + def test_creates_files(self, writer, tmp_path, monkeypatch): + # Pin against the installed dfe-schemas package: a local DFE_SCHEMAS_DIR + # (often from .env) can be an older partial tree and would under-count. + monkeypatch.delenv("DFE_SCHEMAS_DIR", raising=False) written = writer.write_all(tmp_path / "ddl") # Pinned, not derived from generate_all(): deriving both sides would let a # regression that silently drops a profile pass. Every schema version is