From 7621cb5fbc0f509a18f4acdac8ba736d642ffcc1 Mon Sep 17 00:00:00 2001 From: Paul Fidika Date: Mon, 20 Jul 2026 03:45:39 -0600 Subject: [PATCH 1/3] =?UTF-8?q?gw#601:=20generic=20worker-activity=20progr?= =?UTF-8?q?ess=20=E2=80=94=20ActivityUpdate=20envelope,=20evidence-gated?= =?UTF-8?q?=20watchdog=20heartbeat,=20typed=20activity=5Ffailed;=20mint/wa?= =?UTF-8?q?rmup=20phases=20(th#929=20companion)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- proto/worker_scheduler.proto | 51 ++++- src/gen_worker/activity.py | 233 +++++++++++++++++++++ src/gen_worker/executor.py | 28 ++- src/gen_worker/local_cells.py | 9 +- src/gen_worker/pb/worker_scheduler_pb2.py | 184 ++++++++-------- src/gen_worker/pb/worker_scheduler_pb2.pyi | 39 +++- tests/test_activity_gw601.py | 172 +++++++++++++++ 7 files changed, 613 insertions(+), 103 deletions(-) create mode 100644 src/gen_worker/activity.py create mode 100644 tests/test_activity_gw601.py diff --git a/proto/worker_scheduler.proto b/proto/worker_scheduler.proto index a42cc977..0836389e 100644 --- a/proto/worker_scheduler.proto +++ b/proto/worker_scheduler.proto @@ -41,14 +41,15 @@ service WorkerScheduler { message WorkerMessage { oneof msg { - Hello hello = 1; - StateDelta state_delta = 2; - JobAccepted job_accepted = 3; - JobResult job_result = 4; - JobProgress job_progress = 5; - ModelEvent model_event = 6; - FnUnavailable fn_unavailable = 7; - FnDegraded fn_degraded = 8; + Hello hello = 1; + StateDelta state_delta = 2; + JobAccepted job_accepted = 3; + JobResult job_result = 4; + JobProgress job_progress = 5; + ModelEvent model_event = 6; + FnUnavailable fn_unavailable = 7; + FnDegraded fn_degraded = 8; + ActivityUpdate activity_update = 9; } } @@ -575,6 +576,40 @@ enum ModelState { MODEL_STATE_HOST_CAPACITY_PROGRESS = 8; } +// --------------------------------------------------------------------------- +// Worker activity (th#929 / gw#601) +// --------------------------------------------------------------------------- + +// Worker -> orchestrator: generic progress for the worker's current internal +// job (self-mint compile, warmup, ...). Liveness contract: while an activity +// RUNS the worker keeps seq advancing — phase transitions, step advances, or +// watchdog heartbeats — and the hub enforces ONE generic stall rule on +// silence; there are no per-activity wall-clock budgets. A silent death is a +// bug by contract: terminal FAILED carries the exception. Model downloads +// keep their richer ModelEvent byte shape; the hub folds them into the same +// last-progress store. Additive proto3 oneof field: old hubs ignore it. +message ActivityUpdate { + string kind = 1; // "self_mint_compile" | "warmup" | ... + string phase = 2; // e.g. "load", "inductor_compile", "warmup_forward", "seal_publish" + int64 step = 3; // k in "phase k/K"; 0 = not stepwise + int64 total_steps = 4; // K; 0 = unknown + // Monotonic per worker process across all activities. The hub records + // last-progress on every increase; heartbeats bump it without changing + // phase/step. + uint64 seq = 5; + ActivityState state = 6; + string error = 7; // set when state = FAILED (the typed activity_failed terminal) + string detail = 8; // human-readable elaboration + int64 updated_at_unix_ms = 9; // worker clock, informational; hub stamps receipt +} + +enum ActivityState { + ACTIVITY_STATE_UNSPECIFIED = 0; + ACTIVITY_STATE_RUNNING = 1; + ACTIVITY_STATE_COMPLETED = 2; + ACTIVITY_STATE_FAILED = 3; +} + // --------------------------------------------------------------------------- // Availability / drain // --------------------------------------------------------------------------- diff --git a/src/gen_worker/activity.py b/src/gen_worker/activity.py new file mode 100644 index 00000000..5797ca06 --- /dev/null +++ b/src/gen_worker/activity.py @@ -0,0 +1,233 @@ +"""Generic worker-activity progress (gw#601 / th#929). + +The worker reports whatever internal job it is doing — self-mint compile, +warmup — as ActivityUpdate envelopes on the existing worker->hub stream. +Liveness contract: while an activity runs, seq keeps advancing (phase +transitions, step advances, or watchdog heartbeats around long silent calls); +the hub enforces ONE generic stall rule on silence. A silent death is a bug +by contract: terminal FAILED carries the exception. + +Kind/phase strings are wire-shared with tensorhub +(internal/orchestrator/grpc/worker_activity.go) — keep them identical. + +Without a bound transport sink (cozy-local, tests) reports land on the +logger, which IS the local progress UI. +""" + +from __future__ import annotations + +import asyncio +import logging +import threading +import time +from typing import Callable, Optional + +from .pb import worker_scheduler_pb2 as pb + +logger = logging.getLogger(__name__) + +KIND_SELF_MINT_COMPILE = "self_mint_compile" +KIND_WARMUP = "warmup" + +PHASE_LOAD = "load" +PHASE_TRACE_GRAPH = "trace_graph" +PHASE_INDUCTOR_COMPILE = "inductor_compile" +PHASE_WARMUP_FORWARD = "warmup_forward" +PHASE_SEAL_PUBLISH = "seal_publish" + +# Default watchdog cadence; the hub's stall rule (~10 min) tolerates many +# missed beats. +HEARTBEAT_INTERVAL_S = 60.0 +# Minimum evidence advance (process CPU seconds) per interval for a heartbeat: +# a hung (blocked/deadlocked) call stops accruing CPU and the beat stops with +# it, which is exactly the silence the hub enforces on. +_EVIDENCE_EPS = 0.05 + +_lock = threading.Lock() +_seq = 0 +_sink: Optional[Callable[[pb.ActivityUpdate], None]] = None +_current: Optional["Activity"] = None + + +def bind_sink(emit, loop: asyncio.AbstractEventLoop) -> None: + """Route reports onto the worker->hub stream: emit is the async + WorkerMessage sender, loop the transport loop. Thread-safe emission.""" + def sink(update: pb.ActivityUpdate) -> None: + coro = emit(pb.WorkerMessage(activity_update=update)) + try: + running = asyncio.get_running_loop() + except RuntimeError: + running = None + if running is loop: + loop.create_task(coro) + elif not loop.is_closed(): + asyncio.run_coroutine_threadsafe(coro, loop) + else: + coro.close() + global _sink + with _lock: + _sink = sink + + +def _next_seq() -> int: + global _seq + with _lock: + _seq += 1 + return _seq + + +def _emit(update: pb.ActivityUpdate) -> None: + with _lock: + sink = _sink + try: + if sink is not None: + sink(update) + else: + state = pb.ActivityState.Name(update.state) + logger.info( + "[activity] %s %s %s/%s %s %s", update.kind, update.phase, + update.step, update.total_steps, state, update.error or update.detail, + ) + except Exception: # reporting must never break the work it reports on + logger.debug("activity report dropped", exc_info=True) + + +class Activity: + """One running activity. Use begin() / the context manager `running()`.""" + + def __init__(self, kind: str) -> None: + self.kind = kind + self._phase = "" + self._step = 0 + self._total = 0 + self._done = False + + def _report(self, state: int, error: str = "", detail: str = "") -> None: + _emit(pb.ActivityUpdate( + kind=self.kind, phase=self._phase, step=self._step, + total_steps=self._total, seq=_next_seq(), state=state, + error=error, detail=detail, + updated_at_unix_ms=int(time.time() * 1000), + )) + + def phase(self, phase: str, step: int = 0, total: int = 0) -> None: + self._phase, self._step, self._total = phase, step, total + self._report(pb.ActivityState.ACTIVITY_STATE_RUNNING) + + def heartbeat(self) -> None: + """Re-report the current phase with a fresh seq (liveness proof).""" + self._report(pb.ActivityState.ACTIVITY_STATE_RUNNING) + + def completed(self) -> None: + if not self._done: + self._done = True + self._report(pb.ActivityState.ACTIVITY_STATE_COMPLETED) + _end(self) + + def failed(self, exc: BaseException) -> None: + """The typed activity_failed terminal — a silent death is a bug.""" + if not self._done: + self._done = True + self._report( + pb.ActivityState.ACTIVITY_STATE_FAILED, + error=f"{type(exc).__name__}: {exc}"[:2000], + ) + _end(self) + + +def begin(kind: str, phase: str = "") -> Activity: + global _current + act = Activity(kind) + with _lock: + _current = act + act.phase(phase) if phase else act.heartbeat() + return act + + +def current_phase(phase: str, step: int = 0, total: int = 0) -> None: + """Report a phase on the current activity; no-op when none is running. + Setups serialize under the executor load lock, so one current is enough.""" + with _lock: + act = _current + if act is not None and not act._done: + act.phase(phase, step, total) + + +def _end(act: Activity) -> None: + global _current + with _lock: + if _current is act: + _current = None + + +class running: + """Context manager: begin() on enter; COMPLETED on clean exit, FAILED + (carrying the exception) on raise.""" + + def __init__(self, kind: str, phase: str = "") -> None: + self._kind, self._phase = kind, phase + self.activity: Optional[Activity] = None + + def __enter__(self) -> Activity: + self.activity = begin(self._kind, self._phase) + return self.activity + + def __exit__(self, exc_type, exc, tb) -> None: + assert self.activity is not None + if exc is not None: + self.activity.failed(exc) + else: + self.activity.completed() + _end(self.activity) + + +def _process_cpu_evidence() -> float: + return time.process_time() + + +class watchdog: + """Bracket for a long call that may stay wire-silent (inductor compile, + large fuse): a background thread samples an evidence counter every + interval and heartbeats the activity ONLY while evidence advances. A hung + call stops accruing evidence, the beat stops within one interval, and the + hub's stall rule takes it from there. + + Default evidence is process CPU seconds; pass a monotonic callable + (e.g. compile-wall-seconds) for calls with better signals.""" + + def __init__( + self, + act: Activity, + *, + interval_s: float = HEARTBEAT_INTERVAL_S, + evidence: Optional[Callable[[], float]] = None, + ) -> None: + self._act = act + self._interval = interval_s + self._evidence = evidence or _process_cpu_evidence + self._stop = threading.Event() + self._thread = threading.Thread( + target=self._run, name="activity-watchdog", daemon=True, + ) + + def _run(self) -> None: + try: + last = self._evidence() + except Exception: + last = 0.0 + while not self._stop.wait(self._interval): + try: + now = self._evidence() + except Exception: + continue + if now - last >= _EVIDENCE_EPS: + last = now + self._act.heartbeat() + + def __enter__(self) -> "watchdog": + self._thread.start() + return self + + def __exit__(self, *exc) -> None: + self._stop.set() + self._thread.join(timeout=5) diff --git a/src/gen_worker/executor.py b/src/gen_worker/executor.py index be6a28b7..4ec332a8 100644 --- a/src/gen_worker/executor.py +++ b/src/gen_worker/executor.py @@ -29,6 +29,7 @@ import msgspec +from . import activity as activity_mod from .api.binding import ModelRef, wire_ref from .api.errors import ( ArtifactTransferError, @@ -2979,6 +2980,10 @@ async def ensure_setup( if spec.cls is None: return None # function-shaped endpoint: no instance, no setup self.store.bind_loop() + try: + activity_mod.bind_sink(self._send, asyncio.get_running_loop()) + except RuntimeError: + pass rec = self._class_record(spec) async with rec.lock: if rec.ready and not rec.stale: @@ -3061,9 +3066,20 @@ async def ensure_setup( spec.name, ) await self._rollback_failed_setup(rec) + # gw#601: setup+warmup is one reportable activity. The watchdog + # heartbeats through long wire-silent calls (inductor etc.) while + # they provably burn CPU; a hang stops the beat within one + # interval and the hub's generic stall rule owns termination. + act = activity_mod.begin( + activity_mod.KIND_SELF_MINT_COMPILE if spec.compile is not None + else activity_mod.KIND_WARMUP, + activity_mod.PHASE_LOAD, + ) try: - instance = await self._setup_locked(spec, rec, snapshots) + with activity_mod.watchdog(act): + instance = await self._setup_locked(spec, rec, snapshots) except BaseException as exc: + act.failed(exc) # Setup is a transaction: endpoint construction, tenant # setup/warmup, residency registration, and compile-target # publication either all reach READY or all ownership is @@ -3096,6 +3112,7 @@ async def ensure_setup( self.unavailable.pop(s.name, None) rec.instance = instance rec.ready = True + act.completed() self._clear_recovered_compile_failures(rec) self._on_state_change() return instance @@ -3179,7 +3196,9 @@ async def _run_synthesized_warmup( logger.info("boot warmup skipped for %s: %s", skip.spec.name, skip.reason) objects = tuple({id(obj): obj for obj in proof_objects}.values()) evidence = _WarmupEvidence() - for wj in jobs: + for wj_index, wj in enumerate(jobs, start=1): + activity_mod.current_phase( + activity_mod.PHASE_WARMUP_FORWARD, wj_index, len(jobs)) before = { id(obj): ( compile_cache.execution_count(obj), @@ -3463,6 +3482,9 @@ async def run_warmup() -> Tuple[int, Dict[int, set[str]]]: ) return evidence.count, evidence.functions_by_object + activity_mod.current_phase( + activity_mod.PHASE_INDUCTOR_COMPILE if inj.pending_self_mints + else activity_mod.PHASE_WARMUP_FORWARD) compile_seconds_before = ( compile_cache.compile_wall_seconds() if proves_inductor else 0.0) if inj.active_compile_artifacts: @@ -3517,6 +3539,8 @@ async def run_warmup() -> Tuple[int, Dict[int, set[str]]]: # advertising/publishing this boot's capture. from . import fleet_cells as fleet_cells_mod + activity_mod.current_phase( + activity_mod.PHASE_SEAL_PUBLISH) finalized = fleet_cells_mod.finalize_self_mint( pipe, pending_mint) inj.pending_self_mints.pop(id(pipe), None) diff --git a/src/gen_worker/local_cells.py b/src/gen_worker/local_cells.py index 9739d20d..34640127 100644 --- a/src/gen_worker/local_cells.py +++ b/src/gen_worker/local_cells.py @@ -58,6 +58,7 @@ from pathlib import Path from typing import Any, Optional +from . import activity as activity_mod from . import compile_cache as cc logger = logging.getLogger(__name__) @@ -171,7 +172,13 @@ def _mint(pipe: Any, cfg: Any, target: Path, family: str) -> Path: _sweep_stale_mints(root) capture = root / _MINT_DIR / f"{target.name[: -len('.tar.gz')]}-{os.getpid()}" started = time.monotonic() - cc.mint_artifact(pipe, cfg, family, target, capture, say=_say) + # gw#601: the local mint reports the same activity envelope as the fleet + # path — without a transport sink it lands on the logger (the local UI), + # heartbeating through the long silent compile. + with activity_mod.running( + activity_mod.KIND_SELF_MINT_COMPILE, activity_mod.PHASE_INDUCTOR_COMPILE, + ) as act, activity_mod.watchdog(act): + cc.mint_artifact(pipe, cfg, family, target, capture, say=_say) _say( f"compile cell saved: {target} " f"({target.stat().st_size / 1e6:.1f} MB, {time.monotonic() - started:.0f}s total); " diff --git a/src/gen_worker/pb/worker_scheduler_pb2.py b/src/gen_worker/pb/worker_scheduler_pb2.py index 749fd6d8..ee9d363c 100644 --- a/src/gen_worker/pb/worker_scheduler_pb2.py +++ b/src/gen_worker/pb/worker_scheduler_pb2.py @@ -24,7 +24,7 @@ -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x16worker_scheduler.proto\x12\x0e\x63ozy.scheduler\"\xab\x03\n\rWorkerMessage\x12&\n\x05hello\x18\x01 \x01(\x0b\x32\x15.cozy.scheduler.HelloH\x00\x12\x31\n\x0bstate_delta\x18\x02 \x01(\x0b\x32\x1a.cozy.scheduler.StateDeltaH\x00\x12\x33\n\x0cjob_accepted\x18\x03 \x01(\x0b\x32\x1b.cozy.scheduler.JobAcceptedH\x00\x12/\n\njob_result\x18\x04 \x01(\x0b\x32\x19.cozy.scheduler.JobResultH\x00\x12\x33\n\x0cjob_progress\x18\x05 \x01(\x0b\x32\x1b.cozy.scheduler.JobProgressH\x00\x12\x31\n\x0bmodel_event\x18\x06 \x01(\x0b\x32\x1a.cozy.scheduler.ModelEventH\x00\x12\x37\n\x0e\x66n_unavailable\x18\x07 \x01(\x0b\x32\x1d.cozy.scheduler.FnUnavailableH\x00\x12\x31\n\x0b\x66n_degraded\x18\x08 \x01(\x0b\x32\x1a.cozy.scheduler.FnDegradedH\x00\x42\x05\n\x03msg\"\xb0\x02\n\x10SchedulerMessage\x12-\n\thello_ack\x18\x01 \x01(\x0b\x32\x18.cozy.scheduler.HelloAckH\x00\x12)\n\x07run_job\x18\x02 \x01(\x0b\x32\x16.cozy.scheduler.RunJobH\x00\x12/\n\ncancel_job\x18\x03 \x01(\x0b\x32\x19.cozy.scheduler.CancelJobH\x00\x12+\n\x08model_op\x18\x04 \x01(\x0b\x32\x17.cozy.scheduler.ModelOpH\x00\x12&\n\x05\x64rain\x18\x05 \x01(\x0b\x32\x15.cozy.scheduler.DrainH\x00\x12\x35\n\rtoken_refresh\x18\x06 \x01(\x0b\x32\x1c.cozy.scheduler.TokenRefreshH\x00\x42\x05\n\x03msg\"\xa8\x02\n\x05Hello\x12\x39\n\x10protocol_version\x18\x01 \x01(\x0e\x32\x1f.cozy.scheduler.ProtocolVersion\x12\x11\n\tworker_id\x18\x02 \x01(\t\x12\x12\n\nrelease_id\x18\x03 \x01(\t\x12\x32\n\tresources\x18\x04 \x01(\x0b\x32\x1f.cozy.scheduler.WorkerResources\x12)\n\x05state\x18\x05 \x01(\x0b\x32\x1a.cozy.scheduler.StateDelta\x12.\n\x06models\x18\x06 \x03(\x0b\x32\x1e.cozy.scheduler.ModelResidency\x12.\n\tin_flight\x18\x07 \x03(\x0b\x32\x1b.cozy.scheduler.InFlightJob\"\x9b\x02\n\x0fWorkerResources\x12\x11\n\tgpu_count\x18\x01 \x01(\x05\x12\x18\n\x10vram_total_bytes\x18\x02 \x01(\x03\x12\x10\n\x08gpu_name\x18\x03 \x01(\t\x12\x0e\n\x06gpu_sm\x18\x04 \x01(\t\x12\x16\n\x0einstalled_libs\x18\x05 \x03(\t\x12\x14\n\x0cimage_digest\x18\x06 \x01(\t\x12\x12\n\ngit_commit\x18\x07 \x01(\t\x12\x13\n\x0binstance_id\x18\x08 \x01(\t\x12/\n\x0bhost_canary\x18\t \x01(\x0b\x32\x1a.cozy.scheduler.HostCanary\x12\x15\n\rtorch_version\x18\n \x01(\t\x12\x1a\n\x12gen_worker_version\x18\x0b \x01(\t\"\xc9\x01\n\nHostCanary\x12\x13\n\x0bmemcpy_gbps\x18\x01 \x01(\x01\x12\x10\n\x08h2d_gbps\x18\x02 \x01(\x01\x12\x10\n\x08\x64\x32h_gbps\x18\x03 \x01(\x01\x12\x17\n\x0fpinned_alloc_ok\x18\x04 \x01(\x08\x12\x17\n\x0f\x63pu_single_mbps\x18\x05 \x01(\x01\x12\x16\n\x0e\x63pu_multi_mbps\x18\x06 \x01(\x01\x12\r\n\x05vcpus\x18\x07 \x01(\x05\x12\x14\n\x0cram_total_gb\x18\x08 \x01(\x01\x12\x13\n\x0b\x64uration_ms\x18\t \x01(\x03\"\x95\x01\n\x0eModelResidency\x12\x0b\n\x03ref\x18\x01 \x01(\t\x12+\n\x04tier\x18\x02 \x01(\x0e\x32\x1d.cozy.scheduler.ResidencyTier\x12\x12\n\nvram_bytes\x18\x03 \x01(\x03\x12\x17\n\x0fsnapshot_digest\x18\x04 \x01(\t\x12\x1c\n\x14residency_generation\x18\x05 \x01(\x04\"2\n\x0bInFlightJob\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\"\xdd\x01\n\x08HelloAck\x12\x39\n\x10protocol_version\x18\x01 \x01(\x0e\x32\x1f.cozy.scheduler.ProtocolVersion\x12\x15\n\rfile_base_url\x18\x02 \x01(\t\x12\x0c\n\x04keep\x18\x03 \x03(\t\x12\x34\n\x0bresolutions\x18\x04 \x03(\x0b\x32\x1f.cozy.scheduler.ModelResolution\x12;\n\x11\x64\x65sired_residency\x18\x05 \x01(\x0b\x32 .cozy.scheduler.DesiredResidency\"\xf7\x01\n\x10\x44\x65siredResidency\x12\x12\n\ngeneration\x18\x01 \x01(\x04\x12\x11\n\tdisk_refs\x18\x02 \x03(\t\x12,\n\x03hot\x18\x03 \x03(\x0b\x32\x1f.cozy.scheduler.DesiredInstance\x12\x42\n\tsnapshots\x18\x04 \x03(\x0b\x32/.cozy.scheduler.DesiredResidency.SnapshotsEntry\x1aJ\n\x0eSnapshotsEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\'\n\x05value\x18\x02 \x01(\x0b\x32\x18.cozy.scheduler.Snapshot:\x02\x38\x01\"V\n\x0f\x44\x65siredInstance\x12\x15\n\rfunction_name\x18\x01 \x01(\t\x12,\n\x06models\x18\x02 \x03(\x0b\x32\x1c.cozy.scheduler.ModelBinding\"B\n\x0fModelResolution\x12\x0b\n\x03ref\x18\x01 \x01(\t\x12\x14\n\x0cresolved_ref\x18\x02 \x01(\t\x12\x0c\n\x04\x63\x61st\x18\x03 \x01(\t\"\xb3\x02\n\nStateDelta\x12*\n\x05phase\x18\x01 \x01(\x0e\x32\x1b.cozy.scheduler.WorkerPhase\x12\x1b\n\x13\x61vailable_functions\x18\x02 \x03(\t\x12\x19\n\x11loading_functions\x18\x03 \x03(\t\x12\x17\n\x0f\x66ree_vram_bytes\x18\x04 \x01(\x03\x12\x17\n\x0f\x66inalizing_jobs\x18\x05 \x01(\x05\x12%\n\x1dobserved_residency_generation\x18\x06 \x01(\x04\x12\x36\n\x0f\x63ompile_targets\x18\x07 \x03(\x0b\x32\x1d.cozy.scheduler.CompileTarget\x12\x30\n\x0c\x63\x65ll_lookups\x18\x08 \x03(\x0b\x32\x1a.cozy.scheduler.CellLookup\".\n\nCellLookup\x12\x0e\n\x06\x66\x61mily\x18\x01 \x01(\t\x12\x10\n\x08\x63\x65ll_key\x18\x02 \x01(\t\"\xc6\x03\n\rCompileTarget\x12\x16\n\x0eincarnation_id\x18\x01 \x01(\t\x12\x0e\n\x06\x66\x61mily\x18\x02 \x01(\t\x12\x1c\n\x14pipeline_weight_lane\x18\x03 \x01(\t\x12\x13\n\x0blora_bucket\x18\x04 \x01(\x05\x12\x17\n\x0f\x63ontract_digest\x18\x05 \x01(\t\x12\x1a\n\x12\x61\x63tive_compile_ref\x18\x06 \x01(\t\x12&\n\x1e\x61\x63tive_compile_snapshot_digest\x18\x07 \x01(\t\x12\x16\n\x0e\x66unction_names\x18\x08 \x03(\t\x12<\n\x0emodel_bindings\x18\t \x03(\x0b\x32$.cozy.scheduler.CompileTargetBinding\x12\x1a\n\x12requested_cell_key\x18\n \x01(\t\x12Q\n\x13requested_cell_axes\x18\x0b \x03(\x0b\x32\x34.cozy.scheduler.CompileTarget.RequestedCellAxesEntry\x1a\x38\n\x16RequestedCellAxesEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\"J\n\x14\x43ompileTargetBinding\x12\x0c\n\x04slot\x18\x01 \x01(\t\x12\x0b\n\x03ref\x18\x02 \x01(\t\x12\x17\n\x0fsnapshot_digest\x18\x03 \x01(\t\"\x88\x04\n\x06RunJob\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\x12\x15\n\rfunction_name\x18\x03 \x01(\t\x12\x15\n\rinput_payload\x18\x04 \x01(\x0c\x12\x12\n\ntimeout_ms\x18\x05 \x01(\x03\x12\x0e\n\x06tenant\x18\x06 \x01(\t\x12\x12\n\ninvoker_id\x18\x07 \x01(\t\x12\x18\n\x10\x63\x61pability_token\x18\x08 \x01(\t\x12/\n\x0boutput_mode\x18\t \x01(\x0e\x32\x1a.cozy.scheduler.OutputMode\x12\x30\n\x07\x63ompute\x18\n \x01(\x0b\x32\x1f.cozy.scheduler.ResolvedCompute\x12,\n\x06models\x18\x0b \x03(\x0b\x32\x1c.cozy.scheduler.ModelBinding\x12\x38\n\tsnapshots\x18\x0c \x03(\x0b\x32%.cozy.scheduler.RunJob.SnapshotsEntry\x12\x42\n\x10required_compile\x18\r \x01(\x0b\x32(.cozy.scheduler.RequiredCompileExecution\x1aJ\n\x0eSnapshotsEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\'\n\x05value\x18\x02 \x01(\x0b\x32\x18.cozy.scheduler.Snapshot:\x02\x38\x01\"\x82\x01\n\x18RequiredCompileExecution\x12\x1d\n\x15target_incarnation_id\x18\x01 \x01(\t\x12\x10\n\x08\x63\x65ll_ref\x18\x02 \x01(\t\x12\x1c\n\x14\x63\x65ll_snapshot_digest\x18\x03 \x01(\t\x12\x17\n\x0f\x63ontract_digest\x18\x04 \x01(\t\"Y\n\x0fResolvedCompute\x12\x13\n\x0b\x61\x63\x63\x65lerator\x18\x01 \x01(\t\x12\x11\n\tgpu_index\x18\x02 \x01(\x05J\x04\x08\x03\x10\x04J\x04\x08\x04\x10\x05R\tgpu_countR\x07vram_gb\"q\n\x0cModelBinding\x12\x0c\n\x04slot\x18\x01 \x01(\t\x12\x0b\n\x03ref\x18\x02 \x01(\t\x12*\n\x05loras\x18\x03 \x03(\x0b\x32\x1b.cozy.scheduler.LoraOverlay\x12\x1a\n\x12inference_defaults\x18\x04 \x01(\t\"F\n\x0bLoraOverlay\x12\x0b\n\x03ref\x18\x01 \x01(\t\x12\x0e\n\x06weight\x18\x02 \x01(\x01\x12\x1a\n\x12inference_defaults\x18\x03 \x01(\t\"G\n\x08Snapshot\x12\x0e\n\x06\x64igest\x18\x01 \x01(\t\x12+\n\x05\x66iles\x18\x02 \x03(\x0b\x32\x1c.cozy.scheduler.SnapshotFile\"M\n\x0cSnapshotFile\x12\x0c\n\x04path\x18\x01 \x01(\t\x12\x12\n\nsize_bytes\x18\x02 \x01(\x03\x12\x0e\n\x06\x62lake3\x18\x03 \x01(\t\x12\x0b\n\x03url\x18\x04 \x01(\t\"2\n\x0bJobAccepted\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\"\xce\x01\n\tJobResult\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\x12)\n\x06status\x18\x03 \x01(\x0e\x32\x19.cozy.scheduler.JobStatus\x12\x10\n\x06inline\x18\x04 \x01(\x0cH\x00\x12\x12\n\x08\x62lob_ref\x18\x05 \x01(\tH\x00\x12\x14\n\x0csafe_message\x18\x06 \x01(\t\x12+\n\x07metrics\x18\x07 \x01(\x0b\x32\x1a.cozy.scheduler.JobMetricsB\x08\n\x06output\"\xb4\x02\n\nJobMetrics\x12\x12\n\nruntime_ms\x18\x01 \x01(\x03\x12\x10\n\x08queue_ms\x18\x02 \x01(\x03\x12\x18\n\x10rss_at_end_bytes\x18\x03 \x01(\x03\x12\x17\n\x0fpeak_vram_bytes\x18\x04 \x01(\x03\x12\x1c\n\x14\x63oncurrency_at_start\x18\x05 \x01(\x05\x12\x1f\n\x17output_media_duration_s\x18\x06 \x01(\x01\x12\x14\n\x0cinput_tokens\x18\x07 \x01(\x03\x12\x1b\n\x13input_cached_tokens\x18\x08 \x01(\x03\x12\x15\n\routput_tokens\x18\t \x01(\x03\x12\x14\n\x0coutput_count\x18\n \x01(\x03\x12\x14\n\x0cslot_held_ms\x18\x0b \x01(\x03\x12\x18\n\x10\x66inalize_wall_ms\x18\x0c \x01(\x03\"c\n\x0bJobProgress\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\x12\x0b\n\x03seq\x18\x03 \x01(\x03\x12\x0c\n\x04\x64\x61ta\x18\x04 \x01(\x0c\x12\x14\n\x0c\x63ontent_type\x18\x05 \x01(\t\"0\n\tCancelJob\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\"\xa0\x01\n\x07ModelOp\x12\'\n\x02op\x18\x01 \x01(\x0e\x32\x1b.cozy.scheduler.ModelOpKind\x12\x0b\n\x03ref\x18\x02 \x01(\t\x12*\n\x08snapshot\x18\x03 \x01(\x0b\x32\x18.cozy.scheduler.Snapshot\x12\x14\n\x0coperation_id\x18\x04 \x01(\t\x12\x1d\n\x15target_incarnation_id\x18\x05 \x01(\t\"\x9b\x04\n\nModelEvent\x12\x0b\n\x03ref\x18\x01 \x01(\t\x12)\n\x05state\x18\x02 \x01(\x0e\x32\x1a.cozy.scheduler.ModelState\x12\x12\n\nvram_bytes\x18\x03 \x01(\x03\x12\r\n\x05\x65rror\x18\x04 \x01(\t\x12\x12\n\nbytes_done\x18\x05 \x01(\x03\x12\x13\n\x0b\x62ytes_total\x18\x06 \x01(\x03\x12\x13\n\x0b\x64uration_ms\x18\x07 \x01(\x03\x12\x12\n\ncache_hits\x18\x08 \x01(\x03\x12\x14\n\x0c\x63\x61\x63he_misses\x18\t \x01(\x03\x12\x10\n\x08warmup_s\x18\n \x01(\x01\x12\x1f\n\x17host_ram_required_bytes\x18\x0b \x01(\x03\x12\'\n\x1fhost_ram_available_before_bytes\x18\x0c \x01(\x03\x12&\n\x1ehost_ram_available_after_bytes\x18\r \x01(\x03\x12\x1d\n\x15host_ram_evicted_refs\x18\x0e \x03(\t\x12$\n\x1chost_ram_capacity_generation\x18\x0f \x01(\x04\x12\x17\n\x0fsnapshot_digest\x18\x10 \x01(\t\x12\x1c\n\x14residency_generation\x18\x11 \x01(\x04\x12\x14\n\x0coperation_id\x18\x12 \x01(\t\x12\x1d\n\x15target_incarnation_id\x18\x13 \x01(\t\x12\x15\n\rnetwork_bytes\x18\x14 \x01(\x03\"\xaa\x01\n\rFnUnavailable\x12\x15\n\rfunction_name\x18\x01 \x01(\t\x12\x0e\n\x06reason\x18\x02 \x01(\t\x12\x0e\n\x06\x64\x65tail\x18\x03 \x01(\t\x12\x35\n\x04\x61xes\x18\x04 \x03(\x0b\x32\'.cozy.scheduler.FnUnavailable.AxesEntry\x1a+\n\tAxesEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\"\x8d\x01\n\nFnDegraded\x12\x15\n\rfunction_name\x18\x01 \x01(\t\x12\x0e\n\x06wanted\x18\x02 \x01(\t\x12\x0b\n\x03ran\x18\x03 \x01(\t\x12\x0e\n\x06reason\x18\x04 \x01(\t\x12\x1e\n\x16\x65st_latency_multiplier\x18\x05 \x01(\x01\x12\x1b\n\x13recommended_vram_gb\x18\x06 \x01(\x01\"\x1c\n\x05\x44rain\x12\x13\n\x0b\x64\x65\x61\x64line_ms\x18\x01 \x01(\x03\"6\n\x0cTokenRefresh\x12\r\n\x05token\x18\x01 \x01(\t\x12\x17\n\x0f\x65xpires_at_unix\x18\x02 \x01(\x03*Q\n\x0fProtocolVersion\x12 \n\x1cPROTOCOL_VERSION_UNSPECIFIED\x10\x00\x12\x1c\n\x18PROTOCOL_VERSION_CURRENT\x10\x03*y\n\rResidencyTier\x12\x1e\n\x1aRESIDENCY_TIER_UNSPECIFIED\x10\x00\x12\x17\n\x13RESIDENCY_TIER_DISK\x10\x01\x12\x16\n\x12RESIDENCY_TIER_RAM\x10\x02\x12\x17\n\x13RESIDENCY_TIER_VRAM\x10\x03*\xd8\x01\n\x0bWorkerPhase\x12\x1c\n\x18WORKER_PHASE_UNSPECIFIED\x10\x00\x12\x18\n\x14WORKER_PHASE_BOOTING\x10\x01\x12#\n\x1fWORKER_PHASE_DOWNLOADING_MODELS\x10\x02\x12\"\n\x1eWORKER_PHASE_LOADING_PIPELINES\x10\x03\x12\x18\n\x14WORKER_PHASE_WARMING\x10\x04\x12\x16\n\x12WORKER_PHASE_READY\x10\x05\x12\x16\n\x12WORKER_PHASE_ERROR\x10\x06*V\n\nOutputMode\x12\x1b\n\x17OUTPUT_MODE_UNSPECIFIED\x10\x00\x12\x13\n\x0fOUTPUT_MODE_URL\x10\x01\x12\x16\n\x12OUTPUT_MODE_INLINE\x10\x02*\x9b\x01\n\tJobStatus\x12\x1a\n\x16JOB_STATUS_UNSPECIFIED\x10\x00\x12\x11\n\rJOB_STATUS_OK\x10\x01\x12\x16\n\x12JOB_STATUS_INVALID\x10\x02\x12\x18\n\x14JOB_STATUS_RETRYABLE\x10\x03\x12\x14\n\x10JOB_STATUS_FATAL\x10\x04\x12\x17\n\x13JOB_STATUS_CANCELED\x10\x05*\xa7\x01\n\x0bModelOpKind\x12\x1d\n\x19MODEL_OP_KIND_UNSPECIFIED\x10\x00\x12%\n!MODEL_OP_KIND_ADOPT_COMPILE_CACHE\x10\x04\"\x04\x08\x01\x10\x01\"\x04\x08\x02\x10\x02\"\x04\x08\x03\x10\x03*\x16MODEL_OP_KIND_DOWNLOAD*\x12MODEL_OP_KIND_LOAD*\x14MODEL_OP_KIND_UNLOAD*\x82\x02\n\nModelState\x12\x1b\n\x17MODEL_STATE_UNSPECIFIED\x10\x00\x12\x1b\n\x17MODEL_STATE_DOWNLOADING\x10\x01\x12\x17\n\x13MODEL_STATE_ON_DISK\x10\x02\x12\x16\n\x12MODEL_STATE_IN_RAM\x10\x03\x12\x17\n\x13MODEL_STATE_IN_VRAM\x10\x04\x12\x17\n\x13MODEL_STATE_EVICTED\x10\x05\x12\x16\n\x12MODEL_STATE_FAILED\x10\x06\x12\x17\n\x13MODEL_STATE_ADOPTED\x10\x07\x12&\n\"MODEL_STATE_HOST_CAPACITY_PROGRESS\x10\x08\x32\x61\n\x0fWorkerScheduler\x12N\n\x07\x43onnect\x12\x1d.cozy.scheduler.WorkerMessage\x1a .cozy.scheduler.SchedulerMessage(\x01\x30\x01\x42\x61Z_github.com/cozy-creator/tensorhub/internal/orchestrator/grpc/pb/workerscheduler;workerschedulerb\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x16worker_scheduler.proto\x12\x0e\x63ozy.scheduler\"\xe6\x03\n\rWorkerMessage\x12&\n\x05hello\x18\x01 \x01(\x0b\x32\x15.cozy.scheduler.HelloH\x00\x12\x31\n\x0bstate_delta\x18\x02 \x01(\x0b\x32\x1a.cozy.scheduler.StateDeltaH\x00\x12\x33\n\x0cjob_accepted\x18\x03 \x01(\x0b\x32\x1b.cozy.scheduler.JobAcceptedH\x00\x12/\n\njob_result\x18\x04 \x01(\x0b\x32\x19.cozy.scheduler.JobResultH\x00\x12\x33\n\x0cjob_progress\x18\x05 \x01(\x0b\x32\x1b.cozy.scheduler.JobProgressH\x00\x12\x31\n\x0bmodel_event\x18\x06 \x01(\x0b\x32\x1a.cozy.scheduler.ModelEventH\x00\x12\x37\n\x0e\x66n_unavailable\x18\x07 \x01(\x0b\x32\x1d.cozy.scheduler.FnUnavailableH\x00\x12\x31\n\x0b\x66n_degraded\x18\x08 \x01(\x0b\x32\x1a.cozy.scheduler.FnDegradedH\x00\x12\x39\n\x0f\x61\x63tivity_update\x18\t \x01(\x0b\x32\x1e.cozy.scheduler.ActivityUpdateH\x00\x42\x05\n\x03msg\"\xb0\x02\n\x10SchedulerMessage\x12-\n\thello_ack\x18\x01 \x01(\x0b\x32\x18.cozy.scheduler.HelloAckH\x00\x12)\n\x07run_job\x18\x02 \x01(\x0b\x32\x16.cozy.scheduler.RunJobH\x00\x12/\n\ncancel_job\x18\x03 \x01(\x0b\x32\x19.cozy.scheduler.CancelJobH\x00\x12+\n\x08model_op\x18\x04 \x01(\x0b\x32\x17.cozy.scheduler.ModelOpH\x00\x12&\n\x05\x64rain\x18\x05 \x01(\x0b\x32\x15.cozy.scheduler.DrainH\x00\x12\x35\n\rtoken_refresh\x18\x06 \x01(\x0b\x32\x1c.cozy.scheduler.TokenRefreshH\x00\x42\x05\n\x03msg\"\xa8\x02\n\x05Hello\x12\x39\n\x10protocol_version\x18\x01 \x01(\x0e\x32\x1f.cozy.scheduler.ProtocolVersion\x12\x11\n\tworker_id\x18\x02 \x01(\t\x12\x12\n\nrelease_id\x18\x03 \x01(\t\x12\x32\n\tresources\x18\x04 \x01(\x0b\x32\x1f.cozy.scheduler.WorkerResources\x12)\n\x05state\x18\x05 \x01(\x0b\x32\x1a.cozy.scheduler.StateDelta\x12.\n\x06models\x18\x06 \x03(\x0b\x32\x1e.cozy.scheduler.ModelResidency\x12.\n\tin_flight\x18\x07 \x03(\x0b\x32\x1b.cozy.scheduler.InFlightJob\"\x9b\x02\n\x0fWorkerResources\x12\x11\n\tgpu_count\x18\x01 \x01(\x05\x12\x18\n\x10vram_total_bytes\x18\x02 \x01(\x03\x12\x10\n\x08gpu_name\x18\x03 \x01(\t\x12\x0e\n\x06gpu_sm\x18\x04 \x01(\t\x12\x16\n\x0einstalled_libs\x18\x05 \x03(\t\x12\x14\n\x0cimage_digest\x18\x06 \x01(\t\x12\x12\n\ngit_commit\x18\x07 \x01(\t\x12\x13\n\x0binstance_id\x18\x08 \x01(\t\x12/\n\x0bhost_canary\x18\t \x01(\x0b\x32\x1a.cozy.scheduler.HostCanary\x12\x15\n\rtorch_version\x18\n \x01(\t\x12\x1a\n\x12gen_worker_version\x18\x0b \x01(\t\"\xc9\x01\n\nHostCanary\x12\x13\n\x0bmemcpy_gbps\x18\x01 \x01(\x01\x12\x10\n\x08h2d_gbps\x18\x02 \x01(\x01\x12\x10\n\x08\x64\x32h_gbps\x18\x03 \x01(\x01\x12\x17\n\x0fpinned_alloc_ok\x18\x04 \x01(\x08\x12\x17\n\x0f\x63pu_single_mbps\x18\x05 \x01(\x01\x12\x16\n\x0e\x63pu_multi_mbps\x18\x06 \x01(\x01\x12\r\n\x05vcpus\x18\x07 \x01(\x05\x12\x14\n\x0cram_total_gb\x18\x08 \x01(\x01\x12\x13\n\x0b\x64uration_ms\x18\t \x01(\x03\"\x95\x01\n\x0eModelResidency\x12\x0b\n\x03ref\x18\x01 \x01(\t\x12+\n\x04tier\x18\x02 \x01(\x0e\x32\x1d.cozy.scheduler.ResidencyTier\x12\x12\n\nvram_bytes\x18\x03 \x01(\x03\x12\x17\n\x0fsnapshot_digest\x18\x04 \x01(\t\x12\x1c\n\x14residency_generation\x18\x05 \x01(\x04\"2\n\x0bInFlightJob\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\"\xdd\x01\n\x08HelloAck\x12\x39\n\x10protocol_version\x18\x01 \x01(\x0e\x32\x1f.cozy.scheduler.ProtocolVersion\x12\x15\n\rfile_base_url\x18\x02 \x01(\t\x12\x0c\n\x04keep\x18\x03 \x03(\t\x12\x34\n\x0bresolutions\x18\x04 \x03(\x0b\x32\x1f.cozy.scheduler.ModelResolution\x12;\n\x11\x64\x65sired_residency\x18\x05 \x01(\x0b\x32 .cozy.scheduler.DesiredResidency\"\xf7\x01\n\x10\x44\x65siredResidency\x12\x12\n\ngeneration\x18\x01 \x01(\x04\x12\x11\n\tdisk_refs\x18\x02 \x03(\t\x12,\n\x03hot\x18\x03 \x03(\x0b\x32\x1f.cozy.scheduler.DesiredInstance\x12\x42\n\tsnapshots\x18\x04 \x03(\x0b\x32/.cozy.scheduler.DesiredResidency.SnapshotsEntry\x1aJ\n\x0eSnapshotsEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\'\n\x05value\x18\x02 \x01(\x0b\x32\x18.cozy.scheduler.Snapshot:\x02\x38\x01\"V\n\x0f\x44\x65siredInstance\x12\x15\n\rfunction_name\x18\x01 \x01(\t\x12,\n\x06models\x18\x02 \x03(\x0b\x32\x1c.cozy.scheduler.ModelBinding\"B\n\x0fModelResolution\x12\x0b\n\x03ref\x18\x01 \x01(\t\x12\x14\n\x0cresolved_ref\x18\x02 \x01(\t\x12\x0c\n\x04\x63\x61st\x18\x03 \x01(\t\"\xb3\x02\n\nStateDelta\x12*\n\x05phase\x18\x01 \x01(\x0e\x32\x1b.cozy.scheduler.WorkerPhase\x12\x1b\n\x13\x61vailable_functions\x18\x02 \x03(\t\x12\x19\n\x11loading_functions\x18\x03 \x03(\t\x12\x17\n\x0f\x66ree_vram_bytes\x18\x04 \x01(\x03\x12\x17\n\x0f\x66inalizing_jobs\x18\x05 \x01(\x05\x12%\n\x1dobserved_residency_generation\x18\x06 \x01(\x04\x12\x36\n\x0f\x63ompile_targets\x18\x07 \x03(\x0b\x32\x1d.cozy.scheduler.CompileTarget\x12\x30\n\x0c\x63\x65ll_lookups\x18\x08 \x03(\x0b\x32\x1a.cozy.scheduler.CellLookup\".\n\nCellLookup\x12\x0e\n\x06\x66\x61mily\x18\x01 \x01(\t\x12\x10\n\x08\x63\x65ll_key\x18\x02 \x01(\t\"\xc6\x03\n\rCompileTarget\x12\x16\n\x0eincarnation_id\x18\x01 \x01(\t\x12\x0e\n\x06\x66\x61mily\x18\x02 \x01(\t\x12\x1c\n\x14pipeline_weight_lane\x18\x03 \x01(\t\x12\x13\n\x0blora_bucket\x18\x04 \x01(\x05\x12\x17\n\x0f\x63ontract_digest\x18\x05 \x01(\t\x12\x1a\n\x12\x61\x63tive_compile_ref\x18\x06 \x01(\t\x12&\n\x1e\x61\x63tive_compile_snapshot_digest\x18\x07 \x01(\t\x12\x16\n\x0e\x66unction_names\x18\x08 \x03(\t\x12<\n\x0emodel_bindings\x18\t \x03(\x0b\x32$.cozy.scheduler.CompileTargetBinding\x12\x1a\n\x12requested_cell_key\x18\n \x01(\t\x12Q\n\x13requested_cell_axes\x18\x0b \x03(\x0b\x32\x34.cozy.scheduler.CompileTarget.RequestedCellAxesEntry\x1a\x38\n\x16RequestedCellAxesEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\"J\n\x14\x43ompileTargetBinding\x12\x0c\n\x04slot\x18\x01 \x01(\t\x12\x0b\n\x03ref\x18\x02 \x01(\t\x12\x17\n\x0fsnapshot_digest\x18\x03 \x01(\t\"\x88\x04\n\x06RunJob\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\x12\x15\n\rfunction_name\x18\x03 \x01(\t\x12\x15\n\rinput_payload\x18\x04 \x01(\x0c\x12\x12\n\ntimeout_ms\x18\x05 \x01(\x03\x12\x0e\n\x06tenant\x18\x06 \x01(\t\x12\x12\n\ninvoker_id\x18\x07 \x01(\t\x12\x18\n\x10\x63\x61pability_token\x18\x08 \x01(\t\x12/\n\x0boutput_mode\x18\t \x01(\x0e\x32\x1a.cozy.scheduler.OutputMode\x12\x30\n\x07\x63ompute\x18\n \x01(\x0b\x32\x1f.cozy.scheduler.ResolvedCompute\x12,\n\x06models\x18\x0b \x03(\x0b\x32\x1c.cozy.scheduler.ModelBinding\x12\x38\n\tsnapshots\x18\x0c \x03(\x0b\x32%.cozy.scheduler.RunJob.SnapshotsEntry\x12\x42\n\x10required_compile\x18\r \x01(\x0b\x32(.cozy.scheduler.RequiredCompileExecution\x1aJ\n\x0eSnapshotsEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\'\n\x05value\x18\x02 \x01(\x0b\x32\x18.cozy.scheduler.Snapshot:\x02\x38\x01\"\x82\x01\n\x18RequiredCompileExecution\x12\x1d\n\x15target_incarnation_id\x18\x01 \x01(\t\x12\x10\n\x08\x63\x65ll_ref\x18\x02 \x01(\t\x12\x1c\n\x14\x63\x65ll_snapshot_digest\x18\x03 \x01(\t\x12\x17\n\x0f\x63ontract_digest\x18\x04 \x01(\t\"Y\n\x0fResolvedCompute\x12\x13\n\x0b\x61\x63\x63\x65lerator\x18\x01 \x01(\t\x12\x11\n\tgpu_index\x18\x02 \x01(\x05J\x04\x08\x03\x10\x04J\x04\x08\x04\x10\x05R\tgpu_countR\x07vram_gb\"q\n\x0cModelBinding\x12\x0c\n\x04slot\x18\x01 \x01(\t\x12\x0b\n\x03ref\x18\x02 \x01(\t\x12*\n\x05loras\x18\x03 \x03(\x0b\x32\x1b.cozy.scheduler.LoraOverlay\x12\x1a\n\x12inference_defaults\x18\x04 \x01(\t\"F\n\x0bLoraOverlay\x12\x0b\n\x03ref\x18\x01 \x01(\t\x12\x0e\n\x06weight\x18\x02 \x01(\x01\x12\x1a\n\x12inference_defaults\x18\x03 \x01(\t\"G\n\x08Snapshot\x12\x0e\n\x06\x64igest\x18\x01 \x01(\t\x12+\n\x05\x66iles\x18\x02 \x03(\x0b\x32\x1c.cozy.scheduler.SnapshotFile\"M\n\x0cSnapshotFile\x12\x0c\n\x04path\x18\x01 \x01(\t\x12\x12\n\nsize_bytes\x18\x02 \x01(\x03\x12\x0e\n\x06\x62lake3\x18\x03 \x01(\t\x12\x0b\n\x03url\x18\x04 \x01(\t\"2\n\x0bJobAccepted\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\"\xce\x01\n\tJobResult\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\x12)\n\x06status\x18\x03 \x01(\x0e\x32\x19.cozy.scheduler.JobStatus\x12\x10\n\x06inline\x18\x04 \x01(\x0cH\x00\x12\x12\n\x08\x62lob_ref\x18\x05 \x01(\tH\x00\x12\x14\n\x0csafe_message\x18\x06 \x01(\t\x12+\n\x07metrics\x18\x07 \x01(\x0b\x32\x1a.cozy.scheduler.JobMetricsB\x08\n\x06output\"\xb4\x02\n\nJobMetrics\x12\x12\n\nruntime_ms\x18\x01 \x01(\x03\x12\x10\n\x08queue_ms\x18\x02 \x01(\x03\x12\x18\n\x10rss_at_end_bytes\x18\x03 \x01(\x03\x12\x17\n\x0fpeak_vram_bytes\x18\x04 \x01(\x03\x12\x1c\n\x14\x63oncurrency_at_start\x18\x05 \x01(\x05\x12\x1f\n\x17output_media_duration_s\x18\x06 \x01(\x01\x12\x14\n\x0cinput_tokens\x18\x07 \x01(\x03\x12\x1b\n\x13input_cached_tokens\x18\x08 \x01(\x03\x12\x15\n\routput_tokens\x18\t \x01(\x03\x12\x14\n\x0coutput_count\x18\n \x01(\x03\x12\x14\n\x0cslot_held_ms\x18\x0b \x01(\x03\x12\x18\n\x10\x66inalize_wall_ms\x18\x0c \x01(\x03\"c\n\x0bJobProgress\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\x12\x0b\n\x03seq\x18\x03 \x01(\x03\x12\x0c\n\x04\x64\x61ta\x18\x04 \x01(\x0c\x12\x14\n\x0c\x63ontent_type\x18\x05 \x01(\t\"0\n\tCancelJob\x12\x12\n\nrequest_id\x18\x01 \x01(\t\x12\x0f\n\x07\x61ttempt\x18\x02 \x01(\x03\"\xa0\x01\n\x07ModelOp\x12\'\n\x02op\x18\x01 \x01(\x0e\x32\x1b.cozy.scheduler.ModelOpKind\x12\x0b\n\x03ref\x18\x02 \x01(\t\x12*\n\x08snapshot\x18\x03 \x01(\x0b\x32\x18.cozy.scheduler.Snapshot\x12\x14\n\x0coperation_id\x18\x04 \x01(\t\x12\x1d\n\x15target_incarnation_id\x18\x05 \x01(\t\"\x9b\x04\n\nModelEvent\x12\x0b\n\x03ref\x18\x01 \x01(\t\x12)\n\x05state\x18\x02 \x01(\x0e\x32\x1a.cozy.scheduler.ModelState\x12\x12\n\nvram_bytes\x18\x03 \x01(\x03\x12\r\n\x05\x65rror\x18\x04 \x01(\t\x12\x12\n\nbytes_done\x18\x05 \x01(\x03\x12\x13\n\x0b\x62ytes_total\x18\x06 \x01(\x03\x12\x13\n\x0b\x64uration_ms\x18\x07 \x01(\x03\x12\x12\n\ncache_hits\x18\x08 \x01(\x03\x12\x14\n\x0c\x63\x61\x63he_misses\x18\t \x01(\x03\x12\x10\n\x08warmup_s\x18\n \x01(\x01\x12\x1f\n\x17host_ram_required_bytes\x18\x0b \x01(\x03\x12\'\n\x1fhost_ram_available_before_bytes\x18\x0c \x01(\x03\x12&\n\x1ehost_ram_available_after_bytes\x18\r \x01(\x03\x12\x1d\n\x15host_ram_evicted_refs\x18\x0e \x03(\t\x12$\n\x1chost_ram_capacity_generation\x18\x0f \x01(\x04\x12\x17\n\x0fsnapshot_digest\x18\x10 \x01(\t\x12\x1c\n\x14residency_generation\x18\x11 \x01(\x04\x12\x14\n\x0coperation_id\x18\x12 \x01(\t\x12\x1d\n\x15target_incarnation_id\x18\x13 \x01(\t\x12\x15\n\rnetwork_bytes\x18\x14 \x01(\x03\"\xc6\x01\n\x0e\x41\x63tivityUpdate\x12\x0c\n\x04kind\x18\x01 \x01(\t\x12\r\n\x05phase\x18\x02 \x01(\t\x12\x0c\n\x04step\x18\x03 \x01(\x03\x12\x13\n\x0btotal_steps\x18\x04 \x01(\x03\x12\x0b\n\x03seq\x18\x05 \x01(\x04\x12,\n\x05state\x18\x06 \x01(\x0e\x32\x1d.cozy.scheduler.ActivityState\x12\r\n\x05\x65rror\x18\x07 \x01(\t\x12\x0e\n\x06\x64\x65tail\x18\x08 \x01(\t\x12\x1a\n\x12updated_at_unix_ms\x18\t \x01(\x03\"\xaa\x01\n\rFnUnavailable\x12\x15\n\rfunction_name\x18\x01 \x01(\t\x12\x0e\n\x06reason\x18\x02 \x01(\t\x12\x0e\n\x06\x64\x65tail\x18\x03 \x01(\t\x12\x35\n\x04\x61xes\x18\x04 \x03(\x0b\x32\'.cozy.scheduler.FnUnavailable.AxesEntry\x1a+\n\tAxesEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\"\x8d\x01\n\nFnDegraded\x12\x15\n\rfunction_name\x18\x01 \x01(\t\x12\x0e\n\x06wanted\x18\x02 \x01(\t\x12\x0b\n\x03ran\x18\x03 \x01(\t\x12\x0e\n\x06reason\x18\x04 \x01(\t\x12\x1e\n\x16\x65st_latency_multiplier\x18\x05 \x01(\x01\x12\x1b\n\x13recommended_vram_gb\x18\x06 \x01(\x01\"\x1c\n\x05\x44rain\x12\x13\n\x0b\x64\x65\x61\x64line_ms\x18\x01 \x01(\x03\"6\n\x0cTokenRefresh\x12\r\n\x05token\x18\x01 \x01(\t\x12\x17\n\x0f\x65xpires_at_unix\x18\x02 \x01(\x03*Q\n\x0fProtocolVersion\x12 \n\x1cPROTOCOL_VERSION_UNSPECIFIED\x10\x00\x12\x1c\n\x18PROTOCOL_VERSION_CURRENT\x10\x03*y\n\rResidencyTier\x12\x1e\n\x1aRESIDENCY_TIER_UNSPECIFIED\x10\x00\x12\x17\n\x13RESIDENCY_TIER_DISK\x10\x01\x12\x16\n\x12RESIDENCY_TIER_RAM\x10\x02\x12\x17\n\x13RESIDENCY_TIER_VRAM\x10\x03*\xd8\x01\n\x0bWorkerPhase\x12\x1c\n\x18WORKER_PHASE_UNSPECIFIED\x10\x00\x12\x18\n\x14WORKER_PHASE_BOOTING\x10\x01\x12#\n\x1fWORKER_PHASE_DOWNLOADING_MODELS\x10\x02\x12\"\n\x1eWORKER_PHASE_LOADING_PIPELINES\x10\x03\x12\x18\n\x14WORKER_PHASE_WARMING\x10\x04\x12\x16\n\x12WORKER_PHASE_READY\x10\x05\x12\x16\n\x12WORKER_PHASE_ERROR\x10\x06*V\n\nOutputMode\x12\x1b\n\x17OUTPUT_MODE_UNSPECIFIED\x10\x00\x12\x13\n\x0fOUTPUT_MODE_URL\x10\x01\x12\x16\n\x12OUTPUT_MODE_INLINE\x10\x02*\x9b\x01\n\tJobStatus\x12\x1a\n\x16JOB_STATUS_UNSPECIFIED\x10\x00\x12\x11\n\rJOB_STATUS_OK\x10\x01\x12\x16\n\x12JOB_STATUS_INVALID\x10\x02\x12\x18\n\x14JOB_STATUS_RETRYABLE\x10\x03\x12\x14\n\x10JOB_STATUS_FATAL\x10\x04\x12\x17\n\x13JOB_STATUS_CANCELED\x10\x05*\xa7\x01\n\x0bModelOpKind\x12\x1d\n\x19MODEL_OP_KIND_UNSPECIFIED\x10\x00\x12%\n!MODEL_OP_KIND_ADOPT_COMPILE_CACHE\x10\x04\"\x04\x08\x01\x10\x01\"\x04\x08\x02\x10\x02\"\x04\x08\x03\x10\x03*\x16MODEL_OP_KIND_DOWNLOAD*\x12MODEL_OP_KIND_LOAD*\x14MODEL_OP_KIND_UNLOAD*\x82\x02\n\nModelState\x12\x1b\n\x17MODEL_STATE_UNSPECIFIED\x10\x00\x12\x1b\n\x17MODEL_STATE_DOWNLOADING\x10\x01\x12\x17\n\x13MODEL_STATE_ON_DISK\x10\x02\x12\x16\n\x12MODEL_STATE_IN_RAM\x10\x03\x12\x17\n\x13MODEL_STATE_IN_VRAM\x10\x04\x12\x17\n\x13MODEL_STATE_EVICTED\x10\x05\x12\x16\n\x12MODEL_STATE_FAILED\x10\x06\x12\x17\n\x13MODEL_STATE_ADOPTED\x10\x07\x12&\n\"MODEL_STATE_HOST_CAPACITY_PROGRESS\x10\x08*\x84\x01\n\rActivityState\x12\x1e\n\x1a\x41\x43TIVITY_STATE_UNSPECIFIED\x10\x00\x12\x1a\n\x16\x41\x43TIVITY_STATE_RUNNING\x10\x01\x12\x1c\n\x18\x41\x43TIVITY_STATE_COMPLETED\x10\x02\x12\x19\n\x15\x41\x43TIVITY_STATE_FAILED\x10\x03\x32\x61\n\x0fWorkerScheduler\x12N\n\x07\x43onnect\x12\x1d.cozy.scheduler.WorkerMessage\x1a .cozy.scheduler.SchedulerMessage(\x01\x30\x01\x42\x61Z_github.com/cozy-creator/tensorhub/internal/orchestrator/grpc/pb/workerscheduler;workerschedulerb\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) @@ -40,94 +40,98 @@ _globals['_RUNJOB_SNAPSHOTSENTRY']._serialized_options = b'8\001' _globals['_FNUNAVAILABLE_AXESENTRY']._loaded_options = None _globals['_FNUNAVAILABLE_AXESENTRY']._serialized_options = b'8\001' - _globals['_PROTOCOLVERSION']._serialized_start=6210 - _globals['_PROTOCOLVERSION']._serialized_end=6291 - _globals['_RESIDENCYTIER']._serialized_start=6293 - _globals['_RESIDENCYTIER']._serialized_end=6414 - _globals['_WORKERPHASE']._serialized_start=6417 - _globals['_WORKERPHASE']._serialized_end=6633 - _globals['_OUTPUTMODE']._serialized_start=6635 - _globals['_OUTPUTMODE']._serialized_end=6721 - _globals['_JOBSTATUS']._serialized_start=6724 - _globals['_JOBSTATUS']._serialized_end=6879 - _globals['_MODELOPKIND']._serialized_start=6882 - _globals['_MODELOPKIND']._serialized_end=7049 - _globals['_MODELSTATE']._serialized_start=7052 - _globals['_MODELSTATE']._serialized_end=7310 + _globals['_PROTOCOLVERSION']._serialized_start=6470 + _globals['_PROTOCOLVERSION']._serialized_end=6551 + _globals['_RESIDENCYTIER']._serialized_start=6553 + _globals['_RESIDENCYTIER']._serialized_end=6674 + _globals['_WORKERPHASE']._serialized_start=6677 + _globals['_WORKERPHASE']._serialized_end=6893 + _globals['_OUTPUTMODE']._serialized_start=6895 + _globals['_OUTPUTMODE']._serialized_end=6981 + _globals['_JOBSTATUS']._serialized_start=6984 + _globals['_JOBSTATUS']._serialized_end=7139 + _globals['_MODELOPKIND']._serialized_start=7142 + _globals['_MODELOPKIND']._serialized_end=7309 + _globals['_MODELSTATE']._serialized_start=7312 + _globals['_MODELSTATE']._serialized_end=7570 + _globals['_ACTIVITYSTATE']._serialized_start=7573 + _globals['_ACTIVITYSTATE']._serialized_end=7705 _globals['_WORKERMESSAGE']._serialized_start=43 - _globals['_WORKERMESSAGE']._serialized_end=470 - _globals['_SCHEDULERMESSAGE']._serialized_start=473 - _globals['_SCHEDULERMESSAGE']._serialized_end=777 - _globals['_HELLO']._serialized_start=780 - _globals['_HELLO']._serialized_end=1076 - _globals['_WORKERRESOURCES']._serialized_start=1079 - _globals['_WORKERRESOURCES']._serialized_end=1362 - _globals['_HOSTCANARY']._serialized_start=1365 - _globals['_HOSTCANARY']._serialized_end=1566 - _globals['_MODELRESIDENCY']._serialized_start=1569 - _globals['_MODELRESIDENCY']._serialized_end=1718 - _globals['_INFLIGHTJOB']._serialized_start=1720 - _globals['_INFLIGHTJOB']._serialized_end=1770 - _globals['_HELLOACK']._serialized_start=1773 - _globals['_HELLOACK']._serialized_end=1994 - _globals['_DESIREDRESIDENCY']._serialized_start=1997 - _globals['_DESIREDRESIDENCY']._serialized_end=2244 - _globals['_DESIREDRESIDENCY_SNAPSHOTSENTRY']._serialized_start=2170 - _globals['_DESIREDRESIDENCY_SNAPSHOTSENTRY']._serialized_end=2244 - _globals['_DESIREDINSTANCE']._serialized_start=2246 - _globals['_DESIREDINSTANCE']._serialized_end=2332 - _globals['_MODELRESOLUTION']._serialized_start=2334 - _globals['_MODELRESOLUTION']._serialized_end=2400 - _globals['_STATEDELTA']._serialized_start=2403 - _globals['_STATEDELTA']._serialized_end=2710 - _globals['_CELLLOOKUP']._serialized_start=2712 - _globals['_CELLLOOKUP']._serialized_end=2758 - _globals['_COMPILETARGET']._serialized_start=2761 - _globals['_COMPILETARGET']._serialized_end=3215 - _globals['_COMPILETARGET_REQUESTEDCELLAXESENTRY']._serialized_start=3159 - _globals['_COMPILETARGET_REQUESTEDCELLAXESENTRY']._serialized_end=3215 - _globals['_COMPILETARGETBINDING']._serialized_start=3217 - _globals['_COMPILETARGETBINDING']._serialized_end=3291 - _globals['_RUNJOB']._serialized_start=3294 - _globals['_RUNJOB']._serialized_end=3814 - _globals['_RUNJOB_SNAPSHOTSENTRY']._serialized_start=2170 - _globals['_RUNJOB_SNAPSHOTSENTRY']._serialized_end=2244 - _globals['_REQUIREDCOMPILEEXECUTION']._serialized_start=3817 - _globals['_REQUIREDCOMPILEEXECUTION']._serialized_end=3947 - _globals['_RESOLVEDCOMPUTE']._serialized_start=3949 - _globals['_RESOLVEDCOMPUTE']._serialized_end=4038 - _globals['_MODELBINDING']._serialized_start=4040 - _globals['_MODELBINDING']._serialized_end=4153 - _globals['_LORAOVERLAY']._serialized_start=4155 - _globals['_LORAOVERLAY']._serialized_end=4225 - _globals['_SNAPSHOT']._serialized_start=4227 - _globals['_SNAPSHOT']._serialized_end=4298 - _globals['_SNAPSHOTFILE']._serialized_start=4300 - _globals['_SNAPSHOTFILE']._serialized_end=4377 - _globals['_JOBACCEPTED']._serialized_start=4379 - _globals['_JOBACCEPTED']._serialized_end=4429 - _globals['_JOBRESULT']._serialized_start=4432 - _globals['_JOBRESULT']._serialized_end=4638 - _globals['_JOBMETRICS']._serialized_start=4641 - _globals['_JOBMETRICS']._serialized_end=4949 - _globals['_JOBPROGRESS']._serialized_start=4951 - _globals['_JOBPROGRESS']._serialized_end=5050 - _globals['_CANCELJOB']._serialized_start=5052 - _globals['_CANCELJOB']._serialized_end=5100 - _globals['_MODELOP']._serialized_start=5103 - _globals['_MODELOP']._serialized_end=5263 - _globals['_MODELEVENT']._serialized_start=5266 - _globals['_MODELEVENT']._serialized_end=5805 - _globals['_FNUNAVAILABLE']._serialized_start=5808 - _globals['_FNUNAVAILABLE']._serialized_end=5978 - _globals['_FNUNAVAILABLE_AXESENTRY']._serialized_start=5935 - _globals['_FNUNAVAILABLE_AXESENTRY']._serialized_end=5978 - _globals['_FNDEGRADED']._serialized_start=5981 - _globals['_FNDEGRADED']._serialized_end=6122 - _globals['_DRAIN']._serialized_start=6124 - _globals['_DRAIN']._serialized_end=6152 - _globals['_TOKENREFRESH']._serialized_start=6154 - _globals['_TOKENREFRESH']._serialized_end=6208 - _globals['_WORKERSCHEDULER']._serialized_start=7312 - _globals['_WORKERSCHEDULER']._serialized_end=7409 + _globals['_WORKERMESSAGE']._serialized_end=529 + _globals['_SCHEDULERMESSAGE']._serialized_start=532 + _globals['_SCHEDULERMESSAGE']._serialized_end=836 + _globals['_HELLO']._serialized_start=839 + _globals['_HELLO']._serialized_end=1135 + _globals['_WORKERRESOURCES']._serialized_start=1138 + _globals['_WORKERRESOURCES']._serialized_end=1421 + _globals['_HOSTCANARY']._serialized_start=1424 + _globals['_HOSTCANARY']._serialized_end=1625 + _globals['_MODELRESIDENCY']._serialized_start=1628 + _globals['_MODELRESIDENCY']._serialized_end=1777 + _globals['_INFLIGHTJOB']._serialized_start=1779 + _globals['_INFLIGHTJOB']._serialized_end=1829 + _globals['_HELLOACK']._serialized_start=1832 + _globals['_HELLOACK']._serialized_end=2053 + _globals['_DESIREDRESIDENCY']._serialized_start=2056 + _globals['_DESIREDRESIDENCY']._serialized_end=2303 + _globals['_DESIREDRESIDENCY_SNAPSHOTSENTRY']._serialized_start=2229 + _globals['_DESIREDRESIDENCY_SNAPSHOTSENTRY']._serialized_end=2303 + _globals['_DESIREDINSTANCE']._serialized_start=2305 + _globals['_DESIREDINSTANCE']._serialized_end=2391 + _globals['_MODELRESOLUTION']._serialized_start=2393 + _globals['_MODELRESOLUTION']._serialized_end=2459 + _globals['_STATEDELTA']._serialized_start=2462 + _globals['_STATEDELTA']._serialized_end=2769 + _globals['_CELLLOOKUP']._serialized_start=2771 + _globals['_CELLLOOKUP']._serialized_end=2817 + _globals['_COMPILETARGET']._serialized_start=2820 + _globals['_COMPILETARGET']._serialized_end=3274 + _globals['_COMPILETARGET_REQUESTEDCELLAXESENTRY']._serialized_start=3218 + _globals['_COMPILETARGET_REQUESTEDCELLAXESENTRY']._serialized_end=3274 + _globals['_COMPILETARGETBINDING']._serialized_start=3276 + _globals['_COMPILETARGETBINDING']._serialized_end=3350 + _globals['_RUNJOB']._serialized_start=3353 + _globals['_RUNJOB']._serialized_end=3873 + _globals['_RUNJOB_SNAPSHOTSENTRY']._serialized_start=2229 + _globals['_RUNJOB_SNAPSHOTSENTRY']._serialized_end=2303 + _globals['_REQUIREDCOMPILEEXECUTION']._serialized_start=3876 + _globals['_REQUIREDCOMPILEEXECUTION']._serialized_end=4006 + _globals['_RESOLVEDCOMPUTE']._serialized_start=4008 + _globals['_RESOLVEDCOMPUTE']._serialized_end=4097 + _globals['_MODELBINDING']._serialized_start=4099 + _globals['_MODELBINDING']._serialized_end=4212 + _globals['_LORAOVERLAY']._serialized_start=4214 + _globals['_LORAOVERLAY']._serialized_end=4284 + _globals['_SNAPSHOT']._serialized_start=4286 + _globals['_SNAPSHOT']._serialized_end=4357 + _globals['_SNAPSHOTFILE']._serialized_start=4359 + _globals['_SNAPSHOTFILE']._serialized_end=4436 + _globals['_JOBACCEPTED']._serialized_start=4438 + _globals['_JOBACCEPTED']._serialized_end=4488 + _globals['_JOBRESULT']._serialized_start=4491 + _globals['_JOBRESULT']._serialized_end=4697 + _globals['_JOBMETRICS']._serialized_start=4700 + _globals['_JOBMETRICS']._serialized_end=5008 + _globals['_JOBPROGRESS']._serialized_start=5010 + _globals['_JOBPROGRESS']._serialized_end=5109 + _globals['_CANCELJOB']._serialized_start=5111 + _globals['_CANCELJOB']._serialized_end=5159 + _globals['_MODELOP']._serialized_start=5162 + _globals['_MODELOP']._serialized_end=5322 + _globals['_MODELEVENT']._serialized_start=5325 + _globals['_MODELEVENT']._serialized_end=5864 + _globals['_ACTIVITYUPDATE']._serialized_start=5867 + _globals['_ACTIVITYUPDATE']._serialized_end=6065 + _globals['_FNUNAVAILABLE']._serialized_start=6068 + _globals['_FNUNAVAILABLE']._serialized_end=6238 + _globals['_FNUNAVAILABLE_AXESENTRY']._serialized_start=6195 + _globals['_FNUNAVAILABLE_AXESENTRY']._serialized_end=6238 + _globals['_FNDEGRADED']._serialized_start=6241 + _globals['_FNDEGRADED']._serialized_end=6382 + _globals['_DRAIN']._serialized_start=6384 + _globals['_DRAIN']._serialized_end=6412 + _globals['_TOKENREFRESH']._serialized_start=6414 + _globals['_TOKENREFRESH']._serialized_end=6468 + _globals['_WORKERSCHEDULER']._serialized_start=7707 + _globals['_WORKERSCHEDULER']._serialized_end=7804 # @@protoc_insertion_point(module_scope) diff --git a/src/gen_worker/pb/worker_scheduler_pb2.pyi b/src/gen_worker/pb/worker_scheduler_pb2.pyi index 5e20308d..15240d08 100644 --- a/src/gen_worker/pb/worker_scheduler_pb2.pyi +++ b/src/gen_worker/pb/worker_scheduler_pb2.pyi @@ -60,6 +60,13 @@ class ModelState(int, metaclass=_enum_type_wrapper.EnumTypeWrapper): MODEL_STATE_FAILED: _ClassVar[ModelState] MODEL_STATE_ADOPTED: _ClassVar[ModelState] MODEL_STATE_HOST_CAPACITY_PROGRESS: _ClassVar[ModelState] + +class ActivityState(int, metaclass=_enum_type_wrapper.EnumTypeWrapper): + __slots__ = () + ACTIVITY_STATE_UNSPECIFIED: _ClassVar[ActivityState] + ACTIVITY_STATE_RUNNING: _ClassVar[ActivityState] + ACTIVITY_STATE_COMPLETED: _ClassVar[ActivityState] + ACTIVITY_STATE_FAILED: _ClassVar[ActivityState] PROTOCOL_VERSION_UNSPECIFIED: ProtocolVersion PROTOCOL_VERSION_CURRENT: ProtocolVersion RESIDENCY_TIER_UNSPECIFIED: ResidencyTier @@ -93,9 +100,13 @@ MODEL_STATE_EVICTED: ModelState MODEL_STATE_FAILED: ModelState MODEL_STATE_ADOPTED: ModelState MODEL_STATE_HOST_CAPACITY_PROGRESS: ModelState +ACTIVITY_STATE_UNSPECIFIED: ActivityState +ACTIVITY_STATE_RUNNING: ActivityState +ACTIVITY_STATE_COMPLETED: ActivityState +ACTIVITY_STATE_FAILED: ActivityState class WorkerMessage(_message.Message): - __slots__ = ("hello", "state_delta", "job_accepted", "job_result", "job_progress", "model_event", "fn_unavailable", "fn_degraded") + __slots__ = ("hello", "state_delta", "job_accepted", "job_result", "job_progress", "model_event", "fn_unavailable", "fn_degraded", "activity_update") HELLO_FIELD_NUMBER: _ClassVar[int] STATE_DELTA_FIELD_NUMBER: _ClassVar[int] JOB_ACCEPTED_FIELD_NUMBER: _ClassVar[int] @@ -104,6 +115,7 @@ class WorkerMessage(_message.Message): MODEL_EVENT_FIELD_NUMBER: _ClassVar[int] FN_UNAVAILABLE_FIELD_NUMBER: _ClassVar[int] FN_DEGRADED_FIELD_NUMBER: _ClassVar[int] + ACTIVITY_UPDATE_FIELD_NUMBER: _ClassVar[int] hello: Hello state_delta: StateDelta job_accepted: JobAccepted @@ -112,7 +124,8 @@ class WorkerMessage(_message.Message): model_event: ModelEvent fn_unavailable: FnUnavailable fn_degraded: FnDegraded - def __init__(self, hello: _Optional[_Union[Hello, _Mapping]] = ..., state_delta: _Optional[_Union[StateDelta, _Mapping]] = ..., job_accepted: _Optional[_Union[JobAccepted, _Mapping]] = ..., job_result: _Optional[_Union[JobResult, _Mapping]] = ..., job_progress: _Optional[_Union[JobProgress, _Mapping]] = ..., model_event: _Optional[_Union[ModelEvent, _Mapping]] = ..., fn_unavailable: _Optional[_Union[FnUnavailable, _Mapping]] = ..., fn_degraded: _Optional[_Union[FnDegraded, _Mapping]] = ...) -> None: ... + activity_update: ActivityUpdate + def __init__(self, hello: _Optional[_Union[Hello, _Mapping]] = ..., state_delta: _Optional[_Union[StateDelta, _Mapping]] = ..., job_accepted: _Optional[_Union[JobAccepted, _Mapping]] = ..., job_result: _Optional[_Union[JobResult, _Mapping]] = ..., job_progress: _Optional[_Union[JobProgress, _Mapping]] = ..., model_event: _Optional[_Union[ModelEvent, _Mapping]] = ..., fn_unavailable: _Optional[_Union[FnUnavailable, _Mapping]] = ..., fn_degraded: _Optional[_Union[FnDegraded, _Mapping]] = ..., activity_update: _Optional[_Union[ActivityUpdate, _Mapping]] = ...) -> None: ... class SchedulerMessage(_message.Message): __slots__ = ("hello_ack", "run_job", "cancel_job", "model_op", "drain", "token_refresh") @@ -573,6 +586,28 @@ class ModelEvent(_message.Message): network_bytes: int def __init__(self, ref: _Optional[str] = ..., state: _Optional[_Union[ModelState, str]] = ..., vram_bytes: _Optional[int] = ..., error: _Optional[str] = ..., bytes_done: _Optional[int] = ..., bytes_total: _Optional[int] = ..., duration_ms: _Optional[int] = ..., cache_hits: _Optional[int] = ..., cache_misses: _Optional[int] = ..., warmup_s: _Optional[float] = ..., host_ram_required_bytes: _Optional[int] = ..., host_ram_available_before_bytes: _Optional[int] = ..., host_ram_available_after_bytes: _Optional[int] = ..., host_ram_evicted_refs: _Optional[_Iterable[str]] = ..., host_ram_capacity_generation: _Optional[int] = ..., snapshot_digest: _Optional[str] = ..., residency_generation: _Optional[int] = ..., operation_id: _Optional[str] = ..., target_incarnation_id: _Optional[str] = ..., network_bytes: _Optional[int] = ...) -> None: ... +class ActivityUpdate(_message.Message): + __slots__ = ("kind", "phase", "step", "total_steps", "seq", "state", "error", "detail", "updated_at_unix_ms") + KIND_FIELD_NUMBER: _ClassVar[int] + PHASE_FIELD_NUMBER: _ClassVar[int] + STEP_FIELD_NUMBER: _ClassVar[int] + TOTAL_STEPS_FIELD_NUMBER: _ClassVar[int] + SEQ_FIELD_NUMBER: _ClassVar[int] + STATE_FIELD_NUMBER: _ClassVar[int] + ERROR_FIELD_NUMBER: _ClassVar[int] + DETAIL_FIELD_NUMBER: _ClassVar[int] + UPDATED_AT_UNIX_MS_FIELD_NUMBER: _ClassVar[int] + kind: str + phase: str + step: int + total_steps: int + seq: int + state: ActivityState + error: str + detail: str + updated_at_unix_ms: int + def __init__(self, kind: _Optional[str] = ..., phase: _Optional[str] = ..., step: _Optional[int] = ..., total_steps: _Optional[int] = ..., seq: _Optional[int] = ..., state: _Optional[_Union[ActivityState, str]] = ..., error: _Optional[str] = ..., detail: _Optional[str] = ..., updated_at_unix_ms: _Optional[int] = ...) -> None: ... + class FnUnavailable(_message.Message): __slots__ = ("function_name", "reason", "detail", "axes") class AxesEntry(_message.Message): diff --git a/tests/test_activity_gw601.py b/tests/test_activity_gw601.py new file mode 100644 index 00000000..091bf698 --- /dev/null +++ b/tests/test_activity_gw601.py @@ -0,0 +1,172 @@ +"""gw#601 generic worker-activity progress: monotonic phase envelopes over +the real executor setup/warmup path, the evidence-gated watchdog heartbeat +(an induced hang stops the beat within one interval), and the typed +activity_failed terminal (an induced crash never dies silently).""" + +from __future__ import annotations + +import asyncio +import threading +import time +from typing import List + +import msgspec +import pytest + +from gen_worker import activity +from gen_worker.api import Resources, endpoint +from gen_worker.executor import Executor +from gen_worker.pb import worker_scheduler_pb2 as pb +from gen_worker.registry import extract_specs + + +class _In(msgspec.Struct): + prompt: str = "x" + + +class _Out(msgspec.Struct): + y: str + + +@pytest.fixture(autouse=True) +def _reset_activity_sink(): + yield + with activity._lock: + activity._sink = None + activity._current = None + + +def _updates(sent: List[pb.WorkerMessage]) -> List[pb.ActivityUpdate]: + return [m.activity_update for m in sent if m.WhichOneof("msg") == "activity_update"] + + +# --------------------------------------------------------------------------- +# real executor code path: scripted setup+warmup emits monotonic phases +# --------------------------------------------------------------------------- + + +def test_executor_setup_emits_monotonic_activity_phases(): + sent: List[pb.WorkerMessage] = [] + + async def _send(msg: pb.WorkerMessage) -> None: + sent.append(msg) + + @endpoint(resources=Resources(vram_gb=8)) + class Ep: + def setup(self) -> None: + pass + + def generate(self, ctx, payload: _In) -> _Out: + return _Out(y="ok") + + def generate_turbo(self, ctx, payload: _In) -> _Out: + return _Out(y="ok") + + specs = extract_specs(Ep) + ex = Executor(specs, _send) + + async def _go() -> None: + await ex.ensure_setup(specs[0]) + # the sink schedules sends onto this loop; let them flush + for _ in range(10): + await asyncio.sleep(0) + + asyncio.run(_go()) + + ups = _updates(sent) + assert ups, "no activity envelopes emitted" + assert all(u.kind == activity.KIND_WARMUP for u in ups) + seqs = [u.seq for u in ups] + assert seqs == sorted(seqs) and len(set(seqs)) == len(seqs), seqs + phases = [(u.phase, u.step, u.total_steps, u.state) for u in ups] + running = pb.ActivityState.ACTIVITY_STATE_RUNNING + assert phases[0] == (activity.PHASE_LOAD, 0, 0, running) + assert (activity.PHASE_WARMUP_FORWARD, 1, 2, running) in phases + assert (activity.PHASE_WARMUP_FORWARD, 2, 2, running) in phases + assert ups[-1].state == pb.ActivityState.ACTIVITY_STATE_COMPLETED + + +def test_executor_setup_crash_emits_typed_activity_failed(): + sent: List[pb.WorkerMessage] = [] + + async def _send(msg: pb.WorkerMessage) -> None: + sent.append(msg) + + @endpoint(resources=Resources(vram_gb=8)) + class Ep: + def setup(self) -> None: + raise RuntimeError("induced setup crash") + + def generate(self, ctx, payload: _In) -> _Out: # pragma: no cover + return _Out(y="ok") + + specs = extract_specs(Ep) + ex = Executor(specs, _send) + + async def _go() -> None: + with pytest.raises(RuntimeError, match="induced setup crash"): + await ex.ensure_setup(specs[0]) + for _ in range(10): + await asyncio.sleep(0) + + asyncio.run(_go()) + + ups = _updates(sent) + assert ups and ups[-1].state == pb.ActivityState.ACTIVITY_STATE_FAILED + assert "RuntimeError: induced setup crash" in ups[-1].error + + +# --------------------------------------------------------------------------- +# watchdog: heartbeats only while evidence advances +# --------------------------------------------------------------------------- + + +def test_watchdog_heartbeats_while_evidence_advances_then_stops_on_hang(): + reports: List[pb.ActivityUpdate] = [] + with activity._lock: + activity._sink = reports.append + + evidence_val = [0.0] + advancing = threading.Event() + advancing.set() + + def evidence() -> float: + if advancing.is_set(): + evidence_val[0] += 1.0 + return evidence_val[0] + + act = activity.begin(activity.KIND_SELF_MINT_COMPILE, activity.PHASE_INDUCTOR_COMPILE) + baseline = len(reports) + interval = 0.05 + with activity.watchdog(act, interval_s=interval, evidence=evidence): + time.sleep(6 * interval) + beats_while_advancing = len(reports) - baseline + # induced hang: the wrapped call stops accruing evidence + advancing.clear() + time.sleep(2 * interval) # one interval of grace, then silence + beats_at_hang = len(reports) + time.sleep(6 * interval) + beats_after_hang = len(reports) + act.completed() + + assert beats_while_advancing >= 2, "no heartbeats while evidence advanced" + assert beats_after_hang == beats_at_hang, "heartbeat outlived the hang" + assert reports[-1].state == pb.ActivityState.ACTIVITY_STATE_COMPLETED + seqs = [u.seq for u in reports] + assert seqs == sorted(seqs) and len(set(seqs)) == len(seqs) + + +def test_running_context_manager_reports_failed_with_exception(): + reports: List[pb.ActivityUpdate] = [] + with activity._lock: + activity._sink = reports.append + + with pytest.raises(ValueError): + with activity.running(activity.KIND_SELF_MINT_COMPILE, activity.PHASE_LOAD): + raise ValueError("mint exploded") + + assert reports[-1].state == pb.ActivityState.ACTIVITY_STATE_FAILED + assert "ValueError: mint exploded" in reports[-1].error + # terminal cleared the current activity: later phase reports are no-ops + activity.current_phase(activity.PHASE_SEAL_PUBLISH) + assert reports[-1].state == pb.ActivityState.ACTIVITY_STATE_FAILED From 760b47abb4475bb98ae1e6bf829bb124c30e8441 Mon Sep 17 00:00:00 2001 From: Paul Fidika Date: Mon, 20 Jul 2026 03:49:02 -0600 Subject: [PATCH 2/3] gw#601: mypy annotations for activity module --- src/gen_worker/activity.py | 24 +++++++++++++++++++----- 1 file changed, 19 insertions(+), 5 deletions(-) diff --git a/src/gen_worker/activity.py b/src/gen_worker/activity.py index 5797ca06..6ee8291e 100644 --- a/src/gen_worker/activity.py +++ b/src/gen_worker/activity.py @@ -20,7 +20,8 @@ import logging import threading import time -from typing import Callable, Optional +from types import TracebackType +from typing import Any, Callable, Coroutine, Optional from .pb import worker_scheduler_pb2 as pb @@ -49,7 +50,10 @@ _current: Optional["Activity"] = None -def bind_sink(emit, loop: asyncio.AbstractEventLoop) -> None: +def bind_sink( + emit: Callable[["pb.WorkerMessage"], Coroutine[Any, Any, None]], + loop: asyncio.AbstractEventLoop, +) -> None: """Route reports onto the worker->hub stream: emit is the async WorkerMessage sender, loop the transport loop. Thread-safe emission.""" def sink(update: pb.ActivityUpdate) -> None: @@ -102,7 +106,7 @@ def __init__(self, kind: str) -> None: self._total = 0 self._done = False - def _report(self, state: int, error: str = "", detail: str = "") -> None: + def _report(self, state: "pb.ActivityState", error: str = "", detail: str = "") -> None: _emit(pb.ActivityUpdate( kind=self.kind, phase=self._phase, step=self._step, total_steps=self._total, seq=_next_seq(), state=state, @@ -172,7 +176,12 @@ def __enter__(self) -> Activity: self.activity = begin(self._kind, self._phase) return self.activity - def __exit__(self, exc_type, exc, tb) -> None: + def __exit__( + self, + exc_type: Optional[type[BaseException]], + exc: Optional[BaseException], + tb: Optional[TracebackType], + ) -> None: assert self.activity is not None if exc is not None: self.activity.failed(exc) @@ -228,6 +237,11 @@ def __enter__(self) -> "watchdog": self._thread.start() return self - def __exit__(self, *exc) -> None: + def __exit__( + self, + exc_type: Optional[type[BaseException]], + exc: Optional[BaseException], + tb: Optional[TracebackType], + ) -> None: self._stop.set() self._thread.join(timeout=5) From 2e4934e57b35226ba3cce5c4ce2f34a1b2a367a4 Mon Sep 17 00:00:00 2001 From: Paul Fidika Date: Mon, 20 Jul 2026 03:49:39 -0600 Subject: [PATCH 3/3] gw#601: bind_sink accepts Awaitable sender (matches executor _send typing) --- src/gen_worker/activity.py | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/src/gen_worker/activity.py b/src/gen_worker/activity.py index 6ee8291e..beabbf54 100644 --- a/src/gen_worker/activity.py +++ b/src/gen_worker/activity.py @@ -21,7 +21,7 @@ import threading import time from types import TracebackType -from typing import Any, Callable, Coroutine, Optional +from typing import Awaitable, Callable, Optional from .pb import worker_scheduler_pb2 as pb @@ -51,23 +51,22 @@ def bind_sink( - emit: Callable[["pb.WorkerMessage"], Coroutine[Any, Any, None]], + emit: Callable[["pb.WorkerMessage"], Awaitable[None]], loop: asyncio.AbstractEventLoop, ) -> None: """Route reports onto the worker->hub stream: emit is the async WorkerMessage sender, loop the transport loop. Thread-safe emission.""" def sink(update: pb.ActivityUpdate) -> None: - coro = emit(pb.WorkerMessage(activity_update=update)) + async def _ship() -> None: + await emit(pb.WorkerMessage(activity_update=update)) try: running = asyncio.get_running_loop() except RuntimeError: running = None if running is loop: - loop.create_task(coro) + loop.create_task(_ship()) elif not loop.is_closed(): - asyncio.run_coroutine_threadsafe(coro, loop) - else: - coro.close() + asyncio.run_coroutine_threadsafe(_ship(), loop) global _sink with _lock: _sink = sink