diff --git a/src/leapflow/cli/commands/interactive.py b/src/leapflow/cli/commands/interactive.py index ccfbc10..c3f79c9 100644 --- a/src/leapflow/cli/commands/interactive.py +++ b/src/leapflow/cli/commands/interactive.py @@ -806,24 +806,6 @@ async def handle_input(text: str) -> None: _print_execution_result(result) return - if canonical == "arm": - from leapflow.cli.commands.scheduler import cmd_arm - - await cmd_arm( - ctx, - cmd_text.split()[1:] if len(cmd_text.split()) > 1 else [], - ) - return - - if canonical == "task": - from leapflow.cli.commands.scheduler import cmd_tasks - - await cmd_tasks( - ctx, - cmd_text.split()[1:] if len(cmd_text.split()) > 1 else [], - ) - return - # Registered engine command with no in-process fast path (e.g. # /board): route through the shared command_execute contract so a # recognized slash command never leaks into the LLM chat stream. diff --git a/src/leapflow/cli/commands/registry.py b/src/leapflow/cli/commands/registry.py index 59303fa..f700da2 100644 --- a/src/leapflow/cli/commands/registry.py +++ b/src/leapflow/cli/commands/registry.py @@ -141,10 +141,8 @@ def supports_runtime(self, runtime: CommandRuntime) -> bool: CommandDef("plugin enable", "Re-enable a disabled plugin (daemon mode, approval required)", "Skills & Tools", args_hint="", requires_host=True, effect=CommandEffect.SESSION, execution=CommandExecution.SHORT_OPERATION), CommandDef("plugin generate", "Generate and install a plugin from a natural-language description", "Skills & Tools", args_hint="[--preview|--dry-run|--id ] ", requires_host=True, effect=CommandEffect.SESSION, execution=CommandExecution.LONG_RUNNING), - # Scheduler - CommandDef("arm", "Schedule a skill for timed execution", "Scheduler", args_hint=" "), - CommandDef("task", "List scheduled tasks", "Scheduler"), - CommandDef("schedule", "List active scheduled tasks", "Scheduler", aliases=("schedule list",), args_hint="[list|status|history|cancel] ...", effect=CommandEffect.READ_ONLY, execution=CommandExecution.INSTANT), + # Scheduler — a single namespace; the bare command lists scheduled tasks. + CommandDef("schedule", "List scheduled tasks", "Scheduler", aliases=("schedule list",), args_hint="[list|status|history|cancel] ...", effect=CommandEffect.READ_ONLY, execution=CommandExecution.INSTANT), CommandDef("schedule status", "Show detailed status for one task", "Scheduler", args_hint="", effect=CommandEffect.READ_ONLY, execution=CommandExecution.INSTANT), CommandDef("schedule history", "Show recent execution log entries", "Scheduler", args_hint="[task_id]", effect=CommandEffect.READ_ONLY, execution=CommandExecution.INSTANT), CommandDef("schedule cancel", "Cancel/disable a scheduled task", "Scheduler", args_hint="", effect=CommandEffect.SESSION, execution=CommandExecution.SHORT_OPERATION), diff --git a/src/leapflow/cli/commands/scheduler.py b/src/leapflow/cli/commands/scheduler.py deleted file mode 100644 index e21e742..0000000 --- a/src/leapflow/cli/commands/scheduler.py +++ /dev/null @@ -1,289 +0,0 @@ -# Copyright (c) Alibaba, Inc. and its affiliates. -"""Scheduler CLI commands — arm tasks and manage scheduled execution. - -Provides ``leap arm`` and ``leap tasks`` subcommands for the interactive REPL. -""" - -from __future__ import annotations - -import time -from typing import TYPE_CHECKING, List - -if TYPE_CHECKING: - from leapflow.cli.context import Context - - -# --------------------------------------------------------------------------- -# Helpers -# --------------------------------------------------------------------------- - - -def _format_duration(seconds: float) -> str: - """Format a duration in seconds into a human-readable string.""" - if seconds <= 0: - return "now" - if seconds < 60: - return f"{int(seconds)}s" - if seconds < 3600: - return f"in {int(seconds / 60)}m" - if seconds < 86400: - h = int(seconds / 3600) - m = int((seconds % 3600) / 60) - return f"in {h}h{m}m" if m else f"in {h}h" - d = int(seconds / 86400) - return f"in {d}d" - - -def _format_trigger(task) -> str: - """Format trigger info for display.""" - if task.trigger_type == "interval": - sec = task.trigger_config.get("interval_seconds", 0) - if sec < 60: - return f"every {int(sec)}s" - if sec < 3600: - return f"every {int(sec / 60)}m" - if sec < 86400: - return f"every {int(sec / 3600)}h" - return f"every {int(sec / 86400)}d" - if task.trigger_type == "cron": - return task.trigger_config.get("expression", "cron") - if task.trigger_type == "event": - return f"event:{task.trigger_config.get('event_pattern', '?')}" - if task.trigger_type == "condition": - expr = task.trigger_config.get("expression", "?") - return f"cond:{expr[:20]}" - return task.trigger_type - - -def _get_coordinator(ctx: "Context"): - """Get or create a TaskCoordinator from context.""" - from leapflow.scheduler.coordinator import TaskCoordinator - from leapflow.scheduler.local_scheduler import LocalScheduler - from leapflow.scheduler.store import TaskStore - - # Use existing coordinator if available - if hasattr(ctx, "coordinator") and ctx.coordinator is not None: - return ctx.coordinator - - # Build one from settings - store = TaskStore(ctx.settings.duckdb_path) - - # Local scheduler with a simple skill executor - class _SimpleExecutor: - """Minimal skill executor for scheduled tasks.""" - - async def execute(self, skill_name: str, parameters: dict) -> dict: - # Try to execute via session if available - if ctx.session: - try: - result = await ctx.session.execute_skill(skill_name, params=parameters) - return {"ok": True, "output": str(result)[:200]} - except Exception as e: - return {"ok": False, "error": str(e)} - return {"ok": False, "error": "No session available"} - - local_scheduler = LocalScheduler( - store=store, - executor=_SimpleExecutor(), - tick_seconds=ctx.settings.scheduler_tick_seconds, - grace_seconds=ctx.settings.scheduler_grace_seconds, - ) - - # Cloud dispatcher (optional, only if compute backend available) - cloud_dispatcher = None - try: - from leapflow.scheduler.cloud_dispatcher import CloudDispatcher - from leapflow.scheduler.compute.modelscope_studio import ModelScopeStudioBackend - from leapflow.scheduler.worker_packager import WorkerPackager - - backend = ModelScopeStudioBackend() - packager = WorkerPackager() - cloud_dispatcher = CloudDispatcher(backend, packager) - except (ImportError, Exception): - pass # Cloud not available — local-only mode - - coordinator = TaskCoordinator( - store=store, - local_scheduler=local_scheduler, - cloud_dispatcher=cloud_dispatcher, - default_tier=ctx.settings.scheduler_default_tier, - ) - ctx.coordinator = coordinator - return coordinator - - -# --------------------------------------------------------------------------- -# Commands -# --------------------------------------------------------------------------- - - -async def cmd_arm(ctx: "Context", args: List[str]) -> int: - """leap arm --trigger "" [--local|--cloud] [--max-runs N] - - Arm a skill for scheduled execution. - """ - if not args or "--help" in args or "-h" in args: - print("Usage: arm --trigger \"\" [--local|--cloud] [--max-runs N]") - print() - print("Trigger formats:") - print(" 30m / 2h / 1d — interval") - print(" every 5m — interval") - print(" 0 9 * * * — cron (5 fields)") - print(" event:ci.passed — event") - print(" condition:expr > val — condition") - print() - print("Options:") - print(" --local Force local execution") - print(" --cloud Force cloud execution") - print(" --max-runs N Stop after N executions (-1 = unlimited)") - return 0 - - # Parse arguments - skill_name = args[0] - trigger_expr = "" - execution_tier = "auto" - max_runs = -1 - - i = 1 - while i < len(args): - if args[i] == "--trigger" and i + 1 < len(args): - trigger_expr = args[i + 1] - i += 2 - elif args[i] == "--local": - execution_tier = "local" - i += 1 - elif args[i] == "--cloud": - execution_tier = "cloud" - i += 1 - elif args[i] == "--max-runs" and i + 1 < len(args): - try: - max_runs = int(args[i + 1]) - except ValueError: - print(f" Error: --max-runs must be an integer, got '{args[i + 1]}'") - return 1 - i += 2 - else: - # If no --trigger flag, treat remaining as trigger expression - if not trigger_expr: - trigger_expr = " ".join(args[i:]) - break - i += 1 - - if not trigger_expr: - print(" Error: trigger expression required.") - print(" Usage: arm --trigger \"30m\"") - return 1 - - try: - coordinator = _get_coordinator(ctx) - task = await coordinator.arm( - skill_name=skill_name, - trigger_expr=trigger_expr, - execution_tier=execution_tier, - max_runs=max_runs, - ) - now = time.time() - next_in = _format_duration(task.next_due_at - now) - print(f" Armed: {task.task_id[:8]} skill={skill_name} tier={task.execution_tier} next={next_in}") - except ValueError as e: - print(f" Error: {e}") - return 1 - except RuntimeError as e: - print(f" Error: {e}") - return 1 - - return 0 - - -async def cmd_tasks(ctx: "Context", args: List[str]) -> int: - """leap tasks [status | logs | cancel ] - - Manage scheduled tasks. - """ - coordinator = _get_coordinator(ctx) - - if not args: - # List all tasks - tasks = await coordinator.list_tasks() - if not tasks: - print(" No scheduled tasks.") - return 0 - - now = time.time() - print(f" {'ID':<10} {'Skill':<20} {'Trigger':<16} {'Tier':<7} {'State':<10} Next Due") - print(f" {'-' * 80}") - for t in tasks: - tid = t.task_id[:8] - skill = t.skill_name[:18] - trigger = _format_trigger(t)[:14] - tier = t.execution_tier[:5] - state = t.state[:8] - if t.next_due_at > 0: - next_due = _format_duration(t.next_due_at - now) - else: - next_due = "-" - print(f" {tid:<10} {skill:<20} {trigger:<16} {tier:<7} {state:<10} {next_due}") - print(f"\n {len(tasks)} task(s).") - return 0 - - subcmd = args[0].lower() - - if subcmd == "status" and len(args) > 1: - task_id = args[1] - try: - status = await coordinator.status(task_id) - t = status.task - print(f" Task: {t.task_id}") - print(f" Skill: {t.skill_name}") - print(f" Trigger: {_format_trigger(t)}") - print(f" Tier: {t.execution_tier}") - print(f" State: {t.state}") - print(f" Runs: {t.run_count}" + (f" / {t.max_runs}" if t.max_runs > 0 else "")) - print(f" Running: {status.is_running}") - if t.next_due_at > 0: - now = time.time() - print(f" Next due: {_format_duration(t.next_due_at - now)}") - if t.cloud_worker_id: - print(f" Worker: {t.cloud_worker_id}") - except ValueError as e: - print(f" Error: {e}") - return 1 - return 0 - - if subcmd == "logs" and len(args) > 1: - task_id = args[1] - tail = 50 - if len(args) > 2: - try: - tail = int(args[2]) - except ValueError: - pass - try: - logs = await coordinator.logs(task_id, tail=tail) - if not logs: - print(" No logs available.") - else: - for line in logs: - print(f" {line}") - except ValueError as e: - print(f" Error: {e}") - return 1 - return 0 - - if subcmd == "cancel" and len(args) > 1: - task_id = args[1] - try: - await coordinator.cancel(task_id) - print(f" Cancelled: {task_id[:8]}") - except ValueError as e: - print(f" Error: {e}") - return 1 - return 0 - - # Unknown subcommand - print("Usage: tasks [status | logs | cancel ]") - print(" (no args) — List all scheduled tasks") - print(" status — Show detailed status") - print(" logs — Show recent logs") - print(" cancel — Cancel a task") - return 0 diff --git a/src/leapflow/cli/commands/slash_handlers.py b/src/leapflow/cli/commands/slash_handlers.py index f192bb1..286ebb7 100644 --- a/src/leapflow/cli/commands/slash_handlers.py +++ b/src/leapflow/cli/commands/slash_handlers.py @@ -2021,10 +2021,6 @@ async def command_execute( return await _execute_hub(ctx, name, args) if name == "run": return {"ok": True, "stream": True, "prompt": args} - if name == "arm": - return await _execute_scheduler_arm(ctx, args) - if name == "task": - return _execute_scheduler_task(ctx) if name == "schedule" or name.startswith("schedule "): sched_args = name[len("schedule"):].strip() if sched_args: @@ -2112,13 +2108,117 @@ def build_checkpoint_payload(ctx: "Context", args: str = "", session_id: str = " return {"ok": False, "message": f"Unknown checkpoint subcommand: {verb}. Use list or rollback."} +# Lifecycle states in which a task is still firing (or about to). Everything +# else -- paused, suspended, done, failed -- is inert and must not advertise a +# countdown, which is why an inert task renders its next run as an em dash +# rather than a misleading "now". +_SCHEDULE_ACTIVE_STATES: frozenset[str] = frozenset( + {"armed", "watching", "due", "confirming", "executing"} +) + + +def _humanize_interval(seconds: float) -> str: + """Render an interval in seconds as a compact 'every Nx' phrase.""" + sec = int(seconds) + if sec < 60: + return f"every {sec}s" + if sec < 3600: + return f"every {sec // 60}m" + if sec < 86400: + return f"every {sec // 3600}h" + return f"every {sec // 86400}d" + + +def _humanize_trigger(task: Any) -> str: + """Human-readable trigger description for the schedule list. + + Event triggers read 'on event' (with the pattern when known) rather than a + bare 'event', so the list explains *why* such a task has no next-run time. + """ + cfg = task.trigger_config if isinstance(task.trigger_config, dict) else {} + if task.trigger_type == "interval": + return _humanize_interval(cfg.get("interval_seconds", 0)) + if task.trigger_type == "cron": + return f"cron {cfg.get('expression', '?')}" + if task.trigger_type == "event": + pattern = cfg.get("event_pattern") or cfg.get("event") or "" + return f"on event: {pattern}" if pattern else "on event" + if task.trigger_type == "condition": + expr = str(cfg.get("expression", "?")) + return f"when {expr[:24]}" + return str(task.trigger_type) + + +def _humanize_next_run(task: Any, now: float, active: bool) -> str: + """Relative next-run string; inert tasks and event waits render as an em dash. + + An inert (paused/suspended/done/failed) task has no meaningful countdown, + and an event task fires on a signal rather than a clock, so both return the + em dash instead of a time that would imply an imminent clock-driven run. + """ + if not active or task.next_due_at <= 0: + return "—" + delta = task.next_due_at - now + if delta <= 0: + return "now" + if delta < 60: + return f"{int(delta)}s" + if delta < 3600: + return f"in {int(delta / 60)}m" + if delta < 86400: + return f"in {int(delta / 3600)}h" + return f"in {int(delta / 86400)}d" + + +def _schedule_task_entry(task: Any, now: float) -> dict[str, Any]: + """Build one structured schedule-list row from an ArmedTask.""" + params = task.parameters if isinstance(task.parameters, dict) else {} + mode = params.get("execution_mode") or "script" + active = task.state in _SCHEDULE_ACTIVE_STATES + runs = str(task.run_count) + (f"/{task.max_runs}" if task.max_runs > 0 else "") + return { + "task_id": task.task_id, + "short_id": task.task_id[:8], + "skill": task.skill_name, + "source": task.source, + "trigger": _humanize_trigger(task), + "next_run": _humanize_next_run(task, now, active), + "next_due_at": task.next_due_at, + "state": task.state, + "active": active, + "mode": mode, + "runs": runs, + } + + +def _schedule_list_text(entries: list[dict[str, Any]], summary: dict[str, Any]) -> str: + """Plain-text fallback of the schedule list for non-Rich consumers/logs. + + The TUI renders the structured ``view`` as a table; this text form keeps + skill names, state, and timing available wherever only ``message`` is read. + """ + inactive = summary.get("inactive") or {} + census = [f"{summary['active']} active"] + [ + f"{count} {state}" for state, count in sorted(inactive.items()) + ] + lines = [f"Scheduled tasks ({summary['total']}) — {' · '.join(census)}"] + for e in entries: + mode_suffix = " [agent]" if e["mode"] == "agent" else "" + lines.append( + f" {e['short_id']} {e['skill']} source={e['source']} {e['trigger']}" + f" next={e['next_run']} {e['state']}{mode_suffix}" + ) + return "\n".join(lines) + + def build_schedule_payload(ctx: "Context", args: str = "") -> dict[str, Any]: """Handle /schedule list, /schedule history, and /schedule cancel commands.""" from leapflow.scheduler.coordinator import TaskCoordinator from leapflow.scheduler.execution_log import DuckDBExecutionLogStore from leapflow.scheduler.store import TaskStore - # Resolve coordinator from context — same wiring as /arm and /task + # Resolve coordinator from context; fall back to a read-only store so the + # daemon (where no coordinator is wired) can still list tasks from DuckDB. coordinator: TaskCoordinator | None = getattr(ctx, "coordinator", None) task_store: TaskStore | None = None @@ -2138,49 +2238,38 @@ def build_schedule_payload(ctx: "Context", args: str = "") -> dict[str, Any]: # ── /schedule list (default) ───────────────────────────────────── if verb == "list" or not args.strip(): if task_store is None: - return {"ok": True, "message": "No scheduler active."} + return {"ok": True, "view": "schedule", "tasks": [], "message": "No scheduler active."} try: tasks = task_store.load_all() except Exception as exc: return {"ok": False, "message": f"Failed to load tasks: {exc}"} if not tasks: - return {"ok": True, "message": "No scheduled tasks."} + return {"ok": True, "view": "schedule", "tasks": [], "message": "No scheduled tasks."} import time as _time now = _time.time() - lines = ["Active scheduled tasks:"] - for t in tasks: - tid = t.task_id[:8] - trigger = t.trigger_type - if t.trigger_type == "interval": - sec = (t.trigger_config or {}).get("interval_seconds", 0) - if sec < 60: - trigger = f"every {int(sec)}s" - elif sec < 3600: - trigger = f"every {int(sec / 60)}m" - else: - trigger = f"every {int(sec / 3600)}h" - elif t.trigger_type == "cron": - trigger = (t.trigger_config or {}).get("expression", "cron") - if t.next_due_at > 0: - delta = t.next_due_at - now - if delta <= 0: - next_str = "now" - elif delta < 60: - next_str = f"{int(delta)}s" - elif delta < 3600: - next_str = f"{int(delta / 60)}m" - else: - next_str = f"{int(delta / 3600)}h" - else: - next_str = "-" - enabled = t.state not in ("suspended", "done", "failed", "paused") - params = t.parameters if isinstance(t.parameters, dict) else {} - mode = params.get("execution_mode") or "script" - lines.append( - f" {tid} skill={t.skill_name} mode={mode} trigger={trigger}" - f" next={next_str} enabled={enabled}" - ) - return {"ok": True, "message": "\n".join(lines)} + entries = [_schedule_task_entry(t, now) for t in tasks] + # Active tasks first, then soonest next run; inert tasks (no due time) + # sort last so the top of the list is always what fires next. + entries.sort(key=lambda e: ( + not e["active"], + e["next_due_at"] if e["next_due_at"] > 0 else float("inf"), + )) + inactive_counts: dict[str, int] = {} + for e in entries: + if not e["active"]: + inactive_counts[e["state"]] = inactive_counts.get(e["state"], 0) + 1 + summary = { + "total": len(entries), + "active": sum(1 for e in entries if e["active"]), + "inactive": inactive_counts, + } + return { + "ok": True, + "view": "schedule", + "tasks": entries, + "summary": summary, + "message": _schedule_list_text(entries, summary), + } # ── /schedule history [task_id] ────────────────────────────────── if verb == "history": @@ -3512,35 +3601,6 @@ async def _execute_hub(ctx: "Context", name: str, args: str) -> dict[str, Any]: return {"ok": False, "message": f"Hub command '{command}' is not yet implemented in this runtime."} -def _execute_scheduler_task(ctx: "Context") -> dict[str, Any]: - """Execute /task command.""" - scheduler = getattr(ctx, "scheduler", None) - if scheduler is None: - return {"ok": True, "view": "task", "tasks": [], "message": "No scheduler active."} - tasks = scheduler.list_tasks() if hasattr(scheduler, "list_tasks") else [] - entries = [ - {"name": t.name, "schedule": t.schedule, "next_run": str(getattr(t, "next_run", ""))} - for t in tasks - ] - return {"ok": True, "view": "task", "tasks": entries} - - -async def _execute_scheduler_arm(ctx: "Context", args: str) -> dict[str, Any]: - """Execute /arm command.""" - scheduler = getattr(ctx, "scheduler", None) - if scheduler is None: - return {"ok": False, "message": "Scheduler not active in this session."} - tokens = args.strip().split(None, 1) - if len(tokens) < 2: - return {"ok": False, "message": "Usage: /arm "} - skill_name, cron_expr = tokens - try: - task_id = await scheduler.arm(skill_name, cron_expr) - return {"ok": True, "message": f"Armed: {skill_name} → {cron_expr} (id={task_id})"} - except Exception as e: - return {"ok": False, "message": str(e)} - - def render_command_payload(console: "LeapConsole", payload: dict[str, Any]) -> None: """Render a generic command_execute result payload in the TUI.""" view = str(payload.get("view") or "") @@ -3575,8 +3635,8 @@ def render_command_payload(console: "LeapConsole", payload: dict[str, Any]) -> N if view == "skill_show": _render_skill_show_view(console, payload) return - if view == "task": - _render_task_view(console, payload) + if view == "schedule": + _render_schedule_view(console, payload) return if view == "dashboard": _render_dashboard_view(console, payload) @@ -3972,20 +4032,84 @@ def _render_skill_show_view(console: "LeapConsole", payload: dict[str, Any]) -> console.print(Panel(info, title=str(payload.get("name") or "Skill"), border_style="cyan")) -def _render_task_view(console: "LeapConsole", payload: dict[str, Any]) -> None: +# Glyph + Rich style per lifecycle state. The glyph carries the status at a +# glance (green ● = firing, yellow ⏸ = paused, red ✗ = failed), so the list is +# scannable without reading the label; an unknown state falls back to a neutral +# dot rather than crashing on a state added later in the scheduler. +_SCHEDULE_STATE_STYLE: dict[str, tuple[str, str]] = { + "armed": ("●", "green"), + "watching": ("◔", "green"), + "due": ("◆", "yellow"), + "confirming": ("◆", "yellow"), + "executing": ("▸", "cyan"), + "paused": ("⏸", "yellow"), + "suspended": ("⊘", "bright_black"), + "done": ("✓", "bright_black"), + "failed": ("✗", "red"), +} + + +def _render_schedule_view(console: "LeapConsole", payload: dict[str, Any]) -> None: + """Render /schedule list as an aligned, status-coloured table. + + Replaces the previous space-joined ``key=value`` line list, which never + aligned into columns and flattened a nine-state lifecycle into one + enabled=True/False boolean. Columns are chosen so the short id (the handle + every other /schedule subcommand takes) is never cropped. Source identifies + whether the system or a user created each task. The Mode column appears only + when at least one task runs in agent mode -- otherwise it is a constant + ``script`` that adds noise without information. + """ tasks = payload.get("tasks") or [] - msg = payload.get("message") - if msg: - console.system(str(msg)) - return if not tasks: - console.system("No scheduled tasks.") + msg = payload.get("message") + console.system(str(msg) if msg else "No scheduled tasks.") return + from rich.table import Table - table = Table(show_header=True, header_style="bold", border_style="dim") - table.add_column("Name", style="cyan") - table.add_column("Schedule") - table.add_column("Next Run") + from rich.text import Text + + show_mode = any(t.get("mode") == "agent" for t in tasks) + summary = payload.get("summary") or {} + inactive = summary.get("inactive") or {} + census = [f"{summary.get('active', 0)} active"] + [ + f"{count} {state}" for state, count in sorted(inactive.items()) + ] + + table = Table( + title=f"Scheduled tasks ({summary.get('total', len(tasks))})", + title_style="bold", + title_justify="left", + caption=" · ".join(census), + caption_style="dim", + caption_justify="left", + header_style="bold", + border_style="dim", + ) + table.add_column("ID", style="cyan", no_wrap=True) + # "Task" rather than "Skill": the row is a scheduled task and this column + # names what it runs; "Skill" was read as the LeapFlow skill catalog. + table.add_column("Task", overflow="fold") + table.add_column("Source", no_wrap=True) + table.add_column("Trigger") + table.add_column("Next run", justify="right") + table.add_column("Status") + if show_mode: + table.add_column("Mode") + for t in tasks: - table.add_row(str(t.get("name") or ""), str(t.get("schedule") or ""), str(t.get("next_run") or "")) + state = str(t.get("state") or "") + glyph, style = _SCHEDULE_STATE_STYLE.get(state, ("•", "white")) + row: list[Any] = [ + str(t.get("short_id") or ""), + str(t.get("skill") or ""), + str(t.get("source") or "user").title(), + str(t.get("trigger") or ""), + str(t.get("next_run") or ""), + Text(f"{glyph} {state}", style=style), + ] + if show_mode: + row.append("agent" if t.get("mode") == "agent" else "script") + table.add_row(*row) + console.print(table) diff --git a/src/leapflow/monitor/manager.py b/src/leapflow/monitor/manager.py index 38355ab..b07f76c 100644 --- a/src/leapflow/monitor/manager.py +++ b/src/leapflow/monitor/manager.py @@ -36,7 +36,7 @@ from leapflow.scheduler.local_scheduler import LocalScheduler from leapflow.scheduler.store import TaskStore from leapflow.scheduler.triggers.event import EventTrigger -from leapflow.scheduler.types import ArmedTask, TaskState +from leapflow.scheduler.types import ArmedTask, TaskSource, TaskState logger = logging.getLogger(__name__) @@ -237,6 +237,7 @@ async def arm_watch(self, spec: WatchSpec) -> WatchView: execution_tier="local", max_runs=spec.max_runs, parameters=spec.to_task_parameters(), + source=TaskSource.SYSTEM.value, ) # Backfill watch_id into parameters and stamp watch metadata so ticks # and listings can identify the row without re-deriving it. diff --git a/src/leapflow/plugins/tool_plugins/scheduler_tools.py b/src/leapflow/plugins/tool_plugins/scheduler_tools.py index abfbd2b..7ecfdcd 100644 --- a/src/leapflow/plugins/tool_plugins/scheduler_tools.py +++ b/src/leapflow/plugins/tool_plugins/scheduler_tools.py @@ -13,6 +13,7 @@ from typing import Any from leapflow.plugins.protocol import ToolMetadata +from leapflow.scheduler.types import TaskSource logger = logging.getLogger(__name__) @@ -295,6 +296,7 @@ async def _handle_create(self, **kwargs: Any) -> dict[str, Any]: trigger_expr=trigger_expr, parameters=parameters, max_retries=max_retries, + source=TaskSource.USER.value, ) except (ValueError, RuntimeError) as exc: return {"ok": False, "error": "arm_failed", "message": str(exc)} @@ -322,6 +324,7 @@ async def _handle_list(self, **kwargs: Any) -> dict[str, Any]: { "task_id": t.task_id, "skill_name": t.skill_name, + "source": t.source, "state": t.state, "execution_mode": _task_execution_mode(t), "trigger_type": t.trigger_type, diff --git a/src/leapflow/scheduler/__init__.py b/src/leapflow/scheduler/__init__.py index 0c9f936..29b80d8 100644 --- a/src/leapflow/scheduler/__init__.py +++ b/src/leapflow/scheduler/__init__.py @@ -10,6 +10,7 @@ from leapflow.scheduler.types import ( ArmedTask, TaskState, + TaskSource, ExecutionTier, TaskStatus, Trigger, @@ -29,6 +30,7 @@ # Types & enums "ArmedTask", "TaskState", + "TaskSource", "ExecutionTier", "TaskStatus", "Trigger", diff --git a/src/leapflow/scheduler/coordinator.py b/src/leapflow/scheduler/coordinator.py index e959ac7..52d6820 100644 --- a/src/leapflow/scheduler/coordinator.py +++ b/src/leapflow/scheduler/coordinator.py @@ -22,6 +22,7 @@ ArmedTask, ExecutionTier, SchedulerExecutionMode, + TaskSource, TaskState, TaskStatus, ) @@ -204,6 +205,7 @@ async def arm( context_snapshot: Optional[dict] = None, max_retries: Optional[int] = None, retry_backoff_s: Optional[float] = None, + source: str = TaskSource.USER.value, ) -> ArmedTask: """Create and register an armed task. @@ -213,8 +215,13 @@ async def arm( 4. Persist to store 5. Route to local_scheduler.register() or cloud_dispatcher.deploy() """ - # 1. Parse trigger + # 1. Parse trigger and validate the creation source at the boundary. trigger_type, trigger_config = parse_trigger_expression(trigger_expr) + try: + task_source = TaskSource(source).value + except ValueError as exc: + allowed = ", ".join(item.value for item in TaskSource) + raise ValueError(f"Invalid task source '{source}'; expected one of: {allowed}") from exc # 2. Decide tier tier = execution_tier if execution_tier != "auto" else self._default_tier @@ -253,6 +260,7 @@ async def arm( next_due_at=trigger.next_due_at, max_retries=effective_max_retries, retry_backoff_s=effective_backoff, + source=task_source, ) # 7. Persist diff --git a/src/leapflow/scheduler/store.py b/src/leapflow/scheduler/store.py index 32b427c..272663e 100644 --- a/src/leapflow/scheduler/store.py +++ b/src/leapflow/scheduler/store.py @@ -14,7 +14,7 @@ from pathlib import Path from typing import Any, List, Optional, Union -from leapflow.scheduler.types import ArmedTask +from leapflow.scheduler.types import ArmedTask, TaskSource from leapflow.storage.connection import ConnectionHolder, LocalConnectionHolder from leapflow.storage.write_buffer import execute_with_retry @@ -82,10 +82,12 @@ def _ensure_table(self) -> None: metadata TEXT DEFAULT '{}', max_retries INTEGER DEFAULT 0, retry_count INTEGER DEFAULT 0, - retry_backoff_s DOUBLE DEFAULT 60.0 + retry_backoff_s DOUBLE DEFAULT 60.0, + source TEXT DEFAULT 'user' ) """) self._migrate_retry_columns() + self._migrate_source_column() def _migrate_retry_columns(self) -> None: """Idempotent migration: add retry columns to pre-existing tables.""" @@ -101,6 +103,36 @@ def _migrate_retry_columns(self) -> None: except Exception: # noqa: BLE001 — column already exists pass + def _migrate_source_column(self) -> None: + """Add task source and backfill legacy monitor watches as system tasks. + + Before the source field existed, monitor watches were the only + system-created scheduler rows and were already durably tagged with + ``metadata.kind=watch``. Other legacy rows came from the user-facing + scheduler tool and therefore retain the column default of ``user``. + """ + columns = { + str(row[1]) + for row in self._con.execute("PRAGMA table_info('armed_tasks')").fetchall() + } + if "source" in columns: + return + + self._con.execute( + "ALTER TABLE armed_tasks ADD COLUMN source TEXT DEFAULT 'user'" + ) + rows = self._con.execute( + "SELECT task_id, metadata FROM armed_tasks", + ).fetchall() + for task_id, raw_metadata in rows: + if self._safe_json_loads(raw_metadata).get("kind") != "watch": + continue + execute_with_retry( + self._con, + "UPDATE armed_tasks SET source = ? WHERE task_id = ?", + [TaskSource.SYSTEM.value, task_id], + ) + # ------------------------------------------------------------------ # CRUD # ------------------------------------------------------------------ @@ -115,8 +147,8 @@ def save(self, task: ArmedTask) -> None: state, execution_tier, context_snapshot, confidence, created_at, next_due_at, last_run_at, run_count, max_runs, grace_seconds, parameters, cloud_worker_id, metadata, - max_retries, retry_count, retry_backoff_s - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + max_retries, retry_count, retry_backoff_s, source + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, [ task.task_id, @@ -139,6 +171,7 @@ def save(self, task: ArmedTask) -> None: task.max_retries, task.retry_count, task.retry_backoff_s, + task.source, ], ) @@ -283,6 +316,7 @@ def _row_to_task(self, row: tuple) -> ArmedTask: max_retries=row[17] if len(row) > 17 else 0, retry_count=row[18] if len(row) > 18 else 0, retry_backoff_s=row[19] if len(row) > 19 else 60.0, + source=row[20] if len(row) > 20 else TaskSource.USER.value, ) @staticmethod diff --git a/src/leapflow/scheduler/types.py b/src/leapflow/scheduler/types.py index b87198a..30aa34f 100644 --- a/src/leapflow/scheduler/types.py +++ b/src/leapflow/scheduler/types.py @@ -38,6 +38,13 @@ class ExecutionTier(str, Enum): AUTO = "auto" +class TaskSource(str, Enum): + """Who created a scheduled task.""" + + SYSTEM = "system" + USER = "user" + + class SchedulerExecutionMode(str, Enum): """How a scheduled task runs its skill on each trigger. @@ -98,6 +105,7 @@ class ArmedTask: max_retries: int = 0 retry_count: int = 0 retry_backoff_s: float = 60.0 + source: str = TaskSource.USER.value @dataclass diff --git a/tests/test_monitor_subsystem.py b/tests/test_monitor_subsystem.py index 5c8d04d..1dfb5d6 100644 --- a/tests/test_monitor_subsystem.py +++ b/tests/test_monitor_subsystem.py @@ -25,6 +25,7 @@ from leapflow.scheduler.coordinator import parse_trigger_expression from leapflow.scheduler.triggers import create_trigger from leapflow.scheduler.triggers.event import EventTrigger +from leapflow.scheduler.types import TaskSource from leapflow.storage.connection import LocalConnectionHolder @@ -140,6 +141,9 @@ async def test_manager_arm_list_and_state_transitions(tmp_path: Path) -> None: assert view.state == "armed" assert view.client_coupled is False assert view.to_dict()["client_coupled"] is False + stored = manager._task_store.load(view.watch_id) # noqa: SLF001 + assert stored is not None + assert stored.source == TaskSource.SYSTEM.value assert [v.watch_id for v in manager.list_watches()] == [view.watch_id] assert manager.has_active_watches() is True diff --git a/tests/test_scheduler_crud_retry.py b/tests/test_scheduler_crud_retry.py index 4fe1b26..f04893b 100644 --- a/tests/test_scheduler_crud_retry.py +++ b/tests/test_scheduler_crud_retry.py @@ -12,7 +12,7 @@ from leapflow.scheduler.coordinator import TaskCoordinator from leapflow.scheduler.local_scheduler import LocalScheduler from leapflow.scheduler.store import TaskStore -from leapflow.scheduler.types import ArmedTask, TaskState +from leapflow.scheduler.types import ArmedTask, TaskSource, TaskState # --------------------------------------------------------------------------- @@ -193,6 +193,31 @@ def test_retry_fields_roundtrip(self, tmp_store: TaskStore): assert loaded.retry_backoff_s == 30.0 +class TestTaskSource: + """Creation source is explicit and survives persistence.""" + + def test_armed_task_defaults_to_user(self) -> None: + task = ArmedTask( + skill_name="manual", + trigger_type="interval", + trigger_config={"interval_seconds": 60}, + ) + assert task.source == TaskSource.USER.value + + def test_system_source_roundtrip(self, tmp_store: TaskStore) -> None: + task = ArmedTask( + task_id="system_task", + skill_name="health", + trigger_type="interval", + trigger_config={"interval_seconds": 60}, + source=TaskSource.SYSTEM.value, + ) + tmp_store.save(task) + loaded = tmp_store.load(task.task_id) + assert loaded is not None + assert loaded.source == TaskSource.SYSTEM.value + + class TestRetryLogic: """Retry behavior in LocalScheduler._execute_task.""" @@ -438,6 +463,20 @@ async def test_arm_uses_config_defaults(self, tmp_store: TaskStore): task = await coordinator.arm("my_skill", "5m") assert task.max_retries == 5 assert task.retry_backoff_s == 30.0 + assert task.source == TaskSource.USER.value + + @pytest.mark.asyncio + async def test_arm_rejects_unknown_source(self, tmp_store: TaskStore) -> None: + local_sched = AsyncMock() + local_sched.register = AsyncMock() + coordinator = TaskCoordinator( + store=tmp_store, + local_scheduler=local_sched, + default_tier="local", + ) + + with pytest.raises(ValueError, match="Invalid task source"): + await coordinator.arm("my_skill", "5m", source="operator") @pytest.mark.asyncio async def test_arm_per_task_overrides_config(self, tmp_store: TaskStore): @@ -464,7 +503,7 @@ async def test_arm_per_task_overrides_config(self, tmp_store: TaskStore): class TestStoreMigration: - """Idempotent column migration for retry fields.""" + """Idempotent column migration for scheduler fields.""" def test_migration_is_idempotent(self, tmp_path: Path): """Creating TaskStore twice doesn't fail (columns already exist).""" @@ -490,6 +529,31 @@ def test_retry_columns_present_after_migration(self, tmp_store: TaskStore): assert loaded.max_retries == 7 assert loaded.retry_backoff_s == 120.0 + def test_source_migration_backfills_legacy_watches(self, tmp_path: Path) -> None: + """Legacy monitor watches become system tasks; other rows remain user tasks.""" + db = tmp_path / "source_migration.duckdb" + store = TaskStore(db) + store.save(ArmedTask( + task_id="legacy_watch", + skill_name="hardware", + trigger_type="interval", + trigger_config={"interval_seconds": 60}, + metadata={"kind": "watch"}, + )) + store.save(ArmedTask( + task_id="legacy_user", + skill_name="report", + trigger_type="interval", + trigger_config={"interval_seconds": 60}, + )) + store._con.execute("ALTER TABLE armed_tasks DROP COLUMN source") # noqa: SLF001 + store.close() + + migrated = TaskStore(db) + assert migrated.load("legacy_watch").source == TaskSource.SYSTEM.value + assert migrated.load("legacy_user").source == TaskSource.USER.value + migrated.close() + # --------------------------------------------------------------------------- # TUI payload builders diff --git a/tests/test_scheduler_tools.py b/tests/test_scheduler_tools.py index 82f05d7..bdbad2d 100644 --- a/tests/test_scheduler_tools.py +++ b/tests/test_scheduler_tools.py @@ -18,7 +18,7 @@ from leapflow.plugins.protocol import ToolMetadata, ToolPlugin from leapflow.plugins.tool_plugins.scheduler_tools import SchedulerToolsPlugin from leapflow.scheduler.local_scheduler import LocalScheduler -from leapflow.scheduler.types import ArmedTask +from leapflow.scheduler.types import ArmedTask, TaskSource # --------------------------------------------------------------------------- @@ -178,7 +178,11 @@ async def test_cancel_unbound(self, plugin: SchedulerToolsPlugin) -> None: class TestBoundHandlers: @pytest.mark.asyncio - async def test_create_success(self, bound_plugin: SchedulerToolsPlugin) -> None: + async def test_create_success( + self, + bound_plugin: SchedulerToolsPlugin, + mock_coordinator: AsyncMock, + ) -> None: result = await bound_plugin._handle_create( trigger_expression="30m", instruction="daily_report", @@ -186,6 +190,7 @@ async def test_create_success(self, bound_plugin: SchedulerToolsPlugin) -> None: assert result["ok"] is True assert "task_id" in result assert result["state"] == "armed" + assert mock_coordinator.arm.call_args.kwargs["source"] == TaskSource.USER.value @pytest.mark.asyncio async def test_create_missing_fields(self, bound_plugin: SchedulerToolsPlugin) -> None: @@ -217,6 +222,7 @@ async def test_list_success(self, bound_plugin: SchedulerToolsPlugin) -> None: assert result["ok"] is True assert result["count"] == 1 assert result["tasks"][0]["skill_name"] == "report" + assert result["tasks"][0]["source"] == TaskSource.USER.value @pytest.mark.asyncio async def test_status_success(self, bound_plugin: SchedulerToolsPlugin) -> None: @@ -577,9 +583,9 @@ async def test_status_includes_execution_mode(self) -> None: class TestScheduleListStatusFormat: - """Slash command output shows mode column and per-task status detail.""" + """Slash command output reports per-task mode and status detail.""" - def test_list_shows_execution_mode_column(self) -> None: + def test_list_reports_execution_mode_per_task(self) -> None: from leapflow.cli.commands.slash_handlers import build_schedule_payload ctx = MagicMock() @@ -589,6 +595,7 @@ def test_list_shows_execution_mode_column(self) -> None: skill_name="a", trigger_type="interval", trigger_config={"interval_seconds": 60}, task_id="aaaa1111", state="armed", parameters={"execution_mode": "agent"}, + source=TaskSource.SYSTEM.value, ), ArmedTask( skill_name="b", trigger_type="interval", @@ -601,8 +608,90 @@ def test_list_shows_execution_mode_column(self) -> None: ctx.coordinator = coordinator result = build_schedule_payload(ctx, "list") assert result["ok"] is True - assert "mode=agent" in result["message"] - assert "mode=script" in result["message"] + # Structured view drives the table; mode is per-task, not a constant + # "mode=script" repeated on every row. + assert result["view"] == "schedule" + modes = {t["skill"]: t["mode"] for t in result["tasks"]} + assert modes == {"a": "agent", "b": "script"} + sources = {t["skill"]: t["source"] for t in result["tasks"]} + assert sources == {"a": "system", "b": "user"} + assert result["summary"]["total"] == 2 + assert result["summary"]["active"] == 2 + + def test_schedule_table_renders_source_column(self) -> None: + from rich.console import Console + + from leapflow.cli.commands.slash_handlers import _render_schedule_view + + console = Console(record=True, width=120) + _render_schedule_view(console, { + "tasks": [ + { + "short_id": "aaaa1111", "skill": "health", + "source": "system", "trigger": "every 5m", + "next_run": "in 3m", "state": "armed", "mode": "script", + }, + { + "short_id": "bbbb2222", "skill": "report", + "source": "user", "trigger": "cron 0 9 * * *", + "next_run": "in 2h", "state": "armed", "mode": "script", + }, + ], + "summary": {"total": 2, "active": 2, "inactive": {}}, + }) + rendered = console.export_text() + assert "Source" in rendered + assert "System" in rendered + assert "User" in rendered + + def test_list_inert_task_has_no_countdown(self) -> None: + """A paused task must not advertise a next-run time (regression). + + The old key=value list showed ``next=now enabled=False`` for a paused, + overdue task -- a disabled task claiming it fires immediately. An inert + task now renders its next run as an em dash. + """ + from leapflow.cli.commands.slash_handlers import build_schedule_payload + + ctx = MagicMock() + task_store = MagicMock() + task_store.load_all.return_value = [ + ArmedTask( + skill_name="paused_one", trigger_type="interval", + trigger_config={"interval_seconds": 60}, task_id="cccc3333", + state="paused", next_due_at=time.time() - 5, + ), + ] + coordinator = MagicMock() + coordinator._store = task_store + ctx.coordinator = coordinator + result = build_schedule_payload(ctx, "list") + assert result["ok"] is True + entry = result["tasks"][0] + assert entry["active"] is False + assert entry["next_run"] == "\u2014" + assert result["summary"]["active"] == 0 + assert result["summary"]["inactive"] == {"paused": 1} + + def test_list_event_trigger_has_no_countdown(self) -> None: + from leapflow.cli.commands.slash_handlers import build_schedule_payload + + ctx = MagicMock() + task_store = MagicMock() + task_store.load_all.return_value = [ + ArmedTask( + skill_name="on_signal", trigger_type="event", + trigger_config={"event_pattern": "ci.passed"}, task_id="dddd4444", + state="armed", next_due_at=0.0, + ), + ] + coordinator = MagicMock() + coordinator._store = task_store + ctx.coordinator = coordinator + result = build_schedule_payload(ctx, "list") + entry = result["tasks"][0] + assert entry["trigger"] == "on event: ci.passed" + assert entry["next_run"] == "\u2014" def test_status_missing_task_id(self) -> None: from leapflow.cli.commands.slash_handlers import build_schedule_payload diff --git a/tests/test_slash_command_router.py b/tests/test_slash_command_router.py index 9663b42..3434e4a 100644 --- a/tests/test_slash_command_router.py +++ b/tests/test_slash_command_router.py @@ -53,8 +53,8 @@ def test_command_router_all_commands_supported_in_daemon() -> None: "/skill show demo", "/hub search test", "/gateway", - "/arm test_skill 0 * * * *", - "/task", + "/schedule", + "/schedule list", ): inv = daemon_router.parse(cmd_text) assert inv is not None, f"parse failed for {cmd_text}"