From d22a236e9eec1969447a77d6a08fc85a2e9792ea Mon Sep 17 00:00:00 2001 From: duguwanglong Date: Mon, 20 Jul 2026 16:16:09 +0800 Subject: [PATCH] fix(runtime): reduce noisy warnings and harden startup Prefer user-defined plugin resources over built-ins and skip stale workflow trigger configs without alarming startup logs. Classify expected runtime conditions consistently, deduplicate MCP failures, and add regression coverage for startup isolation and notification loading. --- flocks/agent/agent_factory.py | 34 ++-- flocks/channel/builtin/weixin/channel.py | 2 +- flocks/ingest/kafka/manager.py | 38 +++- flocks/ingest/syslog/manager.py | 35 +++- flocks/mcp/client.py | 10 +- flocks/mcp/errors.py | 35 ++++ flocks/mcp/server.py | 12 +- flocks/notifications/service.py | 1 - flocks/server/app.py | 3 +- flocks/server/routes/notifications.py | 2 - flocks/server/routes/session.py | 163 +++++++----------- flocks/server/routes/workflow.py | 11 +- flocks/session/features/memory.py | 2 +- flocks/session/runner.py | 21 +-- flocks/skill/skill.py | 47 +++-- flocks/tool/device/startup.py | 3 + flocks/tool/file/glob.py | 6 +- flocks/tool/registry.py | 16 +- flocks/workflow/center.py | 111 ++++++++++-- flocks/workflow/fs_store.py | 9 +- flocks/workflow/poller_manager.py | 34 +++- tests/agent/test_agent_factory.py | 38 ++++ tests/channel/test_weixin_log_levels.py | 47 +++++ tests/ingest/test_kafka_manager.py | 94 ++++++++++ .../test_syslog_manager_bind_failure.py | 99 +++++++++-- tests/mcp/test_mcp_client_sse.py | 21 +++ tests/mcp/test_mcp_server.py | 67 +++++++ tests/provider/test_api_service_management.py | 52 ++++++ .../routes/test_notifications_routes.py | 6 +- tests/server/test_app_errors.py | 42 +++++ tests/server/test_input_dispatcher.py | 21 ++- tests/session/test_runner_step.py | 19 ++ .../session/test_session_memory_log_levels.py | 40 +++++ tests/skill/test_skill.py | 19 +- tests/tool/test_device_startup_sync.py | 28 +++ tests/tool/test_tools.py | 49 +++++- tests/workflow/test_fs_store.py | 39 +++++ tests/workflow/test_poller_manager.py | 70 +++++++- tests/workflow/test_workflow_paths.py | 98 ++++++++++- webui/src/api/notifications.ts | 6 +- webui/src/components/layout/Layout.test.tsx | 32 +++- webui/src/components/layout/Layout.tsx | 20 +-- 42 files changed, 1248 insertions(+), 254 deletions(-) create mode 100644 flocks/mcp/errors.py create mode 100644 tests/channel/test_weixin_log_levels.py create mode 100644 tests/session/test_session_memory_log_levels.py diff --git a/flocks/agent/agent_factory.py b/flocks/agent/agent_factory.py index d8c08f265..2c2ac9b37 100644 --- a/flocks/agent/agent_factory.py +++ b/flocks/agent/agent_factory.py @@ -223,8 +223,9 @@ def scan_and_load(dirs: Optional[List[Path]] = None) -> Dict[str, AgentInfo]: """ Scan agent directories and load all valid agents. - Scans built-in agents first, then user plugin agents, then project plugin - agents. Name conflicts are skipped with a warning (first wins). + Scans built-in agents first, then user plugin agents, then project-bundled + agents. The first match wins, so user customizations take precedence over + project bundles while core built-ins remain protected. ``native`` is determined by the source directory, not by agent.yaml: - ``_BUILTIN_AGENTS_DIR`` → native=True @@ -238,20 +239,21 @@ def scan_and_load(dirs: Optional[List[Path]] = None) -> Dict[str, AgentInfo]: Returns: Dict mapping agent name → AgentInfo. """ - # (directory, is_native) pairs — order determines first-wins conflict resolution - search_dirs: List[tuple[Path, bool]] = [(_BUILTIN_AGENTS_DIR, True)] + # (directory, is_native, source) tuples — order determines first-wins priority. + search_dirs: List[tuple[Path, bool, str]] = [(_BUILTIN_AGENTS_DIR, True, "builtin")] if _PLUGIN_AGENTS_DIR.exists(): - search_dirs.append((_PLUGIN_AGENTS_DIR, False)) + search_dirs.append((_PLUGIN_AGENTS_DIR, False, "user")) # Project-level plugin agents: resolved at call time because cwd is dynamic project_agents_dir = Path.cwd() / ".flocks" / "plugins" / "agents" if project_agents_dir.exists() and project_agents_dir != _PLUGIN_AGENTS_DIR: - search_dirs.append((project_agents_dir, True)) + search_dirs.append((project_agents_dir, True, "project")) if dirs: - search_dirs.extend((d, False) for d in dirs) + search_dirs.extend((d, False, "extra") for d in dirs) result: Dict[str, AgentInfo] = {} + selected_sources: Dict[str, tuple[str, Path]] = {} - for scan_dir, is_native in search_dirs: + for scan_dir, is_native, source in search_dirs: if not scan_dir.is_dir(): continue for agent_dir in _iter_agent_dirs(scan_dir): @@ -259,13 +261,21 @@ def scan_and_load(dirs: Optional[List[Path]] = None) -> Dict[str, AgentInfo]: if agent is None: continue if agent.name in result: - log.warn("agent.factory.name_conflict", { + selected_source, selected_dir = selected_sources[agent.name] + details = { "name": agent.name, - "existing_source": "previous scan", - "skipped_source": str(agent_dir), - }) + "selected_source": selected_source, + "selected_path": str(selected_dir), + "skipped_source": source, + "skipped_path": str(agent_dir), + } + if source != selected_source and source != "extra": + log.info("agent.factory.lower_priority_skipped", details) + else: + log.warn("agent.factory.name_conflict", details) continue result[agent.name] = agent + selected_sources[agent.name] = (source, agent_dir) log.debug("agent.factory.loaded", { "name": agent.name, "dir": str(agent_dir), diff --git a/flocks/channel/builtin/weixin/channel.py b/flocks/channel/builtin/weixin/channel.py index cb1e30d5e..fb1c51c7b 100644 --- a/flocks/channel/builtin/weixin/channel.py +++ b/flocks/channel/builtin/weixin/channel.py @@ -237,7 +237,7 @@ async def start( "base_url": self._base_url, }) if self._group_policy != "disabled": - log.warning("weixin.group_policy.note", { + log.info("weixin.group_policy.note", { "group_policy": self._group_policy, "note": ( "QR-login connects an iLink bot identity (e.g. ...@im.bot), not a " diff --git a/flocks/ingest/kafka/manager.py b/flocks/ingest/kafka/manager.py index 16acca4b6..d265b0320 100644 --- a/flocks/ingest/kafka/manager.py +++ b/flocks/ingest/kafka/manager.py @@ -297,7 +297,14 @@ async def start_all(self) -> None: if not workflow_id: continue if isinstance(data, dict) and data.get("enabled"): - await self.restart_workflow(workflow_id) + try: + await self.restart_workflow(workflow_id, startup=True) + except Exception as exc: + self._status[workflow_id] = {"state": "failed", "error": str(exc)} + log.warning( + "kafka.start_failed", + {"workflow_id": workflow_id, "error": str(exc)}, + ) async def stop_all(self) -> None: for workflow_id in list(self._tasks.keys()): @@ -360,7 +367,12 @@ async def stop_workflow(self, workflow_id: str) -> None: if workflow_id in self._status: self._status[workflow_id] = {"state": "stopped", "error": None} - async def restart_workflow(self, workflow_id: str) -> Dict[str, Any]: + async def restart_workflow( + self, + workflow_id: str, + *, + startup: bool = False, + ) -> Dict[str, Any]: """Restart the consumer and return its post-connect runtime status. Blocks until the consumer connects, the connection fails, or @@ -377,6 +389,21 @@ async def restart_workflow(self, workflow_id: str) -> Dict[str, Any]: self._status[workflow_id] = {"state": "stopped", "error": None} return {"state": "stopped", "error": None} + # Load and cache the workflow JSON once; avoids a disk read per message. + wf_data = read_workflow_from_fs(workflow_id) + if not wf_data: + err = "workflow_not_found" + if startup: + self._status[workflow_id] = {"state": "stopped", "error": err} + log.info("kafka.workflow_not_found_on_start", { + "workflow_id": workflow_id, + "action": "stale_config_skipped", + }) + return {"state": "stopped", "error": err} + self._status[workflow_id] = {"state": "failed", "error": err} + log.warning("kafka.workflow_not_found", {"workflow_id": workflow_id}) + return {"state": "failed", "error": err} + input_broker = str(data.get("inputBroker") or "").strip() input_topic = str(data.get("inputTopic") or "").strip() if not input_broker or not input_topic: @@ -385,13 +412,6 @@ async def restart_workflow(self, workflow_id: str) -> Dict[str, Any]: log.warning("kafka.config_incomplete", {"workflow_id": workflow_id}) return {"state": "failed", "error": err} - # Load and cache the workflow JSON once; avoids a disk read per message. - wf_data = read_workflow_from_fs(workflow_id) - if not wf_data: - err = "workflow_not_found" - self._status[workflow_id] = {"state": "failed", "error": err} - log.warning("kafka.workflow_not_found_on_start", {"workflow_id": workflow_id}) - return {"state": "failed", "error": err} workflow_json = wf_data.get("workflowJson") if not workflow_json: err = "workflow_json_missing" diff --git a/flocks/ingest/syslog/manager.py b/flocks/ingest/syslog/manager.py index 3fccaed88..8e6913b31 100644 --- a/flocks/ingest/syslog/manager.py +++ b/flocks/ingest/syslog/manager.py @@ -186,7 +186,17 @@ async def start_all(self) -> None: if not workflow_id: continue if isinstance(data, dict) and data.get("enabled"): - await self.restart_workflow(workflow_id) + try: + await self.restart_workflow(workflow_id, startup=True) + except Exception as exc: + self._listener_status[workflow_id] = { + "state": "failed", + "error": str(exc), + } + log.warning( + "syslog.start_failed", + {"workflow_id": workflow_id, "error": str(exc)}, + ) async def stop_all(self) -> None: for workflow_id in list(self._tasks.keys()): @@ -245,7 +255,12 @@ async def stop_workflow(self, workflow_id: str) -> None: if workflow_id in self._listener_status: self._listener_status[workflow_id] = {"state": "stopped", "error": None} - async def restart_workflow(self, workflow_id: str) -> Dict[str, Any]: + async def restart_workflow( + self, + workflow_id: str, + *, + startup: bool = False, + ) -> Dict[str, Any]: """Restart the listener and return its post-bind runtime status. This call blocks until the underlying socket either binds successfully, @@ -268,8 +283,15 @@ async def restart_workflow(self, workflow_id: str) -> Dict[str, Any]: wf_data = read_workflow_from_fs(workflow_id) if not wf_data: err = "workflow_not_found" + if startup: + self._listener_status[workflow_id] = {"state": "stopped", "error": err} + log.info("syslog.workflow_not_found_on_start", { + "workflow_id": workflow_id, + "action": "stale_config_skipped", + }) + return {"state": "stopped", "error": err} self._listener_status[workflow_id] = {"state": "failed", "error": err} - log.warning("syslog.workflow_not_found_on_start", {"workflow_id": workflow_id}) + log.warning("syslog.workflow_not_found", {"workflow_id": workflow_id}) return {"state": "failed", "error": err} workflow_json = wf_data.get("workflowJson") if not workflow_json: @@ -286,6 +308,10 @@ async def restart_workflow(self, workflow_id: str) -> Dict[str, Any]: self._listener_status[workflow_id] = {"state": "failed", "error": err} log.warning("syslog.workflow_plan_failed", {"workflow_id": workflow_id, "error": str(exc)}) return self.get_listener_status(workflow_id) + + host = str(data.get("host") or "0.0.0.0") + port = int(data.get("port") or 5140) + protocol = str(data.get("protocol") or "udp").lower() queue: asyncio.Queue = asyncio.Queue(maxsize=_MAX_QUEUE_SIZE) self._queues[workflow_id] = queue @@ -295,9 +321,6 @@ async def restart_workflow(self, workflow_id: str) -> Dict[str, Any]: ready = asyncio.Event() self._listener_ready[workflow_id] = ready - host = str(data.get("host") or "0.0.0.0") - port = int(data.get("port") or 5140) - protocol = str(data.get("protocol") or "udp").lower() self._listener_status[workflow_id] = { "state": "binding", "error": None, diff --git a/flocks/mcp/client.py b/flocks/mcp/client.py index c2bf675ae..8acdf8e33 100644 --- a/flocks/mcp/client.py +++ b/flocks/mcp/client.py @@ -703,15 +703,7 @@ async def list_resources(self) -> List[McpResource]: Raises: RuntimeError: If not connected """ - try: - result = await self._submit_command("list_resources") - return result - except Exception as exc: - log.error("mcp.client.list_resources_error", { - "server": self.name, - "error": str(exc), - }) - raise + return await self._submit_command("list_resources") async def read_resource(self, uri: str) -> Any: """ diff --git a/flocks/mcp/errors.py b/flocks/mcp/errors.py new file mode 100644 index 000000000..7849bf81a --- /dev/null +++ b/flocks/mcp/errors.py @@ -0,0 +1,35 @@ +"""Helpers for classifying MCP protocol errors.""" + +from __future__ import annotations + +from mcp.types import METHOD_NOT_FOUND + + +def is_method_not_found_error(exc: BaseException) -> bool: + """Return whether an exception chain contains JSON-RPC method-not-found.""" + return _contains_method_not_found(exc, seen=set()) + + +def _contains_method_not_found(exc: BaseException, *, seen: set[int]) -> bool: + exc_id = id(exc) + if exc_id in seen: + return False + seen.add(exc_id) + + error = getattr(exc, "error", None) + if getattr(error, "code", None) == METHOD_NOT_FOUND: + return True + + if isinstance(exc, BaseExceptionGroup): + if any(_contains_method_not_found(child, seen=seen) for child in exc.exceptions): + return True + + cause = exc.__cause__ + if cause is not None and _contains_method_not_found(cause, seen=seen): + return True + context = exc.__context__ + if context is not None and _contains_method_not_found(context, seen=seen): + return True + + message = str(exc).strip().lower() + return message == "method not found" or message.endswith(": method not found") diff --git a/flocks/mcp/server.py b/flocks/mcp/server.py index b7a20cd10..ae6dd0240 100644 --- a/flocks/mcp/server.py +++ b/flocks/mcp/server.py @@ -8,6 +8,7 @@ import time from typing import Dict, Optional, Any, List from flocks.mcp.client import McpClient +from flocks.mcp.errors import is_method_not_found_error from flocks.mcp.types import ( McpStatus, McpStatusInfo, @@ -226,10 +227,13 @@ async def _connect_and_register( resources = await client.list_resources() self._resources_cache[name] = resources except Exception as e: - log.warn("mcp.resources_unavailable", { - "server": name, - "error": str(e) - }) + if is_method_not_found_error(e): + log.info("mcp.resources_unsupported", {"server": name}) + else: + log.warn("mcp.resources_unavailable", { + "server": name, + "error": str(e) + }) self._resources_cache[name] = [] # 5. Register tools diff --git a/flocks/notifications/service.py b/flocks/notifications/service.py index 931fc2fc0..3cf3a02cd 100644 --- a/flocks/notifications/service.py +++ b/flocks/notifications/service.py @@ -206,7 +206,6 @@ async def list_active( *, user_id: str, locale: str | None = None, - current_version: str | None = None, ) -> list[NotificationResponse]: target_locale = _normalize_locale(locale) config_notifications = await cls._load_config_notifications() diff --git a/flocks/server/app.py b/flocks/server/app.py index 78271d130..c71a1bbb6 100644 --- a/flocks/server/app.py +++ b/flocks/server/app.py @@ -1111,7 +1111,8 @@ async def validation_exception_handler(request: Request, exc: RequestValidationE @app.exception_handler(StarletteHTTPException) async def http_exception_handler(request: Request, exc: StarletteHTTPException): """Handle HTTP exceptions""" - log.error("http.error", { + http_log = log.warn if 400 <= exc.status_code < 500 else log.error + http_log("http.error", { "path": request.url.path, "status": exc.status_code, "detail": exc.detail, diff --git a/flocks/server/routes/notifications.py b/flocks/server/routes/notifications.py index 709a9240f..0792abd3b 100644 --- a/flocks/server/routes/notifications.py +++ b/flocks/server/routes/notifications.py @@ -26,13 +26,11 @@ async def list_active_notifications( request: Request, locale: str | None = None, - current_version: str | None = None, ) -> list[NotificationResponse]: user = require_user(request) return await NotificationService.list_active( user_id=user.id, locale=locale, - current_version=current_version, ) diff --git a/flocks/server/routes/session.py b/flocks/server/routes/session.py index 3cf75b890..3a6fc98a5 100644 --- a/flocks/server/routes/session.py +++ b/flocks/server/routes/session.py @@ -9,6 +9,7 @@ import asyncio import json import time +from pathlib import Path from typing import List, Optional, Any, Dict, Literal, Union, Tuple from fastapi import APIRouter, HTTPException, status, Query, Request from fastapi.responses import StreamingResponse @@ -48,6 +49,62 @@ # extension (e.g. a PNG named report.pdf.exe whose tail would otherwise be # ".exe"). _UPLOAD_SAFE_EXTS = frozenset({"png", "jpg", "jpeg", "gif", "webp", "bmp", "pdf"}) +_UPLOAD_EXT_BY_MIME = { + "image/png": ".png", + "image/jpeg": ".jpg", + "image/jpg": ".jpg", + "image/gif": ".gif", + "image/webp": ".webp", + "image/bmp": ".bmp", + "application/pdf": ".pdf", +} + + +def _session_uploads_dir(session_id: str) -> Path: + """Return the application-owned upload directory for one session.""" + from flocks.config.config import Config + + uploads_root = (Config.get_data_path() / "uploads").resolve() + target = (uploads_root / session_id).resolve() + if target == uploads_root or not target.is_relative_to(uploads_root): + raise ValueError(f"Invalid session ID for upload path: {session_id}") + return target + + +def _materialize_data_url_part( + session_id: str, + data_url: str, + mime_hint: str, + filename_hint: Optional[str], + *, + failure_event: str = "session.prompt_queue.materialize_failed", +) -> str: + """Persist a data URL under the application data directory.""" + try: + import base64 + from flocks.utils.id import Identifier + + _header, _sep, encoded = data_url.partition(",") + if not encoded: + return data_url + raw_bytes = base64.b64decode(encoded) + uploads_root = _session_uploads_dir(session_id) + uploads_root.mkdir(parents=True, exist_ok=True) + + ext = _UPLOAD_EXT_BY_MIME.get(mime_hint, "") + if not ext and filename_hint: + _, _, tail = filename_hint.rpartition(".") + if tail.lower() in _UPLOAD_SAFE_EXTS: + ext = "." + tail.lower() + target = uploads_root / f"{Identifier.create('part')}{ext}" + target.write_bytes(raw_bytes) + return target.resolve().as_uri() + except Exception as exc: + log.warn(failure_event, { + "sessionID": session_id, + "error": str(exc), + }) + return data_url def _context_usage_cache_key(session_id: str, session: SessionModel) -> Tuple[str, int]: @@ -782,17 +839,14 @@ async def delete_session(sessionID: str, request: Request) -> bool: await Session.delete(session.project_id, sessionID) # Best-effort cleanup of any image/file uploads materialised for this - # session via ``_materialize_data_url_to_disk`` (see prompt_async). + # session via ``_materialize_data_url_part`` (see prompt_async). # The session DB row is gone, so the on-disk bytes are now orphaned — - # remove them to keep the workspace tidy. We deliberately swallow any + # remove them to keep application data tidy. We deliberately swallow any # filesystem errors: deletion of the session record is the contract, # the upload cleanup is incidental. try: import shutil - from flocks.workspace.manager import WorkspaceManager - - ws = WorkspaceManager.get_instance() - uploads_root = ws.resolve_workspace_path(f"uploads/{sessionID}") + uploads_root = _session_uploads_dir(sessionID) if uploads_root.exists() and uploads_root.is_dir(): shutil.rmtree(uploads_root, ignore_errors=True) log.info("session.uploads.cleaned", { @@ -2764,51 +2818,6 @@ async def _process_session_message( # ------------------------------------------------------------------ from flocks.session.message import FilePart - def _materialize_data_url_to_disk( - data_url: str, mime_hint: str, filename_hint: Optional[str] - ) -> str: - """Decode a ``data:`` URL to ``~/.flocks/workspace/uploads//...``. - - Returns a ``file://`` URL pointing at the persisted file. On failure - the original ``data:`` URL is returned unchanged (older code paths - still cope with that, just with the now-known token-cost penalty). - """ - try: - import base64 - from flocks.workspace.manager import WorkspaceManager - - header, _, encoded = data_url.partition(",") - if not encoded: - return data_url - raw_bytes = base64.b64decode(encoded) - - ws = WorkspaceManager.get_instance() - # Use resolve_workspace_path to guard against path traversal if - # sessionID were ever user-controlled (e.g. ../../../tmp/x). - uploads_root = ws.resolve_workspace_path(f"uploads/{sessionID}") - uploads_root.mkdir(parents=True, exist_ok=True) - - ext_map = { - "image/png": ".png", "image/jpeg": ".jpg", "image/jpg": ".jpg", - "image/gif": ".gif", "image/webp": ".webp", "image/bmp": ".bmp", - "application/pdf": ".pdf", - } - ext = ext_map.get(mime_hint, "") - if not ext and filename_hint: - _, _, tail = filename_hint.rpartition(".") - if tail.lower() in _UPLOAD_SAFE_EXTS: - ext = "." + tail.lower() - unique_name = f"{Identifier.create('part')}{ext}" - target = uploads_root / unique_name - target.write_bytes(raw_bytes) - return target.resolve().as_uri() - except Exception as exc: - log.warn("session.message.file_part.materialize_failed", { - "sessionID": sessionID, - "error": str(exc), - }) - return data_url - for raw_part in request.parts or []: part_type = raw_part.get("type") if part_type == "text": @@ -2824,7 +2833,13 @@ def _materialize_data_url_to_disk( continue # Materialize ``data:`` URLs to disk before persisting the part. if url.startswith("data:"): - url = _materialize_data_url_to_disk(url, mime, raw_part.get("filename")) + url = _materialize_data_url_part( + sessionID, + url, + mime, + raw_part.get("filename"), + failure_event="session.message.file_part.materialize_failed", + ) file_part_id = raw_part.get("id") or Identifier.create("part") file_part = FilePart( id=file_part_id, @@ -3277,50 +3292,6 @@ def _materialize_queued_parts(session_id: str, parts: List[Dict[str, Any]]) -> L return prepared -def _materialize_data_url_part( - session_id: str, - data_url: str, - mime_hint: str, - filename_hint: Optional[str], -) -> str: - try: - import base64 - from flocks.workspace.manager import WorkspaceManager - from flocks.utils.id import Identifier - - _header, _sep, encoded = data_url.partition(",") - if not encoded: - return data_url - raw_bytes = base64.b64decode(encoded) - ws = WorkspaceManager.get_instance() - uploads_root = ws.resolve_workspace_path(f"uploads/{session_id}") - uploads_root.mkdir(parents=True, exist_ok=True) - - ext_map = { - "image/png": ".png", - "image/jpeg": ".jpg", - "image/jpg": ".jpg", - "image/gif": ".gif", - "image/webp": ".webp", - "image/bmp": ".bmp", - "application/pdf": ".pdf", - } - ext = ext_map.get(mime_hint, "") - if not ext and filename_hint: - _, _, tail = filename_hint.rpartition(".") - if tail.lower() in _UPLOAD_SAFE_EXTS: - ext = "." + tail.lower() - target = uploads_root / f"{Identifier.create('part')}{ext}" - target.write_bytes(raw_bytes) - return target.resolve().as_uri() - except Exception as exc: - log.warn("session.prompt_queue.materialize_failed", { - "sessionID": session_id, - "error": str(exc), - }) - return data_url - - def _event_from_queued_prompt(item, working_directory: str): from flocks.input.events import UserInputEvent diff --git a/flocks/server/routes/workflow.py b/flocks/server/routes/workflow.py index 1bf338498..bc7c5cc23 100644 --- a/flocks/server/routes/workflow.py +++ b/flocks/server/routes/workflow.py @@ -329,9 +329,9 @@ def _workflow_integration_config_key(workflow_id: str) -> str: def _read_workflow_from_fs(workflow_id: str) -> Optional[Dict[str, Any]]: """Read workflow data from the filesystem. - Search order (lowest → highest priority), same roots as - resolve_global_workflow_roots / resolve_project_workflow_roots; per-id dir - is ``//`` with ``workflow.json`` inside. + Search order (lowest → highest priority) visits project-bundled roots before + global user roots; per-id dir is ``//`` with ``workflow.json`` + inside. """ return shared_read_workflow_from_fs(workflow_id) @@ -531,9 +531,8 @@ def _scan_workflow_base_dir(base_dir: Path, source: str) -> Dict[str, Dict[str, def _list_workflows_from_fs() -> List[Dict[str, Any]]: """Scan global and project workflow directories and return merged list. - Scan order matches *_all_scan_dirs()* (lowest -> highest priority): each - root from *resolve_global_workflow_roots* then each from - *resolve_project_workflow_roots(workspace)*; under each root, immediate + Scan order matches *_all_scan_dirs()* (lowest -> highest priority): project + bundle roots first, then global user roots; under each root, immediate subdirectories with *workflow.json* are workflows. Later entries override earlier ones when the workflow directory name (*id*) diff --git a/flocks/session/features/memory.py b/flocks/session/features/memory.py index 7c2370bfb..418531b5b 100644 --- a/flocks/session/features/memory.py +++ b/flocks/session/features/memory.py @@ -66,7 +66,7 @@ async def initialize(self) -> bool: memory_config_dict = config.memory if hasattr(config, 'memory') and config.memory else None if not memory_config_dict: - log.warn("session.memory.no_config", {"session_id": self.session_id}) + log.info("session.memory.no_config", {"session_id": self.session_id}) memory_config = MemoryConfig(enabled=True) else: if isinstance(memory_config_dict, dict): diff --git a/flocks/session/runner.py b/flocks/session/runner.py index e9c10b5ae..b605cca4d 100644 --- a/flocks/session/runner.py +++ b/flocks/session/runner.py @@ -1341,18 +1341,15 @@ async def device_asset_prompt_factory() -> Optional[str]: except Exception as e: error_attempt += 1 - log.error("runner.step.error", { - "error": str(e), - "attempt": error_attempt, - }) - + # Convert exception to error dict for retry check error_dict = self._exception_to_error_dict(e) - + # Check if retryable retry_message = SessionRetry.retryable(error_dict) + will_retry = retry_message is not None and error_attempt <= MAX_ERROR_RETRIES - if retry_message is not None and error_attempt <= MAX_ERROR_RETRIES: + if will_retry: # Error is retryable and we have budget left delay_ms = SessionRetry.delay(error_attempt, error_dict) # Always cap the sleep to RETRY_MAX_DELAY_NO_HEADERS so a @@ -1361,8 +1358,9 @@ async def device_asset_prompt_factory() -> Optional[str]: from flocks.session.lifecycle.retry import RETRY_MAX_DELAY_NO_HEADERS delay_ms = min(delay_ms, RETRY_MAX_DELAY_NO_HEADERS) next_retry_time = int(asyncio.get_event_loop().time() * 1000) + delay_ms - - log.info("runner.step.retry", { + + log.warn("runner.step.retry", { + "error": str(e), "attempt": error_attempt, "delay_ms": delay_ms, "reason": retry_message, @@ -1393,7 +1391,10 @@ async def device_asset_prompt_factory() -> Optional[str]: "max_retries": MAX_ERROR_RETRIES, }) else: - log.error("runner.step.not_retryable", {"error": str(e)}) + log.error("runner.step.not_retryable", { + "error": str(e), + "attempt": error_attempt, + }) final_error_message = str(e) if SessionRetry.is_connection_error(error_dict): diff --git a/flocks/skill/skill.py b/flocks/skill/skill.py index db5f2e1e5..3b6be7b08 100644 --- a/flocks/skill/skill.py +++ b/flocks/skill/skill.py @@ -198,10 +198,9 @@ class Skill: Skill discovery and management. Discovers SKILL.md files from (lowest → highest priority): - - .flocks dirs (global + project-level) - - .claude dirs (global ~/.claude + project-level) - - ~/.flocks (global user-level) - - /.flocks (project-level, wins on collision) + - .claude dirs + - global built-in and project-bundled .flocks dirs + - ~/.flocks/plugins (user-level customizations, wins on collision) """ _cache: Optional[Dict[str, SkillInfo]] = None @@ -392,11 +391,21 @@ def _scan_directory( skill_info = cls._parse_skill_md(match, source=source) if skill_info: if skill_info.name in skills: - log.warn("skill.duplicate", { - "name": skill_info.name, - "existing": skills[skill_info.name].location, - "duplicate": match, - }) + existing = skills[skill_info.name] + if existing.source != skill_info.source: + log.info("skill.override", { + "name": skill_info.name, + "selected": match, + "selected_source": skill_info.source, + "replaced": existing.location, + "replaced_source": existing.source, + }) + else: + log.warn("skill.duplicate", { + "name": skill_info.name, + "existing": existing.location, + "duplicate": match, + }) skills[skill_info.name] = skill_info log.debug("skill.found", { @@ -416,9 +425,10 @@ def _discover(cls) -> Dict[str, SkillInfo]: Discover all skills. Last wins on name collision. Scan order (lowest → highest priority): - 1. .claude dirs (global ~/.claude + project-level) - 2. ~/.flocks (global user-level, overrides .claude) - 3. /.flocks (project-level, highest priority) + 1. .claude dirs + 2. ~/.flocks/skill[s] built-ins + 3. /.flocks project-bundled skills + 4. ~/.flocks/plugins user customizations Source labels: "flocks" — built-in skills inside .flocks/skills/ directories @@ -450,16 +460,12 @@ def _discover(cls) -> Dict[str, SkillInfo]: for claude_dir in cls._find_dirs_up(".claude", current_dir, worktree): cls._scan_directory(claude_dir, "skills/**/SKILL.md", skills, source="claude") - # 2) Global ~/.flocks — overrides .claude - # Built-in skills: source="flocks"; user-installed plugins: source="user" + # 2) Global built-in skills — overrides .claude if os.path.isdir(global_flocks): for pattern in builtin_patterns: cls._scan_directory(global_flocks, pattern, skills, source="flocks") - for pattern in plugin_patterns: - cls._scan_directory(global_flocks, pattern, skills, source="user") - # 3) Project-level .flocks — highest priority - # Built-in skills: source="flocks"; project-installed plugins: source="project" + # 3) Project-bundled .flocks skills — built-in baseline for this project for flocks_dir in cls._find_dirs_up(".flocks", current_dir, worktree): if os.path.normpath(flocks_dir) == os.path.normpath(global_flocks): continue @@ -468,6 +474,11 @@ def _discover(cls) -> Dict[str, SkillInfo]: for pattern in plugin_patterns: cls._scan_directory(flocks_dir, pattern, skills, source="project") + # 4) User-installed plugins — explicit customizations win over project bundles + if os.path.isdir(global_flocks): + for pattern in plugin_patterns: + cls._scan_directory(global_flocks, pattern, skills, source="user") + log.info("skill.discovery.complete", {"count": len(skills), "names": list(skills.keys())}) return skills diff --git a/flocks/tool/device/startup.py b/flocks/tool/device/startup.py index 82a32920a..7b1b0454b 100644 --- a/flocks/tool/device/startup.py +++ b/flocks/tool/device/startup.py @@ -7,6 +7,8 @@ """ from __future__ import annotations +import aiosqlite + from flocks.storage.storage import Storage from flocks.utils.log import Log @@ -40,6 +42,7 @@ async def _heal_stale_service_ids() -> None: from flocks.tool.device.store import storage_key_to_service_id async with Storage.connect(Storage.get_db_path()) as db: + db.row_factory = aiosqlite.Row cur = await db.execute("SELECT id, storage_key, service_id FROM device_integrations") rows = await cur.fetchall() updates: list[tuple[str, str]] = [] diff --git a/flocks/tool/file/glob.py b/flocks/tool/file/glob.py index c9f5eb8a0..3a5e9b769 100644 --- a/flocks/tool/file/glob.py +++ b/flocks/tool/file/glob.py @@ -24,6 +24,7 @@ # Constants MAX_FILES = 100 +_ripgrep_fallback_logged = False # Description matching Flocks' glob.txt @@ -195,7 +196,10 @@ async def glob_tool( 'mtime': mtime }) else: - log.warn("glob.ripgrep_not_found", {"fallback": "python_glob"}) + global _ripgrep_fallback_logged + if not _ripgrep_fallback_logged: + _ripgrep_fallback_logged = True + log.info("glob.ripgrep_not_found", {"fallback": "python_glob"}) for filepath in fallback_glob(search_path, pattern): if len(files) >= MAX_FILES: diff --git a/flocks/tool/registry.py b/flocks/tool/registry.py index d2bff0d7a..cb97ca1fb 100644 --- a/flocks/tool/registry.py +++ b/flocks/tool/registry.py @@ -376,6 +376,13 @@ def _normalize_param_key(name: str) -> str: return "".join(ch for ch in str(name).lower() if ch.isalnum()) +def _api_service_storage_key(provider: str) -> str: + """Resolve the single persisted key used for an API-like provider.""" + from flocks.config.api_versioning import versioned_storage_key_for + + return versioned_storage_key_for(provider) or provider + + _SCOPED_SCHEMA_ALIASES: Dict[str, Dict[str, str]] = { # SkyEye historically surfaced "威胁级别", which some callers guessed as # threat_level. Keep the schema canonical on hazard_level, but accept the @@ -516,7 +523,7 @@ async def execute(self, ctx: ToolContext, **kwargs) -> ToolResult: # Validate required parameters for required_param in schema.required: if required_param not in effective_kwargs: - log.error("tool.execute.missing_param", { + log.warn("tool.execute.missing_param", { "tool": self.info.name, "missing": required_param, "provided": list(effective_kwargs.keys()), @@ -1188,7 +1195,7 @@ def _bootstrap_user_api_services(cls) -> None: or not info.enabled ): continue - provider = info.provider + provider = _api_service_storage_key(info.provider) if provider in seen_providers: continue seen_providers.add(provider) @@ -1240,9 +1247,10 @@ def _sync_api_service_states(cls) -> None: disabled_count = 0 restored_count = 0 for tool in cls._tools.values(): - provider = tool.info.provider - if not provider: + configured_provider = tool.info.provider + if not configured_provider: continue + provider = _api_service_storage_key(configured_provider) svc = api_services.get(provider, {}) svc_enabled = svc.get("enabled", False) if not svc_enabled: diff --git a/flocks/workflow/center.py b/flocks/workflow/center.py index ad05aff3b..b0fb594bc 100644 --- a/flocks/workflow/center.py +++ b/flocks/workflow/center.py @@ -100,6 +100,11 @@ def _normalize_workflow_id(path: Path) -> str: return digest[:24] +def _logical_workflow_id(path: Path) -> str: + """Return the filesystem workflow ID shared across discovery roots.""" + return path.parent.name + + def _fingerprint(path: Path) -> str: digest = hashlib.sha256() with path.open("rb") as f: @@ -156,6 +161,17 @@ def resolve_project_workflow_roots(base_dir: Optional[Path] = None) -> list[Path ] +def resolve_workflow_scan_roots( + base_dir: Optional[Path] = None, +) -> list[tuple[Path, str]]: + """Return workflow roots ordered from lowest to highest priority.""" + return [ + (root, "project") for root in resolve_project_workflow_roots(base_dir) + ] + [ + (root, "global") for root in resolve_global_workflow_roots() + ] + + def _resolve_project_workflow_root(base_dir: Optional[Path] = None) -> Path: """/.flocks/plugins/workflows/ — canonical project-level workflow storage.""" root = base_dir or Path.cwd() @@ -256,6 +272,17 @@ async def _read_registry(workflow_id: str) -> Dict[str, Any]: return data +def _load_discoverable_workflow(workflow_path: Path) -> Optional[Dict[str, Any]]: + """Load and validate a workflow, returning None when intentionally hidden.""" + raw = json.loads(workflow_path.read_text(encoding="utf-8")) + meta_path = workflow_path.parent / "meta.json" + meta = json.loads(meta_path.read_text(encoding="utf-8")) if meta_path.is_file() else None + if is_hidden_workflow(raw, meta): + return None + Workflow.from_dict(raw) + return raw + + async def _scan_workflow_dir( workflow_root: Path, source_type: str, @@ -270,12 +297,9 @@ async def _scan_workflow_dir( return for workflow_path in sorted(workflow_root.glob("*/workflow.json")): try: - raw = json.loads(workflow_path.read_text(encoding="utf-8")) - meta_path = workflow_path.parent / "meta.json" - meta = json.loads(meta_path.read_text(encoding="utf-8")) if meta_path.is_file() else None - if is_hidden_workflow(raw, meta): + raw = _load_discoverable_workflow(workflow_path) + if raw is None: continue - Workflow.from_dict(raw) except Exception as exc: log.warning( "workflow.center.scan.skip_invalid", @@ -284,6 +308,7 @@ async def _scan_workflow_dir( continue workflow_id = _normalize_workflow_id(workflow_path) + logical_workflow_id = _logical_workflow_id(workflow_path) fp = _fingerprint(workflow_path) now_ms = _now_ms() existing = await WorkflowStore.kv_get(_registry_key(workflow_id)) or {} @@ -291,6 +316,7 @@ async def _scan_workflow_dir( draft_changed = bool(existing) and existing.get("fingerprint") != fp entry = { "workflowId": workflow_id, + "logicalWorkflowId": logical_workflow_id, "name": raw.get("name") or workflow_path.parent.name, "description": raw.get("description") or "", "sourceType": source_type, @@ -308,7 +334,7 @@ async def _scan_workflow_dir( "serviceUrl": existing.get("serviceUrl"), } await WorkflowStore.kv_put(_registry_key(workflow_id), entry) - by_id[workflow_id] = entry + by_id[logical_workflow_id] = entry async def scan_skill_workflows(base_dir: Optional[Path] = None) -> List[Dict[str, Any]]: @@ -316,21 +342,15 @@ async def scan_skill_workflows(base_dir: Optional[Path] = None) -> List[Dict[str Scan order (lowest → highest priority), see resolve_global_workflow_roots / resolve_project_workflow_roots: - 1. ~/.flocks/plugins/workflow/ (global legacy, sourceType="global") - 2. ~/.flocks/workflow/ (global compat, sourceType="global") - 3. ~/.flocks/plugins/workflows/ (global canonical, sourceType="global") - 4. /.flocks/plugins/workflow/ (project legacy, sourceType="project") - 5. /.flocks/workflow/ (project compat, sourceType="project") - 6. /.flocks/plugins/workflows/ (project canonical, sourceType="project") + 1. /.flocks project-bundled workflows + 2. ~/.flocks user workflows (highest priority) When two directories contain a workflow with the same ID, the later (higher-priority) entry wins. """ by_id: Dict[str, Dict[str, Any]] = {} - for path in resolve_global_workflow_roots(): - await _scan_workflow_dir(path, "global", by_id) - for path in resolve_project_workflow_roots(base_dir): - await _scan_workflow_dir(path, "project", by_id) + for path, source_type in resolve_workflow_scan_roots(base_dir): + await _scan_workflow_dir(path, source_type, by_id) entries = list(by_id.values()) entries.sort(key=lambda item: item.get("updatedAt", 0), reverse=True) @@ -376,16 +396,71 @@ def format_workflow_entries( async def list_registry_entries() -> List[Dict[str, Any]]: """List registered skill workflows.""" keys = await WorkflowStore.kv_list(_REGISTRY_PREFIX) - items: List[Dict[str, Any]] = [] + selected_entries: Dict[str, tuple[tuple[int, int, int], Dict[str, Any]]] = {} for raw_key in keys: key = _key_to_string(raw_key) entry = await WorkflowStore.kv_get(key) if entry: - items.append(entry) + logical_id = entry.get("logicalWorkflowId") + if not logical_id: + workflow_path = entry.get("workflowPath") + logical_id = ( + _logical_workflow_id(Path(str(workflow_path))) + if workflow_path + else str(entry.get("workflowId") or key) + ) + priority = _registry_entry_priority(entry) + if priority[1] >= 0 and not _registry_entry_is_discoverable(entry): + continue + # Only collapse entries that belong to the currently active project + # or user roots. Registrations from another project remain manageable. + selection_key = ( + f"logical:{logical_id}" + if priority[1] >= 0 + else f"registry:{entry.get('workflowId') or key}" + ) + selected = selected_entries.get(selection_key) + if selected is None or priority >= selected[0]: + selected_entries[selection_key] = (priority, entry) + items = [entry for _priority, entry in selected_entries.values()] items.sort(key=lambda item: item.get("updatedAt", 0), reverse=True) return items +def _registry_entry_priority(entry: Dict[str, Any]) -> tuple[int, int, int]: + """Rank registry entries using the same project-then-user scan order.""" + workflow_path = Path(str(entry.get("workflowPath") or "")) + try: + resolved_path = workflow_path.resolve() + except OSError: + resolved_path = workflow_path + + roots = [root for root, _source in resolve_workflow_scan_roots(Path.cwd())] + root_priority = -1 + for index, root in enumerate(roots): + try: + resolved_path.relative_to(root.resolve()) + except (OSError, ValueError): + continue + root_priority = index + break + + return ( + int(resolved_path.is_file()), + root_priority, + int(entry.get("updatedAt") or 0), + ) + + +def _registry_entry_is_discoverable(entry: Dict[str, Any]) -> bool: + """Return whether an active-root registry entry still passes discovery.""" + workflow_path = Path(str(entry.get("workflowPath") or "")) + try: + return _load_discoverable_workflow(workflow_path) is not None + except Exception: + return False + + def _service_release_file(workflow_id: str, release_id: str) -> Path: base = Config.get_data_path() / _SERVICE_DATA_DIR / "releases" / workflow_id base.mkdir(parents=True, exist_ok=True) diff --git a/flocks/workflow/fs_store.py b/flocks/workflow/fs_store.py index 8f1900268..0f3471959 100644 --- a/flocks/workflow/fs_store.py +++ b/flocks/workflow/fs_store.py @@ -9,7 +9,7 @@ from flocks.utils.log import Log -from .center import resolve_global_workflow_roots, resolve_project_workflow_roots +from .center import resolve_workflow_scan_roots log = Log.create(service="workflow.fs-store") @@ -104,12 +104,7 @@ def find_workspace_root() -> Path: def workflow_scan_dirs() -> list[tuple[Path, str]]: """Return all workflow roots ordered from lowest to highest priority.""" - workspace = find_workspace_root() - return [ - (root, "global") for root in resolve_global_workflow_roots() - ] + [ - (root, "project") for root in resolve_project_workflow_roots(workspace) - ] + return resolve_workflow_scan_roots(find_workspace_root()) def read_workflow_dir( diff --git a/flocks/workflow/poller_manager.py b/flocks/workflow/poller_manager.py index c4596286c..179376c48 100644 --- a/flocks/workflow/poller_manager.py +++ b/flocks/workflow/poller_manager.py @@ -191,7 +191,18 @@ async def start_all(self) -> None: if not workflow_id: continue if isinstance(data, dict) and data.get("enabled"): - await self.restart_workflow(workflow_id) + try: + await self.restart_workflow(workflow_id, startup=True) + except Exception as exc: + self._status[workflow_id] = { + **self._base_status(workflow_id), + "state": "failed", + "error": str(exc), + } + log.warning( + "poller.start_failed", + {"workflow_id": workflow_id, "error": str(exc)}, + ) async def stop_all(self) -> None: for workflow_id in list(self._tasks.keys()): @@ -229,7 +240,12 @@ async def stop_workflow(self, workflow_id: str) -> None: self._run_cancel_events.pop(workflow_id, None) self._status[workflow_id] = current - async def restart_workflow(self, workflow_id: str) -> Dict[str, Any]: + async def restart_workflow( + self, + workflow_id: str, + *, + startup: bool = False, + ) -> Dict[str, Any]: await self.stop_workflow(workflow_id) try: stored = await WorkflowStore.get_config(workflow_id, kind="workflow_poller_config") @@ -237,8 +253,7 @@ async def restart_workflow(self, workflow_id: str) -> Dict[str, Any]: log.warning("poller.restart_read_failed", {"workflow_id": workflow_id, "error": str(exc)}) return {"workflowId": workflow_id, "state": "failed", "error": str(exc)} - config = self._normalize_config(workflow_id, stored) - if not config.get("enabled"): + if not isinstance(stored, dict) or not stored.get("enabled"): self._status[workflow_id] = { **self._base_status(workflow_id), "workflowId": workflow_id, @@ -253,11 +268,20 @@ async def restart_workflow(self, workflow_id: str) -> Dict[str, Any]: self._status[workflow_id] = { **self.get_status(workflow_id), "workflowId": workflow_id, - "state": "failed", + "state": "stopped" if startup else "failed", "error": err, } + poller_log = log.info if startup else log.warning + poller_log( + "poller.workflow_not_found_on_start" if startup else "poller.workflow_not_found", + { + "workflow_id": workflow_id, + **({"action": "stale_config_skipped"} if startup else {}), + }, + ) return self.get_status(workflow_id) + config = self._normalize_config(workflow_id, stored) workflow_json = wf_data.get("workflowJson") if not workflow_json: err = "workflow_json_missing" diff --git a/tests/agent/test_agent_factory.py b/tests/agent/test_agent_factory.py index 30ce129e6..78f86fdac 100644 --- a/tests/agent/test_agent_factory.py +++ b/tests/agent/test_agent_factory.py @@ -15,6 +15,7 @@ import textwrap from pathlib import Path +from unittest.mock import MagicMock import pytest @@ -780,6 +781,43 @@ def test_user_plugin_agent_is_not_native( assert "custom-agent" in result assert result["custom-agent"].native is False + def test_user_plugin_agent_wins_over_project_bundle( + self, tmp_path: Path, monkeypatch: pytest.MonkeyPatch + ): + user_agents_dir = tmp_path / "user_plugins" / "agents" + project_dir = tmp_path / "project" + project_agents_dir = project_dir / ".flocks" / "plugins" / "agents" + self._write_agent(user_agents_dir, "shared-agent") + self._write_agent(project_agents_dir, "shared-agent") + (user_agents_dir / "shared-agent" / "agent.yaml").write_text( + "name: shared-agent\ndescription: user customization\nmode: subagent\n", + encoding="utf-8", + ) + (project_agents_dir / "shared-agent" / "agent.yaml").write_text( + "name: shared-agent\ndescription: project bundle\nmode: subagent\n", + encoding="utf-8", + ) + + info_log = MagicMock() + warn_log = MagicMock() + monkeypatch.setattr(_factory_module, "_PLUGIN_AGENTS_DIR", user_agents_dir) + monkeypatch.setattr(_factory_module.log, "info", info_log) + monkeypatch.setattr(_factory_module.log, "warn", warn_log) + monkeypatch.chdir(project_dir) + + result = scan_and_load() + + assert result["shared-agent"].description == "user customization" + assert any( + call.args and call.args[0] == "agent.factory.lower_priority_skipped" + for call in info_log.call_args_list + ) + assert not any( + call.args and call.args[0] == "agent.factory.name_conflict" + and call.args[1].get("name") == "shared-agent" + for call in warn_log.call_args_list + ) + # =========================================================================== # Project-level plugin agent CRUD diff --git a/tests/channel/test_weixin_log_levels.py b/tests/channel/test_weixin_log_levels.py new file mode 100644 index 000000000..cc083af66 --- /dev/null +++ b/tests/channel/test_weixin_log_levels.py @@ -0,0 +1,47 @@ +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock + +import pytest + +import flocks.channel.builtin.weixin.channel as weixin_module + + +@pytest.mark.asyncio +async def test_group_policy_note_logs_info(monkeypatch): + sessions = [MagicMock(), MagicMock()] + info_log = MagicMock() + warn_log = MagicMock() + warning_log = MagicMock() + + monkeypatch.setattr(weixin_module, "AIOHTTP_AVAILABLE", True) + monkeypatch.setattr(weixin_module, "CRYPTO_AVAILABLE", True) + monkeypatch.setattr( + weixin_module, + "aiohttp", + SimpleNamespace( + ClientTimeout=MagicMock(return_value=MagicMock()), + ClientSession=MagicMock(side_effect=sessions), + ), + ) + monkeypatch.setattr(weixin_module.ilink, "make_ssl_connector", MagicMock(return_value=None)) + monkeypatch.setattr(weixin_module, "ContextTokenStore", MagicMock(return_value=MagicMock())) + monkeypatch.setattr(weixin_module, "MessageDedup", MagicMock(return_value=MagicMock())) + monkeypatch.setattr(weixin_module, "MediaCache", MagicMock(return_value=MagicMock())) + monkeypatch.setattr(weixin_module.log, "info", info_log) + monkeypatch.setattr(weixin_module.log, "warn", warn_log) + monkeypatch.setattr(weixin_module.log, "warning", warning_log) + + channel = weixin_module.WeixinChannel() + monkeypatch.setattr(channel, "_poll_loop", AsyncMock(return_value=None)) + monkeypatch.setattr(channel, "_close_sessions", AsyncMock(return_value=None)) + + await channel.start( + {"token": "token", "accountId": "account", "groupPolicy": "all"}, + AsyncMock(), + ) + + assert any(call.args and call.args[0] == "weixin.group_policy.note" for call in info_log.call_args_list) + assert not any( + call.args and call.args[0] == "weixin.group_policy.note" + for call in warn_log.call_args_list + warning_log.call_args_list + ) diff --git a/tests/ingest/test_kafka_manager.py b/tests/ingest/test_kafka_manager.py index 31de29d02..e5803177a 100644 --- a/tests/ingest/test_kafka_manager.py +++ b/tests/ingest/test_kafka_manager.py @@ -25,6 +25,95 @@ from flocks.workflow.triggers.models import TriggerDefinition +@pytest.mark.asyncio +async def test_start_all_skips_stale_workflow_config_as_info( + monkeypatch: pytest.MonkeyPatch, +) -> None: + workflow_id = "removed-workflow" + config = { + "enabled": True, + "inputBroker": "", + "inputTopic": "", + } + info_events: list[tuple[str, dict]] = [] + warning_events: list[str] = [] + + async def _list_configs(*, kind: str): # noqa: ANN202 + assert kind == "workflow_kafka_config" + return [(workflow_id, config)] + + async def _get_config(_workflow_id: str, *, kind: str) -> dict: + assert kind == "workflow_kafka_config" + return config + + monkeypatch.setattr(kafka_manager.WorkflowStore, "list_configs", _list_configs) + monkeypatch.setattr(kafka_manager.WorkflowStore, "get_config", _get_config) + monkeypatch.setattr(kafka_manager, "read_workflow_from_fs", lambda _workflow_id: None) + monkeypatch.setattr( + kafka_manager.log, + "info", + lambda event, data=None: info_events.append((event, data)), + ) + monkeypatch.setattr( + kafka_manager.log, + "warning", + lambda event, _data=None: warning_events.append(event), + ) + + manager = kafka_manager.KafkaManager() + await manager.start_all() + + assert manager.get_consumer_status(workflow_id) == { + "state": "stopped", + "error": "workflow_not_found", + } + assert info_events == [ + ( + "kafka.workflow_not_found_on_start", + {"workflow_id": workflow_id, "action": "stale_config_skipped"}, + ) + ] + assert warning_events == [] + + +@pytest.mark.asyncio +async def test_start_all_continues_after_one_workflow_fails( + monkeypatch: pytest.MonkeyPatch, +) -> None: + restarted: list[str] = [] + warning_events: list[tuple[str, dict]] = [] + + async def _list_configs(*, kind: str): # noqa: ANN202 + assert kind == "workflow_kafka_config" + return [ + ("broken-workflow", {"enabled": True}), + ("healthy-workflow", {"enabled": True}), + ] + + async def _restart(workflow_id: str, *, startup: bool = False) -> dict: + assert startup is True + restarted.append(workflow_id) + if workflow_id == "broken-workflow": + raise ValueError("invalid config") + return {"state": "running", "error": None} + + monkeypatch.setattr(kafka_manager.WorkflowStore, "list_configs", _list_configs) + monkeypatch.setattr(kafka_manager.log, "warning", lambda event, data: warning_events.append((event, data))) + + manager = kafka_manager.KafkaManager() + monkeypatch.setattr(manager, "restart_workflow", _restart) + + await manager.start_all() + + assert restarted == ["broken-workflow", "healthy-workflow"] + assert warning_events == [ + ( + "kafka.start_failed", + {"workflow_id": "broken-workflow", "error": "invalid config"}, + ) + ] + + @pytest.mark.asyncio async def test_worker_pool_bounds_in_flight_dispatches(monkeypatch: pytest.MonkeyPatch) -> None: """The fixed worker pool must cap concurrent ``_trigger_workflow`` calls.""" @@ -190,6 +279,11 @@ async def _fake_get_config(workflow_id: str, *, kind: str) -> dict: return {"enabled": True, "inputBroker": "", "inputTopic": ""} monkeypatch.setattr(kafka_manager.WorkflowStore, "get_config", _fake_get_config) + monkeypatch.setattr( + kafka_manager, + "read_workflow_from_fs", + lambda _workflow_id: {"workflowJson": {}}, + ) status = await manager.restart_workflow("wf-no-broker") assert status["state"] == "failed" diff --git a/tests/ingest/test_syslog_manager_bind_failure.py b/tests/ingest/test_syslog_manager_bind_failure.py index ecd621fcb..75809307d 100644 --- a/tests/ingest/test_syslog_manager_bind_failure.py +++ b/tests/ingest/test_syslog_manager_bind_failure.py @@ -14,6 +14,7 @@ import asyncio import socket +from unittest.mock import AsyncMock, MagicMock import pytest @@ -28,6 +29,84 @@ def _find_busy_udp_port() -> tuple[socket.socket, int]: return sock, port +@pytest.mark.asyncio +async def test_start_all_logs_stale_workflow_config_as_info( + monkeypatch: pytest.MonkeyPatch, +) -> None: + workflow_id = "removed-workflow" + config = {"workflowId": workflow_id, "enabled": True} + info_log = MagicMock() + warning_log = MagicMock() + + monkeypatch.setattr( + syslog_manager.WorkflowStore, + "list_configs", + AsyncMock(return_value=[(workflow_id, config)]), + ) + monkeypatch.setattr( + syslog_manager.WorkflowStore, + "get_config", + AsyncMock(return_value=config), + ) + monkeypatch.setattr(syslog_manager, "read_workflow_from_fs", lambda _workflow_id: None) + monkeypatch.setattr(syslog_manager.log, "info", info_log) + monkeypatch.setattr(syslog_manager.log, "warning", warning_log) + + manager = syslog_manager.SyslogManager() + await manager.start_all() + + assert manager.get_listener_status(workflow_id) == { + "state": "stopped", + "error": "workflow_not_found", + } + info_log.assert_called_once_with( + "syslog.workflow_not_found_on_start", + { + "workflow_id": workflow_id, + "action": "stale_config_skipped", + }, + ) + warning_log.assert_not_called() + + +@pytest.mark.asyncio +async def test_start_all_continues_after_one_workflow_fails( + monkeypatch: pytest.MonkeyPatch, +) -> None: + restarted: list[str] = [] + warning_log = MagicMock() + + monkeypatch.setattr( + syslog_manager.WorkflowStore, + "list_configs", + AsyncMock( + return_value=[ + ("broken-workflow", {"enabled": True}), + ("healthy-workflow", {"enabled": True}), + ] + ), + ) + + async def _restart(workflow_id: str, *, startup: bool = False) -> dict: + assert startup is True + restarted.append(workflow_id) + if workflow_id == "broken-workflow": + raise ValueError("invalid config") + return {"state": "listening", "error": None} + + manager = syslog_manager.SyslogManager() + monkeypatch.setattr(manager, "restart_workflow", _restart) + monkeypatch.setattr(syslog_manager.log, "warning", warning_log) + + await manager.start_all() + + assert restarted == ["broken-workflow", "healthy-workflow"] + warning_log.assert_called_once_with( + "syslog.start_failed", + {"workflow_id": "broken-workflow", "error": "invalid config"}, + ) + + @pytest.mark.asyncio async def test_restart_workflow_reports_failure_on_port_conflict( monkeypatch: pytest.MonkeyPatch, @@ -46,11 +125,6 @@ async def test_restart_workflow_reports_failure_on_port_conflict( "inputKey": "syslog_message", } - async def _fake_storage_read(key: str): # noqa: ANN001 - if key == syslog_manager.SyslogManager._config_key(workflow_id): - return config - return None - def _fake_read_workflow_from_fs(wid: str): # noqa: ANN001 return { "id": wid, @@ -62,7 +136,11 @@ def _fake_read_workflow_from_fs(wid: str): # noqa: ANN001 } # Patch the *module-level* names ``manager.py`` looks up at call time. - monkeypatch.setattr(syslog_manager.Storage, "read", _fake_storage_read) + monkeypatch.setattr( + syslog_manager.WorkflowStore, + "get_config", + AsyncMock(return_value=config), + ) monkeypatch.setattr(syslog_manager, "read_workflow_from_fs", _fake_read_workflow_from_fs) manager = syslog_manager.SyslogManager() @@ -95,10 +173,11 @@ async def test_restart_workflow_returns_stopped_when_disabled( "inputKey": "syslog_message", } - async def _fake_storage_read(key: str): # noqa: ANN001 - return config - - monkeypatch.setattr(syslog_manager.Storage, "read", _fake_storage_read) + monkeypatch.setattr( + syslog_manager.WorkflowStore, + "get_config", + AsyncMock(return_value=config), + ) manager = syslog_manager.SyslogManager() status = await manager.restart_workflow(workflow_id) diff --git a/tests/mcp/test_mcp_client_sse.py b/tests/mcp/test_mcp_client_sse.py index 479a32a44..ed4e9487d 100644 --- a/tests/mcp/test_mcp_client_sse.py +++ b/tests/mcp/test_mcp_client_sse.py @@ -5,9 +5,12 @@ from types import MethodType, SimpleNamespace import pytest +from mcp.shared.exceptions import McpError +from mcp.types import ErrorData, METHOD_NOT_FOUND import flocks.mcp.client as mcp_client_module from flocks.mcp.client import McpClient, _extract_root_cause +from flocks.mcp.errors import is_method_not_found_error def _make_session_class( @@ -374,3 +377,21 @@ class MockRequest: result = _extract_root_cause(exc) assert "401" in result assert "secret" not in result + + def test_method_not_found_detection_uses_structured_error_code(self): + exc = McpError(ErrorData(code=METHOD_NOT_FOUND, message="Unsupported operation")) + assert is_method_not_found_error(exc) is True + + def test_method_not_found_detection_checks_every_exception_group_child(self): + method_error = McpError( + ErrorData(code=METHOD_NOT_FOUND, message="Unsupported operation") + ) + exc = ExceptionGroup( + "request failed", + [RuntimeError("transport closed"), method_error], + ) + assert is_method_not_found_error(exc) is True + + def test_method_not_found_detection_does_not_match_unrelated_phrase(self): + exc = RuntimeError("cache method not found during local lookup") + assert is_method_not_found_error(exc) is False diff --git a/tests/mcp/test_mcp_server.py b/tests/mcp/test_mcp_server.py index b89a6b26d..7e3f6c098 100644 --- a/tests/mcp/test_mcp_server.py +++ b/tests/mcp/test_mcp_server.py @@ -4,7 +4,10 @@ from types import SimpleNamespace import pytest +from mcp.shared.exceptions import McpError +from mcp.types import ErrorData, METHOD_NOT_FOUND +from flocks.mcp import server as mcp_server_module from flocks.mcp.server import McpServerManager from flocks.mcp.types import McpStatus @@ -81,6 +84,70 @@ def __init__(self, **kwargs) -> None: assert manager._status["legacy-demo"].status == McpStatus.CONNECTED +@pytest.mark.asyncio +async def test_connect_treats_missing_resources_method_as_optional(monkeypatch: pytest.MonkeyPatch): + class ToolsOnlyClient(_FakeMcpClient): + async def list_resources(self) -> list: + raise McpError( + ErrorData(code=METHOD_NOT_FOUND, message="Unsupported operation") + ) + + info_events: list[str] = [] + warn_events: list[str] = [] + monkeypatch.setattr(mcp_server_module, "McpClient", ToolsOnlyClient) + monkeypatch.setattr( + mcp_server_module.log, + "info", + lambda event, _data=None: info_events.append(event), + ) + monkeypatch.setattr( + mcp_server_module.log, + "warn", + lambda event, _data=None: warn_events.append(event), + ) + + manager = McpServerManager() + await manager._connect_and_register( + "tools-only", + {"type": "local", "command": ["python", "-m", "demo"]}, + ) + + assert manager._resources_cache["tools-only"] == [] + assert manager._status["tools-only"].status == McpStatus.CONNECTED + assert "mcp.resources_unsupported" in info_events + assert "mcp.resources_unavailable" not in warn_events + + +@pytest.mark.asyncio +async def test_connect_logs_resource_discovery_failure_once(monkeypatch: pytest.MonkeyPatch): + class BrokenResourcesClient(_FakeMcpClient): + async def list_resources(self) -> list: + raise RuntimeError("resource endpoint unavailable") + + warn_events: list[str] = [] + error_events: list[str] = [] + monkeypatch.setattr(mcp_server_module, "McpClient", BrokenResourcesClient) + monkeypatch.setattr( + mcp_server_module.log, + "warn", + lambda event, _data=None: warn_events.append(event), + ) + monkeypatch.setattr( + mcp_server_module.log, + "error", + lambda event, _data=None: error_events.append(event), + ) + + manager = McpServerManager() + await manager._connect_and_register( + "broken-resources", + {"type": "local", "command": ["python", "-m", "demo"]}, + ) + + assert warn_events == ["mcp.resources_unavailable"] + assert error_events == [] + + @pytest.mark.asyncio async def test_init_is_serialized_for_concurrent_callers(monkeypatch: pytest.MonkeyPatch): manager = McpServerManager() diff --git a/tests/provider/test_api_service_management.py b/tests/provider/test_api_service_management.py index 4647fb8df..fbb1d9ebf 100644 --- a/tests/provider/test_api_service_management.py +++ b/tests/provider/test_api_service_management.py @@ -585,6 +585,58 @@ def test_bootstrap_adds_enabled_true_when_key_missing(self): assert written["apiKey"] == "{secret:key}" assert written["base_url"] == "https://api.example.com" + def test_bootstrap_writes_unique_versioned_storage_key(self): + tools = {"tdp_tool": _make_api_tool("tdp_tool", "tdp_api")} + with patch( + "flocks.config.api_versioning.versioned_storage_key_for", + return_value="tdp_api_v3_3_10", + ): + mock_set = self._run_bootstrap(tools, get_raw_side_effect=lambda _: None) + + mock_set.assert_called_once_with("tdp_api_v3_3_10", {"enabled": True}) + + def test_bootstrap_then_sync_reads_the_same_versioned_storage_key(self): + tool = _make_api_tool("tdp_tool", "tdp_api") + tools = {"tdp_tool": tool} + services = { + "tdp_api": {"apiKey": "legacy-copy"}, + "tdp_api_v3_3_10": {"apiKey": "versioned-copy"}, + } + + def get_raw(storage_key: str): + value = services.get(storage_key) + return value.copy() if value is not None else None + + def set_service(storage_key: str, value: dict): + services[storage_key] = value.copy() + + with ( + patch.object(ToolRegistry, "_tools", tools), + patch.object(ToolRegistry, "_enabled_defaults", {"tdp_tool": True}), + patch( + "flocks.config.api_versioning.versioned_storage_key_for", + return_value="tdp_api_v3_3_10", + ), + patch( + "flocks.config.config_writer.ConfigWriter.get_api_service_raw", + side_effect=get_raw, + ), + patch( + "flocks.config.config_writer.ConfigWriter.set_api_service", + side_effect=set_service, + ), + patch( + "flocks.config.config_writer.ConfigWriter.list_api_services_raw", + side_effect=lambda: services, + ), + ): + ToolRegistry._bootstrap_user_api_services() + ToolRegistry._sync_api_service_states() + + assert services["tdp_api_v3_3_10"]["enabled"] is True + assert "enabled" not in services["tdp_api"] + assert tool.info.enabled is True + def test_bootstrap_does_not_overwrite_explicit_disabled(self): """When api_services..enabled is False, bootstrap leaves it untouched.""" existing = {"enabled": False, "apiKey": "{secret:key}"} diff --git a/tests/server/routes/test_notifications_routes.py b/tests/server/routes/test_notifications_routes.py index 1245f8295..a0400e801 100644 --- a/tests/server/routes/test_notifications_routes.py +++ b/tests/server/routes/test_notifications_routes.py @@ -25,7 +25,7 @@ async def test_notifications_require_browser_login(client: AsyncClient): async def test_active_notifications_and_dismiss_forever(client: AsyncClient): response = await client.get( "/api/notifications/active", - params={"locale": "zh-CN", "current_version": "2026.04.27"}, + params={"locale": "zh-CN"}, ) assert response.status_code == 200, response.text items = response.json() @@ -37,14 +37,14 @@ async def test_active_notifications_and_dismiss_forever(client: AsyncClient): response = await client.get( "/api/notifications/active", - params={"locale": "zh-CN", "current_version": "2026.04.27"}, + params={"locale": "zh-CN"}, ) assert response.status_code == 200, response.text assert response.json() == [] response = await client.get( "/api/notifications/active", - params={"locale": "zh-CN", "current_version": "2026.05.01"}, + params={"locale": "zh-CN"}, ) assert response.status_code == 200, response.text assert response.json() == [] diff --git a/tests/server/test_app_errors.py b/tests/server/test_app_errors.py index 673cd89a1..71e97b879 100644 --- a/tests/server/test_app_errors.py +++ b/tests/server/test_app_errors.py @@ -1,6 +1,8 @@ import json +from unittest.mock import MagicMock import pytest +from starlette.exceptions import HTTPException from starlette.requests import Request from flocks.server import app as app_module @@ -31,3 +33,43 @@ async def test_general_exception_response_does_not_expose_traceback(): "error": "InternalServerError", "message": "Internal server error", } + + +@pytest.mark.asyncio +async def test_http_4xx_logs_warning(monkeypatch): + warning = MagicMock() + error = MagicMock() + monkeypatch.setattr(app_module.log, "warn", warning) + monkeypatch.setattr(app_module.log, "error", error) + + response = await app_module.http_exception_handler( + _request("/missing"), + HTTPException(status_code=404, detail="Not found"), + ) + + assert response.status_code == 404 + warning.assert_called_once_with( + "http.error", + {"path": "/missing", "status": 404, "detail": "Not found"}, + ) + error.assert_not_called() + + +@pytest.mark.asyncio +async def test_http_5xx_logs_error(monkeypatch): + warning = MagicMock() + error = MagicMock() + monkeypatch.setattr(app_module.log, "warn", warning) + monkeypatch.setattr(app_module.log, "error", error) + + response = await app_module.http_exception_handler( + _request("/unavailable"), + HTTPException(status_code=503, detail="Unavailable"), + ) + + assert response.status_code == 503 + error.assert_called_once_with( + "http.error", + {"path": "/unavailable", "status": 503, "detail": "Unavailable"}, + ) + warning.assert_not_called() diff --git a/tests/server/test_input_dispatcher.py b/tests/server/test_input_dispatcher.py index ff1c8bae6..17a11d8f0 100644 --- a/tests/server/test_input_dispatcher.py +++ b/tests/server/test_input_dispatcher.py @@ -357,13 +357,13 @@ def test_materialize_queued_data_url_returns_readable_file_uri(self, monkeypatch from flocks.server.routes import session as session_routes from flocks.session.utils.file_extractor import read_file_part_bytes - class FakeWorkspace: - def resolve_workspace_path(self, rel_path: str): - return tmp_path / rel_path - monkeypatch.setattr( "flocks.workspace.manager.WorkspaceManager.get_instance", - lambda: FakeWorkspace(), + lambda: (_ for _ in ()).throw(PermissionError("workspace access denied")), + ) + monkeypatch.setattr( + "flocks.config.config.Config.get_data_path", + staticmethod(lambda: tmp_path), ) data_url = "data:image/png;base64," + base64.b64encode(b"png-bytes").decode() @@ -377,6 +377,17 @@ def resolve_workspace_path(self, rel_path: str): assert url.startswith("file://") assert read_file_part_bytes(url) == b"png-bytes" + def test_session_upload_path_rejects_traversal(self, monkeypatch, tmp_path): + from flocks.server.routes import session as session_routes + + monkeypatch.setattr( + "flocks.config.config.Config.get_data_path", + staticmethod(lambda: tmp_path), + ) + + with pytest.raises(ValueError, match="Invalid session ID"): + session_routes._session_uploads_dir("../outside") + @pytest.mark.asyncio async def test_prompt_async_queues_when_session_running_without_creating_message(self, monkeypatch): from flocks.server.routes import session as session_routes diff --git a/tests/session/test_runner_step.py b/tests/session/test_runner_step.py index 9d67e6565..994ce4b42 100644 --- a/tests/session/test_runner_step.py +++ b/tests/session/test_runner_step.py @@ -2249,6 +2249,8 @@ async def test_process_step_limits_connection_error_retries(monkeypatch): provider.is_configured.return_value = True assistant_msg = SimpleNamespace(id="msg_assistant_connection_error") update_mock = AsyncMock(return_value=None) + warn_log = MagicMock() + error_log = MagicMock() call_count = 0 async def fake_call_llm(*_args, **_kwargs): @@ -2271,6 +2273,8 @@ async def fake_call_llm(*_args, **_kwargs): monkeypatch.setattr(runner_mod.Message, "update", update_mock) monkeypatch.setattr(runner_mod.SessionRetry, "sleep", AsyncMock(return_value=None)) monkeypatch.setattr(runner, "_call_llm", fake_call_llm) + monkeypatch.setattr(runner_mod.log, "warn", warn_log) + monkeypatch.setattr(runner_mod.log, "error", error_log) result = await runner._process_step([last_user], last_user) @@ -2284,6 +2288,21 @@ async def fake_call_llm(*_args, **_kwargs): assert final_update["error"]["data"]["message"] == "Connection error." assert final_update["error"]["data"]["displayMessage"] == runner_mod.CONNECTION_ERROR_DISPLAY_MESSAGE + retry_logs = [ + call for call in warn_log.call_args_list + if call.args and call.args[0] == "runner.step.retry" + ] + max_retry_logs = [ + call for call in error_log.call_args_list + if call.args and call.args[0] == "runner.step.max_retries_exceeded" + ] + assert len(retry_logs) == 3 + assert len(max_retry_logs) == 1 + assert not any( + call.args and call.args[0] == "runner.step.error" + for call in [*warn_log.call_args_list, *error_log.call_args_list] + ) + @pytest.mark.asyncio async def test_process_step_marks_aborted_llm_message_as_error(monkeypatch): diff --git a/tests/session/test_session_memory_log_levels.py b/tests/session/test_session_memory_log_levels.py new file mode 100644 index 000000000..5a0571cf1 --- /dev/null +++ b/tests/session/test_session_memory_log_levels.py @@ -0,0 +1,40 @@ +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock + +import pytest + +import flocks.session.features.memory as memory_module + + +@pytest.mark.asyncio +async def test_missing_memory_config_logs_info(monkeypatch, tmp_path): + manager = SimpleNamespace(initialize=AsyncMock(return_value=None)) + info_log = MagicMock() + warn_log = MagicMock() + + monkeypatch.setattr( + memory_module.Config, + "get", + AsyncMock(return_value=SimpleNamespace(memory=None)), + ) + monkeypatch.setattr( + memory_module.MemoryManager, + "get_instance", + lambda **_kwargs: manager, + ) + monkeypatch.setattr(memory_module.log, "info", info_log) + monkeypatch.setattr(memory_module.log, "warn", warn_log) + + memory = memory_module.SessionMemory( + session_id="session-no-memory-config", + project_id="project", + workspace_dir=str(tmp_path), + enabled=True, + ) + + assert await memory.initialize() is True + info_log.assert_any_call( + "session.memory.no_config", + {"session_id": "session-no-memory-config"}, + ) + assert not any(call.args and call.args[0] == "session.memory.no_config" for call in warn_log.call_args_list) diff --git a/tests/skill/test_skill.py b/tests/skill/test_skill.py index e704e785e..c993c511f 100644 --- a/tests/skill/test_skill.py +++ b/tests/skill/test_skill.py @@ -10,6 +10,7 @@ from pathlib import Path from unittest.mock import patch +import flocks.skill.skill as skill_module from flocks.skill.skill import Skill, SkillInfo, SkillRequires, SkillInstallSpec @@ -392,8 +393,8 @@ async def test_discover_flocks_source_for_builtin_skills(tmp_path): @pytest.mark.asyncio -async def test_discover_project_overrides_global(tmp_path): - """Project-level plugin skill must override global skill of the same name.""" +async def test_discover_user_plugin_overrides_project_bundle(tmp_path): + """User plugin skills must override project-bundled skills of the same name.""" fake_home = tmp_path / "home" project_dir = tmp_path / "myproject" project_dir.mkdir() @@ -409,14 +410,24 @@ async def test_discover_project_overrides_global(tmp_path): patch("os.path.expanduser", return_value=str(fake_home)), patch("flocks.skill.skill.Instance.get_directory", return_value=str(project_dir)), patch("flocks.skill.skill.Instance.get_worktree", return_value=str(project_dir)), + patch.object(skill_module.log, "info") as info_log, + patch.object(skill_module.log, "warn") as warn_log, ): Skill.clear_cache() skills = await Skill.all() skill = next((s for s in skills if s.name == "shared-skill"), None) assert skill is not None - assert skill.source == "project", "project-level skill should override global" - assert "Project version" in skill.description + assert skill.source == "user", "user-level skill should override project bundle" + assert "Global version" in skill.description + assert any( + call.args and call.args[0] == "skill.override" + for call in info_log.call_args_list + ) + assert not any( + call.args and call.args[0] == "skill.duplicate" + for call in warn_log.call_args_list + ) @pytest.mark.asyncio diff --git a/tests/tool/test_device_startup_sync.py b/tests/tool/test_device_startup_sync.py index 37986aa48..0481e22e5 100644 --- a/tests/tool/test_device_startup_sync.py +++ b/tests/tool/test_device_startup_sync.py @@ -15,6 +15,7 @@ """ from __future__ import annotations +import sqlite3 from pathlib import Path from typing import List @@ -24,6 +25,33 @@ from flocks.tool.device import startup +@pytest.mark.asyncio +async def test_heal_stale_service_ids_reads_named_sqlite_rows(tmp_path, monkeypatch): + db_path = tmp_path / "devices.db" + with sqlite3.connect(db_path) as db: + db.execute( + "CREATE TABLE device_integrations (id TEXT, storage_key TEXT, service_id TEXT)" + ) + db.execute( + "INSERT INTO device_integrations VALUES (?, ?, ?)", + ("dev-1", "onesec_api_v2_8_2", "stale_service_id"), + ) + + monkeypatch.setattr(startup.Storage, "get_db_path", staticmethod(lambda: db_path)) + monkeypatch.setattr( + "flocks.tool.device.store.storage_key_to_service_id", + lambda storage_key: "onesec_api" if storage_key == "onesec_api_v2_8_2" else storage_key, + ) + + await startup._heal_stale_service_ids() + + with sqlite3.connect(db_path) as db: + row = db.execute( + "SELECT service_id FROM device_integrations WHERE id = ?", ("dev-1",) + ).fetchone() + assert row == ("onesec_api",) + + @pytest.fixture def isolated_plugins(monkeypatch, tmp_path): """Point plugin discovery at an empty tmp HOME so production plugins diff --git a/tests/tool/test_tools.py b/tests/tool/test_tools.py index e63e97cf7..895cf0f80 100644 --- a/tests/tool/test_tools.py +++ b/tests/tool/test_tools.py @@ -20,6 +20,7 @@ import uuid from pathlib import Path from typing import Dict, Any, List +from unittest.mock import MagicMock # Import the tool system from flocks.tool import ( @@ -35,6 +36,8 @@ ParameterType, ) from flocks.tool.code import bash as bash_module +import flocks.tool.file.glob as glob_module +import flocks.tool.registry as registry_module import flocks.tool.system.question as question_module @@ -579,6 +582,33 @@ async def test_glob_no_matches(self, tool_context, temp_dir): assert result.success assert "No files found" in result.output + @pytest.mark.asyncio + async def test_python_fallback_logs_info_once(self, monkeypatch, tool_context, temp_dir): + monkeypatch.setattr(glob_module, "_ripgrep_fallback_logged", False) + monkeypatch.setattr(glob_module, "find_ripgrep", lambda: None) + info_log = MagicMock() + warn_log = MagicMock() + monkeypatch.setattr(glob_module.log, "info", info_log) + monkeypatch.setattr(glob_module.log, "warn", warn_log) + + for _ in range(2): + result = await glob_module.glob_tool( + tool_context, + pattern="*.txt", + path=temp_dir, + ) + assert result.success + + fallback_calls = [ + call for call in info_log.call_args_list + if call.args and call.args[0] == "glob.ripgrep_not_found" + ] + assert len(fallback_calls) == 1 + assert not any( + call.args and call.args[0] == "glob.ripgrep_not_found" + for call in warn_log.call_args_list + ) + # ============================================================================= # P1 Tools Tests @@ -1030,8 +1060,13 @@ class TestErrorHandling: """Test error handling across tools""" @pytest.mark.asyncio - async def test_missing_required_parameter(self, tool_context): + async def test_missing_required_parameter(self, monkeypatch, tool_context): """Test error when required parameter is missing""" + warn_log = MagicMock() + error_log = MagicMock() + monkeypatch.setattr(registry_module.log, "warn", warn_log) + monkeypatch.setattr(registry_module.log, "error", error_log) + result = await ToolRegistry.execute( "read", ctx=tool_context @@ -1040,6 +1075,18 @@ async def test_missing_required_parameter(self, tool_context): assert not result.success assert "required" in result.error.lower() or "missing" in result.error.lower() + warn_log.assert_called_once_with( + "tool.execute.missing_param", + { + "tool": "read", + "missing": "filePath", + "provided": [], + }, + ) + assert not any( + call.args and call.args[0] == "tool.execute.missing_param" + for call in error_log.call_args_list + ) @pytest.mark.asyncio async def test_nonexistent_tool(self, tool_context): diff --git a/tests/workflow/test_fs_store.py b/tests/workflow/test_fs_store.py index 711f072ee..f5c7b4c31 100644 --- a/tests/workflow/test_fs_store.py +++ b/tests/workflow/test_fs_store.py @@ -53,6 +53,45 @@ def test_read_workflow_from_fs_refreshes_cached_workspace_root( assert fs_store.find_workspace_root() == second_workspace +def test_read_workflow_from_fs_prefers_user_workflow_over_project_bundle( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +): + workflow_id = "shared-workflow" + project_root = tmp_path / "project-workflows" + user_root = tmp_path / "user-workflows" + + for root, name in ( + (project_root, "project bundle"), + (user_root, "user customization"), + ): + workflow_dir = root / workflow_id + workflow_dir.mkdir(parents=True) + (workflow_dir / "workflow.json").write_text( + json.dumps( + { + "name": name, + "start": "n1", + "nodes": [{"id": "n1", "type": "python", "code": "pass"}], + "edges": [], + } + ), + encoding="utf-8", + ) + + monkeypatch.setattr( + fs_store, + "resolve_workflow_scan_roots", + lambda _workspace: [(project_root, "project"), (user_root, "global")], + ) + + workflow = fs_store.read_workflow_from_fs(workflow_id) + + assert workflow is not None + assert workflow["source"] == "global" + assert workflow["workflowJson"]["name"] == "user customization" + + def test_read_workflow_dir_uses_latest_file_mtime_when_meta_is_stale( tmp_path: Path, ): diff --git a/tests/workflow/test_poller_manager.py b/tests/workflow/test_poller_manager.py index 4b7f915df..b105cd121 100644 --- a/tests/workflow/test_poller_manager.py +++ b/tests/workflow/test_poller_manager.py @@ -12,6 +12,52 @@ from flocks.workflow.runner import RunWorkflowResult +@pytest.mark.asyncio +async def test_start_all_skips_stale_workflow_config_as_info( + monkeypatch: pytest.MonkeyPatch, +) -> None: + workflow_id = "removed-workflow" + config = {"enabled": True, "intervalSeconds": "not-a-number"} + info_events: list[tuple[str, dict]] = [] + warning_events: list[str] = [] + + async def _list_configs(*, kind: str): # noqa: ANN202 + assert kind == "workflow_poller_config" + return [(workflow_id, config)] + + async def _get_config(_workflow_id: str, *, kind: str) -> dict[str, Any]: + assert kind == "workflow_poller_config" + return config + + monkeypatch.setattr(poller_manager.WorkflowStore, "list_configs", _list_configs) + monkeypatch.setattr(poller_manager.WorkflowStore, "get_config", _get_config) + monkeypatch.setattr(poller_manager, "read_workflow_from_fs", lambda _workflow_id: None) + monkeypatch.setattr( + poller_manager.log, + "info", + lambda event, data=None: info_events.append((event, data)), + ) + monkeypatch.setattr( + poller_manager.log, + "warning", + lambda event, _data=None: warning_events.append(event), + ) + + manager = poller_manager.WorkflowPollerManager() + await manager.start_all() + + status = manager.get_status(workflow_id) + assert status["state"] == "stopped" + assert status["error"] == "workflow_not_found" + assert info_events == [ + ( + "poller.workflow_not_found_on_start", + {"workflow_id": workflow_id, "action": "stale_config_skipped"}, + ) + ] + assert warning_events == [] + + @pytest.mark.asyncio async def test_restart_disabled_config_reports_stopped(monkeypatch: pytest.MonkeyPatch) -> None: manager = poller_manager.WorkflowPollerManager() @@ -411,22 +457,42 @@ def _fake_run_workflow( # noqa: ANN001 async def test_start_all_only_restarts_enabled_configs(monkeypatch: pytest.MonkeyPatch) -> None: manager = poller_manager.WorkflowPollerManager() restarted: list[str] = [] + warning_events: list[tuple[str, dict[str, str]]] = [] async def _fake_list_configs(*, kind: str) -> list[tuple[str, dict[str, Any]]]: return [ + ("wf-broken", {"enabled": True}), ("wf-enabled", {"enabled": True}), ("wf-disabled", {"enabled": False}), ] - async def _fake_restart(workflow_id: str) -> dict[str, Any]: + async def _fake_restart( + workflow_id: str, + *, + startup: bool = False, + ) -> dict[str, Any]: + assert startup is True restarted.append(workflow_id) + if workflow_id == "wf-broken": + raise ValueError("invalid config") return {"workflowId": workflow_id, "state": "running"} monkeypatch.setattr(poller_manager.WorkflowStore, "list_configs", _fake_list_configs) monkeypatch.setattr(manager, "restart_workflow", _fake_restart) + monkeypatch.setattr( + poller_manager.log, + "warning", + lambda event, data: warning_events.append((event, data)), + ) await manager.start_all() - assert restarted == ["wf-enabled"] + assert restarted == ["wf-broken", "wf-enabled"] + assert warning_events == [ + ( + "poller.start_failed", + {"workflow_id": "wf-broken", "error": "invalid config"}, + ) + ] @pytest.mark.asyncio diff --git a/tests/workflow/test_workflow_paths.py b/tests/workflow/test_workflow_paths.py index ff31e0237..2194b5bfb 100644 --- a/tests/workflow/test_workflow_paths.py +++ b/tests/workflow/test_workflow_paths.py @@ -162,5 +162,99 @@ async def test_new_canonical_path_wins_over_legacy( results = await center.scan_skill_workflows(tmp_path) # The new canonical path (plugins/workflows/) has higher priority and wins - names = [r["name"] for r in results] - assert "shared-wf-new" in names + assert [r["name"] for r in results] == ["shared-wf-new"] + + +@pytest.mark.asyncio +async def test_scan_and_registry_prefer_user_workflow_over_project_bundle( + tmp_path: Path, + isolated_storage, + monkeypatch: pytest.MonkeyPatch, +) -> None: + project_root = tmp_path / "project" + user_root = tmp_path / "user" + for root, name in ( + (project_root, "project bundle"), + (user_root, "user customization"), + ): + workflow_dir = root / "shared-workflow" + workflow_dir.mkdir(parents=True) + (workflow_dir / "workflow.json").write_text( + json.dumps(_workflow_payload(name)), + encoding="utf-8", + ) + + monkeypatch.setattr( + center, + "resolve_project_workflow_roots", + lambda _base: [project_root], + ) + monkeypatch.setattr(center, "resolve_global_workflow_roots", lambda: [user_root]) + + scanned = await center.scan_skill_workflows(tmp_path) + registered = await center.list_registry_entries() + + assert len(scanned) == 1 + assert scanned[0]["name"] == "user customization" + assert scanned[0]["sourceType"] == "global" + matching_registry_entries = [ + entry + for entry in registered + if entry.get("logicalWorkflowId") == "shared-workflow" + and Path(str(entry.get("workflowPath"))).is_relative_to(tmp_path) + ] + assert len(matching_registry_entries) == 1 + assert matching_registry_entries[0]["name"] == "user customization" + + +@pytest.mark.asyncio +@pytest.mark.parametrize("invalid_user_workflow", ["invalid_json", "hidden"]) +async def test_registry_falls_back_when_user_workflow_is_no_longer_discoverable( + tmp_path: Path, + isolated_storage, + monkeypatch: pytest.MonkeyPatch, + invalid_user_workflow: str, +) -> None: + project_root = tmp_path / "project" + user_root = tmp_path / "user" + project_dir = project_root / "shared-workflow" + user_dir = user_root / "shared-workflow" + project_dir.mkdir(parents=True) + user_dir.mkdir(parents=True) + (project_dir / "workflow.json").write_text( + json.dumps(_workflow_payload("project bundle")), + encoding="utf-8", + ) + user_workflow_path = user_dir / "workflow.json" + user_workflow_path.write_text( + json.dumps(_workflow_payload("user customization")), + encoding="utf-8", + ) + + monkeypatch.setattr( + center, + "resolve_project_workflow_roots", + lambda _base: [project_root], + ) + monkeypatch.setattr(center, "resolve_global_workflow_roots", lambda: [user_root]) + + await center.scan_skill_workflows(tmp_path) + if invalid_user_workflow == "invalid_json": + user_workflow_path.write_text("{", encoding="utf-8") + else: + (user_dir / "meta.json").write_text( + json.dumps({"hidden": True}), + encoding="utf-8", + ) + + scanned = await center.scan_skill_workflows(tmp_path) + registered = await center.list_registry_entries() + matching_registry_entries = [ + entry + for entry in registered + if entry.get("logicalWorkflowId") == "shared-workflow" + and Path(str(entry.get("workflowPath"))).is_relative_to(tmp_path) + ] + + assert [entry["name"] for entry in scanned] == ["project bundle"] + assert [entry["name"] for entry in matching_registry_entries] == ["project bundle"] diff --git a/webui/src/api/notifications.ts b/webui/src/api/notifications.ts index 6c8927989..98ac94aa0 100644 --- a/webui/src/api/notifications.ts +++ b/webui/src/api/notifications.ts @@ -32,14 +32,10 @@ export interface NotificationAckStatus { acknowledged: boolean; } -export const getActiveNotifications = async ( - locale?: string, - currentVersion?: string | null, -): Promise => { +export const getActiveNotifications = async (locale?: string): Promise => { const response = await client.get('/api/notifications/active', { params: { ...(locale ? { locale } : {}), - ...(currentVersion ? { current_version: currentVersion } : {}), }, }); return response.data; diff --git a/webui/src/components/layout/Layout.test.tsx b/webui/src/components/layout/Layout.test.tsx index d8c3ab000..f8911592d 100644 --- a/webui/src/components/layout/Layout.test.tsx +++ b/webui/src/components/layout/Layout.test.tsx @@ -482,7 +482,10 @@ describe('Layout onboarding entry', () => { expect(await screen.findByText('admin.roleMember')).toBeInTheDocument(); expect(await screen.findByText('v2026.6.21')).toBeInTheDocument(); expect(screen.queryByRole('link', { name: 'Flocks' })).not.toBeInTheDocument(); - await waitFor(() => expect(checkUpdate).toHaveBeenCalledWith('zh-CN', 'flockspro')); + expect(checkUpdate).not.toHaveBeenCalled(); + await waitFor(() => { + expect(getActiveNotifications).toHaveBeenCalledWith('zh-CN'); + }); const sidebarShell = container.querySelector('aside > div'); const logoRow = sidebarShell?.firstElementChild as HTMLElement | null; @@ -845,6 +848,33 @@ describe('Layout onboarding entry', () => { expect(await screen.findByText('Token 免费期已延长')).toBeInTheDocument(); expect(screen.getByText('Flocks v2026.04.28 更新内容')).toBeInTheDocument(); }); + + it('loads backend notifications without waiting for the update check', async () => { + localStorage.setItem('flocks_onboarding_dismissed', 'true'); + const updateCheck = deferred<{ + has_update: boolean; + latest_version: null; + current_version: string; + error: null; + }>(); + checkUpdate.mockReturnValue(updateCheck.promise); + + renderHomeWithLayout(); + + await waitFor(() => { + expect(getActiveNotifications).toHaveBeenCalledWith('zh-CN'); + }); + expect(getNotificationAckStatus).not.toHaveBeenCalled(); + + await act(async () => { + updateCheck.resolve({ + has_update: false, + latest_version: null, + current_version: '0.2.0', + error: null, + }); + }); + }); }); describe('Layout WebUI contract pages navigation', () => { diff --git a/webui/src/components/layout/Layout.tsx b/webui/src/components/layout/Layout.tsx index f4758e1ca..df4160bfe 100644 --- a/webui/src/components/layout/Layout.tsx +++ b/webui/src/components/layout/Layout.tsx @@ -217,6 +217,8 @@ export default function Layout() { const [flocksproStatusReady, setFlocksproStatusReady] = useState(false); const [flocksproVersion, setFlocksproVersion] = useState(null); const canManageUpdates = user?.role === 'admin'; + const notificationGateReady = flocksproStatusReady + && (!canManageUpdates || hasCompletedUpdateCheck); const canCreateWorkspaceCustomPage = user?.role === 'admin'; const { pages: webuiContractPages, workspaces: webuiContractWorkspaces = [] } = useWebUIContractPages(); const [openWorkspaceMenuId, setOpenWorkspaceMenuId] = useState(null); @@ -253,7 +255,7 @@ export default function Layout() { }, [handleOpenOnboarding]); const refreshUpdateStatus = useCallback(async (bypassMinGap = false) => { - if (!flocksproStatusReady) return; + if (!flocksproStatusReady || !canManageUpdates) return; const now = Date.now(); if (checkingUpdateRef.current) return; @@ -304,7 +306,7 @@ export default function Layout() { }, [canManageUpdates, flocksproStatusReady, i18n.language, isFlocksproActive]); useEffect(() => { - if (!flocksproStatusReady) return undefined; + if (!flocksproStatusReady || !canManageUpdates) return undefined; const initialCheckTimerId = window.setTimeout(() => { refreshUpdateStatus(true); @@ -335,7 +337,7 @@ export default function Layout() { document.removeEventListener('visibilitychange', handleVisibilityChange); window.removeEventListener('focus', handleWindowFocus); }; - }, [flocksproStatusReady, refreshUpdateStatus]); + }, [canManageUpdates, flocksproStatusReady, refreshUpdateStatus]); useEffect(() => { let cancelled = false; @@ -397,16 +399,14 @@ export default function Layout() { lastNotificationFetchKeyRef.current = null; return; } - if (!hasCompletedUpdateCheck) return; - - const fetchKey = `${user.id}:${i18n.language}:${currentVersion ?? 'pending-version'}`; + const fetchKey = `${user.id}:${i18n.language}`; if (lastNotificationFetchKeyRef.current === fetchKey) return; const previousFetchKey = lastNotificationFetchKeyRef.current; lastNotificationFetchKeyRef.current = fetchKey; setBackendNotificationsReady(false); let cancelled = false; - void getActiveNotifications(i18n.language, currentVersion) + void getActiveNotifications(i18n.language) .then((items) => { if (cancelled) return; setNotifications((prev) => { @@ -431,7 +431,7 @@ export default function Layout() { return () => { cancelled = true; }; - }, [currentVersion, hasCompletedUpdateCheck, i18n.language, user?.id]); + }, [i18n.language, user?.id]); useEffect(() => { if (!user?.id) { @@ -439,7 +439,7 @@ export default function Layout() { setUpdateNotificationReady(false); return; } - if (!hasCompletedUpdateCheck) return; + if (!notificationGateReady) return; setUpdateNotificationReady(false); const notification = buildUpdateNotification(updateInfo, i18n.language); @@ -465,7 +465,7 @@ export default function Layout() { return () => { cancelled = true; }; - }, [hasCompletedUpdateCheck, i18n.language, updateInfo, user?.id]); + }, [i18n.language, notificationGateReady, updateInfo, user?.id]); const allNotifications = updateNotification ? [...notifications, updateNotification].sort((a, b) => a.priority - b.priority)