diff --git a/src/recotem/_artifact_identity.py b/src/recotem/_artifact_identity.py index 6945665f..a3355909 100644 --- a/src/recotem/_artifact_identity.py +++ b/src/recotem/_artifact_identity.py @@ -137,10 +137,14 @@ def check_artifact_recipe_hash(header_dict: Any, *, recipe: Any, name: str) -> N retrain, restart serve. The old artifact loads, ``/v1/health`` reports ``ok``, and ``/v1/recipes/{name}`` reports the *artifact's* ``algorithms`` / ``cutoff`` / ``metric``, which now contradict the recipe on - disk with nothing marking them as historical. Every block the running - server holds comes from the body it parsed for the model it is serving, so - the whole response is of a piece: it is the *older* recipe throughout, and - the warning is the only thing that says so. + disk with nothing marking them as historical, and this warning is the only + thing that says so. + + While it is firing the response is not of a piece, and the split is worth + knowing: everything read out of the *artifact* -- the model, and the header + fields ``/v1/recipes/{name}`` reports -- is the older recipe, while + ``item_metadata`` is read from the recipe on disk at the moment the model + was loaded. The two agree again as soon as a retrain silences this warning. **Warn, never refuse.** Unlike ``recipe_name``, a hash difference is not a contradiction -- it is the expected state whenever a recipe is edited in a diff --git a/src/recotem/serving/watcher.py b/src/recotem/serving/watcher.py index 02f6284b..1884dd2e 100644 --- a/src/recotem/serving/watcher.py +++ b/src/recotem/serving/watcher.py @@ -199,17 +199,28 @@ def _sha256_bytes(data: bytes) -> str: class _RecipeWatchState: """Internal state the watcher maintains per recipe.""" + #: The recipe **as it is on disk now**, refreshed by every rescan that + #: re-parses the YAML. ``None`` only for a YAML-parse-failure stub, which + #: never parsed at all. + #: + #: This used to be the body parsed at process start and was never + #: refreshed outside the two YAML-error-recovery branches below, so for a + #: recipe that always parsed cleanly it stayed the startup body for the + #: life of the process. It was read as a proxy for "the body the serving + #: model was built from", which it is only until the first hot-swap: edit + #: a recipe and retrain, and the startup body is neither the body on disk + #: nor the body the new artifact was trained from. The watcher cannot + #: reconstruct the trained-from body -- it holds only its ``recipe_hash`` + #: -- so the body on disk is the only one it can name, and it is the one + #: ``app.py`` already uses on the startup load path. Keeping the two + #: paths on the same body is what makes the hot-swap path and a restart + #: produce the same entry. recipe: Any # Recipe + #: Where the artifact for this recipe currently lives, i.e. + #: ``recipe.output.path``. Kept as a separate field rather than read + #: through ``recipe`` because a re-point has to drop everything the + #: watcher remembers about the previous file; see ``_repoint_artifact``. artifact_path: str - #: The most recently *parsed* body of this recipe's YAML, refreshed on - #: every successful rescan. ``recipe`` above is deliberately NOT - #: refreshed: it is the body the currently-served model was built from, - #: and re-pointing ``output.path`` or ``item_metadata`` underneath a live - #: model is a behaviour change, not a warning fix. This field exists so - #: the drift warning can compare the artifact against the recipe *as it - #: is on disk right now* without altering what is served. ``None`` until - #: the first rescan after startup, where ``recipe`` is still current. - latest_recipe: Any = None last_marker: Any = None last_sha256: str = "" #: Last-known contents of the ``.sha256`` sidecar pointer file. @@ -242,15 +253,14 @@ class _RecipeWatchState: last_attempted_marker: Any = None #: Set to True after the first TypeError from artifact_path + ".sha256" #: so subsequent polls skip the sidecar check rather than flooding logs - #: with the same warning on every poll cycle (M7). + #: with the same warning on every poll cycle (M7). Cleared by + #: ``_scan_recipes_dir`` whenever the recipe YAML is re-parsed, so a + #: changed configuration gets a fresh evaluation (C4). sidecar_unsupported: bool = False - #: The yaml_mtime at which sidecar_unsupported was set. When the recipe - #: YAML changes (yaml_mtime differs from this value) the sidecar_unsupported - #: flag is cleared so the new configuration gets a fresh evaluation (C4). - sidecar_unsupported_at_mtime: float | None = None #: Counter for consecutive transient OSErrors on sidecar reads. After - #: 3 consecutive non-ENOENT OSErrors the watcher skips sidecar checks until - #: the next mtime change to avoid triggering full reloads indefinitely (m7). + #: 3 consecutive non-ENOENT OSErrors the watcher sets ``sidecar_unsupported`` + #: and skips sidecar checks until the recipe YAML is re-parsed, so a + #: persistently unreadable sidecar cannot drive a reload every tick (m7). sidecar_io_error_count: int = 0 #: True while the outstanding ``last_load_error`` on this recipe's registry #: entry is the one the rescan-parse path wrote. Set where that path calls @@ -620,8 +630,10 @@ def _scan_recipes_dir(self) -> None: and cached[0] == current_mtime ): recipe = cached[1] + reparsed = False else: recipe = load_recipe(yaml_file, recipes_root=self._recipes_dir) + reparsed = True if current_mtime is not None: self._yaml_mtime_cache[yaml_file] = (current_mtime, recipe) except Exception as exc: @@ -743,20 +755,37 @@ def _scan_recipes_dir(self) -> None: # parsed successfully so _poll_artifacts uses the correct path. existing_state = self._states[recipe.name] # The YAML parsed, so this is the recipe as it is on disk now. - # Record it for the drift comparison regardless of which - # recovery branch below does or does not fire; without this the - # comparison keeps using the body read at process start and is - # wrong in both directions (it warns after a retrain that made - # the two agree, and stays silent when an edit made them - # disagree). - existing_state.latest_recipe = recipe + # Adopt it regardless of which recovery branch below does or + # does not fire. Without this the watcher keeps the body read + # at process start: the drift comparison is wrong in both + # directions (it warns after a retrain that made the two agree, + # and stays silent when an edit made them disagree), and the + # ``item_metadata`` table ``_build_entry`` joins against is the + # one the startup body named, so a repointed metadata path plus + # a retrain serves the new model on the old table with nothing + # warning -- the hashes legitimately agree. + existing_state.recipe = recipe + if reparsed: + # C4: the recipe changed, so anything latched against the + # previous configuration gets a fresh evaluation. The + # sidecar path is ``artifact_path + ".sha256"`` and + # ``artifact_path`` is ``output.path``, so an edit can move + # the sidecar out from under a latch set on the old one. + # This is the re-parse the latch's own comment asks about; + # it used to be looked for in ``_check_sidecar_changed`` + # through ``getattr(recipe, "_yaml_path", None)``, an + # attribute ``Recipe`` does not define and, with + # ``extra="forbid"`` and no private attributes, cannot + # carry -- so the flag was never cleared and the sidecar + # backstop stayed off for the life of the process. + existing_state.sidecar_unsupported = False + existing_state.sidecar_io_error_count = 0 if not existing_state.artifact_path and recipe.output.path: logger.info( "recipe_yaml_failure_recovered", name=recipe.name, ) existing_state.artifact_path = recipe.output.path - existing_state.recipe = recipe # Reset last_marker so the next tick triggers a fresh load. existing_state.last_marker = None elif existing_state.yaml_rescan_error: @@ -775,15 +804,26 @@ def _scan_recipes_dir(self) -> None: # is the one *this* path wrote; an artifact-load failure # never sets the flag and so is never cleared here. self._registry.set_load_error(recipe.name, None) - # Adopt the corrected recipe body (and the output path it - # derives) so the repair actually takes effect. - existing_state.recipe = recipe + # The corrected body is already adopted above; take the + # output path it derives so the repair takes effect. existing_state.artifact_path = recipe.output.path # Reset last_marker so the next tick reloads. As in the # discovery branch above, the load itself stays on the # _poll_artifacts thread pool rather than blocking the # scan (M-1). existing_state.last_marker = None + elif existing_state.artifact_path != recipe.output.path: + # Neither recovery branch fired, so this is a healthy + # recipe whose ``output.path`` was edited under the running + # watcher. Follow it: ``output.path`` is where ``train`` + # writes, and a watcher that keeps polling the path it read + # at discovery ignores the edit for the life of the process + # with nothing to show for it -- the drift warning is + # emitted from the load path, and no load ever happens at + # the new location. + self._repoint_artifact( + recipe.name, existing_state, recipe.output.path + ) existing_state.yaml_rescan_error = False for gone in current_names - found_names: @@ -825,6 +865,48 @@ def _scan_recipes_dir(self) -> None: if found_names != current_names: _metrics.set_active_recipes(self._registry.loaded_count()) + def _repoint_artifact( + self, name: str, state: _RecipeWatchState, new_path: str + ) -> None: + """Follow an edited ``output.path`` to the file it now names. + + Everything the watcher remembers about this recipe's artifact -- the + change marker, the payload digest, the ``.sha256`` sidecar contents, + the load backoff and the repeat-suppression signatures -- describes the + *previous* file and says nothing at all about the new one, so all of it + is dropped here rather than carried across. Leaving ``last_sha256`` + behind in particular would let a new path holding byte-identical + content take the "unchanged bytes" short-circuit in ``_load_recipe`` + and never rebuild the entry, so the served entry would keep naming the + old path. + + The model currently in the registry is deliberately left serving. When + the new path does not exist yet -- the ordinary case, an operator who + edited the recipe and has not retrained -- the next poll records + "artifact missing or unreadable" against the entry and + ``/v1/health/details`` degrades, which is the same M-2 contract every + other missing artifact gets. The entry's own ``artifact_path`` keeps + naming the file it was last loaded from, which is what that field + documents; the ``recipe_output_path_changed`` line below is what ties + the two together in a log. + """ + logger.info( + "recipe_output_path_changed", + name=name, + previous_path=str(state.artifact_path), + path=str(new_path), + ) + state.artifact_path = new_path + state.last_marker = None + state.last_sha256 = "" + state.last_sidecar_contents = None + state.sidecar_io_error_count = 0 + state.sidecar_unsupported = False + state._last_stat_error_class = None + state._last_failure_signature = None + state.stat_error_outstanding = False + self._clear_load_backoff(state) + # ------------------------------------------------------------------ # Artifact polling # ------------------------------------------------------------------ @@ -1086,13 +1168,7 @@ def _load_recipe( return try: - entry = self._build_entry( - name, - state.recipe, - data, - artifact_path, - drift_recipe=state.latest_recipe or state.recipe, - ) + entry = self._build_entry(name, state.recipe, data, artifact_path) except ArtifactError as exc: kid_log, kid_reason = _extract_kid_safe(data) if kid_reason is not None: @@ -1197,16 +1273,16 @@ def _build_entry( recipe: Any, data: bytes, artifact_path: str, - *, - drift_recipe: Any = None, ) -> ModelEntry: """Parse, verify, deserialize data and return a fresh ModelEntry. - *recipe* is the body this model is being built from -- it decides - ``item_metadata`` and everything else that reaches the entry. - *drift_recipe* is compared against the artifact's ``recipe_hash`` and - is only ever read for that warning; it defaults to *recipe* so the - startup path, which has no separate on-disk body yet, is unchanged. + *recipe* is the recipe as it is on disk: it decides ``item_metadata`` + and everything else that reaches the entry, and it is the body compared + against the artifact's ``recipe_hash`` for the drift warning. Both + callers pass a freshly-parsed body -- ``app.py`` because it has just + read the directory, the watcher because ``_scan_recipes_dir`` refreshes + ``state.recipe`` on every re-parse -- so a hot-swap and a restart build + the same entry from the same bytes. """ from recotem.artifact.format import parse_header_from_bytes from recotem.artifact.signing import unpickle_payload, verify_hmac @@ -1255,11 +1331,7 @@ def _build_entry( # arrives by hot-swap after a recipe edit is the same staleness, # and a check wired into only one of the two load paths is the # divergence #270 had to correct. - check_artifact_recipe_hash( - header_dict, - recipe=drift_recipe if drift_recipe is not None else recipe, - name=name, - ) + check_artifact_recipe_hash(header_dict, recipe=recipe, name=name) # Preflight the irspack version before deserializing: a skewed artifact # fails inside the C++ __setstate__ with an error that names neither @@ -1677,30 +1749,13 @@ def _check_sidecar_changed(state: _RecipeWatchState) -> bool: """ artifact_path = state.artifact_path - # Short-circuit: if a previous poll already determined that sidecar - # construction is unsupported for this path, re-evaluate only if the - # recipe YAML mtime has changed since the flag was set (C4). + # Short-circuit: a previous poll already determined that the sidecar is + # unusable for this path. The flag is cleared in ``_scan_recipes_dir`` + # when the recipe YAML is re-parsed (C4) -- that is the watcher's own + # record of "the configuration changed", and it is where the sidecar path + # itself can move, since the path is derived from ``output.path``. if state.sidecar_unsupported: - import os as _os - - yaml_mtime: float | None = None - try: - recipe_yaml = getattr(state.recipe, "_yaml_path", None) or getattr( - state.recipe, "yaml_path", None - ) - if recipe_yaml is not None: - yaml_mtime = _os.stat(recipe_yaml).st_mtime - except OSError: - pass - if ( - yaml_mtime is None - or state.sidecar_unsupported_at_mtime is None - or yaml_mtime == state.sidecar_unsupported_at_mtime - ): - return False - # YAML mtime changed — clear the flag and re-evaluate. - state.sidecar_unsupported = False - state.sidecar_unsupported_at_mtime = None + return False # Only meaningful for local-FS paths where we can form a sibling sidecar. # For remote URIs (s3://, gs://) this is a no-op; the marker comparison @@ -1713,19 +1768,7 @@ def _check_sidecar_changed(state: _RecipeWatchState) -> bool: path=str(artifact_path), exc_type=type(exc).__name__, ) - import os as _os2 - - yaml_mtime2: float | None = None - try: - recipe_yaml2 = getattr(state.recipe, "_yaml_path", None) or getattr( - state.recipe, "yaml_path", None - ) - if recipe_yaml2 is not None: - yaml_mtime2 = _os2.stat(recipe_yaml2).st_mtime - except OSError: - pass state.sidecar_unsupported = True - state.sidecar_unsupported_at_mtime = yaml_mtime2 return False try: @@ -1794,21 +1837,10 @@ def _check_sidecar_changed(state: _RecipeWatchState) -> bool: if state.sidecar_io_error_count >= 3: # After 3 consecutive non-ENOENT errors, stop triggering full # reloads on every tick to avoid a reload storm from a - # persistently unreadable sidecar (m7). The flag is cleared - # when the next yaml_mtime change is detected (C4 logic above). - import os as _os3 - - yaml_mtime3: float | None = None - try: - recipe_yaml3 = getattr(state.recipe, "_yaml_path", None) or getattr( - state.recipe, "yaml_path", None - ) - if recipe_yaml3 is not None: - yaml_mtime3 = _os3.stat(recipe_yaml3).st_mtime - except OSError: - pass + # persistently unreadable sidecar (m7). The flag is cleared on + # the next re-parse of the recipe YAML (C4, in + # ``_scan_recipes_dir``). state.sidecar_unsupported = True - state.sidecar_unsupported_at_mtime = yaml_mtime3 state.sidecar_io_error_count = 0 logger.warning( "sidecar_io_errors_suppressed", diff --git a/tests/unit/test_serving_watcher.py b/tests/unit/test_serving_watcher.py index f635b3b1..e2e6cda5 100644 --- a/tests/unit/test_serving_watcher.py +++ b/tests/unit/test_serving_watcher.py @@ -5126,55 +5126,67 @@ def test_sidecar_enoent_still_returns_false(tmp_path: Path) -> None: # --------------------------------------------------------------------------- -# C4: sidecar_unsupported resets when recipe YAML mtime changes +# C4: sidecar_unsupported resets when the recipe YAML is re-parsed # --------------------------------------------------------------------------- -def test_sidecar_unsupported_clears_on_yaml_mtime_change( +def test_sidecar_unsupported_is_never_cleared_by_the_sidecar_check_itself( tmp_path: Path, ) -> None: - """When sidecar_unsupported=True and the recipe YAML mtime changes, - _check_sidecar_changed must clear the flag and re-evaluate (C4).""" - from unittest.mock import MagicMock, patch + """C4 recovery belongs to the rescan, and must not move back in here. + + This test used to build the state around a ``MagicMock`` carrying a + ``_yaml_path`` attribute and assert that ``_check_sidecar_changed`` cleared + the latch once that file's mtime moved. It passed, and it proved nothing: + ``Recipe`` is a pydantic model with ``extra="forbid"`` and no private + attributes, so no real recipe has ever carried ``_yaml_path`` or + ``yaml_path``. The ``getattr`` chain answered ``None`` for every recipe + that has ever existed, the guard it fed always declined, and the recovery + never ran once in a live server -- the double was the only object in + existence for which the code worked. + + The clearing now happens where the watcher already decides the YAML + changed, in ``_scan_recipes_dir``; + ``tests/unit/test_watcher_live_recipe_body.py`` covers the recovery itself + against a real recipe. Pinned here is the other half of that split: with a + real recipe, this function declines and leaves the latch alone however the + YAML moves underneath it, so a reintroduced mtime probe would fail here. + """ + import os + from recotem.recipe.loader import load_recipe from recotem.serving.watcher import _check_sidecar_changed, _RecipeWatchState - yaml_path = tmp_path / "recipe.yaml" - yaml_path.write_text("name: test\n") - artifact_path = str(tmp_path / "model.recotem") - - recipe = MagicMock() - recipe.name = "c4_test" - recipe._yaml_path = yaml_path + recipes_dir = tmp_path / "recipes" + recipes_dir.mkdir() + artifact_path = tmp_path / "model.recotem" + yaml_path = _write_recipe_yaml(recipes_dir, "c4_test", artifact_path) - initial_mtime = yaml_path.stat().st_mtime + # A sidecar that exists and has genuinely changed, so a "False" below can + # only be the latch and not an absent or unchanged file. + sidecar_path = Path(str(artifact_path) + ".sha256") + sidecar_path.write_text("sha_v2\n") state = _RecipeWatchState( - recipe=recipe, - artifact_path=artifact_path, + recipe=load_recipe(yaml_path), + artifact_path=str(artifact_path), + last_sidecar_contents="sha_v1\n", sidecar_unsupported=True, - sidecar_unsupported_at_mtime=initial_mtime, ) - # With same mtime, still unsupported — returns False immediately. - result = _check_sidecar_changed(state) - assert result is False, "No mtime change → sidecar_unsupported stays True" + assert _check_sidecar_changed(state) is False assert state.sidecar_unsupported is True - # Simulate mtime change by patching os.stat to return a newer mtime. - new_mtime = initial_mtime + 1.0 - - class _FakeStat: - st_mtime = new_mtime + bumped = os.stat(yaml_path).st_mtime + 10 + os.utime(yaml_path, (bumped, bumped)) - with patch("os.stat", return_value=_FakeStat()): - result2 = _check_sidecar_changed(state) - - # After mtime change, sidecar_unsupported must be cleared. - assert state.sidecar_unsupported is False, ( - "sidecar_unsupported must be cleared when recipe YAML mtime changes" + assert _check_sidecar_changed(state) is False, ( + "the sidecar check has no view of the recipe file and must not grow one" + ) + assert state.sidecar_unsupported is True, ( + "only the rescan clears the latch; clearing it here would need a recipe " + "attribute that no Recipe can carry" ) - assert state.sidecar_unsupported_at_mtime is None # --------------------------------------------------------------------------- diff --git a/tests/unit/test_watcher_live_recipe_body.py b/tests/unit/test_watcher_live_recipe_body.py new file mode 100644 index 00000000..7e0153c7 --- /dev/null +++ b/tests/unit/test_watcher_live_recipe_body.py @@ -0,0 +1,496 @@ +"""What the watcher reads out of a recipe once the file on disk has changed. + +``_RecipeWatchState.recipe`` was the body parsed at process start. Nothing +refreshed it: the two assignments that existed sat inside YAML-error-recovery +branches, so a recipe that always parsed cleanly kept its startup body for the +life of the process. Three things read that body and wanted the current one: + +* ``item_metadata`` -- ``_build_entry`` joins the response against the table the + *startup* body named, so an operator who repointed ``item_metadata.path`` and + retrained got the new model joined onto the old table, silently; +* ``output.path`` -- captured once into ``state.artifact_path`` at discovery, so + an edited output path was ignored and the watcher kept polling the old file + forever; and +* the ``.sha256`` sidecar suppression latch, whose documented "clear it when the + recipe changes" recovery was reached through + ``getattr(recipe, "_yaml_path", None)`` -- an attribute ``Recipe`` does not + have and, with ``extra="forbid"`` and no private attributes, cannot have. + +All three shipped in 2.0.0. These tests edit the YAML *under a running +watcher*, which is the only way the startup body and the current one differ. +""" + +from __future__ import annotations + +import errno +import os +import time +from pathlib import Path +from unittest.mock import patch + +import structlog.testing + +from recotem.artifact.signing import KeyRing +from recotem.config import ServeConfig +from recotem.serving.registry import ModelEntry, ModelRegistry +from recotem.serving.watcher import ( + ArtifactWatcher, + _check_sidecar_changed, + _RecipeWatchState, +) +from tests.conftest import ACTIVE_KEY_HEX, build_raw_artifact + +# --------------------------------------------------------------------------- +# Fixtures +# --------------------------------------------------------------------------- + + +def _write_artifact( + path: Path, recipe_name: str, tag: str, recipe_hash: str | None = None +) -> None: + import pickle # noqa: S403 # test fixture: payload built locally + + payload = pickle.dumps({"tag": tag}, protocol=4) # noqa: S301 + header: dict = { + "recipe_name": recipe_name, + "best_class": "TopPop", + "trained_at": "2026-01-01T00:00:00Z", + } + if recipe_hash is not None: + header["recipe_hash"] = recipe_hash + path.write_bytes( + build_raw_artifact( + kid="active", + key_hex=ACTIVE_KEY_HEX, + header_dict=header, + payload_bytes=payload, + ) + ) + + +def _recipe_text( + *, output_path: Path | str, metadata_path: Path | str | None = None +) -> str: + metadata_block = ( + "" + if metadata_path is None + else f"""\ +item_metadata: + type: csv + path: {metadata_path} + fields: [title] + on_field_missing: error +""" + ) + return f"""\ +name: demo +source: + type: csv + path: /tmp/data.csv +schema: + user_column: user_id + item_column: item_id +training: + algorithms: [TopPop] + n_trials: 1 +{metadata_block}output: + path: {output_path} +""" + + +def _rewrite(yaml_path: Path, text: str) -> None: + """Rewrite the recipe and push its mtime forward. + + The watcher re-parses a recipe only when ``st_mtime`` differs from the value + it cached (``_yaml_mtime_cache``). Two writes inside one mtime tick are + invisible to it; local filesystems here resolve consecutive writes but a CI + runner's need not, and a test that leans on that resolution fails there and + nowhere else. Bump it explicitly. + """ + yaml_path.write_text(text) + bumped = os.stat(yaml_path).st_mtime + 10 + os.utime(yaml_path, (bumped, bumped)) + + +def _make_serve_config() -> ServeConfig: + cfg = ServeConfig() + cfg.signing_keys_raw = f"active:{ACTIVE_KEY_HEX}" + cfg.watch_interval = 0.05 + cfg.max_artifact_bytes = 100 * 1024 * 1024 + return cfg + + +def _build_watcher(tmp_path: Path, yaml_text: str, artifact_path: Path): + """A watcher over one recipe, primed to load on its first tick.""" + from recotem.recipe.loader import load_recipe + + recipes_dir = tmp_path / "recipes" + recipes_dir.mkdir(exist_ok=True) + yaml_path = recipes_dir / "demo.yaml" + yaml_path.write_text(yaml_text) + + registry = ModelRegistry() + registry.replace( + "demo", + ModelEntry( + name="demo", + recommender=None, + header={}, + kid="", + artifact_path=str(artifact_path), + loaded=False, + ), + ) + states = { + "demo": _RecipeWatchState( + recipe=load_recipe(yaml_path), + artifact_path=str(artifact_path), + last_sha256="", + last_marker=None, + ) + } + watcher = ArtifactWatcher( + registry=registry, + recipes_dir=recipes_dir, + serve_config=_make_serve_config(), + key_ring=KeyRing(f"active:{ACTIVE_KEY_HEX}"), + initial_states=states, + ) + return watcher, registry, states, yaml_path + + +def _wait_until(predicate, timeout: float = 5.0) -> bool: + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + if predicate(): + return True + time.sleep(0.05) + return False + + +def _titles(registry: ModelRegistry) -> list: + entry = registry.get("demo") + if entry is None or entry.metadata_index is None: + return [] + return [v.get("title") for v in entry.metadata_index.values()] + + +# --------------------------------------------------------------------------- +# item_metadata +# --------------------------------------------------------------------------- + + +def test_item_metadata_follows_the_edited_recipe_after_a_retrain( + tmp_path: Path, +) -> None: + """Repoint ``item_metadata.path``, retrain, and the join must follow. + + This is the ordinary operator move -- edit the recipe and retrain -- and it + is exactly the case where the startup body is provably the wrong one: the + artifact that swaps in was trained FROM the edited recipe, so the recipe on + disk and the body the model was built from are the same body, and the + startup body is neither. The hashes agree, so the drift warning stays + silent and nothing else marks the mismatch: the response is the new model + joined onto the old table. + """ + from recotem._recipe_hash import compute_recipe_hash + from recotem.recipe.loader import load_recipe + + old_meta = tmp_path / "meta_old.csv" + old_meta.write_text("item_id,title\ni1,OLD_TITLE\n") + new_meta = tmp_path / "meta_new.csv" + new_meta.write_text("item_id,title\ni1,NEW_TITLE\n") + + artifact_path = tmp_path / "model.recotem" + watcher, registry, _, yaml_path = _build_watcher( + tmp_path, + _recipe_text(output_path=artifact_path, metadata_path=old_meta), + artifact_path, + ) + startup_hash = compute_recipe_hash(load_recipe(yaml_path)) + _write_artifact( + artifact_path, recipe_name="demo", tag="v1", recipe_hash=startup_hash + ) + + with structlog.testing.capture_logs() as logs: + watcher.start() + try: + assert _wait_until( + lambda: (e := registry.get("demo")) is not None and e.loaded + ), "precondition: the artifact must load" + assert _titles(registry) == ["OLD_TITLE"], ( + "precondition: the startup body's metadata table must be the one " + f"loaded; got {_titles(registry)}" + ) + + # The operator repoints item_metadata AND retrains, so the new + # artifact carries the hash of the recipe as edited. + _rewrite( + yaml_path, + _recipe_text(output_path=artifact_path, metadata_path=new_meta), + ) + edited_hash = compute_recipe_hash(load_recipe(yaml_path)) + _n = len(logs) + assert _wait_until( + lambda: any(r.get("event") == "recipe_loaded" for r in logs[_n:]) + ), "the watcher must re-scan the edited YAML" + + _write_artifact( + artifact_path, recipe_name="demo", tag="v2", recipe_hash=edited_hash + ) + assert _wait_until( + lambda: ( + (e := registry.get("demo")) is not None + and e.recommender == {"tag": "v2"} + ) + ), "the retrained artifact must hot-swap in" + finally: + watcher.stop() + watcher.join(timeout=5.0) + + assert _titles(registry) == ["NEW_TITLE"], ( + "the artifact that just swapped in was trained FROM the edited recipe, " + "yet the metadata join still uses the table the STARTUP body named. The " + f"response is the new model on the old table. Titles: {_titles(registry)}" + ) + warned = [r for r in logs if r.get("event") == "artifact_recipe_hash_mismatch"] + assert not warned, ( + "control: artifact and recipe agree here, so nothing may warn -- the " + f"stale join is silent, which is the point. Emitted: {warned}" + ) + + +def test_item_metadata_is_not_reloaded_without_a_hot_swap(tmp_path: Path) -> None: + """Control for the test above: an edit alone must not move the join. + + Metadata is read where the model is built, so an edit with no retrain + leaves the old model AND the old table in place -- which is the coherent + answer, and the drift warning is what says the recipe has moved on. Without + this arm, "the join followed the edit" above could as easily mean "the join + re-reads on every tick", which would be a different behaviour entirely. + """ + from recotem._recipe_hash import compute_recipe_hash + from recotem.recipe.loader import load_recipe + + old_meta = tmp_path / "meta_old.csv" + old_meta.write_text("item_id,title\ni1,OLD_TITLE\n") + new_meta = tmp_path / "meta_new.csv" + new_meta.write_text("item_id,title\ni1,NEW_TITLE\n") + + artifact_path = tmp_path / "model.recotem" + watcher, registry, _, yaml_path = _build_watcher( + tmp_path, + _recipe_text(output_path=artifact_path, metadata_path=old_meta), + artifact_path, + ) + startup_hash = compute_recipe_hash(load_recipe(yaml_path)) + _write_artifact( + artifact_path, recipe_name="demo", tag="v1", recipe_hash=startup_hash + ) + + with structlog.testing.capture_logs() as logs: + watcher.start() + try: + assert _wait_until( + lambda: (e := registry.get("demo")) is not None and e.loaded + ), "precondition: the artifact must load" + + _rewrite( + yaml_path, + _recipe_text(output_path=artifact_path, metadata_path=new_meta), + ) + _n = len(logs) + assert _wait_until( + lambda: any(r.get("event") == "recipe_loaded" for r in logs[_n:]) + ), "the watcher must re-scan the edited YAML" + # Several further ticks, with no new artifact. + time.sleep(0.5) + finally: + watcher.stop() + watcher.join(timeout=5.0) + + assert _titles(registry) == ["OLD_TITLE"], ( + "no artifact changed, so the served model is still the old one and its " + f"metadata must be too; got {_titles(registry)}" + ) + + +# --------------------------------------------------------------------------- +# output.path +# --------------------------------------------------------------------------- + + +def test_an_edited_output_path_is_followed_to_the_new_artifact( + tmp_path: Path, +) -> None: + """``output.path`` is where train writes; serve must poll where it points now. + + ``state.artifact_path`` was taken from ``recipe.output.path`` once, at + discovery. Editing it under a running server was ignored for the life of + the process -- and silently: the drift warning cannot cover this, because it + is emitted from the load path and no load ever happens at the new location. + """ + old_path = tmp_path / "old.recotem" + new_path = tmp_path / "new.recotem" + _write_artifact(old_path, recipe_name="demo", tag="v1") + + watcher, registry, states, yaml_path = _build_watcher( + tmp_path, _recipe_text(output_path=old_path), old_path + ) + + watcher.start() + try: + assert _wait_until( + lambda: ( + (e := registry.get("demo")) is not None + and e.recommender == {"tag": "v1"} + ) + ), "precondition: the artifact at the original path must load" + + _write_artifact(new_path, recipe_name="demo", tag="NEWPATH") + _rewrite(yaml_path, _recipe_text(output_path=new_path)) + + followed = _wait_until( + lambda: ( + (e := registry.get("demo")) is not None + and e.recommender == {"tag": "NEWPATH"} + ) + ) + finally: + watcher.stop() + watcher.join(timeout=5.0) + + assert followed, ( + "output.path was repointed under a running watcher and the artifact at " + "the new path never loaded: the watcher is still polling " + f"{states['demo'].artifact_path!r}, the path captured at discovery." + ) + assert states["demo"].artifact_path == str(new_path) + assert registry.get("demo").artifact_path == str(new_path) + + +def test_repointing_output_path_at_a_missing_file_keeps_the_model_serving( + tmp_path: Path, +) -> None: + """The half of following an edit that must not become an outage. + + An operator repoints ``output.path`` and has not retrained yet, so nothing + is there. Following the edit must degrade health -- the recipe now names a + file that does not exist -- without evicting the model that is serving, + which is the same contract every other missing artifact gets (M-2). + """ + old_path = tmp_path / "old.recotem" + missing = tmp_path / "not-trained-yet.recotem" + _write_artifact(old_path, recipe_name="demo", tag="v1") + + watcher, registry, states, yaml_path = _build_watcher( + tmp_path, _recipe_text(output_path=old_path), old_path + ) + + watcher.start() + try: + assert _wait_until( + lambda: ( + (e := registry.get("demo")) is not None + and e.recommender == {"tag": "v1"} + ) + ), "precondition: the artifact at the original path must load" + + _rewrite(yaml_path, _recipe_text(output_path=missing)) + + degraded = _wait_until( + lambda: (e := registry.get("demo")) is not None and e.last_load_error + ) + finally: + watcher.stop() + watcher.join(timeout=5.0) + + assert degraded, ( + "output.path now names a file that does not exist and nothing was " + "recorded against the entry: the watcher is still polling the old path " + f"({states['demo'].artifact_path!r}) and reporting success for it." + ) + entry = registry.get("demo") + assert entry.loaded is True, "the model that is serving must not be evicted" + assert entry.recommender == {"tag": "v1"} + + +# --------------------------------------------------------------------------- +# the .sha256 sidecar suppression latch +# --------------------------------------------------------------------------- + + +def _latch_sidecar_unsupported(state: _RecipeWatchState) -> None: + """Drive the real latching path: three consecutive non-ENOENT read errors.""" + perm_error = PermissionError("permission denied") + perm_error.errno = errno.EACCES + with patch.object(Path, "read_text", side_effect=perm_error): + for _ in range(3): + _check_sidecar_changed(state) + + +def test_sidecar_suppression_clears_when_the_recipe_yaml_is_reparsed( + tmp_path: Path, +) -> None: + """The documented C4 recovery, against a real ``Recipe``. + + Three unreadable sidecar reads latch ``sidecar_unsupported`` so a broken + sidecar cannot drive a reload every tick. Nothing else clears it, so the + latch is permanent -- and the sidecar is the only backstop the watcher has + when the change marker cannot discriminate (an ``append_sha`` pointer file + is a constant 22 bytes, so ``(mtime, size)`` degenerates to mtime alone). + The documented escape is "the recipe YAML changed, re-evaluate", reached via + ``getattr(state.recipe, "_yaml_path", None)``, which is ``None`` for every + ``Recipe`` that has ever existed. + """ + artifact_path = tmp_path / "model.recotem" + _write_artifact(artifact_path, recipe_name="demo", tag="v1") + sidecar = Path(str(artifact_path) + ".sha256") + sidecar.write_text("sha_v1\n") + + watcher, _registry, states, yaml_path = _build_watcher( + tmp_path, _recipe_text(output_path=artifact_path), artifact_path + ) + state = states["demo"] + state.last_sidecar_contents = "sha_v1\n" + # Warm the YAML mtime cache the way a running watcher's first tick does, so + # the scans below are the steady-state "nothing changed" case. + watcher._scan_recipes_dir() + + _latch_sidecar_unsupported(state) + assert state.sidecar_unsupported is True, ( + "precondition: three EACCES sidecar reads must latch the suppression" + ) + + # Control: scans that do NOT re-parse (the YAML has not changed) must leave + # the latch alone -- otherwise every tick would clear it and the + # reload-storm guard the latch exists for would be gone. + watcher._scan_recipes_dir() + watcher._scan_recipes_dir() + assert state.sidecar_unsupported is True, ( + "an unchanged YAML must not clear the latch; the guard would be useless" + ) + sidecar.write_text("sha_v2\n") + assert _check_sidecar_changed(state) is False, ( + "control: while latched, a readable and genuinely CHANGED sidecar is " + "ignored -- this is the backstop that is off" + ) + + # The operator edits the recipe. That is the documented signal to + # re-evaluate -- and after it the sidecar path itself may have moved, + # because it is derived from output.path. + _rewrite(yaml_path, _recipe_text(output_path=artifact_path) + "# edited\n") + watcher._scan_recipes_dir() + + assert state.sidecar_unsupported is False, ( + "the recipe YAML was re-parsed and the suppression was not cleared: the " + "documented recovery reads Recipe._yaml_path, which does not exist, so " + "the latch is permanent and the sidecar backstop is off until restart." + ) + assert _check_sidecar_changed(state) is True, ( + "with the latch cleared, the sidecar change made above must be seen" + ) + assert _check_sidecar_changed(state) is False, ( + "control: and the call after it must answer False -- the True above is " + "a change signal, not an unconditional yes" + )