|
9 | 9 | import subprocess |
10 | 10 |
|
11 | 11 | from sqlmesh.core import dialect as d |
| 12 | +from sqlmesh.core import selector as selector_module |
12 | 13 | from sqlmesh.core.audit import StandaloneAudit |
13 | 14 | from sqlmesh.core.environment import Environment |
14 | | -from sqlmesh.core.model import Model, SqlModel |
| 15 | +from sqlmesh.core.model import ExternalModel, Model, SqlModel, create_external_model |
15 | 16 | from sqlmesh.core.model.common import ParsableSql |
16 | 17 | from sqlmesh.core.selector import NativeSelector |
17 | | -from sqlmesh.core.snapshot import SnapshotChangeCategory |
| 18 | +from sqlmesh.core.snapshot import Snapshot, SnapshotChangeCategory |
| 19 | +from sqlmesh.core.snapshot.cache import SnapshotCache |
| 20 | +from sqlmesh.core.snapshot.definition import Node |
18 | 21 | from sqlmesh.utils import UniqueKeyDict |
19 | 22 | from sqlmesh.utils.date import now_timestamp |
20 | 23 | from sqlmesh.utils.git import GitClient |
21 | 24 |
|
22 | 25 |
|
| 26 | +def _external_schema_refresh_case( |
| 27 | + mocker: MockerFixture, tmp_path: Path |
| 28 | +) -> t.Tuple[ |
| 29 | + ExternalModel, |
| 30 | + ExternalModel, |
| 31 | + SqlModel, |
| 32 | + t.Dict[str, Snapshot], |
| 33 | + UniqueKeyDict[str, Model], |
| 34 | + NativeSelector, |
| 35 | +]: |
| 36 | + external_a = create_external_model("db.external_a", columns={"id": "int"}) |
| 37 | + external_b = create_external_model("db.external_b", columns={"id": "int"}) |
| 38 | + child = SqlModel( |
| 39 | + name="db.child", |
| 40 | + query=d.parse_one( |
| 41 | + "SELECT a.id FROM db.external_a AS a JOIN db.external_b AS b ON a.id = b.id" |
| 42 | + ), |
| 43 | + ) |
| 44 | + |
| 45 | + # Inject inconsistent persisted state for this invariant test: the stored external model |
| 46 | + # has optimize_query=False, but its fingerprint is seeded from the model with None. This |
| 47 | + # synthetic mismatch does not imply that an upstream migration produces this state. |
| 48 | + stored_b = t.cast(ExternalModel, external_b.copy(update={"optimize_query": False})) |
| 49 | + stored_b._data_hash = external_b.data_hash |
| 50 | + stored_b._metadata_hash = external_b.metadata_hash |
| 51 | + |
| 52 | + old_nodes: t.Dict[str, Node] = { |
| 53 | + external_a.fqn: external_a, |
| 54 | + stored_b.fqn: stored_b, |
| 55 | + child.fqn: child, |
| 56 | + } |
| 57 | + old_snapshots = { |
| 58 | + name: Snapshot.from_node(node, nodes=old_nodes) for name, node in old_nodes.items() |
| 59 | + } |
| 60 | + for snapshot in old_snapshots.values(): |
| 61 | + snapshot.categorize_as(SnapshotChangeCategory.BREAKING) |
| 62 | + |
| 63 | + current_models: UniqueKeyDict[str, Model] = UniqueKeyDict("models") |
| 64 | + current_models[external_a.fqn] = external_a |
| 65 | + current_models[external_b.fqn] = external_b |
| 66 | + current_models[child.fqn] = child |
| 67 | + |
| 68 | + state_reader = mocker.Mock() |
| 69 | + state_reader.get_environment.return_value = Environment( |
| 70 | + name="prod", |
| 71 | + snapshots=[snapshot.table_info for snapshot in old_snapshots.values()], |
| 72 | + start_at="2023-01-01", |
| 73 | + end_at="2023-01-02", |
| 74 | + plan_id="test_plan", |
| 75 | + ) |
| 76 | + |
| 77 | + snapshot_cache = SnapshotCache(tmp_path / "snapshot_cache") |
| 78 | + |
| 79 | + def load_snapshots(infos): |
| 80 | + snapshot_ids = {info.snapshot_id for info in infos} |
| 81 | + snapshots, _ = snapshot_cache.get_or_load( |
| 82 | + snapshot_ids, |
| 83 | + lambda ids: [ |
| 84 | + snapshot.copy(deep=True) |
| 85 | + for snapshot in old_snapshots.values() |
| 86 | + if snapshot.snapshot_id in ids |
| 87 | + ], |
| 88 | + ) |
| 89 | + return snapshots |
| 90 | + |
| 91 | + state_reader.get_snapshots.side_effect = load_snapshots |
| 92 | + |
| 93 | + selector = NativeSelector( |
| 94 | + state_reader, |
| 95 | + current_models, |
| 96 | + context_path=tmp_path, |
| 97 | + cache_dir=tmp_path, |
| 98 | + ) |
| 99 | + return external_a, external_b, child, old_snapshots, current_models, selector |
| 100 | + |
| 101 | + |
| 102 | +def _current_snapshots(models: UniqueKeyDict[str, Model]) -> t.Dict[str, Snapshot]: |
| 103 | + nodes = t.cast(t.Dict[str, Node], dict(models)) |
| 104 | + return {model.fqn: Snapshot.from_node(model, nodes=nodes) for model in models.values()} |
| 105 | + |
| 106 | + |
| 107 | +def test_unselected_external_coparent_is_not_directly_modified_by_schema_refresh( |
| 108 | + mocker: MockerFixture, tmp_path: Path |
| 109 | +) -> None: |
| 110 | + mocker.patch("sqlmesh.core.constants.MAX_FORK_WORKERS", 2) |
| 111 | + external_a, external_b, child, old_snapshots, _, selector = _external_schema_refresh_case( |
| 112 | + mocker, tmp_path |
| 113 | + ) |
| 114 | + |
| 115 | + schema_refresh = mocker.spy(selector_module, "update_model_schemas") |
| 116 | + models, selected_fqns = selector.select_models([external_a.fqn], "prod") |
| 117 | + current = _current_snapshots(models) |
| 118 | + |
| 119 | + assert external_a.fqn in selected_fqns |
| 120 | + schema_refresh.assert_called_once() |
| 121 | + assert not current[external_b.fqn].is_directly_modified(old_snapshots[external_b.fqn]) |
| 122 | + assert ( |
| 123 | + current[external_b.fqn].fingerprint.metadata_hash |
| 124 | + == old_snapshots[external_b.fqn].fingerprint.metadata_hash |
| 125 | + ) |
| 126 | + assert not current[child.fqn].is_indirectly_modified(old_snapshots[child.fqn]) |
| 127 | + |
| 128 | + |
| 129 | +def test_selected_external_change_remains_directly_modified( |
| 130 | + mocker: MockerFixture, tmp_path: Path |
| 131 | +) -> None: |
| 132 | + mocker.patch("sqlmesh.core.constants.MAX_FORK_WORKERS", 2) |
| 133 | + _, external_b, _, old_snapshots, current_models, selector = _external_schema_refresh_case( |
| 134 | + mocker, tmp_path |
| 135 | + ) |
| 136 | + changed_b = create_external_model("db.external_b", columns={"id": "int", "new": "text"}) |
| 137 | + current_models.update({changed_b.fqn: changed_b}) |
| 138 | + |
| 139 | + models, selected_fqns = selector.select_models([external_b.fqn], "prod") |
| 140 | + current = _current_snapshots(models) |
| 141 | + |
| 142 | + assert external_b.fqn in selected_fqns |
| 143 | + assert current[external_b.fqn].is_directly_modified(old_snapshots[external_b.fqn]) |
| 144 | + |
| 145 | + |
| 146 | +def test_selected_dependency_schema_change_remains_visible( |
| 147 | + mocker: MockerFixture, tmp_path: Path |
| 148 | +) -> None: |
| 149 | + mocker.patch("sqlmesh.core.constants.MAX_FORK_WORKERS", 2) |
| 150 | + external_a, external_b, child, old_snapshots, current_models, selector = ( |
| 151 | + _external_schema_refresh_case(mocker, tmp_path) |
| 152 | + ) |
| 153 | + changed_a = create_external_model("db.external_a", columns={"id": "int", "new": "text"}) |
| 154 | + current_models.update({changed_a.fqn: changed_a}) |
| 155 | + |
| 156 | + models, selected_fqns = selector.select_models([external_a.fqn], "prod") |
| 157 | + current = _current_snapshots(models) |
| 158 | + |
| 159 | + assert external_a.fqn in selected_fqns |
| 160 | + assert current[external_a.fqn].is_directly_modified(old_snapshots[external_a.fqn]) |
| 161 | + assert not current[external_b.fqn].is_directly_modified(old_snapshots[external_b.fqn]) |
| 162 | + assert current[child.fqn].is_indirectly_modified(old_snapshots[child.fqn]) |
| 163 | + |
| 164 | + |
23 | 165 | @pytest.mark.parametrize( |
24 | 166 | "default_catalog", |
25 | 167 | [ |
|
0 commit comments