Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 12 additions & 2 deletions openapi-spec/openapi.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": [
{
Expand Down Expand Up @@ -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": {
Expand Down Expand Up @@ -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"
Expand All @@ -38113,6 +38122,7 @@
"current",
"selected",
"versions",
"current_table_topic_type",
"version"
],
"title": "SourceVersionGetDetailResponse",
Expand Down
15 changes: 10 additions & 5 deletions src/dfe_engine/api/v1/sources.py
Original file line number Diff line number Diff line change
Expand Up @@ -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),
)

Expand Down Expand Up @@ -1832,10 +1833,14 @@ 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``.

``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.
Expand Down Expand Up @@ -2318,7 +2323,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 []),
)

Expand Down
63 changes: 56 additions & 7 deletions src/dfe_engine/source/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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.
Expand All @@ -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]
Expand Down Expand Up @@ -1156,11 +1171,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
Expand Down Expand Up @@ -1425,12 +1465,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
Expand Down Expand Up @@ -1515,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``"
)
Expand Down Expand Up @@ -1556,8 +1604,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(
Expand Down
4 changes: 4 additions & 0 deletions src/dfe_engine/source/registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
Loading
Loading