mirror of
https://github.com/rookiestar28/ComfyUI-OpenClaw.git
synced 2026-08-14 00:48:07 +00:00
refactor: add delta-first event and task cursors
This commit is contained in:
@@ -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,
|
||||
}
|
||||
)
|
||||
|
||||
+24
-1
@@ -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)
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"]),
|
||||
|
||||
@@ -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)
|
||||
|
||||
+24
-11
@@ -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() {
|
||||
|
||||
@@ -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() }
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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" }),
|
||||
]);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user