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
179 changes: 179 additions & 0 deletions core/continuum-core/src/inference/llama_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -911,6 +911,21 @@ pub const PROVIDER_ID: &str = "llama-server";
/// process lifetime (the daemon owns it), so the only non-ready resolution is
/// the timeout.
pub async fn await_ready_serving(timeout: Duration) -> Option<ServingSnapshot> {
// Lane-source-agnostic readiness (the misfit / grid design): if the operator
// pinned an EXTERNAL OpenAI-compatible endpoint via `LLAMA_SERVER_BASE_URL`,
// this node does NOT own a local GPU serving lane — it hosts personas against
// the pinned endpoint (a co-located engine like a K3 llama-server, or later a
// reachable grid peer). Probe that endpoint directly and synthesize the ready
// snapshot instead of waiting on the LOCAL serving daemon's `SERVING_STATE`.
// Both the persona-host gate AND the persona adapter factory call THIS one
// function, so a single change makes the whole hosting path source-agnostic —
// no persona ever requires this box to own the GPU. The endpoint is held to
// the IDENTICAL decode bar as a local lane (a real multi-token generation), so
// a compute-wedged endpoint is rejected, never faked
// ([[fallbacks-are-illegal-fail-loud]]).
if external_serving_pin().is_some() {
return probe_external_serving(timeout).await;
}
let mut rx = SERVING_STATE.get()?.clone();
{
// Fast path: already ready, no await.
Expand All @@ -928,6 +943,170 @@ pub async fn await_ready_serving(timeout: Duration) -> Option<ServingSnapshot> {
}
}

/// The operator-pinned EXTERNAL OpenAI-compatible serving endpoint ROOT
/// (`http://host:port`, no `/v1`), if `LLAMA_SERVER_BASE_URL` is set in the
/// config (read via [`crate::config_env::read`] — the `~/.continuum/config.env`
/// FILE, not a process env var). When set, this node does not own a local GPU
/// serving lane; persona hosting adopts the pinned endpoint. `serving_root()`
/// already resolves to this same value — this predicate just answers "is one
/// pinned", the seam the hosting path branches on.
pub fn external_serving_pin() -> Option<String> {
let raw = crate::config_env::read("LLAMA_SERVER_BASE_URL")?;
let trimmed = raw.trim().trim_end_matches('/');
let root = trimmed.strip_suffix("/v1").unwrap_or(trimmed);
(!root.is_empty()).then(|| root.to_string())
}

/// Read the resident model id an external endpoint reports at `/v1/models`,
/// accepting BOTH the canonical OpenAI shape (`data[].id`) and the
/// llama.cpp/Ollama-compat shape some builds answer (`models[].model|name` — a
/// K3 `llama-server` does exactly this). Returns `None` only if neither shape
/// names a model; the caller substitutes a stable sentinel (a single-resident
/// endpoint ignores the request `model` field anyway).
/// Reachability gate for a pinned EXTERNAL endpoint: `/health` answers 200.
///
/// Why NOT the local-lane [`decode_smoke_ok`] (a 5+ token generation): that bar
/// exists to catch an *untrusted orphan* whose GPU compute path is wedged. It
/// doubles as a speed test — a legitimately slow external lane (a CPU-offloaded
/// MoE like K3 at ~0.03 tok/s under load, or a distant grid peer) cannot finish a
/// multi-token generation inside any sane probe deadline, and worse, a
/// client-side timeout does NOT cancel the server-side generation, so a probe
/// storm monopolizes the endpoint's slot and NOTHING ever hosts. The operator
/// DELIBERATELY pinned this endpoint, so the adopt question is "is it reachable +
/// serving a model" (proved here + by the `/props` window read + the `/v1/models`
/// model read in [`probe_external_serving`]); the decode path is exercised by the
/// persona's real turns, where a genuine wedge (every turn 500s) surfaces LOUD on
/// the first turn — not silently faked. Cheap, off the HTTP layer, so no slot
/// contention and no probe storm.
async fn external_health_ok(root: &str, client: &reqwest::Client) -> bool {
let url = format!("{root}/health");
matches!(
client.get(&url).timeout(PROBE_TIMEOUT).send().await,
Ok(resp) if resp.status().is_success()
)
}

async fn external_active_model(v1_url: &str, client: &reqwest::Client) -> Option<String> {
let url = format!("{v1_url}/models");
let body: serde_json::Value = client
.get(&url)
.timeout(PROBE_TIMEOUT)
.send()
.await
.ok()?
.json()
.await
.ok()?;
// OpenAI: { "data": [ { "id": "..." } ] }
if let Some(id) = body
.get("data")
.and_then(|d| d.as_array())
.and_then(|a| a.first())
.and_then(|m| m.get("id"))
.and_then(|v| v.as_str())
{
return Some(id.to_string());
}
// llama.cpp/Ollama-compat: { "models": [ { "model": "...", "name": "..." } ] }
body.get("models")
.and_then(|d| d.as_array())
.and_then(|a| a.first())
.and_then(|m| m.get("model").or_else(|| m.get("name")))
.and_then(|v| v.as_str())
.map(|s| s.to_string())
}

/// Probe a pinned EXTERNAL OpenAI-compatible endpoint for decode-readiness and,
/// on success, synthesize the ready [`ServingSnapshot`] persona hosting binds
/// against — the SAME shape the local serving daemon publishes, so every reader
/// ([`await_ready_serving`], the persona adapter factory) stays source-agnostic.
///
/// The endpoint is held to the IDENTICAL bar as a local lane: [`decode_smoke_ok`]
/// runs a real multi-token generation, so an endpoint that answers `/v1/models`
/// 200 but 500s (or dribbles ~2 tokens on) every decode is rejected, never
/// adopted ([[fallbacks-are-illegal-fail-loud]]). `served_context_window` comes
/// from the endpoint's own `/props` (`n_ctx`) so personas budget prompts to the
/// truth. Returns `None` if unset, unreachable, wedged, or `/props` won't name a
/// window (a decode-ready endpoint that can't report its context is not safely
/// adoptable — never a guessed window).
pub async fn probe_external_serving(timeout: Duration) -> Option<ServingSnapshot> {
let _ = timeout; // reachability gate is control-plane (fast); no generation to bound.
let root = external_serving_pin()?;
let v1_url = format!("{root}/v1");
let client = reqwest::Client::new();

// Reachability gate (see `external_health_ok` for why not a generation probe
// against a deliberately-pinned, possibly-very-slow endpoint).
if !external_health_ok(&root, &client).await {
crate::probe!(
class = "serving.external_probe",
endpoint = root.as_str(),
"pinned external serving endpoint is unreachable (/health) — \
not adopting; will retry on the next serving edge",
);
return None;
}

// The authoritative per-slot window from the endpoint's own `/props`. A
// decode-ready endpoint that won't name its window is not safely adoptable.
let served_context_window = match LlamaServerProcess::with_root(root.clone())
.served_context_window()
.await
{
Ok(w) if w > 0 => w,
_ => {
crate::probe!(
class = "serving.external_probe",
endpoint = root.as_str(),
"pinned external endpoint decodes but `/props` did not report an n_ctx window — \
refusing to adopt with a guessed window",
);
return None;
}
};

let active_model = external_active_model(&v1_url, &client)
.await
.unwrap_or_else(|| "external-serving-lane".to_string());

crate::probe!(
class = "serving.external_adopt",
endpoint = root.as_str(),
model = active_model.as_str(),
window = served_context_window,
"external serving endpoint adopted as the persona lane — this node hosts personas \
without owning a local GPU lane (misfit / grid serving)",
);

Some(ServingSnapshot {
active_model: Some(active_model),
ready: true,
// Personas point their inference adapter here; `serving_v1_url()` already
// resolves to the pinned endpoint (operator-honored verbatim).
base_url: serving_v1_url(),
adapters: Vec::new(),
served_context_window,
// One shared lane. A future refinement can read `/props.total_slots`.
lanes: 1,
degraded_reason: None,
// NOT `Some(now)`. This path ADOPTS an endpoint the operator pinned; we
// have confirmed it answers, not that it has ever delivered a token to
// us. `ready_verified_at_ms` means exactly the latter, and stamping it
// here would tell every downstream reader that a lane we have never
// pulled a token from was confirmed this instant — which is the wedged-
// slot failure the field was added to expose, manufactured by hand.
// `None` = never confirmed, which is true and which readers already
// handle.
ready_verified_at_ms: None,
// An external OpenAI-compatible endpoint does not tell us whether it can
// see, and we have not asked. Absence of knowledge, reported as absence
// — not as `false` meaning "we checked and it cannot".
vision_ready: false,
vision_base_url: None,
vision_model: None,
})
}

/// Wait until the served window has SETTLED after a serving-mode change that may have triggered
/// a grow-relaunch (e.g. an exam declaring Ludicrous/Performance). The caller must be able to
/// pin the lane WITHOUT bouncing in-flight work, so it must know the relaunch is DONE — not
Expand Down
9 changes: 8 additions & 1 deletion core/continuum-core/src/ipc/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2388,7 +2388,14 @@ pub fn start_server(
.as_ref()
.map(|p| p.fits_on_gpu)
.unwrap_or(false);
if plan_ready {
// Lane-source-agnostic hosting (misfit / grid design): when the
// operator pinned an EXTERNAL OpenAI-compatible endpoint, there is
// no local lane to "fit on GPU" — enter the ready-check regardless
// of the local plan. `await_ready_serving` short-circuits to a
// direct decode-probe of the pinned endpoint (K3, a grid peer).
let external_lane =
crate::inference::llama_server::external_serving_pin().is_some();
if plan_ready || external_lane {
// `fits_on_gpu` is a RESOURCE decision (the model fits VRAM) — it does
// NOT prove the lane can DECODE. A lane can fit yet fail EVERY
// generation with `500 "Compute error."` while `/health` still answers
Expand Down
31 changes: 31 additions & 0 deletions core/continuum-core/src/modules/serving_daemon.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1211,6 +1211,37 @@ impl ServingDaemonModule {
/// Already serving the desired model & ready → no-op. A reconcile already
/// in flight → skip (the gate). Otherwise spawn the reconcile.
fn reconcile_to_plan(&self) -> Option<JoinHandle<()>> {
// External serving pin (misfit / grid design): when the operator pinned an
// EXTERNAL OpenAI-compatible endpoint via `LLAMA_SERVER_BASE_URL`, this node
// does NOT own a local GPU serving lane — it ADOPTS the pinned endpoint. We
// spawn / reclaim NOTHING (that would fight a co-located engine, e.g. a K3
// llama-server, for its port), but we MUST publish the endpoint's ready
// ServingSnapshot to `SERVING_STATE` so every consumer sees a ready lane:
// - `await_ready_serving` (persona-host gate + adapter factory),
// - the adapter's pre-generate model-guard via `current_serving()` — which
// refuses to generate unless the request's model == the resident model,
// read straight from the published snapshot (an empty snapshot → every
// turn "model is not the active served model", the bug this fixes).
// Trust once-ready (a genuine wedge surfaces LOUD on a real turn); re-probe
// only while not-yet-ready. The reachability probe is control-plane-fast, so
// this publishes within a tick and never overlaps.
if crate::inference::llama_server::external_serving_pin().is_some() {
if self.serving_tx.borrow().ready {
return None;
}
let serving_tx = self.serving_tx.clone();
let bus = self.bus.get().cloned();
return Some(tokio::spawn(async move {
if let Some(snap) = crate::inference::llama_server::probe_external_serving(
crate::inference::llama_server::DEFAULT_SERVING_WAIT,
)
.await
{
Self::emit_serving(bus.as_ref(), &snap);
let _ = serving_tx.send_replace(snap);
}
}));
}
// Pull the desired model id, the host-fit PER-LANE served window, AND
// the lane count out of the plan in one borrow — both are the planner's
// single source of truth (task #50). We carry them on the ServingTarget
Expand Down
Loading