diff --git a/mkdocs/docs/api.md b/mkdocs/docs/api.md index 1e17e64f42..8afbca02da 100644 --- a/mkdocs/docs/api.md +++ b/mkdocs/docs/api.md @@ -1405,6 +1405,9 @@ tbl.overwrite(df, snapshot_properties={"abc": "def"}) assert tbl.metadata.snapshots[-1].summary["abc"] == "def" ``` +New snapshot summaries automatically include `engine-name` (`pyiceberg`) and `engine-version` +(the installed PyIceberg version). These values override same-named entries in `snapshot_properties`. + ## Snapshot Management Manage snapshots with operations through the `Table` API: diff --git a/pyiceberg/environment_context.py b/pyiceberg/environment_context.py new file mode 100644 index 0000000000..0a4f4eccaf --- /dev/null +++ b/pyiceberg/environment_context.py @@ -0,0 +1,30 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +from pyiceberg import __version__ + + +class EnvironmentContext: + """Environment context carrying the engine name and version for snapshot summaries.""" + + @staticmethod + def get() -> dict[str, str]: + """Return a new dictionary containing only the engine name and version.""" + return { + "engine-name": "pyiceberg", + "engine-version": __version__, + } diff --git a/pyiceberg/table/snapshots.py b/pyiceberg/table/snapshots.py index 5e9e519a01..ee439c45ee 100644 --- a/pyiceberg/table/snapshots.py +++ b/pyiceberg/table/snapshots.py @@ -25,6 +25,7 @@ from pydantic import Field, PrivateAttr, model_serializer +from pyiceberg.environment_context import EnvironmentContext from pyiceberg.io import FileIO from pyiceberg.manifest import DataFile, DataFileContent, ManifestFile, _manifests from pyiceberg.partitioning import UNPARTITIONED_PARTITION_SPEC, PartitionSpec @@ -409,6 +410,11 @@ def _update_totals(total_property: str, added_property: str, removed_property: s removed_property=REMOVED_EQUALITY_DELETES, ) + if context := EnvironmentContext.get(): + # Defensively select only engine fields so future context additions cannot overwrite snapshot metadata. + summary["engine-name"] = context["engine-name"] + summary["engine-version"] = context["engine-version"] + return summary diff --git a/tests/conftest.py b/tests/conftest.py index cca9146c16..c070e81eac 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -48,6 +48,7 @@ from pytest_lazy_fixtures import lf from pyiceberg.catalog import Catalog, load_catalog +from pyiceberg.environment_context import EnvironmentContext from pyiceberg.expressions import BoundReference from pyiceberg.io import ( ADLS_ACCOUNT_KEY, @@ -104,12 +105,29 @@ from pyiceberg.io.pyarrow import PyArrowFileIO +_original_environment_context_get = EnvironmentContext.get + + def pytest_collection_modifyitems(items: list[pytest.Item]) -> None: for item in items: if not any(item.iter_markers()): item.add_marker("unmarked") +@pytest.fixture(autouse=True, scope="session") +def _disable_environment_context() -> Generator[None, None, None]: + """Disable engine metadata for existing tests, including session-scoped fixtures.""" + with pytest.MonkeyPatch.context() as monkeypatch: + monkeypatch.setattr(EnvironmentContext, "get", staticmethod(dict)) + yield + + +@pytest.fixture +def enable_environment_context(monkeypatch: pytest.MonkeyPatch) -> None: + """Restore real engine metadata for tests that explicitly request it.""" + monkeypatch.setattr(EnvironmentContext, "get", staticmethod(_original_environment_context_get)) + + @pytest.fixture(autouse=True, scope="session") def _isolate_pyiceberg_config() -> None: """Make test runs ignore your local PyIceberg config. diff --git a/tests/table/test_snapshots.py b/tests/table/test_snapshots.py index 0f72b08087..294defdcca 100644 --- a/tests/table/test_snapshots.py +++ b/tests/table/test_snapshots.py @@ -24,6 +24,7 @@ import pyarrow as pa import pytest +from pyiceberg import __version__ from pyiceberg.catalog import Catalog from pyiceberg.exceptions import ValidationException from pyiceberg.io.pyarrow import _dataframe_to_data_files @@ -667,6 +668,47 @@ def overwrite_table(catalog: Catalog, arrow_table_simple: pa.Table) -> Table: return table +def test_snapshot_writes_include_engine_metadata( + enable_environment_context: None, catalog: Catalog, arrow_table_simple: pa.Table +) -> None: + catalog.create_namespace("default") + table = catalog.create_table("default.engine_metadata", arrow_table_simple.schema) + table.append(arrow_table_simple) + table.overwrite(arrow_table_simple) + table.delete() + + table.refresh() + snapshots = table.snapshots() + assert snapshots + for snapshot in snapshots: + assert snapshot.summary is not None + assert snapshot.summary["engine-name"] == "pyiceberg" + assert snapshot.summary["engine-version"] == __version__ + + +def test_snapshot_engine_metadata_overrides_snapshot_properties( + enable_environment_context: None, catalog: Catalog, arrow_table_simple: pa.Table +) -> None: + catalog.create_namespace("default") + table = catalog.create_table("default.engine_metadata", arrow_table_simple.schema) + table.append( + arrow_table_simple, + snapshot_properties={ + "engine-name": "custom-engine", + "engine-version": "custom-version", + "job-id": "snapshot-job", + }, + ) + + table.refresh() + snapshot = table.current_snapshot() + assert snapshot is not None + assert snapshot.summary is not None + assert snapshot.summary["engine-name"] == "pyiceberg" + assert snapshot.summary["engine-version"] == __version__ + assert snapshot.summary["job-id"] == "snapshot-job" + + def _write_data_file(table: Table, rows: pa.Table) -> DataFile: return next( iter( diff --git a/tests/test_environment_context.py b/tests/test_environment_context.py new file mode 100644 index 0000000000..09eea1c86d --- /dev/null +++ b/tests/test_environment_context.py @@ -0,0 +1,32 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +from pyiceberg import __version__ +from pyiceberg.environment_context import EnvironmentContext + + +def test_get_returns_fresh_engine_metadata(enable_environment_context: None) -> None: + first = EnvironmentContext.get() + second = EnvironmentContext.get() + assert first is not second + + first.clear() + + assert second == { + "engine-name": "pyiceberg", + "engine-version": __version__, + } + assert EnvironmentContext.get() == second