From 55b1324136056da9408adc3754934dda0628985a Mon Sep 17 00:00:00 2001 From: BANADDA MUBARAKA <83466862+BANADDA@users.noreply.github.com> Date: Wed, 30 Sep 2026 23:20:39 +0300 Subject: [PATCH] serving: inference miners serve full systems, validators probe them on the trace, and systems are held to a measured latency ceiling --- docs/full_system_playbook.md | 2 + docs/inference_miner.md | 41 ++++++++++++++++ docs/miner_setup.md | 11 +++++ microtensor/cli/operator.py | 45 ++++++++++++++++-- microtensor/harness/engines/gguf.py | 7 ++- microtensor/serving/agent.py | 4 ++ microtensor/serving/archived.py | 72 ++++++++++++++++++++++++++++- microtensor/serving/loop.py | 44 ++++++++++++++++-- microtensor/serving/probe.py | 44 ++++++++++++++++++ microtensor/validator/evaluate.py | 12 ++++- microtensor/validator/live.py | 20 +++++++- 11 files changed, 287 insertions(+), 15 deletions(-) diff --git a/docs/full_system_playbook.md b/docs/full_system_playbook.md index 090bbff..2df8881 100644 --- a/docs/full_system_playbook.md +++ b/docs/full_system_playbook.md @@ -24,6 +24,7 @@ The router is what turns confidence into money. It sends only the requests the s 6. **Calibrate and quantise.** Fold temperature into the file (`scripts/fold_temperature.py`) so the confidence validators recompute is already calibrated. Quantise with the output layer kept at Q8 (`--output-tensor-type q8_0`, or `--token-embedding-type q8_0` for tied embeddings). 7. **Simulate the whole system.** `mt miner simulate --escalation-url URL` runs the small model, harness, router and escalation locally over the train split and prints every trace with the scores validators will compute: end to end quality, the small model alone, escalation rate, waste, misses, calibration and cost. 8. **Package and host.** Write `system.json` at schema version 2 (model, harness with its runtime, router with its features, escalation model pinned to a revision, endpoint), then serve it through the dial out agent. No public IP is needed. +9. **Serve it to customers.** Once certified, inference miners serve your system under the arena's catalogue name on GPUs (`inference_miner.md`, section 7a), and we can serve it from the archive after you stop. ## What validators check @@ -32,6 +33,7 @@ The router is what turns confidence into money. It sends only the requests the s - **Your escalation is honest.** Only your declared, allowlisted model, charged at its published price. - **The archive reproduces the live run.** Your archived harness, router and small model are rerun on sampled tasks; prompts, tokens, decisions and final answers must match. - **Your harness stays inside its package.** No URLs, no network or process modules, no `eval`. +- **Your system is fast enough.** Validators time every request themselves; a system whose end to end p95 is over the arena's latency ceiling earns nothing that round. A system that fails any check is not certified and earns nothing that round. diff --git a/docs/inference_miner.md b/docs/inference_miner.md index 5bcc7c7..138eeb4 100644 --- a/docs/inference_miner.md +++ b/docs/inference_miner.md @@ -373,6 +373,33 @@ up. Reconnect is exponential to a 60 second ceiling. --- +## 7a · Serving full systems + +A certified model can be a full system: a small specialist, its harness, a +router and an escalation model. You serve it exactly like a model, under its +catalogue name, and every request runs the whole system on your machine. + +```bash +mt operator run \ + --serve mt/invoice-4g=sha256:@system:/srv/systems/invoice-4g \ + --escalation-url http://127.0.0.1:18090 \ + --gpu-layers -1 +``` + +- `system:` points at the certified system, either an archive record or an + artifact directory with its `manifest.json`. It is checked against the + manifest before it serves anything. +- `--escalation-url` is an OpenAI compatible server holding our mirror of the + system's escalation model (the name is `microtensor-archive/`). Run it + on the same card with SGLang, or on another card in the same box. +- `--gpu-layers -1` puts the small model on the GPU. Validators replay it on CPU; + the tolerance covers the difference. + +Each response carries the answer, the usage including escalation tokens, and a +trace signed by your hotkey: the small model's answer, confidence and tokens, +the router's features and decision, harness steps, and the escalation answer if +there was one. Clients see a single answer. + ## 8 · How you are verified You send tokens. You build no proof and compute no commitment. The validator @@ -398,6 +425,20 @@ sampled statistics that disagree. It is counted and never gates. Admission is a sequential test, not a fixed count. A clean miner is usually admitted in well under a hundred probes. One validator rejecting keeps you out. +**Full systems are judged on the trace.** The validator checks that the trace is +signed by you and names the certified system, that the answer you returned is +the traced final answer, and then: + +| check | how | +|---|---| +| small model | the traced tokens are replayed on the certified archive on CPU; a token more than 0.5 logits below the model's own choice, or a confidence off by more than 0.02, is a cheat | +| router | its features are recomputed from the replayed small model and its declared rule must give the same decision | +| escalation | only the system's declared model | +| harness | the certified harness is rerun on the same request and must return the same answer | + +Serving other weights under a certified system's name fails the first check at +once: a swapped small model shows margins of tens of logits. + ```bash mt operator status ``` diff --git a/docs/miner_setup.md b/docs/miner_setup.md index 9984517..3a70e5c 100644 --- a/docs/miner_setup.md +++ b/docs/miner_setup.md @@ -27,6 +27,16 @@ artifact, and commit a pointer on chain once per round. Validators fetch it and run it on their own certified hardware, so your machine is busy only while you are training. +In the full system arenas you submit more than a model. A system is four parts +built for one task: the specialist small model with calibrated confidence, the +harness around it (prompts, tools, checks and output templates), a router that +reads the small model's confidence and decides whether to answer, and an +escalation model from the arena's allowlist of open models that takes over only +when the router sends a request up. You host the system through the dial out +agent, validators test it live on withheld tasks, and it is ranked on end to end +quality against total cost within a latency ceiling. Start with +[full_system_playbook.md](full_system_playbook.md). + **You have a GPU and want it earning continuously without training.** You are an inference miner. You post collateral, run a stock engine on certified artifacts, and answer live traffic. You do not choose which models exist and you @@ -148,6 +158,7 @@ and nothing more. You start the next one clean. | | | |---|---| | [system_miner.md](system_miner.md) | Train, package, publish and commit, round by round | +| [full_system_playbook.md](full_system_playbook.md) | Build a full system: small model, harness, router and escalation | | [inference_miner.md](inference_miner.md) | Register, post collateral, serve and get verified | | [compute_miner.md](compute_miner.md) | Enrol a GPU machine in the pool | | [validator_setup.md](validator_setup.md) | Evaluate submissions and verify serving | diff --git a/microtensor/cli/operator.py b/microtensor/cli/operator.py index f9b46de..40bb1ab 100644 --- a/microtensor/cli/operator.py +++ b/microtensor/cli/operator.py @@ -22,7 +22,15 @@ from microtensor.serving import audit as serving_audit from microtensor.serving import client, plan, supervise from microtensor.serving import loop as probe_loop -from microtensor.serving.agent import AgentError, Pool, Served, Settings, run, served +from microtensor.serving.agent import ( + AgentError, + HttpEngine, + Pool, + Served, + Settings, + run, + served, +) from microtensor.serving.client import ServerError log = logging.getLogger("microtensor.cli.operator") @@ -253,7 +261,7 @@ def _run(args: argparse.Namespace) -> int: except (OSError, ValueError) as exc: return fail(f"could not read the pool: {exc}") - pool = Pool(settings.serves) + pool = Pool(settings.serves, build=_engines(args, wallet)) try: asyncio.run(_serve(settings, pool, systems)) except KeyboardInterrupt: @@ -263,6 +271,33 @@ def _run(args: argparse.Namespace) -> int: return 0 +SYSTEM_SCHEME = "system:" + + +def _engines(args: argparse.Namespace, wallet: Any) -> Any: + def build(url: str) -> Any: + if not url.startswith(SYSTEM_SCHEME): + return HttpEngine(url) + from microtensor.harness.sdk import openai_escalation + from microtensor.miner.host import wallet_signer + from microtensor.serving.archived import SystemEngine, open_system, runtime_for + + if not args.escalation_url: + raise AgentError("a served system escalates to our mirrors; pass --escalation-url") + restored = open_system(Path(url[len(SYSTEM_SCHEME) :]), Path(args.restore_dir)) + system = restored.manifest.system + if system is None or system.escalation is None: + raise AgentError(f"{url} is not a full system") + mirrored = f"{args.mirror_org}/{system.escalation.model.split('/', 1)[-1]}" + escalate = openai_escalation(args.escalation_url, mirrored) + runtime = runtime_for( + restored, escalate, hotkey=hotkey_address(wallet), gpu_layers=args.gpu_layers + ) + return SystemEngine(runtime, wallet_signer(wallet), url) + + return build + + def _archived(args: argparse.Namespace, wallet: Any) -> dict[str, Any]: from microtensor.harness.sdk import openai_escalation from microtensor.miner.host import wallet_signer @@ -311,10 +346,12 @@ def _verify(args: argparse.Namespace) -> int: calibrations = probe_loop.load_calibrations(args.calibrations) if not artifacts: return fail(f"no artifact paths in {args.artifacts}") - if not calibrations: + if not calibrations and not all(probe_loop.is_system(p) for p in artifacts.values()): return fail(f"no calibrations in {args.calibrations}") - missing = sorted(set(artifacts) - set(calibrations)) + missing = sorted( + m for m in set(artifacts) - set(calibrations) if not probe_loop.is_system(artifacts[m]) + ) if missing: return fail(f"no calibration for {', '.join(missing)}; an uncalibrated model cannot judge") diff --git a/microtensor/harness/engines/gguf.py b/microtensor/harness/engines/gguf.py index a65a6bc..28ab718 100644 --- a/microtensor/harness/engines/gguf.py +++ b/microtensor/harness/engines/gguf.py @@ -239,9 +239,12 @@ def __call__(self, input_ids: Any, logits: Any) -> bool: class GgufEngine: format = ArtifactFormat.GGUF - def __init__(self, *, threads: int = THREADS, validate: bool = True) -> None: + def __init__( + self, *, threads: int = THREADS, validate: bool = True, gpu_layers: int = GPU_LAYERS + ) -> None: self._threads = threads self._validate = validate + self._gpu_layers = gpu_layers self._model: Any = None self._manifest: LoadManifest | None = None self._answer_ids: dict[str, int] = {} @@ -267,7 +270,7 @@ def load(self, artifact: Path, manifest: LoadManifest) -> None: n_ctx=context, n_threads=self._threads, n_threads_batch=self._threads, - n_gpu_layers=GPU_LAYERS, + n_gpu_layers=self._gpu_layers, seed=SEED, logits_all=False, embedding=False, diff --git a/microtensor/serving/agent.py b/microtensor/serving/agent.py index 2c43e94..08998bb 100644 --- a/microtensor/serving/agent.py +++ b/microtensor/serving/agent.py @@ -527,6 +527,10 @@ async def serve_taken( answer["answered_by"] = str(found.get("answered_by", "front")) if found.get("router_features"): answer["router_features"] = dict(found["router_features"]) + if found.get("trace"): + answer["trace"] = dict(found["trace"]) + if found.get("usage"): + answer["usage"] = dict(found["usage"]) state.release(model, ok=True) return answer except asyncio.CancelledError: diff --git a/microtensor/serving/archived.py b/microtensor/serving/archived.py index 904fa5b..954b7c4 100644 --- a/microtensor/serving/archived.py +++ b/microtensor/serving/archived.py @@ -12,6 +12,7 @@ from microtensor.harness.engines.router import load_router from microtensor.harness.sdk import EscalationModel, Runtime, engine_small from microtensor.registry.manifest import ArtifactManifest, verify_tree +from microtensor.serving.agent import Engine INTAKE = "intake.json" @@ -65,13 +66,80 @@ def restore(record: Path, workdir: Path) -> Restored: return Restored(root=target, manifest=manifest, record=intake) -def runtime_for(restored: Restored, escalate: EscalationModel, *, hotkey: str) -> Runtime: +def open_system(path: Path, workdir: Path) -> Restored: + if (path / INTAKE).is_file(): + return restore(path, workdir) + manifest = ArtifactManifest.from_json((path / "manifest.json").read_bytes()) + ok, reason = verify_tree(path, manifest) + if not ok: + raise ArchiveError(f"the system files do not match their manifest: {reason}") + if manifest.system is None or not manifest.system.full: + raise ArchiveError(f"{path} is not a full system") + return Restored(root=path, manifest=manifest, record={}) + + +def prompt_of(request: Mapping[str, Any]) -> str: + prompt = str(request.get("prompt", "")) + if prompt: + return prompt + for message in reversed(list(request.get("messages") or [])): + if isinstance(message, Mapping) and message.get("role") == "user": + return str(message.get("content", "")) + raise ArchiveError("the request carries no prompt") + + +class SystemEngine(Engine): + def __init__( + self, runtime: Runtime, sign: Callable[[Mapping[str, Any]], str], url: str = "system:" + ) -> None: + super().__init__(url) + self.runtime = runtime + self.sign = sign + + async def generate(self, request: Mapping[str, Any]) -> dict[str, Any]: + import asyncio + + trace = await asyncio.to_thread( + self.runtime.run, + 0, + str(request.get("request_id") or "request"), + prompt_of(request), + dict(request.get("inputs") or {}), + ) + signed = trace.signed_with(self.sign(trace.body())) + escalation = trace.escalation + return { + "text": str(trace.final), + "prompt_tokens": [], + "completion_tokens": list(trace.small.tokens), + "finish_reason": "stop", + "escalated": trace.escalated, + "answered_by": "escalation" if trace.escalated else "small", + "router_features": dict(trace.router.features), + "trace": signed.to_dict(), + "usage": { + "prompt_tokens": trace.small.prompt_tokens + + (escalation.prompt_tokens if escalation else 0), + "completion_tokens": len(trace.small.tokens) + + (escalation.completion_tokens if escalation else 0), + }, + } + + async def stream(self, request: Mapping[str, Any], on_delta: Any) -> dict[str, Any]: + found = await self.generate(request) + await on_delta(found["text"]) + return found + + +def runtime_for( + restored: Restored, escalate: EscalationModel, *, hotkey: str, gpu_layers: int = 0 +) -> Runtime: from microtensor.harness.engines.gguf import GgufEngine system = restored.manifest.system if system is None or system.harness is None or system.escalation is None: raise ArchiveError("the archived submission is not a full system") - engine = GgufEngine() + engine = GgufEngine(gpu_layers=gpu_layers) engine.load(restored.root, restored.manifest.load) return Runtime( restored.root / system.harness.path, diff --git a/microtensor/serving/loop.py b/microtensor/serving/loop.py index b2174f6..cf2d403 100644 --- a/microtensor/serving/loop.py +++ b/microtensor/serving/loop.py @@ -8,6 +8,7 @@ from pathlib import Path from typing import Any +from microtensor.core.system import FULL_SYSTEM from microtensor.serving import probe, verify from microtensor.serving.client import ServerError @@ -94,11 +95,12 @@ def cycle( continue calibration = calibrations.get(wanted) artifact = artifacts.get(wanted) - if calibration is None or artifact is None: + system = artifact is not None and is_system(artifact) + if artifact is None or (calibration is None and not system): log.info("no calibrated artifact for %s; skipping", wanted) continue if wanted not in opened: - opened[wanted] = probe.artifact_model(artifact) + opened[wanted] = open_system(artifact) if system else probe.artifact_model(artifact) for _ in range(per_operator): verdict = _one( @@ -112,6 +114,7 @@ def cycle( model=wanted, prompt=probe.prompt_for(chance), tally=tally, + system=system, ) if verdict is None: break @@ -131,11 +134,12 @@ def _one( credential: str, wallet: Any, engine: Any, - calibration: verify.Calibration, + calibration: verify.Calibration | None, hotkey: str, model: str, prompt: str, tally: Counted, + system: bool = False, ) -> dict[str, Any] | None: try: answered = probe.ask(gateway, credential, hotkey=hotkey, model=model, prompt=prompt) @@ -144,7 +148,10 @@ def _one( log.info("operator %s did not answer: %s", hotkey[:12], exc) return None - found = probe.judge(engine, answered, calibration) + if system or calibration is None: + found = probe.judge_full_system(engine, answered, prompt) + else: + found = probe.judge(engine, answered, calibration) tally.probed += 1 if found.verdict == verify.CHEAT: @@ -191,6 +198,35 @@ def _one( return reported +def is_system(path: Path) -> bool: + if (path / "intake.json").is_file(): + return True + manifest = path / "manifest.json" + if not manifest.is_file(): + return False + try: + return ( + int( + dict(json.loads(manifest.read_text(encoding="utf-8")).get("system") or {}).get( + "schema_version", 1 + ) + ) + >= FULL_SYSTEM + ) + except (ValueError, TypeError): + return False + + +def open_system(path: Path) -> tuple[Any, Any]: + from microtensor.harness.engines.gguf import GgufEngine + from microtensor.serving.archived import open_system as restore + + restored = restore(path, path.parent / f".{path.name}-restored") + engine = GgufEngine() + engine.load(restored.root, restored.manifest.load) + return restored, engine + + def load_calibrations(path: Path) -> dict[str, verify.Calibration]: if not path.exists(): return {} diff --git a/microtensor/serving/probe.py b/microtensor/serving/probe.py index b7a983f..a23543d 100644 --- a/microtensor/serving/probe.py +++ b/microtensor/serving/probe.py @@ -61,6 +61,7 @@ class Answered: escalated: bool = False answered_by: str = "front" router_features: dict[str, float] = field(default_factory=dict) + trace: dict[str, Any] = field(default_factory=dict) @property def empty(self) -> bool: @@ -144,9 +145,52 @@ def ask( escalated=bool(found.get("escalated")), answered_by=str(found.get("answered_by", "front")), router_features={k: float(v) for k, v in dict(found.get("router_features") or {}).items()}, + trace=dict(found.get("trace") or {}), ) +def judge_full_system(held: Any, answered: Answered, prompt: str) -> verify.Judgement: + from microtensor.core.escalation import EscalationModel + from microtensor.core.trace import TraceError, read + from microtensor.core.tracks import get_track + from microtensor.validator.live import chain_verifier + from microtensor.validator.verify_system import verify_system + + restored, engine = held + system = restored.manifest.system + if not answered.trace: + return verify.Judgement(verify.CHEAT, 1.0, 0.0, "a system answer carried no trace") + try: + trace = read(answered.trace, chain_verifier()) + except TraceError as exc: + return verify.Judgement(verify.CHEAT, 1.0, 0.0, str(exc)) + if trace.hotkey != answered.hotkey or trace.system_digest != system.digest(): + reason = "the trace names another operator or system" + return verify.Judgement(verify.CHEAT, 1.0, 0.0, reason) + if str(trace.final).strip() != answered.text.strip(): + return verify.Judgement(verify.CHEAT, 1.0, 0.0, "the answer is not the traced final answer") + pinned = system.escalation + allowed = { + pinned.key: EscalationModel( + model=pinned.model, revision=pinned.revision, usd_per_mtok_in=0.0, usd_per_mtok_out=0.0 + ) + } + verdict = verify_system( + restored.root, + system, + engine, + [trace], + {trace.task_ref: (prompt, {})}, + seed=trace.task_ref, + allowlist=allowed, + chat=get_track(restored.manifest.track).chat, + size=1, + ) + if verdict.certified: + return verify.Judgement(verify.PASS, 0.0, 0.0, "") + return verify.Judgement(verify.CHEAT, 1.0, 0.0, "; ".join(verdict.reasons)) + + def judge(model: Any, answered: Answered, calibration: verify.Calibration) -> verify.Judgement: if answered.empty: return verify.Judgement( diff --git a/microtensor/validator/evaluate.py b/microtensor/validator/evaluate.py index 6606cac..f537d7a 100644 --- a/microtensor/validator/evaluate.py +++ b/microtensor/validator/evaluate.py @@ -675,6 +675,15 @@ def _evaluate_full( if not live.traces: log.info("%s scored zero: no task produced a verified trace", participant.hotkey) return _evaluation(participant, tasks, measured=measured) + p95 = live.p95_ms() + if hardware.max_p95_ms and p95 > hardware.max_p95_ms: + log.info( + "%s scored zero: end to end p95 %.0f ms is over the %d ms ceiling", + participant.hotkey, + p95, + hardware.max_p95_ms, + ) + return _evaluation(participant, tasks, measured=measured) track = get_track(tasks.track) result = run_jailed( @@ -731,6 +740,7 @@ def _evaluate_full( "profile": task.profile, "escalated": trace.escalated, "decided_at_ms": round(trace.router.at_ms, 1), + "measured_ms": round(live.latency_ms.get(task.ref, 0.0), 1), "total_ms": round(trace.total_ms, 1), "trigger": triggers.get(task.ref, ""), "features": {k: round(v, 4) for k, v in trace.router.features.items()}, @@ -789,7 +799,7 @@ def _evaluate_full( n_fixed=n_fixed, n_novel=n_novel, front_only=score.small_quality, - calibration={**score.to_dict(), "live": rows}, + calibration={**score.to_dict(), "p95_ms": round(p95, 1), "live": rows}, expected_ms=score.cost_usd * COST_UNITS_PER_USD, ) diff --git a/microtensor/validator/live.py b/microtensor/validator/live.py index b34eeba..b690501 100644 --- a/microtensor/validator/live.py +++ b/microtensor/validator/live.py @@ -1,9 +1,11 @@ from __future__ import annotations import json +import math +import time import urllib.request from collections.abc import Mapping, Sequence -from dataclasses import dataclass +from dataclasses import dataclass, field from typing import Any, Protocol from microtensor.core.system import SystemManifest @@ -74,6 +76,13 @@ def task(self, name: str, request: Mapping[str, Any]) -> Mapping[str, Any]: class LiveRun: traces: tuple[Trace, ...] failures: tuple[tuple[str, str], ...] + latency_ms: dict[str, float] = field(default_factory=dict) + + def p95_ms(self) -> float: + found = sorted(self.latency_ms.values()) + if not found: + return 0.0 + return found[min(len(found) - 1, math.ceil(0.95 * len(found)) - 1)] @property def answered(self) -> dict[str, Trace]: @@ -103,12 +112,15 @@ def run_live( digest = system.digest() traces: list[Trace] = [] failures: list[tuple[str, str]] = [] + latency: dict[str, float] = {} for task in tasks: + started = time.perf_counter() try: trace = read(client.task(system.endpoint.name, request_for(round_index, task)), verify) except (TraceError, LiveError, OSError, ValueError) as exc: failures.append((task.ref, str(exc))) continue + latency[task.ref] = (time.perf_counter() - started) * 1000.0 wrong = [ what for what, ok in ( @@ -123,7 +135,11 @@ def run_live( failures.append((task.ref, f"the trace names a different {', '.join(wrong)}")) continue traces.append(trace) - return LiveRun(traces=tuple(traces), failures=tuple(failures)) + return LiveRun( + traces=tuple(traces), + failures=tuple(failures), + latency_ms={ref: ms for ref, ms in latency.items() if ref in {t.task_ref for t in traces}}, + ) def chain_verifier() -> Verifier: