From 7ccfe68f3f06ecb5b10ca711434af86e75a1a625 Mon Sep 17 00:00:00 2001 From: Xi Zhang <106144707+X-iZhang@users.noreply.github.com> Date: Thu, 25 Jun 2026 17:31:23 +0100 Subject: [PATCH] feat: add scheduler functionality with cron-style task management (#306) * feat: add scheduler functionality with cron-style task management - Implemented a new scheduler subagent to automate recurring tasks using cron expressions. - Enhanced the subagent factory to include the skill manager and auxiliary chat model for the scheduler. - Created a YAML configuration for the scheduler with a detailed system prompt and toolset. - Updated README files to include documentation on scheduled tasks and usage examples. - Added tests for the scheduler, including command execution, scheduling tools, and middleware integration. - Introduced new dependencies for timezone handling and ensured compatibility in the project configuration. * fix(async-notifier): ensure fallback hint is used for unknown notification kinds * feat: enhance scheduling functionality and improve system message handling --- EvoScientist/EvoScientist.py | 17 +- EvoScientist/cli/async_notifier.py | 12 +- .../commands/implementation/__init__.py | 13 +- .../commands/implementation/schedule.py | 234 +++++++++++++++++ EvoScientist/config/settings.py | 13 +- EvoScientist/cron/__init__.py | 23 ++ EvoScientist/cron/schedule.py | 118 +++++++++ EvoScientist/langgraph_dev/graphs.py | 5 +- EvoScientist/langgraph_dev/langgraph.json | 1 + EvoScientist/langgraph_dev/manager.py | 2 + EvoScientist/middleware/__init__.py | 6 + EvoScientist/middleware/memory.py | 19 +- EvoScientist/middleware/scheduler.py | 222 ++++++++++++++++ EvoScientist/middleware/utils.py | 18 ++ EvoScientist/subagents/_factory.py | 22 +- EvoScientist/subagents/scheduler.yaml | 24 ++ README.md | 30 ++- README.zh-CN.md | 30 ++- pyproject.toml | 1 + tests/test_async_notifier.py | 17 ++ tests/test_cli_serve.py | 1 + tests/test_config.py | 16 ++ tests/test_cron_schedule.py | 136 ++++++++++ tests/test_langgraph_manager.py | 13 +- tests/test_profile_memory_middleware.py | 2 +- tests/test_schedule_command.py | 236 ++++++++++++++++++ tests/test_scheduler.py | 98 ++++++++ tests/test_scheduler_tools.py | 124 +++++++++ uv.lock | 2 + 29 files changed, 1414 insertions(+), 41 deletions(-) create mode 100644 EvoScientist/commands/implementation/schedule.py create mode 100644 EvoScientist/cron/__init__.py create mode 100644 EvoScientist/cron/schedule.py create mode 100644 EvoScientist/middleware/scheduler.py create mode 100644 EvoScientist/subagents/scheduler.yaml create mode 100644 tests/test_cron_schedule.py create mode 100644 tests/test_schedule_command.py create mode 100644 tests/test_scheduler.py create mode 100644 tests/test_scheduler_tools.py diff --git a/EvoScientist/EvoScientist.py b/EvoScientist/EvoScientist.py index 3c979f5..926192e 100644 --- a/EvoScientist/EvoScientist.py +++ b/EvoScientist/EvoScientist.py @@ -668,6 +668,7 @@ def _get_default_middleware( create_memory_lifecycle_middleware, create_memory_middleware, create_runtime_context_middleware, + create_scheduler_middleware, create_tool_selector_middleware, load_fallback_chain, ) @@ -734,8 +735,10 @@ def _get_default_middleware( timeout=cfg.code_interpreter_timeout, max_result_chars=cfg.code_interpreter_max_result_chars, ), - create_runtime_context_middleware(), ] + if cfg.enable_scheduler and not for_async_subagent: + mw.append(create_scheduler_middleware()) + mw.append(create_runtime_context_middleware()) if memory_controls.memory_enabled: mw.append(memory_middleware) if memory_controls.worker_needed(worker_target): @@ -807,7 +810,11 @@ def _get_default_agent(): if not cfg.auto_approve: mw.append( HumanInTheLoopMiddleware( - interrupt_on={"execute": True, "run_in_background": True} + interrupt_on={ + "execute": True, + "run_in_background": True, + "schedule_task": True, + } ) ) @@ -961,7 +968,11 @@ def create_cli_agent( if not cfg.auto_approve: mw.append( HumanInTheLoopMiddleware( - interrupt_on={"execute": True, "run_in_background": True} + interrupt_on={ + "execute": True, + "run_in_background": True, + "schedule_task": True, + } ) ) diff --git a/EvoScientist/cli/async_notifier.py b/EvoScientist/cli/async_notifier.py index 5389c99..cec1ef2 100644 --- a/EvoScientist/cli/async_notifier.py +++ b/EvoScientist/cli/async_notifier.py @@ -480,13 +480,17 @@ def format_notification_lines( """ if not notifs: return [] - tasks = [n for n in notifs if n.kind != "bg-process"] + tasks = [n for n in notifs if n.kind == "agent"] shell = [n for n in notifs if n.kind == "bg-process"] + unknown = [n for n in notifs if n.kind not in {"agent", "bg-process"}] lines: list[tuple[str, str]] = [] if tasks: lines += _render_notification_group(tasks, " ✦ Agent Teams ✦ ", "Task") if shell: lines += _render_notification_group(shell, " ✦ Background ✦ ", "Cmd") + if unknown: + # Fallback so a future kind is never silently dropped from the display. + lines += _render_notification_group(unknown, " ✦ Updates ✦ ", "Task") return lines @@ -515,12 +519,14 @@ def format_batch_message(notifs: list[AsyncTaskNotification]) -> str: ) # bg-process is inspected with check_process; sub-agents with check_async_task. hints: list[str] = [] - if any(n.kind != "bg-process" for n in notifs): + if any(n.kind == "agent" for n in notifs): hints.append("check_async_task (sub-agents)") if any(n.kind == "bg-process" for n in notifs): hints.append("check_process (background processes)") + # Fallback when a batch has only unrecognized kinds (hints empty). + hint_text = " or ".join(hints) if hints else "the appropriate status tool" lines.append( - f"(Signal only — fetch full result via {' or '.join(hints)} if relevant to " + f"(Signal only — fetch full result via {hint_text} if relevant to " "the current step, else acknowledge & continue.)" ) return "\n".join(lines) diff --git a/EvoScientist/commands/implementation/__init__.py b/EvoScientist/commands/implementation/__init__.py index 1333598..0cdab5f 100644 --- a/EvoScientist/commands/implementation/__init__.py +++ b/EvoScientist/commands/implementation/__init__.py @@ -1,5 +1,14 @@ from __future__ import annotations -from . import channel, general, mcp, model, model_fallback, session, skills +from . import channel, general, mcp, model, model_fallback, schedule, session, skills -__all__ = ["channel", "general", "mcp", "model", "model_fallback", "session", "skills"] +__all__ = [ + "channel", + "general", + "mcp", + "model", + "model_fallback", + "schedule", + "session", + "skills", +] diff --git a/EvoScientist/commands/implementation/schedule.py b/EvoScientist/commands/implementation/schedule.py new file mode 100644 index 0000000..230ce2a --- /dev/null +++ b/EvoScientist/commands/implementation/schedule.py @@ -0,0 +1,234 @@ +from __future__ import annotations + +import asyncio +import re +from typing import ClassVar + +from rich.table import Table + +from ..base import Command, CommandContext, SubCommand +from ..manager import manager + + +class ScheduleCommand(Command): + """Manage scheduled (cron) tasks.""" + + name = "/schedule" + description = "Manage scheduled (cron) tasks" + subcommands: ClassVar[list[SubCommand]] = [ + SubCommand("add", 'Add: /schedule add ""'), + SubCommand("list", "List scheduled tasks"), + SubCommand("remove", "Remove a schedule by id"), + SubCommand("run", "Run a schedule's prompt once now (test)"), + SubCommand("pause", "Disable a schedule by id"), + SubCommand("resume", "Enable a schedule by id"), + ] + + async def execute(self, ctx: CommandContext, args: list[str]) -> None: + """Dispatch to the appropriate /schedule subcommand.""" + cfg = getattr(ctx, "config", None) + if cfg is not None and not getattr(cfg, "enable_scheduler", True): + ctx.ui.append_system( + "Scheduled tasks are disabled (`enable_scheduler` is off).", + style="yellow", + ) + return + + from ...cron import schedule as crons + + # Cron SDK calls are sync HTTP; offload to a thread so backend latency + # can never freeze the interactive event loop. + if not await asyncio.to_thread(crons.is_available): + ctx.ui.append_system( + "Scheduler unavailable: the langgraph dev backend is not running.", + style="yellow", + ) + return + + if not args or args[0].lower() == "list": + await self._list(ctx, crons) + return + + sub = args[0].lower() + rest = args[1:] + if sub == "add": + await self._add(ctx, crons, rest) + elif sub == "remove": + await self._remove(ctx, crons, rest[0] if rest else "") + elif sub == "run": + await self._run(ctx, crons, rest[0] if rest else "") + elif sub in ("pause", "resume"): + await self._set_enabled( + ctx, crons, rest[0] if rest else "", sub == "resume" + ) + else: + ctx.ui.append_system("Schedule commands:", style="bold") + for s in self.subcommands: + ctx.ui.append_system( + f" /schedule {s.name:<8} {s.description}", style="dim" + ) + + async def _add(self, ctx: CommandContext, crons, rest: list[str]) -> None: + # Cron may arrive as 5 separate tokens (unquoted) or 1 token (shlex-quoted). + # Split on any whitespace so extra spaces don't break detection; the + # backend rejects genuinely malformed expressions. + if rest and len(rest[0].split()) == 5: # quoted 5-field cron + schedule, prompt_tokens = " ".join(rest[0].split()), rest[1:] + elif len(rest) >= 5: # 5 separate cron fields + schedule, prompt_tokens = " ".join(rest[:5]), rest[5:] + else: + ctx.ui.append_system( + 'Usage: /schedule add "" ""', style="yellow" + ) + return + prompt = " ".join(prompt_tokens).strip().strip('"').strip("'") + if not prompt: + ctx.ui.append_system("A task prompt is required.", style="yellow") + return + # B3: strip unsafe chars; keep only alphanumerics + hyphens (kebab-case). + raw = prompt[:48].lower() + name = re.sub(r"[^a-z0-9]+", "-", raw).strip("-")[:32] or "task" + try: + rec = await asyncio.to_thread( + crons.create_schedule, name=name, schedule=schedule, prompt=prompt + ) + except Exception as exc: + ctx.ui.append_system(f"Error: {exc}", style="red") + return + ctx.ui.append_system( + f"Scheduled '{name}' [{schedule}] — id {rec.get('cron_id')}. " + "Runs unattended in the background.", + style="green", + ) + + async def _list(self, ctx: CommandContext, crons) -> None: + # B1: guard SDK call — backend may die after the is_available() check. + try: + rows = await asyncio.to_thread(crons.list_schedules) + except Exception as exc: + ctx.ui.append_system(f"Error: {exc}", style="red") + return + if not rows: + ctx.ui.append_system( + "No scheduled tasks. Add one: /schedule add ...", style="dim" + ) + return + table = Table(title="Scheduled Tasks", show_header=True) + table.add_column("ID", style="cyan") + table.add_column("Name", style="magenta") + table.add_column("Schedule", style="green") + table.add_column("Enabled", style="yellow") + table.add_column("Next run (UTC)", style="white") + for r in rows: + meta = r.get("metadata") or {} + table.add_row( + str(r.get("cron_id", ""))[:8], + str(meta.get("name", "")), + str(r.get("schedule", "")), + "yes" if r.get("enabled", True) else "no", + str(r.get("next_run_date", "")), + ) + ctx.ui.mount_renderable(table) + + _AMBIGUOUS = object() # B2: sentinel returned when multiple crons match a prefix + _BACKEND_ERROR = object() # sentinel returned when list_schedules() raises + + async def _resolve(self, crons, prefix: str): + """Return the unique matching record, _AMBIGUOUS if >1 match, _BACKEND_ERROR on error, or None.""" + # B1: guard SDK call — backend may die after is_available() check. + try: + all_rows = await asyncio.to_thread(crons.list_schedules) + except Exception as exc: + # Store the exception text so _resolve_or_report can surface it. + self._last_backend_exc = exc + return self._BACKEND_ERROR + # B2: collect ALL matches; ambiguous prefix → sentinel so callers can warn. + matches = [r for r in all_rows if str(r.get("cron_id", "")).startswith(prefix)] + if len(matches) > 1: + return self._AMBIGUOUS + return matches[0] if matches else None + + async def _resolve_or_report(self, ctx: CommandContext, crons, prefix: str): + """Resolve prefix → record, emit UI error on ambiguity/miss/error, return None on failure.""" + match = await self._resolve(crons, prefix) + if match is self._BACKEND_ERROR: + exc = getattr(self, "_last_backend_exc", None) + ctx.ui.append_system( + f"Error: scheduler backend unavailable ({exc})", + style="red", + ) + return None + if match is self._AMBIGUOUS: + ctx.ui.append_system( + f"Multiple schedules match '{prefix}' — use a longer id.", + style="yellow", + ) + return None + if not match: + ctx.ui.append_system(f"No schedule matching {prefix}.", style="yellow") + return None + return match + + async def _remove(self, ctx: CommandContext, crons, prefix: str) -> None: + if not prefix: + ctx.ui.append_system("Usage: /schedule remove ", style="yellow") + return + match = await self._resolve_or_report(ctx, crons, prefix) + if match is None: + return + cron_id = str(match.get("cron_id", "")) + try: + await asyncio.to_thread(crons.delete_schedule, cron_id) + except Exception as exc: + ctx.ui.append_system(f"Error: {exc}", style="red") + return + ctx.ui.append_system(f"Removed schedule {cron_id}.", style="green") + + async def _run(self, ctx: CommandContext, crons, prefix: str) -> None: + if not prefix: + ctx.ui.append_system("Usage: /schedule run ", style="yellow") + return + match = await self._resolve_or_report(ctx, crons, prefix) + if match is None: + return + prompt = (match.get("metadata") or {}).get("prompt", "") + if not str(prompt).strip(): + ctx.ui.append_system( + f"Schedule {prefix} has no stored prompt — cannot run it.", + style="yellow", + ) + return + try: + rec = await asyncio.to_thread(crons.run_now, prompt) + except Exception as exc: + ctx.ui.append_system(f"Error: {exc}", style="red") + return + # Don't promise a location; the task's own prompt decides where output goes. + ctx.ui.append_system( + f"Fired schedule {prefix} once now (run {rec.get('run_id')}). " + "Any output goes wherever the task's instruction specifies.", + style="green", + ) + + async def _set_enabled( + self, ctx: CommandContext, crons, prefix: str, enabled: bool + ) -> None: + if not prefix: + ctx.ui.append_system("Usage: /schedule pause|resume ", style="yellow") + return + match = await self._resolve_or_report(ctx, crons, prefix) + if match is None: + return + cron_id = str(match.get("cron_id", "")) + try: + await asyncio.to_thread(crons.set_enabled, cron_id, enabled) + except Exception as exc: + ctx.ui.append_system(f"Error: {exc}", style="red") + return + ctx.ui.append_system( + f"{'Resumed' if enabled else 'Paused'} schedule {cron_id}.", style="green" + ) + + +# Register schedule command +manager.register(ScheduleCommand()) diff --git a/EvoScientist/config/settings.py b/EvoScientist/config/settings.py index 647f353..cb8e4b2 100644 --- a/EvoScientist/config/settings.py +++ b/EvoScientist/config/settings.py @@ -100,7 +100,7 @@ class EvoScientistConfig: provider: Default LLM provider ('anthropic', 'openai', 'google-genai', or 'nvidia'). model: Default model name (short name or full ID). auxiliary_provider: Provider for auxiliary_model (empty = use main provider). - auxiliary_model: Model for memory workers + tool selector (empty = use main model). + auxiliary_model: Model for memory workers + tool selector + scheduler (empty = use main model). default_mode: Default workspace mode ('daemon' or 'run'). default_workdir: Default workspace directory (empty = use current working directory). show_thinking: Whether to show thinking panels in CLI. @@ -165,6 +165,15 @@ class EvoScientistConfig: # its own port (langgraph_dev_port); this is just the browser server. webui_port: int = 4716 + # --- Scheduled tasks (cron) --- + # Master switch for scheduled tasks (/schedule, NL tools, scheduler context). Defaults + # True so the feature is available out-of-the-box; set False to disable. + enable_scheduler: bool = True + # Default IANA timezone for cron schedules created without an explicit tz. + # Empty string => the host's local IANA zone (resolved via tzlocal), falling + # back to UTC if it can't be determined; set e.g. "Europe/London" to pin one. + scheduler_default_timezone: str = "" + # Whether langgraph dev persists its runtime state to .langgraph_api/ next # to the subprocess cwd. True (default) keeps async-task, scheduler, and # Store API state across subprocess restarts — useful for future @@ -660,6 +669,8 @@ _ENV_MAPPINGS = { "enable_async_subagents": "EVOSCIENTIST_ENABLE_ASYNC_SUBAGENTS", "langgraph_dev_port": "EVOSCIENTIST_LANGGRAPH_DEV_PORT", "webui_port": "EVOSCIENTIST_WEBUI_PORT", + "enable_scheduler": "EVOSCIENTIST_ENABLE_SCHEDULER", + "scheduler_default_timezone": "EVOSCIENTIST_SCHEDULER_DEFAULT_TIMEZONE", "code_interpreter_timeout": "EVOSCIENTIST_CODE_INTERPRETER_TIMEOUT", "code_interpreter_max_result_chars": "EVOSCIENTIST_CODE_INTERPRETER_MAX_RESULT_CHARS", "sandbox_execute_timeout": "EVOSCIENTIST_SANDBOX_EXECUTE_TIMEOUT", diff --git a/EvoScientist/cron/__init__.py b/EvoScientist/cron/__init__.py new file mode 100644 index 0000000..f4a89b6 --- /dev/null +++ b/EvoScientist/cron/__init__.py @@ -0,0 +1,23 @@ +"""Scheduled tasks (cron) for EvoScientist — thin layer over langgraph crons.""" + +from .schedule import ( + SCHEDULED_RUN_KIND, + SCHEDULER_GRAPH_ID, + create_schedule, + delete_schedule, + is_available, + list_schedules, + run_now, + set_enabled, +) + +__all__ = [ + "SCHEDULED_RUN_KIND", + "SCHEDULER_GRAPH_ID", + "create_schedule", + "delete_schedule", + "is_available", + "list_schedules", + "run_now", + "set_enabled", +] diff --git a/EvoScientist/cron/schedule.py b/EvoScientist/cron/schedule.py new file mode 100644 index 0000000..6bc2789 --- /dev/null +++ b/EvoScientist/cron/schedule.py @@ -0,0 +1,118 @@ +"""Thin wrapper over the langgraph dev built-in cron API (langgraph_sdk). + +EvoScientist scheduled tasks ARE langgraph crons targeting the ``scheduler`` +graph. This module is the single choke-point so the ``/schedule`` command and the +NL ``schedule_task`` tool share one implementation. + +Isolation is **process-level**, not data-level: EvoScientist's manager.py restarts +langgraph dev when the active workspace changes, so each workspace gets its own +langgraph-dev process and its own ``.langgraph_api`` cron store. If you point +multiple clients at one hand-started server they will share the same cron store. +""" + +from __future__ import annotations + +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from langgraph_sdk.schema import Cron, Run + +SCHEDULER_GRAPH_ID = "scheduler" +SCHEDULED_RUN_KIND = "scheduled_task" + + +def _scheduler_url() -> str: + from ..EvoScientist import _ensure_config + + cfg = _ensure_config() + port = int(getattr(cfg, "langgraph_dev_port", 6174)) + return f"http://localhost:{port}" + + +def _client(): + from langgraph_sdk import get_sync_client + + return get_sync_client(url=_scheduler_url(), headers={"x-auth-scheme": "langsmith"}) + + +def _default_timezone() -> str | None: + from ..EvoScientist import _ensure_config + + tz = getattr(_ensure_config(), "scheduler_default_timezone", "") or "" + if tz: + return tz + # Resolve the host's real IANA zone (e.g. "Asia/Shanghai") so absolute-time + # schedules fire in local time and track DST. Falls back to None (-> UTC in + # the cron backend) when the local zone can't be determined. + try: + from tzlocal import get_localzone_name + + return get_localzone_name() + except Exception: + return None + + +def is_available() -> bool: + """True when the langgraph dev backend (which fires crons) is reachable.""" + from ..langgraph_dev.manager import is_langgraph_dev_running + + return bool(is_langgraph_dev_running(base_url=_scheduler_url())) + + +def create_schedule( + *, name: str, schedule: str, prompt: str, timezone: str | None = None +) -> Cron: + """Create a recurring scheduled task on the scheduler graph.""" + # Crons are stored in the langgraph-dev process's .langgraph_api store, not + # tagged by workspace. Isolation is process-level (see module docstring). + return _client().crons.create( + assistant_id=SCHEDULER_GRAPH_ID, + schedule=schedule, + input={"messages": [{"role": "user", "content": prompt}]}, + metadata={"run_kind": SCHEDULED_RUN_KIND, "name": name, "prompt": prompt}, + timezone=timezone or _default_timezone(), + ) + + +def list_schedules() -> list[Cron]: + """Return only EvoScientist scheduled tasks. + + Filtered server-side by ``run_kind`` metadata (the cron backend matches by + metadata containment), so we never page through unrelated crons; ``limit`` is + a ceiling on OUR schedules (far below 1000 in practice). We filter on metadata + rather than ``assistant_id`` because the stored ``assistant_id`` is a resolved + UUID, not the ``scheduler`` graph name we create with. + """ + return _client().crons.search( + metadata={"run_kind": SCHEDULED_RUN_KIND}, + limit=1000, + ) + + +def delete_schedule(cron_id: str) -> None: + """Delete a scheduled task by cron id.""" + _client().crons.delete(cron_id) + + +def set_enabled(cron_id: str, enabled: bool) -> Cron: + """Enable or disable a scheduled task by cron id.""" + return _client().crons.update(cron_id, enabled=enabled) + + +def run_now(prompt: str) -> Run: + """Fire a one-off scheduler run immediately (for ``/schedule run``). + + Output goes wherever the task's prompt specifies; there is no push notification. + """ + client = _client() + thread = client.threads.create(graph_id=SCHEDULER_GRAPH_ID) + return client.runs.create( + thread_id=str(thread["thread_id"]), + assistant_id=SCHEDULER_GRAPH_ID, + input={"messages": [{"role": "user", "content": prompt}]}, + metadata={ + "run_kind": SCHEDULED_RUN_KIND, + "name": "manual-run", + "prompt": prompt, + }, + ) diff --git a/EvoScientist/langgraph_dev/graphs.py b/EvoScientist/langgraph_dev/graphs.py index 5e259d6..835e283 100644 --- a/EvoScientist/langgraph_dev/graphs.py +++ b/EvoScientist/langgraph_dev/graphs.py @@ -11,11 +11,11 @@ To add a new async sub-agent: 1. Set ``async: true`` in ``EvoScientist/subagents/.yaml``. 2. Add a one-line binding here:: - _agent = build_async_subagent_graph("") + = build_async_subagent_graph("") 3. Register it in ``EvoScientist/langgraph_dev/langgraph.json``:: - "": "EvoScientist.langgraph_dev.graphs:_agent" + "": "EvoScientist.langgraph_dev.graphs:" The deployed main agent (``EvoScientist_agent``) lives in ``main_graph.py`` because it follows a different mechanism (re-exporting a lazily-constructed @@ -30,5 +30,6 @@ from EvoScientist.subagents._factory import build_async_subagent_graph writing_agent = build_async_subagent_graph("writing-agent") data_analysis_agent = build_async_subagent_graph("data-analysis-agent") +scheduler = build_async_subagent_graph("scheduler") evomemory_subagent_worker = build_memory_worker_graph(MemoryLifecycleRole.SUBAGENT) evomemory_turn_worker = build_memory_worker_graph(MemoryLifecycleRole.TURN) diff --git a/EvoScientist/langgraph_dev/langgraph.json b/EvoScientist/langgraph_dev/langgraph.json index 5c568c4..6cc690c 100644 --- a/EvoScientist/langgraph_dev/langgraph.json +++ b/EvoScientist/langgraph_dev/langgraph.json @@ -4,6 +4,7 @@ "EvoScientist": "EvoScientist.langgraph_dev.main_graph:EvoScientist_agent", "writing-agent": "EvoScientist.langgraph_dev.graphs:writing_agent", "data-analysis-agent": "EvoScientist.langgraph_dev.graphs:data_analysis_agent", + "scheduler": "EvoScientist.langgraph_dev.graphs:scheduler", "evomemory-subagent-worker": "EvoScientist.langgraph_dev.graphs:evomemory_subagent_worker", "evomemory-turn-worker": "EvoScientist.langgraph_dev.graphs:evomemory_turn_worker" }, diff --git a/EvoScientist/langgraph_dev/manager.py b/EvoScientist/langgraph_dev/manager.py index fb0552c..b419892 100644 --- a/EvoScientist/langgraph_dev/manager.py +++ b/EvoScientist/langgraph_dev/manager.py @@ -92,6 +92,8 @@ def needs_langgraph_dev(config: EvoScientistConfig) -> bool: """Return whether this config needs the background langgraph dev server.""" if config.enable_async_subagents: return True + if config.enable_scheduler: + return True memory_controls = MemoryControls.from_config(config) return memory_controls.worker_needed( MemoryObservationTarget.TURN_WORKER diff --git a/EvoScientist/middleware/__init__.py b/EvoScientist/middleware/__init__.py index 0ad56dc..972a68a 100644 --- a/EvoScientist/middleware/__init__.py +++ b/EvoScientist/middleware/__init__.py @@ -29,6 +29,10 @@ from .memory_lifecycle import ( ) from .model_fallback import ModelFallbackMiddleware, load_fallback_chain from .runtime_context import RuntimeContextMiddleware, create_runtime_context_middleware +from .scheduler import ( + SchedulerMiddleware, + create_scheduler_middleware, +) from .tool_error_handler import ToolErrorHandlerMiddleware from .tool_selector import create_tool_selector_middleware from .utils import disable_thinking @@ -46,6 +50,7 @@ __all__ = [ "ModelFallbackMiddleware", "Question", "RuntimeContextMiddleware", + "SchedulerMiddleware", "ToolErrorHandlerMiddleware", "compute_context_editing_trigger", "create_code_interpreter_middleware", @@ -53,6 +58,7 @@ __all__ = [ "create_memory_lifecycle_middleware", "create_memory_middleware", "create_runtime_context_middleware", + "create_scheduler_middleware", "create_tool_selector_middleware", "disable_thinking", "load_fallback_chain", diff --git a/EvoScientist/middleware/memory.py b/EvoScientist/middleware/memory.py index 39dd318..1fb8cc5 100644 --- a/EvoScientist/middleware/memory.py +++ b/EvoScientist/middleware/memory.py @@ -24,7 +24,6 @@ from langchain.agents.middleware.types import ( ModelRequest, ModelResponse, ) -from langchain_core.messages import SystemMessage from .. import paths as _paths from ..memory import ( @@ -35,6 +34,7 @@ from ..memory import ( create_record_observation_tool, create_search_observations_tool, ) +from .utils import append_to_system_message logger = logging.getLogger(__name__) @@ -162,21 +162,6 @@ Notes about this workspace: conventions, commands, tests, and traps. } -def _append_to_system_message( - system_message: SystemMessage | None, - text: str, -) -> SystemMessage: - """Append text to a system message while preserving existing metadata.""" - existing_blocks = list(system_message.content_blocks) if system_message else [] - new_blocks = [ - *existing_blocks, - {"type": "text", "text": text}, - ] - if system_message is None: - return SystemMessage(content=new_blocks) - return system_message.model_copy(update={"content": new_blocks}) - - def _short_hash(text: str, *, n: int = 16) -> str: """Return a deterministic hash fragment for generated profile paths.""" import hashlib @@ -782,7 +767,7 @@ class EvoMemoryMiddleware(AgentMiddleware): observation_index_context=observation_index_context, profile_content=profile_content, ) - new_system = _append_to_system_message(request.system_message, injection) + new_system = append_to_system_message(request.system_message, injection) return request.override(system_message=new_system) def _profile_context_for_request(self) -> str: diff --git a/EvoScientist/middleware/scheduler.py b/EvoScientist/middleware/scheduler.py new file mode 100644 index 0000000..9307c40 --- /dev/null +++ b/EvoScientist/middleware/scheduler.py @@ -0,0 +1,222 @@ +"""Scheduling middleware and tools for the main agent. + +Bundles three NL scheduling tools (schedule_task, list_scheduled_tasks, +cancel_scheduled_task) together with the SchedulerMiddleware that: +- Injects scheduling guidance + a live ```` block into every + system prompt (static→dynamic, mirroring EvoMemoryMiddleware). +- Contributes the three tools via ``self.tools`` so agent wiring only needs to + append this middleware — no separate tool-list management required. +""" + +from __future__ import annotations + +import asyncio +import time +from collections.abc import Awaitable, Callable + +from langchain.agents.middleware.types import ( + AgentMiddleware, + ModelRequest, + ModelResponse, +) +from langchain_core.tools import tool + +from .utils import append_to_system_message + +_CACHE_TTL_SECONDS = 15.0 + +# Static guidance injected into the system prompt; the dynamic +# list follows it (static→dynamic, mirroring EvoMemoryMiddleware's +# + ordering). +_SCHEDULING_INSTRUCTIONS = """ +When the user asks for something on a recurring schedule ("every 10 minutes", +"each morning at 7"), translate the timing into a 5-field cron expression and +call `schedule_task(name, cron, prompt, timezone)`. The `prompt` must be a +complete, self-contained instruction. Spell out the task fully AND name an +explicit destination for its result, e.g. "...and save the summary to +`/memories/daily-papers.md`" or "...and append the status to `experiment_log.json`" +or "...and update my research memory". It can use skills (e.g. +`paper-navigator`) and the workspace; a task that never says where to put its +output leaves no trace. Use `list_scheduled_tasks` / `cancel_scheduled_task` to +manage them. +""" + + +# --------------------------------------------------------------------------- +# NL scheduling tools +# --------------------------------------------------------------------------- + + +@tool +def schedule_task(name: str, cron: str, prompt: str, timezone: str = "") -> str: + """Create a recurring scheduled task that runs unattended in the background. + + Translate the user's natural-language timing into a standard 5-field cron + expression yourself before calling (e.g. 'every 10 minutes' -> '*/10 * * * *', + 'every day at 7am' -> '0 7 * * *', 'every Monday 9am' -> '0 9 * * 1'). + + Args: + name: short human label for the task (e.g. "uk-weather"). + cron: 5-field cron expression. + prompt: the full instruction the background scheduler runs each time. + timezone: optional IANA tz (e.g. "Europe/London"); empty = host local zone. + """ + from ..cron import schedule as crons + + if not crons.is_available(): + return "Scheduler unavailable: the langgraph dev backend is not running." + try: + rec = crons.create_schedule( + name=name, schedule=cron, prompt=prompt, timezone=timezone or None + ) + except Exception as e: + return f"Error: {e}" + return ( + f"Scheduled '{name}' [{cron}] — id {rec.get('cron_id')}. It runs unattended in the " + "background; output goes wherever the task's prompt specifies. Use list_scheduled_tasks to review." + ) + + +@tool +def list_scheduled_tasks() -> str: + """List the user's recurring scheduled tasks (id, name, schedule, enabled).""" + from ..cron import schedule as crons + + if not crons.is_available(): + return "Scheduler unavailable: the langgraph dev backend is not running." + try: + rows = crons.list_schedules() + except Exception as e: + return f"Error: {e}" + if not rows: + return "No scheduled tasks." + lines = [] + for r in rows: + meta = r.get("metadata") or {} + lines.append( + f"- {str(r.get('cron_id', ''))[:8]} | {meta.get('name', '')} | " + f"{r.get('schedule', '')} | {'on' if r.get('enabled', True) else 'off'}" + ) + return "\n".join(lines) + + +@tool +def cancel_scheduled_task(cron_id: str) -> str: + """Cancel (delete) a scheduled task. Pass the id (or its prefix) shown by list_scheduled_tasks.""" + from ..cron import schedule as crons + + if not crons.is_available(): + return "Scheduler unavailable: the langgraph dev backend is not running." + if not cron_id.strip(): + # Empty prefix would match (and delete) the only cron — refuse it. + return "Provide the id (or a prefix) of the task to cancel." + try: + rows = crons.list_schedules() + # B2: collect ALL prefix matches before acting to detect ambiguity. + matches = [r for r in rows if str(r.get("cron_id", "")).startswith(cron_id)] + if not matches: + return f"No scheduled task matching '{cron_id}'." + if len(matches) > 1: + ids = ", ".join(str(r.get("cron_id", ""))[:8] for r in matches) + return f"Multiple schedules match '{cron_id}' ({ids}) — use a longer id." + target = str(matches[0]["cron_id"]) + crons.delete_schedule(target) + except Exception as e: + return f"Error: {e}" + return f"Cancelled scheduled task {target}." + + +# --------------------------------------------------------------------------- +# Middleware +# --------------------------------------------------------------------------- + + +class SchedulerMiddleware(AgentMiddleware): + """Inject scheduling awareness + contribute scheduling tools to the main agent.""" + + name = "scheduler" + + def __init__(self) -> None: + super().__init__() + self._cache: str | None = None + self._cache_at: float = 0.0 + self.tools = [schedule_task, list_scheduled_tasks, cancel_scheduled_task] + + def _schedules_block(self) -> str: + """Build the dynamic ```` block (empty if none / down).""" + from ..cron import schedule as crons + + def _clean(value: object) -> str: + # Flatten to one line + drop angle brackets so a task's own name or + # prompt can't break the block or inject pseudo-tags. + return " ".join(str(value).split()).replace("<", "").replace(">", "") + + try: + if not crons.is_available(): + return "" + rows = crons.list_schedules() + except Exception: + return "" + if not rows: + return "" + lines = [ + "", + "Background cron tasks currently scheduled (they run unattended; you " + "need do nothing — this is for your awareness, e.g. to answer " + "questions about them or avoid creating duplicates):", + ] + for r in rows: + meta = r.get("metadata") or {} + lines.append( + f"- id={str(r.get('cron_id', ''))[:8]} | {_clean(meta.get('name', ''))} | " + f"{_clean(r.get('schedule', ''))} | " + f"{'on' if r.get('enabled', True) else 'off'} | " + f"{_clean(meta.get('prompt', ''))[:80]}" + ) + lines.append("") + return "\n".join(lines) + + def _cached_schedules_block(self) -> str: + now = time.monotonic() + if self._cache is None or (now - self._cache_at) > _CACHE_TTL_SECONDS: + self._cache = self._schedules_block() + self._cache_at = now + return self._cache + + def _injection(self, schedules_block: str) -> str: + """Static instructions, then the dynamic list (static→dynamic, like memory).""" + parts = [_SCHEDULING_INSTRUCTIONS] + if schedules_block: + parts.append(schedules_block) + return "\n\n".join(parts) + + def modify_request(self, request: ModelRequest) -> ModelRequest: + injection = self._injection(self._cached_schedules_block()) + new_system = append_to_system_message(request.system_message, injection) + return request.override(system_message=new_system) + + def wrap_model_call( + self, + request: ModelRequest, + handler: Callable[[ModelRequest], ModelResponse], + ) -> ModelResponse: + return handler(self.modify_request(request)) + + async def awrap_model_call( + self, + request: ModelRequest, + handler: Callable[[ModelRequest], Awaitable[ModelResponse]], + ) -> ModelResponse: + # Offload the (SDK-touching) list build to a thread so we never block the + # event loop or trip langgraph dev's blockbuster detector. + block = await asyncio.to_thread(self._cached_schedules_block) + injection = self._injection(block) + request = request.override( + system_message=append_to_system_message(request.system_message, injection) + ) + return await handler(request) + + +def create_scheduler_middleware() -> SchedulerMiddleware: + """Factory for the scheduler middleware (main agent only).""" + return SchedulerMiddleware() diff --git a/EvoScientist/middleware/utils.py b/EvoScientist/middleware/utils.py index c884f64..bd61d41 100644 --- a/EvoScientist/middleware/utils.py +++ b/EvoScientist/middleware/utils.py @@ -9,6 +9,7 @@ from __future__ import annotations from typing import Any from langchain_core.language_models import BaseChatModel +from langchain_core.messages import SystemMessage def disable_thinking(model: BaseChatModel) -> BaseChatModel: @@ -42,3 +43,20 @@ def disable_thinking(model: BaseChatModel) -> BaseChatModel: # Fallback for non-Pydantic or unusual model classes # Note: bind() may not effectively override first-class Pydantic fields return model.bind(**updates) + + +def append_to_system_message( + system_message: SystemMessage | None, text: str +) -> SystemMessage: + """Append a text block to a system message, preserving its metadata. + + Used by the memory and scheduler middleware. Unlike building a fresh + ``SystemMessage``, ``model_copy`` keeps ``additional_kwargs`` (e.g. + ``cache_control`` prompt-cache breakpoints), ``id``, ``name`` and + ``response_metadata`` from the original message. + """ + existing_blocks = list(system_message.content_blocks) if system_message else [] + new_blocks = [*existing_blocks, {"type": "text", "text": text}] + if system_message is None: + return SystemMessage(content=new_blocks) + return system_message.model_copy(update={"content": new_blocks}) diff --git a/EvoScientist/subagents/_factory.py b/EvoScientist/subagents/_factory.py index af91522..21aeaca 100644 --- a/EvoScientist/subagents/_factory.py +++ b/EvoScientist/subagents/_factory.py @@ -40,13 +40,14 @@ def build_async_subagent_graph(name: str) -> Any: from EvoScientist.config import apply_config_to_env, get_effective_config from EvoScientist.EvoScientist import ( SUBAGENTS_CONFIG, + _ensure_auxiliary_chat_model, _ensure_chat_model, _ensure_general_purpose_subagent, _get_default_backend, _get_default_middleware, _inject_subagent_middleware, ) - from EvoScientist.tools import tavily_search, think_tool + from EvoScientist.tools import skill_manager, tavily_search, think_tool from EvoScientist.utils import load_subagents # Surface API keys as env vars so downstream SDKs (openai, anthropic, …) @@ -55,7 +56,7 @@ def build_async_subagent_graph(name: str) -> Any: apply_config_to_env(cfg) # Mirror the tool registry constructed in EvoScientist._build_base_kwargs. - tool_registry = {"think_tool": think_tool} + tool_registry = {"think_tool": think_tool, "skill_manager": skill_manager} if os.environ.get("TAVILY_API_KEY"): tool_registry["tavily_search"] = tavily_search @@ -102,16 +103,23 @@ def build_async_subagent_graph(name: str) -> Any: _ensure_general_purpose_subagent(subagents) _inject_subagent_middleware(subagents) + middleware = _get_default_middleware( + for_async_subagent=True, + memory_source_agent=name, + ) + + # Scheduler is an unattended timer task → use the cheaper auxiliary model. + model = ( + _ensure_auxiliary_chat_model() if name == "scheduler" else _ensure_chat_model() + ) + return create_deep_agent( name=name, - model=_ensure_chat_model(), + model=model, system_prompt=spec.get("system_prompt", ""), tools=spec.get("tools", []) + agent_mcp_tools, skills=spec.get("skills"), backend=_get_default_backend(), - middleware=_get_default_middleware( - for_async_subagent=True, - memory_source_agent=name, - ), + middleware=middleware, subagents=subagents, ).with_config({"recursion_limit": cfg.recursion_limit}) diff --git a/EvoScientist/subagents/scheduler.yaml b/EvoScientist/subagents/scheduler.yaml new file mode 100644 index 0000000..1626e8a --- /dev/null +++ b/EvoScientist/subagents/scheduler.yaml @@ -0,0 +1,24 @@ +scheduler: + description: "Unattended background executor for a single scheduled (cron) task." + # write_file / read_file / edit_file / ls / execute are DeepAgents built-ins (always present). + # tavily_search requires TAVILY_API_KEY; it is silently omitted when the key is absent. + tools: [think_tool, tavily_search, skill_manager] + skills: ["/skills/"] + # Fired by langgraph-dev crons (client.crons.create). One run = one task. + # Deployment graph: EvoScientist.langgraph_dev.graphs:scheduler + async: true + system_prompt: | + You are the scheduler: an UNATTENDED background worker that runs ONE + scheduled task to completion on a timer. No human is present — never ask + questions; make reasonable assumptions and proceed. + + Your task is the user message. Do it end-to-end, exactly as instructed. + The instruction is the single source of truth: if it tells you to write + files, write them to the paths it specifies (create directories as needed) — + there is no default output directory. If it does not ask for output, just + complete the task. + + Rules: + - Keep outputs concise. Never fabricate data or citations. + - The current date is already in your runtime context. Use `execute` + (e.g. `date`) when you need the exact wall-clock time. diff --git a/README.md b/README.md index fbcce16..71f6d66 100644 --- a/README.md +++ b/README.md @@ -176,6 +176,7 @@ Moving beyond traditional human-in-the-loop systems, EvoScientist adopts a human - [📦 Installation](#-installation) - [🔑 Configuration](#-configuration) - [⚡ Quick Start](#-quick-start) +- [⏰ Scheduled Tasks](#-scheduled-tasks) - [🍪 Examples & Recipes](#-examples--recipes) - [🔌 MCP Integration](#-mcp-integration) - [📱 Channels](#-channels) @@ -535,6 +536,31 @@ for state in EvoScientist_agent.stream(

🔝Back to top

+## ⏰ Scheduled Tasks + +Automate recurring research tasks with cron-style schedules. + +```bash +# Add a schedule (cron expression required for /schedule add) +/schedule add "0 9 * * 1-5" "Summarise the latest ML papers from arXiv with the paper-navigator skill, and save the summary to /memories/daily-papers.md" +/schedule add "*/10 * * * *" "Check my running experiment's status and append the result to experiment_log.json" + +# Manage schedules +/schedule list # list active schedules +/schedule remove # delete a schedule +/schedule run # fire a schedule immediately +/schedule pause # pause without deleting +/schedule resume # resume a paused schedule +``` + +Note: `/schedule add` requires a cron expression (5 fields, e.g. `*/10 * * * *`). To schedule with natural language ("every 10 minutes"), just ask in chat — the agent translates it via the `schedule_task` tool. + +Output goes wherever the task's prompt tells it to write — there is no enforced output directory, so make the prompt specific about file locations. Run `/schedule list` to review schedules; the agent is also made aware of the active schedules via a `` context block, so you can just ask it what's scheduled. + +> **Cost note:** each scheduled run consumes LLM tokens. Delete unused schedules with `/schedule remove` to avoid accumulating charges. + +

🔝Back to top

+ ## 🍪 Examples & Recipes A curated collection of official examples, advanced usage patterns, and community-contributed recipes to help you get the most out of EvoScientist. @@ -608,9 +634,9 @@ Coming soon: - [x] 📑 Technical report on the way - [x] 🔐 OAuth sign-in (CLI coding agent subscribers) - [x] 📺 Web app with workspace UI +- [x] ⏰ Scheduled tasks (cron-style, via `/schedule`) - [ ] 📹 Demo and tutorial in the works - [ ] 📊 Benchmark suite to be released -- [ ] ⏰ Scheduled tasks for the core system planned Stay tuned — more features are on the way! @@ -625,7 +651,7 @@ Stay tuned — more features are on the way! - Xi Zhang
diff --git a/README.zh-CN.md b/README.zh-CN.md index b66e0ec..aaed297 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -185,6 +185,7 @@ EvoScientist 超越了传统的人在回路(Human-in-the-Loop)模式,采 - [📦 安装](#-安装) - [🔑 配置](#-配置) - [⚡ 快速上手](#-快速上手) +- [⏰ 定时任务](#-定时任务) - [🍪 示例与实践](#-示例与实践) - [🔌 MCP 集成](#-mcp-集成) - [📱 渠道接入](#-渠道接入) @@ -544,6 +545,31 @@ for state in EvoScientist_agent.stream(

🔝回到顶部

+## ⏰ 定时任务 + +用 cron 风格的计划任务自动化重复性研究工作。 + +```bash +# 添加计划任务(/schedule add 需要 cron 表达式) +/schedule add "0 9 * * 1-5" "用 paper-navigator 技能总结 arXiv 上最新的 ML 论文,并把摘要保存到 /memories/daily-papers.md" +/schedule add "*/10 * * * *" "检查我正在运行的实验状态,并把结果追加到 experiment_log.json" + +# 管理计划任务 +/schedule list # 列出活跃的计划任务 +/schedule remove # 删除一个计划任务 +/schedule run # 立即触发一次 +/schedule pause # 暂停但不删除 +/schedule resume # 恢复已暂停的计划任务 +``` + +说明:`/schedule add` 需要 cron 表达式(5 字段,例如 `*/10 * * * *`)。想用自然语言("每 10 分钟")排程,直接在对话里说即可——智能体会通过 `schedule_task` 工具自动翻译。 + +输出写到任务 prompt 指定的位置——没有强制的输出目录,所以请在 prompt 里写明文件位置。用 `/schedule list` 查看计划任务;智能体也会通过 `` 上下文块感知当前的计划任务,所以你也可以直接问它有哪些任务。 + +> **成本提示:** 每次计划任务运行都会消耗 LLM token。不用的任务请用 `/schedule remove` 删除,避免持续计费。 + +

🔝回到顶部

+ ## 🍪 示例与实践 收集了一些官方示例、进阶用法和社区贡献的实践方案,帮助你更好地使用 EvoScientist。 @@ -617,9 +643,9 @@ channel_enabled: "telegram,slack,feishu,qq" - [x] 📑 技术报告已发布 - [x] 🔐 OAuth 登录(CLI 编程智能体订阅用户) - [x] 📺 带工作区的 Web 应用界面(beta) +- [x] ⏰ 定时任务(cron 风格,通过 /schedule) - [ ] 📹 Demo 与教程正在制作中 - [ ] 📊 基准测试套件即将推出 -- [ ] ⏰ 核心系统定时任务规划中 敬请期待——更多功能正在路上! @@ -634,7 +660,7 @@ channel_enabled: "telegram,slack,feishu,qq" - Xi Zhang
diff --git a/pyproject.toml b/pyproject.toml index f541088..33de965 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -39,6 +39,7 @@ dependencies = [ "lazy-loader>=0.5", "markdownify>=1.2", "nest-asyncio>=1.6", + "tzlocal>=5.0", "langchain-mcp-adapters>=0.2", # <8.2.7: 8.2.7 Kitty "report-all-keys" breaks CJK input on iTerm2 "textual>=8.0,<8.2.7", diff --git a/tests/test_async_notifier.py b/tests/test_async_notifier.py index 309b746..b7c7769 100644 --- a/tests/test_async_notifier.py +++ b/tests/test_async_notifier.py @@ -431,6 +431,23 @@ def test_format_batch_message_multiple(): assert "check_async_task" in msg.lower() # hint to LLM +def test_format_batch_message_unknown_kind_uses_fallback_hint(): + """A batch with only unrecognized kinds must not emit an empty hint join.""" + notifs = [ + async_notifier.AsyncTaskNotification( + task_id="t1", + agent_name="x", + status="success", + received_at="2026-05-07T12:00:00Z", + prompt="A", + kind="weird", + ) + ] + msg = format_batch_message(notifs) + assert "the appropriate status tool" in msg + assert "via " not in msg # no empty " or ".join(hints) + + def test_dedup_preserves_order(): """dedup_notifications preserves the original order of notifications.""" notifs = [ diff --git a/tests/test_cli_serve.py b/tests/test_cli_serve.py index 4eb7642..89d75e4 100644 --- a/tests/test_cli_serve.py +++ b/tests/test_cli_serve.py @@ -31,6 +31,7 @@ def _make_config( enable_ask_user=enable_ask_user, dangerous_mode=dangerous_mode, enable_async_subagents=False, + enable_scheduler=False, memory_profile_enabled=True, memory_observations_enabled=True, memory_observation_writer=MemoryObservationWriter.ALL, diff --git a/tests/test_config.py b/tests/test_config.py index f10f7c2..abec7ca 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -738,3 +738,19 @@ class TestAuxiliaryModelConfig: config = get_effective_config() assert config.auxiliary_model == "gpt-5.5" assert config.auxiliary_provider == "openai" + + +def test_scheduler_config_defaults_and_env(monkeypatch): + from EvoScientist.config.settings import EvoScientistConfig, get_effective_config + + c = EvoScientistConfig() + assert c.enable_scheduler is True + assert c.scheduler_default_timezone == "" + + monkeypatch.setenv("EVOSCIENTIST_ENABLE_SCHEDULER", "false") + eff = get_effective_config({}) + assert eff.enable_scheduler is False + + monkeypatch.setenv("EVOSCIENTIST_SCHEDULER_DEFAULT_TIMEZONE", "America/New_York") + eff2 = get_effective_config({}) + assert eff2.scheduler_default_timezone == "America/New_York" diff --git a/tests/test_cron_schedule.py b/tests/test_cron_schedule.py new file mode 100644 index 0000000..9b3fb69 --- /dev/null +++ b/tests/test_cron_schedule.py @@ -0,0 +1,136 @@ +"""Unit tests for the cron wrapper (mock langgraph_sdk client).""" + +from unittest.mock import MagicMock + + +def _patch_client(monkeypatch): + from EvoScientist.cron import schedule as crons + + fake = MagicMock() + fake.crons.create.return_value = {"cron_id": "c-1", "schedule": "*/10 * * * *"} + fake.crons.search.return_value = [ + { + "cron_id": "c-1", + "schedule": "*/10 * * * *", + "metadata": {"run_kind": "scheduled_task", "name": "weather"}, + }, + ] + monkeypatch.setattr(crons, "_client", lambda: fake) + monkeypatch.setattr(crons, "_default_timezone", lambda: "Europe/London") + return crons, fake + + +def test_create_schedule_targets_scheduler(monkeypatch): + crons, fake = _patch_client(monkeypatch) + rec = crons.create_schedule( + name="weather", + schedule="*/10 * * * *", + prompt="search uk weather and summarize", + ) + assert rec["cron_id"] == "c-1" + kw = fake.crons.create.call_args.kwargs + assert kw["assistant_id"] == crons.SCHEDULER_GRAPH_ID + assert kw["schedule"] == "*/10 * * * *" + assert kw["input"] == { + "messages": [{"role": "user", "content": "search uk weather and summarize"}] + } + assert kw["metadata"]["run_kind"] == crons.SCHEDULED_RUN_KIND + assert kw["metadata"]["name"] == "weather" + assert kw["timezone"] == "Europe/London" + + +def test_list_schedules_uses_server_side_filter(monkeypatch): + crons, fake = _patch_client(monkeypatch) + out = crons.list_schedules() + assert [c["cron_id"] for c in out] == ["c-1"] + # Filtered server-side by run_kind metadata (no client filter); high limit so + # users with >10 schedules still see them all. + fake.crons.search.assert_called_once_with( + metadata={"run_kind": crons.SCHEDULED_RUN_KIND}, + limit=1000, + ) + + +def test_delete_and_set_enabled(monkeypatch): + crons, fake = _patch_client(monkeypatch) + crons.delete_schedule("c-1") + fake.crons.delete.assert_called_once_with("c-1") + crons.set_enabled("c-1", False) + assert fake.crons.update.call_args.kwargs["enabled"] is False + + +def test_run_now_dispatches_thread_then_run(monkeypatch): + crons, fake = _patch_client(monkeypatch) + fake.threads.create.return_value = {"thread_id": "t-1"} + fake.runs.create.return_value = {"run_id": "r-1"} + rec = crons.run_now("do the thing") + assert rec["run_id"] == "r-1" + fake.threads.create.assert_called_once_with(graph_id=crons.SCHEDULER_GRAPH_ID) + run_kw = fake.runs.create.call_args.kwargs + assert run_kw["thread_id"] == "t-1" + assert run_kw["assistant_id"] == crons.SCHEDULER_GRAPH_ID + assert run_kw["input"] == { + "messages": [{"role": "user", "content": "do the thing"}] + } + assert run_kw["metadata"]["run_kind"] == crons.SCHEDULED_RUN_KIND + assert run_kw["metadata"]["prompt"] == "do the thing" + + +# --------------------------------------------------------------------------- +# Task 3: scheduler registration tests +# --------------------------------------------------------------------------- + + +def test_scheduler_yaml_loads_as_async(): + """scheduler.yaml must define scheduler with async:true and a system_prompt.""" + from pathlib import Path + + import yaml + + import EvoScientist + + subagents_dir = Path(EvoScientist.__file__).parent / "subagents" + yaml_path = subagents_dir / "scheduler.yaml" + assert yaml_path.exists(), f"Missing {yaml_path}" + data = yaml.safe_load(yaml_path.read_text()) + assert "scheduler" in data, "Top-level key must be 'scheduler'" + spec = data["scheduler"] + assert spec.get("async") is True, "scheduler must have 'async: true'" + assert spec.get("system_prompt"), "scheduler must have a non-empty system_prompt" + # description drives task-tool routing — must be a non-empty string. + assert spec.get("description"), "scheduler must have a non-empty description" + + +def test_langgraph_json_registers_scheduler(): + """langgraph.json must map scheduler to the graphs:scheduler binding.""" + import json + from pathlib import Path + + import EvoScientist + + root = Path(EvoScientist.__file__).parent + manifest = json.loads((root / "langgraph_dev" / "langgraph.json").read_text()) + assert manifest["graphs"]["scheduler"] == ( + "EvoScientist.langgraph_dev.graphs:scheduler" + ) + + +def test_scheduler_graph_id_matches_registration(): + """SCHEDULER_GRAPH_ID must match both langgraph.json and scheduler.yaml top-level key.""" + import json + from pathlib import Path + + import yaml + + import EvoScientist + from EvoScientist.cron import schedule as crons + + root = Path(EvoScientist.__file__).parent + manifest = json.loads((root / "langgraph_dev" / "langgraph.json").read_text()) + assert crons.SCHEDULER_GRAPH_ID in manifest["graphs"], ( + f"SCHEDULER_GRAPH_ID={crons.SCHEDULER_GRAPH_ID!r} not found in langgraph.json graphs" + ) + spec = yaml.safe_load((root / "subagents" / "scheduler.yaml").read_text()) + assert crons.SCHEDULER_GRAPH_ID in spec, ( + f"SCHEDULER_GRAPH_ID={crons.SCHEDULER_GRAPH_ID!r} is not the top-level key in scheduler.yaml" + ) diff --git a/tests/test_langgraph_manager.py b/tests/test_langgraph_manager.py index 2eef08f..0a73113 100644 --- a/tests/test_langgraph_manager.py +++ b/tests/test_langgraph_manager.py @@ -240,10 +240,11 @@ class TestEnsureLanggraphDev: def test_skips_when_async_and_memory_workers_disabled( self, tmp_path, runtime_paths ): - """No background server is needed without async subagents or workers.""" + """No background server is needed without async subagents, workers, or scheduler.""" cfg = EvoScientistConfig() cfg.enable_async_subagents = False cfg.memory_workers_enabled = False + cfg.enable_scheduler = False # scheduler crons also require the backend cfg.langgraph_dev_port = 6174 cfg.langgraph_dev_file_persistence = True with ( @@ -265,6 +266,16 @@ class TestEnsureLanggraphDev: start.assert_not_called() assert manager.is_async_subagents_available() is False + def test_needs_langgraph_dev_for_scheduler_only(self): + """enable_scheduler alone requires the backend (crons fire inside it).""" + cfg = EvoScientistConfig() + cfg.enable_async_subagents = False + cfg.memory_workers_enabled = False + cfg.enable_scheduler = True + assert manager.needs_langgraph_dev(cfg) is True + cfg.enable_scheduler = False + assert manager.needs_langgraph_dev(cfg) is False + def test_reuses_existing_healthy_subprocess(self, tmp_path, runtime_paths): """When the subprocess is already running, no new Popen call.""" cfg = EvoScientistConfig() diff --git a/tests/test_profile_memory_middleware.py b/tests/test_profile_memory_middleware.py index fb9e3eb..2179978 100644 --- a/tests/test_profile_memory_middleware.py +++ b/tests/test_profile_memory_middleware.py @@ -84,7 +84,7 @@ def test_append_to_system_message_preserves_metadata(): response_metadata={"provider": "test"}, ) - updated = memory_module._append_to_system_message( + updated = memory_module.append_to_system_message( system_message, "memory context", ) diff --git a/tests/test_schedule_command.py b/tests/test_schedule_command.py new file mode 100644 index 0000000..d3ea9ac --- /dev/null +++ b/tests/test_schedule_command.py @@ -0,0 +1,236 @@ +"""Tests for the /schedule command.""" + +from unittest.mock import MagicMock, patch + +from tests.conftest import run_async as _run + + +def _ctx(): + from EvoScientist.commands.base import CommandContext + + ui = MagicMock() + return CommandContext(agent=None, thread_id="tid", ui=ui), ui + + +def test_list_when_backend_down(): + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + ctx, ui = _ctx() + with patch("EvoScientist.cron.schedule.is_available", return_value=False): + _run(ScheduleCommand().execute(ctx, ["list"])) + msgs = [c.args[0] for c in ui.append_system.call_args_list] + assert any("unavailable" in m.lower() for m in msgs) + + +def test_add_parses_five_field_cron_and_prompt(): + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + ctx, _ui = _ctx() + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch( + "EvoScientist.cron.schedule.create_schedule", + return_value={"cron_id": "c-9"}, + ) as mk, + ): + _run( + ScheduleCommand().execute( + ctx, ["add", "*/10", "*", "*", "*", "*", "search", "uk", "weather"] + ) + ) + kw = mk.call_args.kwargs + assert kw["schedule"] == "*/10 * * * *" + assert kw["prompt"] == "search uk weather" + + +def test_list_renders_table(): + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + ctx, ui = _ctx() + rows = [ + { + "cron_id": "c-1", + "schedule": "0 9 * * *", + "enabled": True, + "next_run_date": "2026-06-25T09:00:00+00:00", + "metadata": {"name": "daily"}, + } + ] + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch("EvoScientist.cron.schedule.list_schedules", return_value=rows), + ): + _run(ScheduleCommand().execute(ctx, ["list"])) + ui.mount_renderable.assert_called_once() + + +def test_add_parses_quoted_cron_and_prompt(): + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + ctx, _ui = _ctx() + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch( + "EvoScientist.cron.schedule.create_schedule", + return_value={"cron_id": "c-9"}, + ) as mk, + ): + _run( + ScheduleCommand().execute(ctx, ["add", "*/10 * * * *", "search uk weather"]) + ) + kw = mk.call_args.kwargs + assert kw["schedule"] == "*/10 * * * *" + assert kw["prompt"] == "search uk weather" + + +def test_run_with_matching_prefix_fires_matched_prompt(): + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + ctx, _ui = _ctx() + rows = [{"cron_id": "c-12345", "metadata": {"prompt": "do the thing"}}] + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch("EvoScientist.cron.schedule.list_schedules", return_value=rows), + patch( + "EvoScientist.cron.schedule.run_now", + return_value={"run_id": "r-1"}, + ) as rn, + ): + _run(ScheduleCommand().execute(ctx, ["run", "c-123"])) + rn.assert_called_once_with("do the thing") + + +def test_run_with_no_match_reports(): + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + ctx, ui = _ctx() + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch("EvoScientist.cron.schedule.list_schedules", return_value=[]), + patch("EvoScientist.cron.schedule.run_now") as rn, + ): + _run(ScheduleCommand().execute(ctx, ["run", "nope"])) + rn.assert_not_called() + msgs = [c.args[0] for c in ui.append_system.call_args_list] + assert any("No schedule matching" in m for m in msgs) + + +def test_pause_resume_set_enabled_with_resolved_id(): + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + rows = [{"cron_id": "c-abcdef", "metadata": {"name": "t"}}] + for sub, expected in (("pause", False), ("resume", True)): + ctx, _ui = _ctx() + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch("EvoScientist.cron.schedule.list_schedules", return_value=rows), + patch("EvoScientist.cron.schedule.set_enabled") as se, + ): + _run(ScheduleCommand().execute(ctx, [sub, "c-abc"])) + se.assert_called_once_with("c-abcdef", expected) + + +# --------------------------------------------------------------------------- +# B1: _list error handling when backend dies after is_available() +# --------------------------------------------------------------------------- + + +def test_list_error_shows_red_message_no_exception(): + """B1: list_schedules raising after is_available() shows a red error, not a traceback.""" + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + ctx, ui = _ctx() + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch( + "EvoScientist.cron.schedule.list_schedules", + side_effect=RuntimeError("backend gone"), + ), + ): + _run(ScheduleCommand().execute(ctx, ["list"])) + msgs = [c.args[0] for c in ui.append_system.call_args_list] + assert any("Error:" in m for m in msgs) + # Verify no exception escaped (test would have raised above otherwise) + + +# --------------------------------------------------------------------------- +# B2: ambiguous prefix detection in /schedule command +# --------------------------------------------------------------------------- + + +def test_remove_ambiguous_prefix_aborts_without_deleting(): + """B2: two crons sharing a prefix → ambiguity message, delete NOT called.""" + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + ctx, ui = _ctx() + rows = [ + {"cron_id": "abc-111", "metadata": {}}, + {"cron_id": "abc-222", "metadata": {}}, + ] + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch("EvoScientist.cron.schedule.list_schedules", return_value=rows), + patch("EvoScientist.cron.schedule.delete_schedule") as mk, + ): + _run(ScheduleCommand().execute(ctx, ["remove", "abc"])) + mk.assert_not_called() + msgs = [c.args[0] for c in ui.append_system.call_args_list] + assert any("Multiple" in m for m in msgs) + + +# --------------------------------------------------------------------------- +# B3: sanitized name from prompt +# --------------------------------------------------------------------------- + + +def test_remove_backend_error_shows_red_error_not_no_match(): + """FIX 1: list_schedules() crashing in _resolve → red 'Error:' message, not 'No schedule matching'.""" + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + ctx, ui = _ctx() + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch( + "EvoScientist.cron.schedule.list_schedules", + side_effect=RuntimeError("boom"), + ), + patch("EvoScientist.cron.schedule.delete_schedule") as mk, + ): + _run(ScheduleCommand().execute(ctx, ["remove", "abc"])) + mk.assert_not_called() + msgs = [c.args[0] for c in ui.append_system.call_args_list] + assert any("Error:" in m for m in msgs), f"Expected red Error: message, got: {msgs}" + assert not any("No schedule matching" in m for m in msgs), ( + f"Backend crash must not look like a name miss, got: {msgs}" + ) + # The error message must mention unavailable/backend to distinguish from other errors. + assert any("unavailable" in m.lower() or "backend" in m.lower() for m in msgs), ( + f"Error message should mention backend/unavailable, got: {msgs}" + ) + + +def test_add_name_sanitized_from_nasty_prompt(): + """B3: prompt with newline / slashes / special chars → clean kebab-case name.""" + import re + + from EvoScientist.commands.implementation.schedule import ScheduleCommand + + ctx, _ui = _ctx() + # prompt with newline, slashes, and dots + nasty_prompt = "Search /tmp/foo\nand summarize! Latest.Papers." + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch( + "EvoScientist.cron.schedule.create_schedule", + return_value={"cron_id": "c-x"}, + ) as mk, + ): + _run(ScheduleCommand().execute(ctx, ["add", "*/5 * * * *", nasty_prompt])) + name = mk.call_args.kwargs["name"] + # Must be non-empty, no spaces, no newlines, no slashes + assert name + assert " " not in name + assert "\n" not in name + assert "/" not in name + # Only safe chars: lowercase alphanumeric and hyphens + assert re.fullmatch(r"[a-z0-9][a-z0-9\-]*", name), f"Unexpected name: {name!r}" diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py new file mode 100644 index 0000000..96ff698 --- /dev/null +++ b/tests/test_scheduler.py @@ -0,0 +1,98 @@ +"""Tests for SchedulerMiddleware (scheduling guide + injection).""" + +from unittest.mock import MagicMock + + +def _mw(): + from EvoScientist.middleware.scheduler import SchedulerMiddleware + + m = SchedulerMiddleware() + m._cache = None + m._cache_at = 0.0 + return m + + +def test_schedules_block_lists_active_crons(monkeypatch): + from EvoScientist.cron import schedule as crons + + monkeypatch.setattr(crons, "is_available", lambda: True) + monkeypatch.setattr( + crons, + "list_schedules", + lambda: [ + { + "cron_id": "abc12345-xyz", + "schedule": "*/10 * * * *", + "enabled": True, + "metadata": {"name": "weather", "prompt": "search uk weather"}, + } + ], + ) + block = _mw()._schedules_block() + assert "" in block + assert "" in block + assert "weather" in block + assert "*/10 * * * *" in block + assert "abc12345" in block + + +def test_schedules_block_empty_when_unavailable(monkeypatch): + from EvoScientist.cron import schedule as crons + + monkeypatch.setattr(crons, "is_available", lambda: False) + assert _mw()._schedules_block() == "" + + +def test_schedules_block_empty_on_error(monkeypatch): + from EvoScientist.cron import schedule as crons + + monkeypatch.setattr(crons, "is_available", lambda: True) + + def _boom(): + raise RuntimeError("backend died") + + monkeypatch.setattr(crons, "list_schedules", _boom) + assert _mw()._schedules_block() == "" + + +def test_modify_request_injects_guide_even_with_no_schedules(monkeypatch): + # The static scheduling guide is injected on every call; the dynamic + # block is added only when there are active schedules. + m = _mw() + monkeypatch.setattr(m, "_cached_schedules_block", lambda: "") + req = MagicMock() + req.system_message = None + m.modify_request(req) + req.override.assert_called_once() + sysmsg = req.override.call_args.kwargs["system_message"] + text = " ".join(str(b) for b in sysmsg.content_blocks) + assert "scheduling_instructions" in text + # The guide mentions in prose; the actual dynamic block + # (identified by its closing tag) must be absent when there are no schedules. + assert "" not in text + + +def test_modify_request_injects_guide_and_schedules(monkeypatch): + m = _mw() + monkeypatch.setattr( + m, "_cached_schedules_block", lambda: "\nx\n" + ) + req = MagicMock() + req.system_message = None + m.modify_request(req) + req.override.assert_called_once() + sysmsg = req.override.call_args.kwargs["system_message"] + text = " ".join(str(b) for b in sysmsg.content_blocks) + assert "scheduling_instructions" in text + assert "" in text + + +def test_middleware_exposes_three_tools(): + """SchedulerMiddleware.tools must contain the three scheduling tools.""" + m = _mw() + tool_names = [t.name for t in m.tools] + assert tool_names == [ + "schedule_task", + "list_scheduled_tasks", + "cancel_scheduled_task", + ] diff --git a/tests/test_scheduler_tools.py b/tests/test_scheduler_tools.py new file mode 100644 index 0000000..91975e4 --- /dev/null +++ b/tests/test_scheduler_tools.py @@ -0,0 +1,124 @@ +"""Tests for the natural-language scheduling tools.""" + +from unittest.mock import patch + + +def test_schedule_task_translates_and_creates(): + from EvoScientist.middleware.scheduler import schedule_task + + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch( + "EvoScientist.cron.schedule.create_schedule", + return_value={"cron_id": "c-7"}, + ) as mk, + ): + out = schedule_task.invoke( + { + "name": "weather", + "cron": "*/10 * * * *", + "prompt": "search uk weather and summarize", + "timezone": "", + } + ) + assert "c-7" in out + assert mk.call_args.kwargs["schedule"] == "*/10 * * * *" + assert mk.call_args.kwargs["name"] == "weather" + + +def test_schedule_task_reports_backend_down(): + from EvoScientist.middleware.scheduler import schedule_task + + with patch("EvoScientist.cron.schedule.is_available", return_value=False): + out = schedule_task.invoke( + {"name": "x", "cron": "* * * * *", "prompt": "do x", "timezone": ""} + ) + assert "unavailable" in out.lower() + + +def test_cancel_scheduled_task(): + from EvoScientist.middleware.scheduler import cancel_scheduled_task + + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch("EvoScientist.cron.schedule.list_schedules", return_value=[]), + patch("EvoScientist.cron.schedule.delete_schedule") as mk, + ): + out = cancel_scheduled_task.invoke({"cron_id": "c-7"}) + mk.assert_not_called() + assert "No scheduled task matching" in out + + +def test_cancel_scheduled_task_prefix_match(): + from EvoScientist.middleware.scheduler import cancel_scheduled_task + + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch( + "EvoScientist.cron.schedule.list_schedules", + return_value=[{"cron_id": "c-7-abc"}], + ), + patch("EvoScientist.cron.schedule.delete_schedule") as mk, + ): + out = cancel_scheduled_task.invoke({"cron_id": "c-7"}) + mk.assert_called_once_with("c-7-abc") + assert "c-7-abc" in out + + +def test_list_scheduled_tasks_formats_rows(): + from EvoScientist.middleware.scheduler import list_scheduled_tasks + + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch( + "EvoScientist.cron.schedule.list_schedules", + return_value=[ + { + "cron_id": "c-1-xyz", + "schedule": "0 9 * * *", + "enabled": True, + "metadata": {"name": "daily"}, + } + ], + ), + ): + out = list_scheduled_tasks.invoke({}) + assert "daily" in out + assert "0 9 * * *" in out + + +# --------------------------------------------------------------------------- +# B2: ambiguous prefix in cancel_scheduled_task tool +# --------------------------------------------------------------------------- + + +def test_cancel_ambiguous_prefix_aborts_without_deleting(): + """B2: two crons sharing a prefix → returns ambiguity message, delete NOT called.""" + from EvoScientist.middleware.scheduler import cancel_scheduled_task + + rows = [{"cron_id": "abc-111"}, {"cron_id": "abc-222"}] + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch("EvoScientist.cron.schedule.list_schedules", return_value=rows), + patch("EvoScientist.cron.schedule.delete_schedule") as mk, + ): + out = cancel_scheduled_task.invoke({"cron_id": "abc"}) + mk.assert_not_called() + assert "Multiple" in out + + +def test_cancel_empty_cron_id_refuses_without_deleting(): + """Empty cron_id would match (and delete) the only cron — must refuse early.""" + from EvoScientist.middleware.scheduler import cancel_scheduled_task + + with ( + patch("EvoScientist.cron.schedule.is_available", return_value=True), + patch( + "EvoScientist.cron.schedule.list_schedules", + return_value=[{"cron_id": "only-one"}], + ), + patch("EvoScientist.cron.schedule.delete_schedule") as mk, + ): + out = cancel_scheduled_task.invoke({"cron_id": " "}) + mk.assert_not_called() + assert "Provide" in out diff --git a/uv.lock b/uv.lock index 830a85c..9214c81 100644 --- a/uv.lock +++ b/uv.lock @@ -972,6 +972,7 @@ dependencies = [ { name = "tavily-python" }, { name = "textual" }, { name = "typer" }, + { name = "tzlocal" }, ] [package.optional-dependencies] @@ -1093,6 +1094,7 @@ requires-dist = [ { name = "tavily-python", specifier = ">=0.7" }, { name = "textual", specifier = ">=8.0,<8.2.7" }, { name = "typer", specifier = ">=0.24" }, + { name = "tzlocal", specifier = ">=5.0" }, ] provides-extras = ["dev", "telegram", "discord", "slack", "wechat", "feishu", "qq", "stt", "oauth", "all-channels"]