From d44165860b605dc49a39bf442b873478d222b845 Mon Sep 17 00:00:00 2001 From: Hogne <227774406+hognek@users.noreply.github.com> Date: Sat, 18 Jul 2026 11:54:16 +0200 Subject: [PATCH 1/8] feat(gpu-queue): op kinds, shadow queue dict, position, cancel_op on GpuArbiter (#1864) --- tests/test_gpu_arbiter_queue_ops.py | 112 +++++++++++++++++++++++++++ tinyagentos/scheduler/gpu_arbiter.py | 73 +++++++++++++---- 2 files changed, 171 insertions(+), 14 deletions(-) create mode 100644 tests/test_gpu_arbiter_queue_ops.py diff --git a/tests/test_gpu_arbiter_queue_ops.py b/tests/test_gpu_arbiter_queue_ops.py new file mode 100644 index 000000000..28426ddc7 --- /dev/null +++ b/tests/test_gpu_arbiter_queue_ops.py @@ -0,0 +1,112 @@ +"""Tests for GPU arbiter queue ops — op shape, position, cancel, snapshot (taOS #1864 A2).""" + +import asyncio +import pytest +from tinyagentos.scheduler.gpu_arbiter import GpuArbiter +from tinyagentos.scheduler.types import Capability, Priority, Task +from tinyagentos.vram_reservation import VramReservationManager + + +def _mgr(free_mb: int, total_mb: int = 16384) -> VramReservationManager: + return VramReservationManager(probe=lambda: (free_mb, total_mb)) + + +def _task(priority=Priority.BACKGROUND, submitter="t"): + async def payload(_res): + await asyncio.sleep(0.05) + return "ok" + return Task(capability=Capability.LLM_CHAT, payload=payload, + preferred_resources=[], priority=priority, submitter=submitter) + + +@pytest.mark.asyncio +async def test_submit_gpu_defaults_backward_compatible(): + arbiter = GpuArbiter(vram_reservation=_mgr(8192)) + result = await arbiter.submit_gpu(_task(), required_vram_mb=1024) + assert result == "ok" # old call shape, no new kwargs + + +@pytest.mark.asyncio +async def test_queue_position_global_for_loads(): + arbiter = GpuArbiter(vram_reservation=_mgr(0)) # everything queues + t1, t2, t3 = _task(), _task(), _task() + f1 = asyncio.ensure_future(arbiter.submit_gpu( + t1, required_vram_mb=1024, op="load", model="a", backend_name="b1")) + f2 = asyncio.ensure_future(arbiter.submit_gpu( + t2, required_vram_mb=1024, op="load", model="b", backend_name="b1")) + f3 = asyncio.ensure_future(arbiter.submit_gpu( + t3, required_vram_mb=1024, op="load", model="c", backend_name="b1")) + await asyncio.sleep(0.05) # let them enqueue + assert arbiter.queue_position(t1.id) == 1 + assert arbiter.queue_position(t2.id) == 2 + assert arbiter.queue_position(t3.id) == 3 + for f in (f1, f2, f3): + f.cancel() + + +@pytest.mark.asyncio +async def test_queue_position_per_model_for_inference(): + arbiter = GpuArbiter(vram_reservation=_mgr(0)) + ta, tb, ta2 = _task(), _task(), _task() + fs = [asyncio.ensure_future(arbiter.submit_gpu( + t, required_vram_mb=1024, op="inference", model=m, backend_name="b1")) + for t, m in ((ta, "m-a"), (tb, "m-b"), (ta2, "m-a"))] + await asyncio.sleep(0.05) + assert arbiter.queue_position(ta.id) == 1 + assert arbiter.queue_position(tb.id) == 1 # only m-b entries count + assert arbiter.queue_position(ta2.id) == 2 # behind ta on m-a + for f in fs: + f.cancel() + + +@pytest.mark.asyncio +async def test_queue_snapshot_non_destructive_and_shaped(): + arbiter = GpuArbiter(vram_reservation=_mgr(0)) + t1 = _task(submitter="pull:x") + f = asyncio.ensure_future(arbiter.submit_gpu( + t1, required_vram_mb=1024, op="load", model="qwen", backend_name="b1")) + await asyncio.sleep(0.05) + snap1 = arbiter.queue_snapshot() + snap2 = arbiter.queue_snapshot() + entry = snap1[0] + assert entry["op"] == "load" and entry["model"] == "qwen" + assert entry["backend_name"] == "b1" and entry["submitter"] == "pull:x" + assert entry["position"] == 1 + assert [e["task_id"] for e in snap1] == [e["task_id"] for e in snap2] + stats = await arbiter.stats() + assert stats["queue_depth"] == 1 # snapshot did not drain + f.cancel() + + +@pytest.mark.asyncio +async def test_cancel_queued_op_removes_and_cancels_future(): + arbiter = GpuArbiter(vram_reservation=_mgr(0)) + t1 = _task() + f = asyncio.ensure_future(arbiter.submit_gpu( + t1, required_vram_mb=1024, op="load", model="m", backend_name="b1")) + await asyncio.sleep(0.05) + assert await arbiter.cancel_op(t1.id) is True + with pytest.raises(asyncio.CancelledError): + await f + assert arbiter.queue_position(t1.id) is None + assert await arbiter.cancel_op(t1.id) is False # idempotent-ish: gone + + +@pytest.mark.asyncio +async def test_cancel_running_op_delegates_to_evict(): + mgr = _mgr(8192) + arbiter = GpuArbiter(vram_reservation=mgr) + started = asyncio.Event() + + async def payload(_res): + started.set() + await asyncio.sleep(30) + + t1 = Task(capability=Capability.LLM_CHAT, payload=payload, + preferred_resources=[], priority=Priority.BACKGROUND, submitter="t") + f = asyncio.ensure_future(arbiter.submit_gpu(t1, required_vram_mb=1024)) + await started.wait() + assert await arbiter.cancel_op(t1.id) is True + await asyncio.sleep(0.05) + assert mgr.reserved_vram_mb == 0 # reservation released + f.cancel() diff --git a/tinyagentos/scheduler/gpu_arbiter.py b/tinyagentos/scheduler/gpu_arbiter.py index a01a1bb10..1467bb650 100644 --- a/tinyagentos/scheduler/gpu_arbiter.py +++ b/tinyagentos/scheduler/gpu_arbiter.py @@ -55,6 +55,9 @@ class _QueuedGpuTask: evictable: bool = field(compare=False) required_gpu_arch: str | None = field(default=None, compare=False) queued_at: float = field(default_factory=time.time, compare=False) + op: str = field(default="inference", compare=False) + model: str | None = field(default=None, compare=False) + backend_name: str | None = field(default=None, compare=False) @dataclass @@ -96,6 +99,8 @@ def __init__( self._running: dict[str, tuple[Task, str | None, int, int]] = {} self._running_tasks: dict[str, asyncio.Task] = {} self._running_lock = asyncio.Lock() + self._queued_entries: dict[str, _QueuedGpuTask] = {} + self._cancelled_ids: set[str] = set() # --- Single VRAM authority (taOS #185) --- # The arbiter does NOT keep its own VRAM ledger. It reserves against a @@ -305,6 +310,9 @@ async def submit_gpu( self, task: Task, required_vram_mb: int = 0, evictable: bool = False, resource_id: str | None = None, required_gpu_arch: str | None = None, + op: str = "inference", + model: str | None = None, + backend_name: str | None = None, ) -> object: """Submit a GPU task with optional hardware-architecture requirements. @@ -314,6 +322,9 @@ async def submit_gpu( evictable: Whether lower-priority tasks can be evicted for this. resource_id: Specific cluster resource to target. required_gpu_arch: CUDA compute capability required (e.g. ``"sm_86"``). + op: Operation kind — ``"inference"`` or ``"load"``. + model: Backend model name (e.g. ``"qwen2.5:7b"``). + backend_name: Config backend name (e.g. ``"local-ollama"``). """ self._submitted += 1 @@ -345,8 +356,10 @@ async def submit_gpu( priority=int(task.priority), seq=self._seq, task=task, required_vram_mb=required_vram_mb, evictable=evictable, required_gpu_arch=required_gpu_arch, + op=op, model=model, backend_name=backend_name, ) await self._queue.put(entry) + self._queued_entries[task.id] = entry self._queued += 1 loop = asyncio.get_running_loop() done: asyncio.Future = loop.create_future() @@ -354,6 +367,7 @@ async def submit_gpu( try: return await done except asyncio.CancelledError: + self._queued_entries.pop(task.id, None) self._evicted += 1 raise # Admitted — reservation already held by _reserve_and_check, so @@ -559,6 +573,10 @@ async def _drain_queue(self) -> None: drained = False while not self._queue.empty() and not drained: entry = self._queue.get_nowait() + if entry.task.id in self._cancelled_ids: + self._cancelled_ids.discard(entry.task.id) + continue + self._queued_entries.pop(entry.task.id, None) admission = await self._reserve_and_check(entry.task.id, entry.required_vram_mb) if not admission.admitted: # Try eviction-to-make-room for higher-priority queued tasks. @@ -603,6 +621,7 @@ def _propagate(ct: asyncio.Task, f: asyncio.Future = future) -> None: for entry in retry: if not self._queue.full(): self._queue.put_nowait(entry) + self._queued_entries[entry.task.id] = entry else: self._dropped += 1 future = getattr(entry.task, "_arbiter_future", None) @@ -631,20 +650,46 @@ async def running_tasks(self) -> list[dict]: for tid, (task, lid, pri, vram) in self._running.items() ] - def queue_snapshot(self) -> list[dict]: - items: list[_QueuedGpuTask] = [] - while not self._queue.empty(): - try: - items.append(self._queue.get_nowait()) - except asyncio.QueueEmpty: + def _ordered_queued(self) -> list[_QueuedGpuTask]: + return sorted(self._queued_entries.values(), key=lambda e: (e.priority, e.seq)) + + def queue_position(self, task_id: str) -> int | None: + entry = self._queued_entries.get(task_id) + if entry is None: + return None + ahead = 0 + for other in self._ordered_queued(): + if other.task.id == task_id: break - result = [ + if entry.op == "inference": + if other.model == entry.model: + ahead += 1 + else: + ahead += 1 + return ahead + 1 + + async def cancel_op(self, task_id: str) -> bool: + entry = self._queued_entries.pop(task_id, None) + if entry is not None: + self._cancelled_ids.add(task_id) + future = getattr(entry.task, "_arbiter_future", None) + if future is not None and not future.done(): + future.cancel() + return True + async with self._running_lock: + running = task_id in self._running + if running: + return bool(await self._evict_task(task_id)) + return False + + def queue_snapshot(self) -> list[dict]: + now = time.time() + return [ {"task_id": e.task.id, "capability": e.task.capability.value, - "priority": e.priority, "vram_mb": e.required_vram_mb, - "queued_seconds": time.time() - e.queued_at} - for e in items + "op": e.op, "model": e.model, "backend_name": e.backend_name, + "submitter": e.task.submitter, "priority": e.priority, + "vram_mb": e.required_vram_mb, + "queued_seconds": now - e.queued_at, + "position": self.queue_position(e.task.id)} + for e in self._ordered_queued() ] - for e in items: - if not self._queue.full(): - self._queue.put_nowait(e) - return result From 43745888ffca9ce6e6b0481dd27bb6d69c6f9657 Mon Sep 17 00:00:00 2001 From: Hogne <227774406+hognek@users.noreply.github.com> Date: Sat, 18 Jul 2026 12:10:06 +0200 Subject: [PATCH 2/8] fix(gpu-arbiter): close shadow-dict drain window, Cancel orphan, validate op - _drain_queue: keep entry in _queued_entries during admission check so concurrent cancel_op finds it (was popped too early before asyncio.to_thread). Pop from shadow dict only after admission is confirmed (line 621). - submit_gpu CancelledError handler: add task_id to _cancelled_ids so the drain loop skips the orphaned queue entry instead of admitting it. - submit_gpu: validate op against {'inference','load'}, raising ValueError for unknown values so typos don't silently fall through to global branch. Kilo review items on PR #1984 (GPU queue A2). Targeted: 29/29 pass. --- tinyagentos/scheduler/gpu_arbiter.py | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/tinyagentos/scheduler/gpu_arbiter.py b/tinyagentos/scheduler/gpu_arbiter.py index 1467bb650..471fc9bd2 100644 --- a/tinyagentos/scheduler/gpu_arbiter.py +++ b/tinyagentos/scheduler/gpu_arbiter.py @@ -328,6 +328,12 @@ async def submit_gpu( """ self._submitted += 1 + # Validate op kind — only 'inference' and 'load' are known. + if op not in ("inference", "load"): + raise ValueError( + f"unknown gpu_op: {op!r} (expected 'inference' or 'load')" + ) + # Hardware architecture check (taOS #796) if required_gpu_arch: arch_ok, arch_reason = self._check_gpu_arch_compatibility( @@ -368,6 +374,7 @@ async def submit_gpu( return await done except asyncio.CancelledError: self._queued_entries.pop(task.id, None) + self._cancelled_ids.add(task.id) self._evicted += 1 raise # Admitted — reservation already held by _reserve_and_check, so @@ -576,7 +583,6 @@ async def _drain_queue(self) -> None: if entry.task.id in self._cancelled_ids: self._cancelled_ids.discard(entry.task.id) continue - self._queued_entries.pop(entry.task.id, None) admission = await self._reserve_and_check(entry.task.id, entry.required_vram_mb) if not admission.admitted: # Try eviction-to-make-room for higher-priority queued tasks. @@ -612,6 +618,7 @@ def _propagate(ct: asyncio.Task, f: asyncio.Future = future) -> None: t.add_done_callback(_propagate) + self._queued_entries.pop(entry.task.id, None) # promoted → _running drained = True else: # Not admitted — release reservation (idempotent) and retry later. From dc0a6311832202c73cff8bdfcb61a2b5ed71bcae Mon Sep 17 00:00:00 2001 From: Hogne <227774406+hognek@users.noreply.github.com> Date: Sat, 18 Jul 2026 12:28:17 +0200 Subject: [PATCH 3/8] fix(gpu-arbiter): close cancel_op/_drain_queue race by re-checking _cancelled_ids after await MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit cancel_op pops from _queued_entries and cancels the future while _drain_queue is suspended on _reserve_and_check — the entry is still admitted despite cancellation. Add a second _cancelled_ids check after the await window so a concurrent cancel_op is caught before the task is spawned. Kilo SUGGESTION on PR #1984 — task t_27501d36 --- tinyagentos/scheduler/gpu_arbiter.py | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/tinyagentos/scheduler/gpu_arbiter.py b/tinyagentos/scheduler/gpu_arbiter.py index 471fc9bd2..a87d37e0c 100644 --- a/tinyagentos/scheduler/gpu_arbiter.py +++ b/tinyagentos/scheduler/gpu_arbiter.py @@ -594,6 +594,16 @@ async def _drain_queue(self) -> None: if evicted > 0: admission = await self._reserve_and_check(entry.task.id, entry.required_vram_mb) if admission.admitted: + # Re-check cancelled_ids after the await window above + # (reserve_and_check / eviction). cancel_op may have fired + # concurrently while we were suspended, removing the entry + # from _queued_entries and cancelling its future. Without + # this re-check the cancelled task is still admitted and + # runs to completion despite the cancellation. + if entry.task.id in self._cancelled_ids: + self._cancelled_ids.discard(entry.task.id) + self._release_reservation(entry.task.id) + continue future = getattr(entry.task, "_arbiter_future", None) # Spawn as background task so drain doesn't block and # eviction-to-make-room stays responsive on subsequent ticks. From 53a7f53844c2627d5924dbe84333dadf58d5cd22 Mon Sep 17 00:00:00 2001 From: Hogne <227774406+hognek@users.noreply.github.com> Date: Sat, 18 Jul 2026 12:34:12 +0200 Subject: [PATCH 4/8] fix(gpu-arbiter): skip cancelled entries in retry/re-queue loop Same race as the admission-path fix: cancel_op may fire during the admission-await window (failed-admission path). The retry list re-queues cancelled entries because it never checks _cancelled_ids. Add the guard so cancelled tasks aren't resurrected. CodeRabbit finding on PR #1984 --- tinyagentos/scheduler/gpu_arbiter.py | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/tinyagentos/scheduler/gpu_arbiter.py b/tinyagentos/scheduler/gpu_arbiter.py index a87d37e0c..1ef25ca16 100644 --- a/tinyagentos/scheduler/gpu_arbiter.py +++ b/tinyagentos/scheduler/gpu_arbiter.py @@ -636,6 +636,11 @@ def _propagate(ct: asyncio.Task, f: asyncio.Future = future) -> None: retry.append(entry) # Re-queue tasks that still can't be admitted for entry in retry: + # Skip entries cancelled during the admission-await window + # (same race: cancel_op fired while we awaited reserve_and_check). + if entry.task.id in self._cancelled_ids: + self._cancelled_ids.discard(entry.task.id) + continue if not self._queue.full(): self._queue.put_nowait(entry) self._queued_entries[entry.task.id] = entry From c41ac97b7b76fcb689b30f16ee677a3a9b4ed85f Mon Sep 17 00:00:00 2001 From: Hogne <227774406+hognek@users.noreply.github.com> Date: Sat, 18 Jul 2026 12:52:14 +0200 Subject: [PATCH 5/8] fix(gpu-arbiter): make cancellation atomic across physical queue, shadow queue, tombstones - cancel_op: release VRAM reservation immediately so capacity is freed rather than leaking until _drain_queue eventually skips the entry. - retry/re-queue loop: re-check _cancelled_ids after put_nowait+dict insert to prevent TOCTOU resurrection of cancelled shadow entries. CodeRabbit finding on PR #1984 (GPU queue A2). --- tinyagentos/scheduler/gpu_arbiter.py | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/tinyagentos/scheduler/gpu_arbiter.py b/tinyagentos/scheduler/gpu_arbiter.py index 1ef25ca16..965dfec1d 100644 --- a/tinyagentos/scheduler/gpu_arbiter.py +++ b/tinyagentos/scheduler/gpu_arbiter.py @@ -644,6 +644,14 @@ def _propagate(ct: asyncio.Task, f: asyncio.Future = future) -> None: if not self._queue.full(): self._queue.put_nowait(entry) self._queued_entries[entry.task.id] = entry + # Re-check _cancelled_ids after insertion — cancel_op may have + # fired between the check above and the put_nowait+dict insert, + # which would resurrect the shadow entry and leak a queue slot. + if entry.task.id in self._cancelled_ids: + self._queued_entries.pop(entry.task.id, None) + self._cancelled_ids.discard(entry.task.id) + # The physical PriorityQueue entry can't be removed but + # _drain_queue skips cancelled entries at dequeue time. else: self._dropped += 1 future = getattr(entry.task, "_arbiter_future", None) @@ -694,6 +702,10 @@ async def cancel_op(self, task_id: str) -> bool: entry = self._queued_entries.pop(task_id, None) if entry is not None: self._cancelled_ids.add(task_id) + # Release the VRAM reservation immediately so capacity is freed + # rather than leaking until _drain_queue eventually skips the entry. + # _release_reservation is idempotent — safe if no reservation exists. + self._release_reservation(task_id) future = getattr(entry.task, "_arbiter_future", None) if future is not None and not future.done(): future.cancel() From 1331ea00e30aa513145c5f43d28971fb94f13bff Mon Sep 17 00:00:00 2001 From: Hogne <227774406+hognek@users.noreply.github.com> Date: Sat, 18 Jul 2026 14:01:50 +0200 Subject: [PATCH 6/8] =?UTF-8?q?fix(gpu-arbiter):=20fix=204=20pre-deploy=20?= =?UTF-8?q?defects=20=E2=80=94=20queue=20resilience,=20None=20guard,=20res?= =?UTF-8?q?ource=5Fid=20threading,=20lease=20renewal?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - _process_queue: wrap _drain_queue in try/except so one bad task cannot permanently kill the queue processor loop (defect 1). - _check_cluster_admission: skip workers whose free_vram_mb is None (non-NVIDIA workers) to prevent TypeErrors (defect 2). - _drain_queue: thread resource_id through _QueuedGpuTask so drained tasks hold a cluster lease on the correct resource (defect 3). - _run_gpu_task: add _renew_lease_loop background coroutine that periodically calls renew_lease every 200 s so the lease never expires mid-run (defect 4). Fixes: #1864 --- tinyagentos/scheduler/gpu_arbiter.py | 62 +++++++++++++++++++++++++++- 1 file changed, 60 insertions(+), 2 deletions(-) diff --git a/tinyagentos/scheduler/gpu_arbiter.py b/tinyagentos/scheduler/gpu_arbiter.py index 965dfec1d..45aec0940 100644 --- a/tinyagentos/scheduler/gpu_arbiter.py +++ b/tinyagentos/scheduler/gpu_arbiter.py @@ -58,6 +58,7 @@ class _QueuedGpuTask: op: str = field(default="inference", compare=False) model: str | None = field(default=None, compare=False) backend_name: str | None = field(default=None, compare=False) + resource_id: str | None = field(default=None, compare=False) @dataclass @@ -363,6 +364,7 @@ async def submit_gpu( required_vram_mb=required_vram_mb, evictable=evictable, required_gpu_arch=required_gpu_arch, op=op, model=model, backend_name=backend_name, + resource_id=resource_id, ) await self._queue.put(entry) self._queued_entries[task.id] = entry @@ -396,6 +398,8 @@ def _check_cluster_admission(self, required_vram_mb: int) -> GpuAdmission: for worker in self._cluster_manager.get_workers(): if worker.status != "online": continue + if worker.free_vram_mb is None: + continue # non-NVIDIA worker — no VRAM probe worker_leases = sum( l.required_vram_mb for l in leases if self._resource_on_worker(l.resource_id, worker.name) and l.required_vram_mb > 0 @@ -412,6 +416,37 @@ def _check_cluster_admission(self, required_vram_mb: int) -> GpuAdmission: reason=f"no cluster worker with {required_vram_mb} MiB free VRAM", ) + async def _renew_lease_loop(self, lease_id: str, stop_event: asyncio.Event) -> None: + """Periodically renew a lease until *stop_event* is set. + + Renews every 200 s so the lease never expires mid-run (the claim + TTL is 300 s). Returns silently when the lease expires or + renewal fails — the caller's finally path still releases it. + """ + RENEW_INTERVAL = 200 + while not stop_event.is_set(): + try: + await asyncio.wait_for(stop_event.wait(), timeout=RENEW_INTERVAL) + return # stop_event set — exit cleanly + except asyncio.TimeoutError: + pass # it's time to renew + try: + if self._cluster_manager is not None: + renewed = await self._cluster_manager.renew_lease( + lease_id, ttl_seconds=300, + ) + if renewed is None: + logger.warning( + "gpu-arbiter: lease %s expired mid-renewal", lease_id, + ) + return + logger.debug("gpu-arbiter: renewed lease %s", lease_id) + except Exception: + logger.exception( + "gpu-arbiter: failed to renew lease %s", lease_id, + ) + return + async def _run_gpu_task( self, task: Task, required_vram_mb: int, evictable: bool, resource_id: str | None, ) -> object: @@ -426,6 +461,8 @@ async def _run_gpu_task( claim_lease fails (taOS #1705 — reservation leak fix). """ lease_id: str | None = None + renew_task: asyncio.Task | None = None + renew_stop: asyncio.Event | None = None try: if self._cluster_manager is not None and resource_id is not None: lease = await self._cluster_manager.claim_lease( @@ -437,6 +474,13 @@ async def _run_gpu_task( f"GPU lease claim failed for {resource_id} (task {task.id})" ) lease_id = lease.lease_id + # Start periodic lease renewal so the lease doesn't expire + # mid-run (taOS #1864 defect 4 — fixed 300 s TTL). + renew_stop = asyncio.Event() + renew_task = asyncio.create_task( + self._renew_lease_loop(lease_id, renew_stop), + name=f"gpu-arbiter-renew-{task.id}", + ) current = asyncio.current_task() async with self._running_lock: self._running[task.id] = (task, lease_id, int(task.priority), required_vram_mb) @@ -453,6 +497,15 @@ async def _run_gpu_task( task.id, task.priority, required_vram_mb) raise finally: + # Stop lease renewal (taOS #1864 defect 4). + if renew_stop is not None: + renew_stop.set() + if renew_task is not None: + renew_task.cancel() + try: + await renew_task + except asyncio.CancelledError: + pass # Release the VRAM reservation whether we completed, errored, # or were cancelled. _evict_task handles its own reservation # release so idempotency matters. @@ -559,7 +612,12 @@ async def _process_queue(self) -> None: while True: await asyncio.sleep(2) if not self._paused: # taOS #796: skip drain while paused - await self._drain_queue() + try: + await self._drain_queue() + except Exception: + logger.exception( + "gpu-arbiter: _drain_queue raised — continuing loop" + ) except asyncio.CancelledError: raise @@ -608,7 +666,7 @@ async def _drain_queue(self) -> None: # Spawn as background task so drain doesn't block and # eviction-to-make-room stays responsive on subsequent ticks. t = asyncio.create_task( - self._run_gpu_task(entry.task, entry.required_vram_mb, entry.evictable, None), + self._run_gpu_task(entry.task, entry.required_vram_mb, entry.evictable, entry.resource_id), name=f"gpu-arbiter-drain-{entry.task.id}", ) From 9680011064b74db98e326c5d823a31aaf0ca651a Mon Sep 17 00:00:00 2001 From: Hogne <227774406+hognek@users.noreply.github.com> Date: Sat, 18 Jul 2026 14:18:02 +0200 Subject: [PATCH 7/8] fix(gpu-arbiter): cancel task on lease-loss and harden _drain_queue error handling MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two Kilo SUGGESTIONS from PR #1984 review: 1. gpu_arbiter.py:438 — When renew_lease returns None, the renew loop silently returns but _run_gpu_task keeps executing GPU work with no active lease. Fix: _renew_lease_loop now accepts an optional task_to_cancel parameter and cancels the parent asyncio Task on renewal failure or expiry. 2. gpu_arbiter.py:617 — Bare except Exception around _drain_queue() swallows programming bugs (KeyError, TypeError, etc.) and retries forever in a tight 2s loop. Fix: narrow to transient errors (NoResourceAvailableError, TimeoutError, OSError), add exponential backoff (2s→60s), and stop queue processor after 10 consecutive failures. --- tinyagentos/scheduler/gpu_arbiter.py | 43 ++++++++++++++++++++++------ 1 file changed, 35 insertions(+), 8 deletions(-) diff --git a/tinyagentos/scheduler/gpu_arbiter.py b/tinyagentos/scheduler/gpu_arbiter.py index 45aec0940..9fcdd173c 100644 --- a/tinyagentos/scheduler/gpu_arbiter.py +++ b/tinyagentos/scheduler/gpu_arbiter.py @@ -416,12 +416,17 @@ def _check_cluster_admission(self, required_vram_mb: int) -> GpuAdmission: reason=f"no cluster worker with {required_vram_mb} MiB free VRAM", ) - async def _renew_lease_loop(self, lease_id: str, stop_event: asyncio.Event) -> None: + async def _renew_lease_loop( + self, lease_id: str, stop_event: asyncio.Event, + task_to_cancel: asyncio.Task | None = None, + ) -> None: """Periodically renew a lease until *stop_event* is set. Renews every 200 s so the lease never expires mid-run (the claim - TTL is 300 s). Returns silently when the lease expires or - renewal fails — the caller's finally path still releases it. + TTL is 300 s). When *task_to_cancel* is provided and a renewal + failure or expiry is detected the referenced task is cancelled so + it doesn't keep executing GPU work without a valid lease (taOS + #1984 — lease-loss detection). """ RENEW_INTERVAL = 200 while not stop_event.is_set(): @@ -437,14 +442,20 @@ async def _renew_lease_loop(self, lease_id: str, stop_event: asyncio.Event) -> N ) if renewed is None: logger.warning( - "gpu-arbiter: lease %s expired mid-renewal", lease_id, + "gpu-arbiter: lease %s expired mid-renewal — " + "cancelling running task", lease_id, ) + if task_to_cancel is not None and not task_to_cancel.done(): + task_to_cancel.cancel() return logger.debug("gpu-arbiter: renewed lease %s", lease_id) except Exception: logger.exception( - "gpu-arbiter: failed to renew lease %s", lease_id, + "gpu-arbiter: failed to renew lease %s — " + "cancelling running task", lease_id, ) + if task_to_cancel is not None and not task_to_cancel.done(): + task_to_cancel.cancel() return async def _run_gpu_task( @@ -477,8 +488,9 @@ async def _run_gpu_task( # Start periodic lease renewal so the lease doesn't expire # mid-run (taOS #1864 defect 4 — fixed 300 s TTL). renew_stop = asyncio.Event() + current = asyncio.current_task() renew_task = asyncio.create_task( - self._renew_lease_loop(lease_id, renew_stop), + self._renew_lease_loop(lease_id, renew_stop, current), name=f"gpu-arbiter-renew-{task.id}", ) current = asyncio.current_task() @@ -609,15 +621,30 @@ async def _evict_task(self, task_id: str) -> int: async def _process_queue(self) -> None: try: + consecutive_failures = 0 + MAX_CONSECUTIVE_FAILURES = 10 while True: await asyncio.sleep(2) if not self._paused: # taOS #796: skip drain while paused try: await self._drain_queue() - except Exception: + consecutive_failures = 0 # reset on success + except (NoResourceAvailableError, asyncio.TimeoutError, OSError): + consecutive_failures += 1 + backoff = min(2 * (2 ** consecutive_failures), 60) logger.exception( - "gpu-arbiter: _drain_queue raised — continuing loop" + "gpu-arbiter: _drain_queue raised (consecutive=%d/%d) — " + "backing off %ds", + consecutive_failures, MAX_CONSECUTIVE_FAILURES, backoff, ) + if consecutive_failures >= MAX_CONSECUTIVE_FAILURES: + logger.critical( + "gpu-arbiter: _drain_queue failed %d consecutive " + "times — stopping queue processor", + consecutive_failures, + ) + return + await asyncio.sleep(backoff) except asyncio.CancelledError: raise From cad6bd0131b8c48d21eb74efd388dc030ece0960 Mon Sep 17 00:00:00 2001 From: Hogne <227774406+hognek@users.noreply.github.com> Date: Sat, 18 Jul 2026 14:43:58 +0200 Subject: [PATCH 8/8] fix(gpu-arbiter): restart queue processor after MAX_CONSECUTIVE_FAILURES instead of dying permanently Previously _process_queue returned after 10 consecutive transient failures (NoResourceAvailableError, TimeoutError, OSError), permanently killing the queue processor with no restart logic. A temporary cluster outage or GPU restart would kill the queue forever until manual intervention. Wrap the inner drain loop in an outer restart loop that waits a 300s cooldown after hitting the failure limit, then re-enters the inner loop with a fresh failure counter. CancelledError propagation is preserved. Kilo WARNING #1984 --- tinyagentos/scheduler/gpu_arbiter.py | 49 ++++++++++++++++------------ 1 file changed, 28 insertions(+), 21 deletions(-) diff --git a/tinyagentos/scheduler/gpu_arbiter.py b/tinyagentos/scheduler/gpu_arbiter.py index 9fcdd173c..dd4240210 100644 --- a/tinyagentos/scheduler/gpu_arbiter.py +++ b/tinyagentos/scheduler/gpu_arbiter.py @@ -621,30 +621,37 @@ async def _evict_task(self, task_id: str) -> int: async def _process_queue(self) -> None: try: - consecutive_failures = 0 MAX_CONSECUTIVE_FAILURES = 10 + COOLDOWN_BACKOFF = 300 # seconds to wait before restarting the processor while True: - await asyncio.sleep(2) - if not self._paused: # taOS #796: skip drain while paused - try: - await self._drain_queue() - consecutive_failures = 0 # reset on success - except (NoResourceAvailableError, asyncio.TimeoutError, OSError): - consecutive_failures += 1 - backoff = min(2 * (2 ** consecutive_failures), 60) - logger.exception( - "gpu-arbiter: _drain_queue raised (consecutive=%d/%d) — " - "backing off %ds", - consecutive_failures, MAX_CONSECUTIVE_FAILURES, backoff, - ) - if consecutive_failures >= MAX_CONSECUTIVE_FAILURES: - logger.critical( - "gpu-arbiter: _drain_queue failed %d consecutive " - "times — stopping queue processor", - consecutive_failures, + consecutive_failures = 0 + while True: + await asyncio.sleep(2) + if not self._paused: # taOS #796: skip drain while paused + try: + await self._drain_queue() + consecutive_failures = 0 # reset on success + except (NoResourceAvailableError, asyncio.TimeoutError, OSError): + consecutive_failures += 1 + backoff = min(2 * (2 ** consecutive_failures), 60) + logger.exception( + "gpu-arbiter: _drain_queue raised (consecutive=%d/%d) — " + "backing off %ds", + consecutive_failures, MAX_CONSECUTIVE_FAILURES, backoff, ) - return - await asyncio.sleep(backoff) + if consecutive_failures >= MAX_CONSECUTIVE_FAILURES: + logger.critical( + "gpu-arbiter: _drain_queue failed %d consecutive " + "times — restarting queue processor after %ds cooldown", + consecutive_failures, COOLDOWN_BACKOFF, + ) + break + await asyncio.sleep(backoff) + # Outer restart loop: wait cooldown, then restart the inner loop + await asyncio.sleep(COOLDOWN_BACKOFF) + logger.warning( + "gpu-arbiter: restarting queue processor after cooldown" + ) except asyncio.CancelledError: raise