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
12 changes: 8 additions & 4 deletions src/recotem/_artifact_identity.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
222 changes: 127 additions & 95 deletions src/recotem/serving/watcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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:
Expand All @@ -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:
Expand Down Expand Up @@ -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
# ------------------------------------------------------------------
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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:
Expand Down Expand Up @@ -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",
Expand Down
Loading