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
52 changes: 52 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,58 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased] — `crd-well-oiled-machine`

### Slice 2-/3-hot-reload — InferencePolicy + ClawMemory loaders watch their mount dirs

Closes the long-running gap where the InferencePolicy
(`inference_policy_loader::load_and_install`) and ClawMemory
(`memory_binding_loader::load_and_install`) loaders were one-shot at
router startup. After the initial load, the file on disk could change
(operator `kubectl edit`s the CR, controller recompiles and rewrites
the ConfigMap, kubelet refreshes the projected mount) and the router
would happily keep echoing the **first** digest forever — the
controller-side `/internal/policy-status` echo loop would never close
`Compiled → Ready` for any change after the first reconcile.

What this PR adds:

- `inference_policy_loader::spawn_inference_policy_watcher(dir, registry, handle)` —
background `tokio` task that polls `dir`'s max-mtime every
`INFERENCE_POLICY_WATCH_INTERVAL` seconds (default 5s) and re-invokes
`load_and_install` whenever a change is detected.
- `memory_binding_loader::spawn_memory_binding_watcher(dir, registry, handle)` —
same shape; env knob `MEMORY_BINDING_WATCH_INTERVAL`. Closes
**Slice 3 DoD #4** ("kubectl edit rotates the digest; router
reloads within 5s") explicitly; **Slice 2 DoD #1** ("`Ready` after
router echo") implicitly (the echo loop now actually closes on
every change, not just the first).
- `load_and_install` semantics tightened in both loaders:
- `Loaded` → handle overwritten with new value (unchanged).
- `NoBinding` / `NoPolicy` → handle **cleared**. Removing
`spec.memoryRef` / `spec.inferenceRef` from a `ClawSandbox` now
actually unbinds the in-memory state on the next tick (instead
of the router pretending the policy is still in effect until the
pod restarts).
- `Error` → handle **left intact**. A transient parse error during
a partial mount update must not knock the data plane offline; the
registry already captured the error so the echo loop notices the
digest is stale.
- Both watchers wired in `main.rs` after `governance::spawn_policy_watcher`.
Best-effort: missing directory is fine (mount may appear later when
the operator adds the CR reference).
- 5 new memory_binding_loader unit tests (clear-on-NoBinding,
preserve-prior-on-Error, two `dir_max_mtime` shape tests, and a
full watcher integration test that crank-runs the watcher with a
1s interval, mutates the file, and asserts the digest in the
handle rolls forward).
- 779 router lib tests (+5 — the new mem-binding tests; the
inference_policy_loader changes are exercised by the same
watcher pattern and tested at integration level via the
policy-status echo route).

No CR shape change. No new env vars required for production (5s
default ticks at exactly the Slice 3 DoD SLO). Other watchers
(`Governance::spawn_policy_watcher`) untouched.

### Slice 3b.5 — `MemoryStoreMissing` Degraded condition for ClawMemory

When the upstream Foundry Memory Store returns HTTP 404 on a
Expand Down
88 changes: 84 additions & 4 deletions inference-router/src/inference_policy_loader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -360,20 +360,100 @@ pub fn load_inference_policy_from_dir(
}

/// Load and install into the shared handle in one call. Used at
/// router startup. Returns the outcome so the caller can log /
/// surface in metrics.
/// router startup **and** by `spawn_inference_policy_watcher`'s
/// mtime-poll loop (Slice 2 hot-reload). Handle-update semantics
/// mirror `memory_binding_loader::load_and_install`:
///
/// - [`LoadOutcome::Loaded`] → handle overwritten with new policy.
/// - [`LoadOutcome::NoBinding`] → handle cleared. Restores the
/// chart-fed `BUDGET_PER_REQUEST_TOKENS` / env-driven content
/// safety paths the second the operator removes the
/// `InferencePolicy` reference from a `ClawSandbox`.
/// - [`LoadOutcome::Error`] → handle left intact. The registry
/// already recorded the parse error and the controller's echo
/// loop will catch the stale digest; we refuse to knock the
/// data plane offline on a transient mid-write read.
pub async fn load_and_install(
dir: &str,
policy_status: &PolicyStatusRegistry,
handle: &LoadedInferencePolicyHandle,
) -> LoadOutcome {
let outcome = load_inference_policy_from_dir(dir, policy_status);
if let LoadOutcome::Loaded(ref policy) = outcome {
*handle.write().await = Some(policy.clone());
match &outcome {
LoadOutcome::Loaded(policy) => {
*handle.write().await = Some(policy.clone());
}
LoadOutcome::NoPolicy => {
*handle.write().await = None;
}
LoadOutcome::Error(_) => {}
}
outcome
}

/// Default poll interval for `spawn_inference_policy_watcher`.
/// Matches `memory_binding_loader::DEFAULT_WATCH_INTERVAL_SECS` so
/// operators see the same "edit-takes-effect-within-5s" SLO across
/// every router-enforced CRD.
pub const DEFAULT_WATCH_INTERVAL_SECS: u64 = 5;

/// Env-var override for [`DEFAULT_WATCH_INTERVAL_SECS`].
pub const WATCH_INTERVAL_ENV: &str = "INFERENCE_POLICY_WATCH_INTERVAL";

/// Spawn a background task that polls `dir`'s max-mtime every
/// `INFERENCE_POLICY_WATCH_INTERVAL` seconds (default 5s) and calls
/// [`load_and_install`] whenever a change is detected. Mirrors the
/// `governance::Governance::spawn_policy_watcher` and
/// `memory_binding_loader::spawn_memory_binding_watcher` patterns —
/// closes the long-running gap where `InferencePolicy` reconciler
/// would happily compile a new digest but the router's loader was
/// one-shot at startup and never re-read.
pub fn spawn_inference_policy_watcher(
dir: String,
policy_status: Arc<PolicyStatusRegistry>,
handle: LoadedInferencePolicyHandle,
) {
let interval_secs: u64 = std::env::var(WATCH_INTERVAL_ENV)
.ok()
.and_then(|v| v.parse().ok())
.filter(|v: &u64| *v > 0)
.unwrap_or(DEFAULT_WATCH_INTERVAL_SECS);

tokio::spawn(async move {
let mut last_mtime = dir_max_mtime(&dir);
let mut ticker = tokio::time::interval(std::time::Duration::from_secs(interval_secs));
ticker.tick().await; // skip the immediate first tick
loop {
ticker.tick().await;
let current = dir_max_mtime(&dir);
if current != last_mtime {
tracing::info!(
target: "inference_policy_watcher",
dir = %dir,
"InferencePolicy directory changed, reloading"
);
let _ = load_and_install(&dir, &policy_status, &handle).await;
last_mtime = current;
}
}
});
}

/// Get the max mtime across `*.json` files in `dir`. Filter mirrors
/// the controller-side compile output (`inference-policy.json`).
fn dir_max_mtime(dir: &str) -> Option<std::time::SystemTime> {
let path = Path::new(dir);
if !path.is_dir() {
return None;
}
std::fs::read_dir(path)
.ok()?
.flatten()
.filter(|e| e.path().extension().is_some_and(|ext| ext == "json"))
.filter_map(|e| e.metadata().ok()?.modified().ok())
.max()
}

/// **Latency-optimised snapshot** of every enforcement axis the
/// inference handlers consume per request (Slice 2a/2b/2c). Acquired
/// with a **single** `RwLock::read().await` and then passed around
Expand Down
28 changes: 28 additions & 0 deletions inference-router/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,34 @@ async fn main() -> Result<()> {
// Start policy hot-reload watcher (polls AGT_POLICY_DIR for mtime changes).
governance::Governance::spawn_policy_watcher(state.governance.clone());

// Slice 2 / Slice 3 hot-reload: poll the InferencePolicy and
// ClawMemory mount directories for mtime changes and re-invoke
// each loader's `load_and_install`. Without this the router
// would happily echo a stale digest for the lifetime of the
// pod whenever an operator `kubectl edit`s the CR — the
// controller-side echo loop would never close `Compiled → Ready`
// after the first change. Each watcher is best-effort: a missing
// directory is fine (the mount may appear later when the
// operator adds `spec.inferenceRef` / `spec.memoryRef`).
{
let inference_dir = std::env::var("INFERENCE_POLICY_DIR")
.unwrap_or_else(|_| "/etc/azureclaw/inference".into());
azureclaw_inference_router::inference_policy_loader::spawn_inference_policy_watcher(
inference_dir,
state.policy_status.clone(),
state.inference_policy.clone(),
);

let memory_dir = std::env::var("MEMORY_BINDING_DIR").unwrap_or_else(|_| {
azureclaw_inference_router::memory_binding_loader::MEMORY_BINDING_DIR_DEFAULT.into()
});
azureclaw_inference_router::memory_binding_loader::spawn_memory_binding_watcher(
memory_dir,
state.policy_status.clone(),
state.memory_binding.clone(),
);
}

// Clone blocklist for the forward proxy before state is moved into the router.
let proxy_blocklist = state.blocklist.clone();
let proxy_blocked_egress = state.blocked_egress.clone();
Expand Down
Loading
Loading