From fd0967ff78ae993fa8a10a1fea0da198b36c04c1 Mon Sep 17 00:00:00 2001 From: rookiestar28 Date: Fri, 27 Mar 2026 15:29:37 +0800 Subject: [PATCH] refactor: add delta-first event and task cursors --- api/events.py | 20 ++++ api/model_manager.py | 25 ++++- services/model_manager.py | 11 ++ services/model_manager_tasks.py | 83 ++++++++++++++- services/model_manager_transfer.py | 6 ++ tests/test_model_manager_service.py | 105 +++++++++++++++++++ tests/test_r95_management_query_contracts.py | 6 ++ web/admin_console_app.js | 35 +++++-- web/openclaw_api.js | 1 + web/tabs/model_manager_tab.js | 57 ++++++++-- web/tests/unit/model_manager_tab.test.js | 21 +++- 11 files changed, 346 insertions(+), 24 deletions(-) diff --git a/api/events.py b/api/events.py index 232eaac..2358415 100644 --- a/api/events.py +++ b/api/events.py @@ -323,6 +323,26 @@ async def events_poll_handler(request: web.Request) -> web.Response: "cursor_status": cursor_status, "warnings": page.warnings, }, + "delta": { + "cursor_key": "since", + "requested_since_seq": since_requested, + "effective_since_seq": since_effective, + "next_since_seq": (events[-1].seq if events else since_effective), + "latest_seq": latest_seq, + "earliest_retained_seq": scan.get("earliest_retained_seq"), + "latest_retained_seq": scan.get("latest_retained_seq"), + "cursor_status": cursor_status, + "snapshot": since_requested == 0, + "truncated": bool( + scan.get("truncated") + or ( + events + and isinstance(scan.get("latest_retained_seq"), int) + and int(scan.get("latest_retained_seq")) > int(events[-1].seq) + ) + ), + "warnings": page.warnings, + }, "scan": scan, } ) diff --git a/api/model_manager.py b/api/model_manager.py index 7c8c6ec..f1e597d 100644 --- a/api/model_manager.py +++ b/api/model_manager.py @@ -17,6 +17,7 @@ try: RoutePlane, endpoint_metadata, ) + from ..services.management_query import normalize_cursor_limit from ..services.model_manager import ModelManagerError, model_manager from ..services.tenant_context import TenantBoundaryError, request_tenant_scope except ImportError: # pragma: no cover @@ -30,6 +31,7 @@ except ImportError: # pragma: no cover RoutePlane, endpoint_metadata, ) + from services.management_query import normalize_cursor_limit # type: ignore from services.model_manager import ModelManagerError, model_manager # type: ignore from services.tenant_context import ( # type: ignore TenantBoundaryError, @@ -159,6 +161,22 @@ async def model_download_list_handler(request: web.Request) -> web.Response: if deny: return deny token_info = resolve_token_info(request) + since_seq = None + delta_warnings = [] + if "since_seq" in request.query: + page = normalize_cursor_limit( + request.query, + cursor_key="since_seq", + default_cursor=0, + min_cursor=0, + default_limit=100, + max_limit=200, + ) + since_seq = int(page.cursor or 0) + delta_warnings = list(page.warnings) + limit = page.limit + else: + limit = _parse_int(request.query.get("limit"), 100, 1, 200) try: with request_tenant_scope( request=request, token_info=token_info, allow_default_when_missing=True @@ -166,9 +184,14 @@ async def model_download_list_handler(request: web.Request) -> web.Response: result = model_manager.list_download_tasks( tenant_id=tenant.tenant_id, state=request.query.get("state", ""), - limit=_parse_int(request.query.get("limit"), 100, 1, 200), + limit=limit, offset=_parse_int(request.query.get("offset"), 0, 0, 10_000), + since_seq=since_seq, ) + if since_seq is not None: + result.setdefault("pagination", {})["warnings"] = delta_warnings + if "delta" in result: + result["delta"]["warnings"] = delta_warnings return _json({"ok": True, **result}) except TenantBoundaryError as exc: return _json({"ok": False, "error": exc.code, "detail": str(exc)}, 403) diff --git a/services/model_manager.py b/services/model_manager.py index 0b7adec..6948c8a 100644 --- a/services/model_manager.py +++ b/services/model_manager.py @@ -138,6 +138,7 @@ class DownloadTask: resume_status: str = "not_started" recovery_attempts: int = 0 last_checkpoint_at: float = 0.0 + change_seq: int = 0 def is_terminal(self) -> bool: return self.state in {"completed", "failed", "cancelled"} @@ -174,6 +175,7 @@ class DownloadTask: "resume_status": self.resume_status, "recovery_attempts": self.recovery_attempts, "last_checkpoint_at": self.last_checkpoint_at, + "change_seq": self.change_seq, } @classmethod @@ -211,6 +213,7 @@ class DownloadTask: resume_status=str(payload.get("resume_status") or "not_started"), recovery_attempts=max(0, int(payload.get("recovery_attempts") or 0)), last_checkpoint_at=float(payload.get("last_checkpoint_at") or 0.0), + change_seq=max(0, int(payload.get("change_seq") or 0)), ) @@ -433,6 +436,7 @@ class ModelManager: self._tasks: Dict[str, DownloadTask] = {} self._futures: Dict[str, Future] = {} self._cancel_events: Dict[str, threading.Event] = {} + self._task_change_seq = 0 self._last_tasks_persist_at = 0.0 self._download_task_cls = DownloadTask self._download_cancelled_cls = DownloadCancelled @@ -619,6 +623,11 @@ class ModelManager: event_store_getter=get_job_event_store, ) + def _bump_task_change_seq_locked(self, task: DownloadTask) -> int: + self._task_change_seq += 1 + task.change_seq = self._task_change_seq + return task.change_seq + def _load_installations(self) -> List[Dict[str, Any]]: return _load_installations_impl(installations_path=self.installations_path) @@ -762,6 +771,7 @@ class ModelManager: state: str = "", limit: int = 100, offset: int = 0, + since_seq: Optional[int] = None, ) -> Dict[str, Any]: return _list_download_tasks_impl( manager=self, @@ -769,6 +779,7 @@ class ModelManager: state=state, limit=limit, offset=offset, + since_seq=since_seq, ) def get_download_task( diff --git a/services/model_manager_tasks.py b/services/model_manager_tasks.py index 441c990..8ef46bf 100644 --- a/services/model_manager_tasks.py +++ b/services/model_manager_tasks.py @@ -51,6 +51,10 @@ def load_tasks_from_disk( if not task.task_id: continue manager._tasks[task.task_id] = task + manager._task_change_seq = max( + int(getattr(manager, "_task_change_seq", 0)), + int(getattr(task, "change_seq", 0)), + ) if not task.is_terminal(): manager._cancel_events[task.task_id] = ( manager._threading_event_factory() @@ -72,6 +76,7 @@ def recover_incomplete_tasks(*, manager: Any) -> None: task.error = "restart_recovery_pending" task.recovery_attempts += 1 task.resume_status = "restart_recovering" + manager._bump_task_change_seq_locked(task) replayable = recoverable[: manager.recovery_replay_limit] overflow = recoverable[manager.recovery_replay_limit :] for task in overflow: @@ -80,6 +85,7 @@ def recover_incomplete_tasks(*, manager: Any) -> None: task.resume_status = "recovery_replay_limit_exceeded" task.finished_at = now task.updated_at = now + manager._bump_task_change_seq_locked(task) manager._emit(task) for task in replayable: event = manager._cancel_events.setdefault( @@ -91,6 +97,7 @@ def recover_incomplete_tasks(*, manager: Any) -> None: task.error = "" task.updated_at = now task.resume_status = "restart_replay_queued" + manager._bump_task_change_seq_locked(task) manager._futures[task.task_id] = manager._executor.submit( manager._run_task, task.task_id ) @@ -194,6 +201,7 @@ def set_resume_status(*, manager: Any, task_id: str, status: str) -> None: return current.resume_status = str(status or "not_started")[:120] current.updated_at = time.time() + manager._bump_task_change_seq_locked(current) manager._persist_tasks_locked(force=True) @@ -266,6 +274,7 @@ def progress(*, manager: Any, task_id: str, downloaded: int, total: int) -> None else 0.0 ) task.updated_at = time.time() + manager._bump_task_change_seq_locked(task) manager._emit(task) manager._persist_tasks_locked(force=False) @@ -277,12 +286,14 @@ def list_download_tasks( state: str = "", limit: int = 100, offset: int = 0, + since_seq: Optional[int] = None, ) -> Dict[str, Any]: limit = max(1, min(200, int(limit))) offset = max(0, int(offset)) state_filter = str(state or "").strip().lower() with manager._lock: tasks = list(manager._tasks.values()) + latest_change_seq = int(getattr(manager, "_task_change_seq", 0)) out = [] for task in tasks: if not manager._tenant_ok(task.tenant_id, tenant_id): @@ -290,13 +301,75 @@ def list_download_tasks( if state_filter and task.state != state_filter: continue out.append(task) - out.sort(key=lambda x: x.created_at, reverse=True) + if since_seq is None: + out.sort(key=lambda x: x.created_at, reverse=True) + total = len(out) + page = [item.to_dict() for item in out[offset : offset + limit]] + return { + "tasks": page, + "pagination": {"limit": limit, "offset": offset, "total": total}, + "filters": {"state": state_filter or None}, + } + + requested_since_seq = max(0, int(since_seq)) + effective_since_seq = requested_since_seq + cursor_status = "ok" + if requested_since_seq > latest_change_seq: + cursor_status = "future_cursor_reset" + effective_since_seq = latest_change_seq + + available_change_seqs = sorted( + int(getattr(task, "change_seq", 0)) + for task in out + if int(getattr(task, "change_seq", 0)) > 0 + ) + earliest_available_seq = available_change_seqs[0] if available_change_seqs else None + latest_available_seq = available_change_seqs[-1] if available_change_seqs else None + if ( + earliest_available_seq is not None + and effective_since_seq != 0 + and effective_since_seq < (earliest_available_seq - 1) + ): + cursor_status = "stale_cursor_reset" + effective_since_seq = max(0, earliest_available_seq - 1) + + out = [ + task + for task in out + if int(getattr(task, "change_seq", 0)) > effective_since_seq + ] + out.sort(key=lambda x: (int(getattr(x, "change_seq", 0)), x.created_at)) total = len(out) - page = [item.to_dict() for item in out[offset : offset + limit]] + page_items = out[:limit] + next_since_seq = ( + int(getattr(page_items[-1], "change_seq", 0)) + if page_items + else effective_since_seq + ) + truncated = bool( + total > len(page_items) + or ( + isinstance(latest_available_seq, int) + and latest_available_seq > next_since_seq + ) + ) return { - "tasks": page, - "pagination": {"limit": limit, "offset": offset, "total": total}, + "tasks": [item.to_dict() for item in page_items], + "pagination": {"limit": limit, "offset": 0, "total": total}, "filters": {"state": state_filter or None}, + "delta": { + "cursor_key": "since_seq", + "requested_since_seq": requested_since_seq, + "effective_since_seq": effective_since_seq, + "next_since_seq": next_since_seq, + "latest_change_seq": latest_change_seq, + "earliest_available_seq": earliest_available_seq, + "latest_available_seq": latest_available_seq, + "cursor_status": cursor_status, + "snapshot": False, + "truncated": truncated, + "warnings": [], + }, } @@ -333,12 +406,14 @@ def cancel_download_task( return task.to_dict() task.cancel_requested = True task.updated_at = time.time() + manager._bump_task_change_seq_locked(task) event.set() if task.state == "queued" and future is not None and future.cancel(): task.state = "cancelled" task.error = "cancelled_before_start" task.finished_at = time.time() task.updated_at = task.finished_at + manager._bump_task_change_seq_locked(task) manager._emit(task) manager._persist_tasks_locked(force=True) return task.to_dict() diff --git a/services/model_manager_transfer.py b/services/model_manager_transfer.py index cf01eea..b5dba40 100644 --- a/services/model_manager_transfer.py +++ b/services/model_manager_transfer.py @@ -126,6 +126,7 @@ def create_download_task( ) with manager._lock: manager._tasks[task.task_id] = task + manager._bump_task_change_seq_locked(task) manager._cancel_events[task.task_id] = manager._threading_event_factory() manager._futures[task.task_id] = manager._executor.submit( manager._run_task, task.task_id @@ -145,6 +146,7 @@ def run_task(*, manager: Any, task_id: str) -> None: task.started_at = time.time() task.updated_at = task.started_at task.resume_status = task.resume_status or "running" + manager._bump_task_change_seq_locked(task) manager._emit(task) manager._persist_tasks_locked(force=True) try: @@ -159,6 +161,7 @@ def run_task(*, manager: Any, task_id: str) -> None: current.progress = 1.0 current.staged_path = staged_path current.computed_sha256 = digest + manager._bump_task_change_seq_locked(current) manager._emit(current) manager._persist_tasks_locked(force=True) except manager._download_cancelled_cls: @@ -170,6 +173,7 @@ def run_task(*, manager: Any, task_id: str) -> None: current.error = "cancelled" current.updated_at = time.time() current.finished_at = current.updated_at + manager._bump_task_change_seq_locked(current) manager._emit(current) manager._persist_tasks_locked(force=True) except Exception as exc: @@ -181,6 +185,7 @@ def run_task(*, manager: Any, task_id: str) -> None: current.error = str(exc) current.updated_at = time.time() current.finished_at = current.updated_at + manager._bump_task_change_seq_locked(current) manager._emit(current) manager._persist_tasks_locked(force=True) @@ -482,6 +487,7 @@ def import_downloaded_model( current.installation_path = rec["installation_path"] current.installation_record_id = rec["id"] current.updated_at = time.time() + manager._bump_task_change_seq_locked(current) manager._emit(current) manager._persist_tasks_locked(force=True) return rec diff --git a/tests/test_model_manager_service.py b/tests/test_model_manager_service.py index 6378c1c..70c926b 100644 --- a/tests/test_model_manager_service.py +++ b/tests/test_model_manager_service.py @@ -228,6 +228,111 @@ class TestModelManagerService(unittest.TestCase): self.manager.import_downloaded_model(task_id=task.task_id) self.assertEqual(ctx.exception.code, "sha256_mismatch") + def test_list_download_tasks_delta_cursor_contract(self): + first = DownloadTask( + task_id="task-1", + model_id="model-1", + name="Model 1", + model_type="checkpoint", + source="catalog", + source_label="Catalog", + download_url="https://example.com/model-1.safetensors", + destination_subdir="checkpoints", + filename="model-1.safetensors", + expected_sha256="a" * 64, + provenance={ + "publisher": "OpenClaw", + "license": "OpenRAIL", + "source_url": "https://example.com/model-1", + }, + tenant_id="default", + created_at=10.0, + change_seq=4, + ) + second = DownloadTask( + task_id="task-2", + model_id="model-2", + name="Model 2", + model_type="checkpoint", + source="catalog", + source_label="Catalog", + download_url="https://example.com/model-2.safetensors", + destination_subdir="checkpoints", + filename="model-2.safetensors", + expected_sha256="b" * 64, + provenance={ + "publisher": "OpenClaw", + "license": "OpenRAIL", + "source_url": "https://example.com/model-2", + }, + tenant_id="default", + created_at=11.0, + change_seq=5, + ) + self.manager._tasks[first.task_id] = first + self.manager._tasks[second.task_id] = second + self.manager._task_change_seq = 5 + + result = self.manager.list_download_tasks(limit=10, since_seq=4) + self.assertEqual([row["task_id"] for row in result["tasks"]], ["task-2"]) + self.assertEqual(result["delta"]["requested_since_seq"], 4) + self.assertEqual(result["delta"]["effective_since_seq"], 4) + self.assertEqual(result["delta"]["next_since_seq"], 5) + self.assertEqual(result["delta"]["cursor_status"], "ok") + self.assertFalse(result["delta"]["truncated"]) + + def test_list_download_tasks_delta_resets_stale_cursor(self): + first = DownloadTask( + task_id="task-1", + model_id="model-1", + name="Model 1", + model_type="checkpoint", + source="catalog", + source_label="Catalog", + download_url="https://example.com/model-1.safetensors", + destination_subdir="checkpoints", + filename="model-1.safetensors", + expected_sha256="a" * 64, + provenance={ + "publisher": "OpenClaw", + "license": "OpenRAIL", + "source_url": "https://example.com/model-1", + }, + tenant_id="default", + created_at=10.0, + change_seq=7, + ) + second = DownloadTask( + task_id="task-2", + model_id="model-2", + name="Model 2", + model_type="checkpoint", + source="catalog", + source_label="Catalog", + download_url="https://example.com/model-2.safetensors", + destination_subdir="checkpoints", + filename="model-2.safetensors", + expected_sha256="b" * 64, + provenance={ + "publisher": "OpenClaw", + "license": "OpenRAIL", + "source_url": "https://example.com/model-2", + }, + tenant_id="default", + created_at=11.0, + change_seq=8, + ) + self.manager._tasks[first.task_id] = first + self.manager._tasks[second.task_id] = second + self.manager._task_change_seq = 8 + + result = self.manager.list_download_tasks(limit=1, since_seq=1) + self.assertEqual([row["task_id"] for row in result["tasks"]], ["task-1"]) + self.assertEqual(result["delta"]["effective_since_seq"], 6) + self.assertEqual(result["delta"]["next_since_seq"], 7) + self.assertEqual(result["delta"]["cursor_status"], "stale_cursor_reset") + self.assertTrue(result["delta"]["truncated"]) + @patch( "services.model_manager.validate_outbound_url", return_value=("https", "example.com", 443, ["1.1.1.1"]), diff --git a/tests/test_r95_management_query_contracts.py b/tests/test_r95_management_query_contracts.py index e50b535..d8b0a38 100644 --- a/tests/test_r95_management_query_contracts.py +++ b/tests/test_r95_management_query_contracts.py @@ -144,6 +144,10 @@ class TestR95EventsApi(unittest.IsolatedAsyncioTestCase): self.assertEqual(payload["pagination"]["cursor_status"], "stale_cursor_reset") self.assertEqual(payload["pagination"]["since_requested"], 3) self.assertEqual(payload["pagination"]["since_effective"], 49) + self.assertEqual(payload["delta"]["requested_since_seq"], 3) + self.assertEqual(payload["delta"]["effective_since_seq"], 49) + self.assertEqual(payload["delta"]["next_since_seq"], 51) + self.assertTrue(payload["delta"]["truncated"]) self.assertEqual([e["seq"] for e in payload["events"]], [50, 51]) self.assertEqual(len(store.calls), 2) @@ -181,6 +185,8 @@ class TestR95EventsApi(unittest.IsolatedAsyncioTestCase): payload = fake_web.json_response.call_args.args[0] self.assertEqual(payload["pagination"]["cursor_status"], "future_cursor_reset") self.assertEqual(payload["pagination"]["since_effective"], 10) + self.assertEqual(payload["delta"]["next_since_seq"], 10) + self.assertFalse(payload["delta"]["truncated"]) codes = {w["code"] for w in payload["pagination"]["warnings"]} self.assertIn("R95_INVALID_LIMIT", codes) self.assertIn("R95_STALE_CURSOR_FUTURE", codes) diff --git a/web/admin_console_app.js b/web/admin_console_app.js index 369faf7..a2b0531 100644 --- a/web/admin_console_app.js +++ b/web/admin_console_app.js @@ -71,15 +71,18 @@ export function mountAdminConsole(root = document) { elements.token.value = api.getToken(); function appendEvent(obj) { + const seq = Number(obj?.seq); + if (!Number.isNaN(seq)) { + if (seq <= api.state.lastSeq) { + return; + } + api.state.lastSeq = seq; + } const line = `[${now()}] ${JSON.stringify(obj)}`; const lines = elements.eventsBox.textContent ? elements.eventsBox.textContent.split("\n") : []; lines.push(line); elements.eventsBox.textContent = lines.slice(-120).join("\n"); elements.eventsBox.scrollTop = elements.eventsBox.scrollHeight; - const seq = Number(obj?.seq); - if (!Number.isNaN(seq) && seq > api.state.lastSeq) { - api.state.lastSeq = seq; - } } async function loadDashboard() { @@ -179,14 +182,24 @@ export function mountAdminConsole(root = document) { } async function pollEvents() { - const response = await api.request(`/events?since=${encodeURIComponent(String(api.state.lastSeq))}&limit=50`); - if (!response.ok) { - setStatus(elements.eventsStatus, `Events poll failed: ${response.error || "unknown"}`, "err"); - return; + let polls = 0; + let totalEvents = 0; + let cursor = api.state.lastSeq; + while (polls < 4) { + const response = await api.request(`/events?since=${encodeURIComponent(String(cursor))}&limit=50`); + if (!response.ok) { + setStatus(elements.eventsStatus, `Events poll failed: ${response.error || "unknown"}`, "err"); + return; + } + const events = Array.isArray(response.data?.events) ? response.data.events : []; + events.forEach(appendEvent); + totalEvents += events.length; + const nextSinceSeq = Number(response.data?.delta?.next_since_seq); + cursor = Number.isFinite(nextSinceSeq) ? nextSinceSeq : api.state.lastSeq; + polls += 1; + if (!response.data?.delta?.truncated) break; } - const events = response.data?.events || []; - events.forEach(appendEvent); - setStatus(elements.eventsStatus, `Polled ${events.length} events`, "ok"); + setStatus(elements.eventsStatus, `Polled ${totalEvents} events`, "ok"); } function disconnectSse() { diff --git a/web/openclaw_api.js b/web/openclaw_api.js index 709ae66..6c77e5e 100644 --- a/web/openclaw_api.js +++ b/web/openclaw_api.js @@ -813,6 +813,7 @@ export class OpenClawAPI { if (params.state) qs.set("state", String(params.state)); if (params.limit != null) qs.set("limit", String(params.limit)); if (params.offset != null) qs.set("offset", String(params.offset)); + if (params.since_seq != null) qs.set("since_seq", String(params.since_seq)); const suffix = qs.toString() ? `?${qs.toString()}` : ""; return this.fetch(`${this._path("/models/downloads")}${suffix}`, { headers: { ...this._adminTokenHeaders() } diff --git a/web/tabs/model_manager_tab.js b/web/tabs/model_manager_tab.js index 270a810..5168df1 100644 --- a/web/tabs/model_manager_tab.js +++ b/web/tabs/model_manager_tab.js @@ -128,6 +128,40 @@ function buildInstallationCard(item) { `; } +function sortTasksForDisplay(tasks) { + return [...tasks].sort((left, right) => Number(right?.created_at || 0) - Number(left?.created_at || 0)); +} + +export function mergeTaskDelta(existingTasks, deltaTasks) { + const merged = new Map(); + for (const task of Array.isArray(existingTasks) ? existingTasks : []) { + if (task?.task_id) { + merged.set(task.task_id, task); + } + } + for (const task of Array.isArray(deltaTasks) ? deltaTasks : []) { + if (task?.task_id) { + merged.set(task.task_id, task); + } + } + return sortTasksForDisplay([...merged.values()]); +} + +function resolveNextTaskCursor(payload, fallbackCursor = 0) { + const delta = payload?.delta || {}; + const nextSinceSeq = Number(delta?.next_since_seq); + if (Number.isFinite(nextSinceSeq) && nextSinceSeq >= 0) { + return nextSinceSeq; + } + const taskSeqs = Array.isArray(payload?.tasks) + ? payload.tasks.map((task) => Number(task?.change_seq)).filter((seq) => Number.isFinite(seq) && seq >= 0) + : []; + if (!taskSeqs.length) { + return fallbackCursor; + } + return Math.max(fallbackCursor, ...taskSeqs); +} + export const ModelManagerTab = { id: "model-manager", title: "Model Manager", @@ -211,6 +245,7 @@ export const ModelManagerTab = { tasks: [], installations: [], pollingTimer: null, + taskSinceSeq: 0, }; const modelManagerAction = { label: "Open Model Manager", @@ -296,12 +331,20 @@ export const ModelManagerTab = { renderResults(); }; - const loadTasks = async () => { - const res = await openclawApi.listModelDownloadTasks({ limit: 100, offset: 0 }); + const loadTasks = async ({ delta = false } = {}) => { + const params = { limit: 100, offset: 0 }; + if (delta && state.taskSinceSeq > 0) { + params.since_seq = state.taskSinceSeq; + } + const res = await openclawApi.listModelDownloadTasks(params); if (!res.ok) { throw new Error(res.error || "tasks_list_failed"); } - state.tasks = Array.isArray(res.data?.tasks) ? res.data.tasks : []; + const incomingTasks = Array.isArray(res.data?.tasks) ? res.data.tasks : []; + state.tasks = params.since_seq != null + ? mergeTaskDelta(state.tasks, incomingTasks) + : sortTasksForDisplay(incomingTasks); + state.taskSinceSeq = resolveNextTaskCursor(res.data, state.taskSinceSeq); renderTasks(); }; @@ -358,7 +401,7 @@ export const ModelManagerTab = { dedupeKey: "model-manager:queue-success", action: modelManagerAction, }); - await loadTasks(); + await loadTasks({ delta: false }); }; const cancelTask = async (taskId) => { @@ -368,7 +411,7 @@ export const ModelManagerTab = { reportIssue(`cancel failed: ${res.error || "request_failed"}`, "model-manager:cancel"); return; } - await loadTasks(); + await loadTasks({ delta: false }); }; const importTask = async (taskId) => { @@ -388,7 +431,7 @@ export const ModelManagerTab = { dedupeKey: "model-manager:import-success", action: modelManagerAction, }); - await loadTasks(); + await loadTasks({ delta: false }); await loadInstallations(); await loadSearch(); }; @@ -429,7 +472,7 @@ export const ModelManagerTab = { } const pane = container.closest(".openclaw-tab-pane"); if (pane && !pane.classList.contains("active")) return; - loadTasks().catch(() => { + loadTasks({ delta: true }).catch(() => { // Poll refresh errors are surfaced by explicit refresh actions. }); }, 3000); diff --git a/web/tests/unit/model_manager_tab.test.js b/web/tests/unit/model_manager_tab.test.js index e50de95..ad2cc0b 100644 --- a/web/tests/unit/model_manager_tab.test.js +++ b/web/tests/unit/model_manager_tab.test.js @@ -22,7 +22,7 @@ vi.mock("../../openclaw_api.js", () => ({ vi.mock("../../openclaw_utils.js", () => utilsMock); -import { ModelManagerTab } from "../../tabs/model_manager_tab.js"; +import { ModelManagerTab, mergeTaskDelta } from "../../tabs/model_manager_tab.js"; describe("model_manager_tab", () => { beforeEach(() => { @@ -69,4 +69,23 @@ describe("model_manager_tab", () => { }) ); }); + + it("merges task deltas without duplicating existing rows", () => { + const merged = mergeTaskDelta( + [ + { task_id: "task-1", state: "running", created_at: 10, change_seq: 3 }, + { task_id: "task-2", state: "queued", created_at: 11, change_seq: 4 }, + ], + [ + { task_id: "task-1", state: "completed", created_at: 10, change_seq: 5 }, + { task_id: "task-3", state: "queued", created_at: 12, change_seq: 6 }, + ] + ); + + expect(merged).toEqual([ + expect.objectContaining({ task_id: "task-3", state: "queued" }), + expect.objectContaining({ task_id: "task-2", state: "queued" }), + expect.objectContaining({ task_id: "task-1", state: "completed" }), + ]); + }); });