diff --git a/src/aipass/daemon/.seedgo/bypass.json b/src/aipass/daemon/.seedgo/bypass.json index 5e6468e1..8b79759f 100644 --- a/src/aipass/daemon/.seedgo/bypass.json +++ b/src/aipass/daemon/.seedgo/bypass.json @@ -107,6 +107,18 @@ "reason": "Authorized cross-branch wake_branch import (DPLAN-0204 §2.8 path A). ai_mail exposes wake_branch only via its handler — ai_mail's own modules and @trigger import it identically; no module entry point exists to use. Root fix (an ai_mail module wrapper) is tracked separately.", "pattern": "Handler imported directly" }, + { + "file": "apps/handlers/schedule/telegram_notifier.py", + "standard": "encapsulation", + "reason": "Authorized cross-branch import (TDPLAN-0008 contract). daemon emits lifecycle pings via skills telegram notifier — the only send path; no module wrapper exists.", + "pattern": "Handler imported directly" + }, + { + "file": "apps/handlers/schedule/telegram_notifier.py", + "standard": "handlers", + "reason": "Authorized cross-branch import (TDPLAN-0008 contract). send_telegram_notification lives in @skills handler — daemon's notifier wraps it fail-soft.", + "pattern": "Cross-handler imports" + }, { "file": "apps/modules/run.py", "standard": "introspection", @@ -149,6 +161,30 @@ "reason": "Test file — lives in tests/ by convention, not in 3-layer apps/ structure", "pattern": "File not in standard 3-layer structure" }, + { + "file": "tests/test_run_module.py", + "standard": "architecture", + "reason": "Test file — lives in tests/ by convention, not in 3-layer apps/ structure", + "pattern": "File not in standard 3-layer structure" + }, + { + "file": "tests/test_run_module.py", + "standard": "documentation", + "reason": "Test methods are self-documenting via descriptive names — pytest convention", + "pattern": "missing docstrings" + }, + { + "file": "tests/test_scheduler_bot.py", + "standard": "architecture", + "reason": "Test file — lives in tests/ by convention, not in 3-layer apps/ structure", + "pattern": "File not in standard 3-layer structure" + }, + { + "file": "tests/test_scheduler_bot.py", + "standard": "encapsulation", + "reason": "Test file imports handlers directly to unit-test them — standard pytest pattern", + "pattern": "Handler imported directly" + }, { "file": "apps/modules/timer_install.py", "standard": "introspection", diff --git a/src/aipass/daemon/README.md b/src/aipass/daemon/README.md index 7d189162..4f7a585b 100644 --- a/src/aipass/daemon/README.md +++ b/src/aipass/daemon/README.md @@ -110,9 +110,10 @@ drone @daemon --help | Module | Description | Status | |--------|-------------|--------| | `update` | Status digest of DAEMON activity | *(partial)* — reads inbox/sessions but data_loader paths return empty | -| `schedule` | Fire-and-forget scheduled follow-ups and task management | Operational | +| `queue` | Unified job queue view — Rich table or `--json` (frozen schema for @skills bot) | Operational | +| `schedule` | *(retired)* Fire-and-forget follow-ups — superseded by `.daemon/schedule.json` | Retired | | `activity_report` | Branch activity reports: `activity`, `activity-report`, `branch-health` | Operational | -| `actions` | Action registry CLI — list, toggle, info, set reminder, set schedule, migrate | Operational | +| `actions` | *(retired)* Action registry — superseded by `.daemon/schedule.json` | Retired | | `scheduler_ops` | Scheduler cron operations facade for scheduler_cron.py | Operational | | `wakeup_ops` | Wake-up cron operations facade for daemon_wakeup.py | Operational | | `timer_install` | Idempotent systemd user timer installer for daemon scheduler | Operational | diff --git a/src/aipass/daemon/apps/daemon.py b/src/aipass/daemon/apps/daemon.py index a921048c..43a17c02 100644 --- a/src/aipass/daemon/apps/daemon.py +++ b/src/aipass/daemon/apps/daemon.py @@ -24,7 +24,7 @@ from aipass.prax.apps.modules.logger import system_logger as logger # Console from aipass.cli.apps.modules import console, error from aipass.daemon.apps.handlers.json import json_handler -from aipass.daemon.apps.modules import update, schedule, activity_report, actions, run, timer_install +from aipass.daemon.apps.modules import update, schedule, activity_report, actions, run, timer_install, queue def _header(text): @@ -46,7 +46,7 @@ def get_modules() -> List[Any]: List of module objects with handle_command function """ modules = [] - for mod in [update, schedule, activity_report, actions, run, timer_install]: + for mod in [update, schedule, activity_report, actions, run, timer_install, queue]: if hasattr(mod, "handle_command"): modules.append(mod) return modules @@ -136,14 +136,15 @@ def print_help(modules: List[Any]): # Show actual routable commands, not module names _COMMAND_HELP = [ ("update", "Returns digest of DAEMON activity for check-ins."), - ("schedule", "CLI interface for fire-and-forget scheduled follow-ups."), + ("queue", "Unified job queue view (--json for frozen schema)."), ("activity", "Quick 24-hour activity summary."), ("activity-report", "Full detailed activity report (--json for raw)."), ("branch-health", "Single branch deep dive (e.g., branch-health DAEMON)."), - ("actions", "CLI interface for the numbered action registry."), ("run", "One scheduler tick: discover .daemon/ jobs, fire due ones."), ("install-timer", "Install + enable daemon-tick systemd user timer (~2 min)."), ("uninstall-timer", "Stop + remove daemon-tick systemd user timer."), + ("schedule", "(retired) Use .daemon/schedule.json — see run --help."), + ("actions", "(retired) Use .daemon/schedule.json — see run --help."), ] for cmd_name, desc in _COMMAND_HELP: diff --git a/src/aipass/daemon/apps/handlers/actions/actions_registry.py b/src/aipass/daemon/apps/handlers/actions/actions_registry.py deleted file mode 100644 index 8581484d..00000000 --- a/src/aipass/daemon/apps/handlers/actions/actions_registry.py +++ /dev/null @@ -1,550 +0,0 @@ -# =================== AIPass ==================== -# Name: actions_registry.py -# Description: Numbered Action Registry -# Version: 1.0.0 -# Created: 2026-03-02 -# Modified: 2026-03-02 -# ============================================= - -""" -Numbered Action Registry — DPLAN-043 - -Central registry for all scheduled actions. Each action gets a sequential -numeric ID (0001, 0002, ...) and can be individually toggled on/off. - -Replaces the old all-or-nothing daemon + kill switch model with granular -per-action control. - -Action types: - - plugin: Backed by a plugin file in apps/plugins/ (migrated from existing system) - - schedule: Custom recurring action (dispatches via wake.py) - - reminder: One-shot action that auto-completes after firing -""" - -import json -from datetime import datetime, timedelta -from pathlib import Path -from typing import Optional - -from aipass.prax import logger -from aipass.daemon.apps.handlers.json import json_handler - -# logger imported from aipass.prax - -# Paths -_DAEMON_ROOT = Path(__file__).resolve().parents[3] # src/aipass/daemon/ -REGISTRY_FILE = _DAEMON_ROOT / "daemon_json" / "actions_registry.json" -PLUGINS_DIR = _DAEMON_ROOT / "apps" / "plugins" - - -def _empty_registry() -> dict: - """Return a fresh empty registry structure (avoids shared mutable state).""" - return {"version": 1, "next_id": 1, "actions": []} - - -# ============================================= -# STORAGE -# ============================================= - - -def load_registry() -> dict: - """Load the actions registry from disk. Returns empty registry if missing.""" - if not REGISTRY_FILE.exists(): - return _empty_registry().copy() - try: - with open(REGISTRY_FILE, "r", encoding="utf-8") as f: - data = json.load(f) - if "actions" not in data: - data["actions"] = [] - if "next_id" not in data: - data["next_id"] = 1 - return data - except (json.JSONDecodeError, OSError) as e: - logger.error("[actions_registry] Failed to load: %s", e) - return _empty_registry().copy() - - -def save_registry(data: dict) -> bool: - """Save the actions registry to disk. Returns True on success.""" - try: - REGISTRY_FILE.parent.mkdir(parents=True, exist_ok=True) - with open(REGISTRY_FILE, "w", encoding="utf-8") as f: - json.dump(data, f, indent=2) - f.write("\n") - return True - except OSError as e: - logger.error("[actions_registry] Failed to save: %s", e) - return False - - -# ============================================= -# ID GENERATION -# ============================================= - - -def _get_next_id(registry: dict) -> str: - """Get next sequential ID as 4-digit string. Advances next_id.""" - next_num = registry.get("next_id", 1) - action_id = f"{next_num:04d}" - registry["next_id"] = next_num + 1 - return action_id - - -# ============================================= -# CRUD OPERATIONS -# ============================================= - - -def create_action( - name: str, - action_type: str, - schedule_type: str, - target_branch: str = "", - prompt: str = "", - time: Optional[str] = None, - interval_minutes: Optional[int] = None, - due_date: Optional[str] = None, - fresh: bool = True, - max_turns: int = 50, - enabled: bool = True, - self_dispatch: bool = False, - plugin_file: Optional[str] = None, -) -> dict: - """ - Create a new action and save to registry. - - Args: - name: Human-readable action name (e.g., "daily_audit") - action_type: "plugin" | "schedule" | "reminder" - schedule_type: "daily" | "hourly" | "interval" | "once" - target_branch: Target branch email (e.g., "@seedgo") - prompt: What the dispatched agent should do - time: For daily: "HH:MM", for hourly: "MM" - interval_minutes: For interval schedule type - due_date: For reminder (once) type, ISO date string - fresh: Start fresh session (True) or resume (False) - max_turns: Max agent turns - enabled: Active by default - self_dispatch: Plugin handles its own dispatch - plugin_file: Plugin filename (without .py) for plugin-backed actions - - Returns: - The created action dict - """ - registry = load_registry() - action_id = _get_next_id(registry) - - action = { - "id": action_id, - "name": name, - "type": action_type, - "schedule_type": schedule_type, - "time": time, - "interval_minutes": interval_minutes, - "due_date": due_date, - "target_branch": target_branch, - "prompt": prompt, - "fresh": fresh, - "max_turns": max_turns, - "enabled": enabled, - "self_dispatch": self_dispatch, - "plugin_file": plugin_file, - "last_run": None, - "next_run": None, - "created": datetime.now().isoformat(), - "completed": None, - } - - registry["actions"].append(action) - save_registry(registry) - - json_handler.log_operation("action_registry_modified", {"action": name}) - logger.info("[actions_registry] Created action %s: %s (%s)", action_id, name, action_type) - return action - - -def get_action(action_id: str) -> Optional[dict]: - """Get a single action by ID. Returns None if not found.""" - registry = load_registry() - for action in registry["actions"]: - if action["id"] == action_id: - return action - return None - - -def list_actions(include_completed: bool = False) -> list: - """ - List all actions. - - Args: - include_completed: If True, include completed reminders. - - Returns: - List of action dicts. - """ - registry = load_registry() - actions = registry["actions"] - if not include_completed: - actions = [a for a in actions if a.get("completed") is None] - return actions - - -def toggle_action(action_id: str, enabled: bool) -> bool: - """Toggle an action on or off. Returns True if found and updated.""" - registry = load_registry() - for action in registry["actions"]: - if action["id"] == action_id: - action["enabled"] = enabled - save_registry(registry) - state = "enabled" if enabled else "disabled" - logger.info("[actions_registry] Action %s %s: %s", action_id, state, action["name"]) - return True - return False - - -def delete_action(action_id: str) -> bool: - """Delete an action by ID. Returns True if found and removed.""" - registry = load_registry() - original_len = len(registry["actions"]) - registry["actions"] = [a for a in registry["actions"] if a["id"] != action_id] - if len(registry["actions"]) < original_len: - save_registry(registry) - logger.info("[actions_registry] Deleted action %s", action_id) - return True - return False - - -def update_last_run(action_id: str, timestamp: Optional[str] = None) -> bool: - """Update last_run timestamp for an action. Returns True if found.""" - if timestamp is None: - timestamp = datetime.now().isoformat() - registry = load_registry() - for action in registry["actions"]: - if action["id"] == action_id: - action["last_run"] = timestamp - action["next_run"] = calc_next_run(action) - save_registry(registry) - return True - return False - - -def mark_reminder_completed(action_id: str) -> bool: - """Mark a reminder as completed (one-shot). Returns True if found.""" - registry = load_registry() - for action in registry["actions"]: - if action["id"] == action_id: - action["completed"] = datetime.now().isoformat() - action["enabled"] = False - save_registry(registry) - logger.info("[actions_registry] Reminder %s completed: %s", action_id, action["name"]) - return True - return False - - -# ============================================= -# DUE CHECKING -# ============================================= - - -def _already_ran_today(action: dict, now: datetime) -> bool: - """Check if a daily action already ran today.""" - last_run = action.get("last_run") - if not last_run: - return False - try: - last_dt = datetime.fromisoformat(last_run) - return last_dt.date() == now.date() - except (ValueError, TypeError) as e: - logger.info("[actions_registry] Daily last_run parse failed: %s", e) - return False - - -def _already_ran_this_hour(action: dict, now: datetime) -> bool: - """Check if an hourly action already ran this hour.""" - last_run = action.get("last_run") - if not last_run: - return False - try: - last_dt = datetime.fromisoformat(last_run) - return last_dt.hour == now.hour and last_dt.date() == now.date() - except (ValueError, TypeError) as e: - logger.info("[actions_registry] Hourly last_run parse failed: %s", e) - return False - - -def _is_daily_due(action: dict, now: datetime) -> bool: - """Check if a daily action is due.""" - target_time = action.get("time", "00:00") - try: - target_h, target_m = map(int, target_time.split(":")) - except (ValueError, AttributeError) as e: - logger.info("[actions_registry] Daily time parse failed for %r: %s", target_time, e) - return False - current_minutes = now.hour * 60 + now.minute - target_minutes = target_h * 60 + target_m - minutes_diff = abs(current_minutes - target_minutes) - minutes_diff = min(minutes_diff, 1440 - minutes_diff) - if minutes_diff > 15: - return False - return not _already_ran_today(action, now) - - -def _is_hourly_due(action: dict, now: datetime) -> bool: - """Check if an hourly action is due.""" - target_m_str = action.get("time", "0") - try: - target_m = int(target_m_str) - except (ValueError, TypeError) as e: - logger.info("[actions_registry] Hourly time parse failed for %r: %s", target_m_str, e) - return False - minutes_diff = abs(now.minute - target_m) - minutes_diff = min(minutes_diff, 60 - minutes_diff) - if minutes_diff > 15: - return False - return not _already_ran_this_hour(action, now) - - -def _is_interval_due(action: dict, now: datetime) -> bool: - """Check if an interval action is due.""" - interval = action.get("interval_minutes", 60) - last_run = action.get("last_run") - if not last_run: - return True - try: - last_dt = datetime.fromisoformat(last_run) - elapsed = (now - last_dt).total_seconds() / 60 - return elapsed >= interval - except (ValueError, TypeError) as e: - logger.info("[actions_registry] Interval last_run parse failed: %s", e) - return True - - -def _is_once_due(action: dict, now: datetime) -> bool: - """Check if a one-shot reminder action is due.""" - due_date = action.get("due_date") - if not due_date: - return False - try: - due_dt = ( - datetime.fromisoformat(due_date).date() - if "T" in due_date - else datetime.strptime(due_date, "%Y-%m-%d").date() - ) - return now.date() >= due_dt - except (ValueError, TypeError) as e: - logger.info("[actions_registry] Once due_date parse failed for %r: %s", due_date, e) - return False - - -def is_action_due(action: dict) -> bool: - """ - Check if an action should run now. - - For daily: matches current hour:minute, hasn't run today - For hourly: matches current minute, hasn't run this hour - For interval: enough time has elapsed since last run - For once (reminder): due_date <= today, not completed - """ - if not action.get("enabled", False): - return False - - if action.get("completed"): - return False - - now = datetime.now() - schedule_type = action.get("schedule_type", "") - - _due_checkers = { - "daily": _is_daily_due, - "hourly": _is_hourly_due, - "interval": _is_interval_due, - "once": _is_once_due, - } - checker = _due_checkers.get(schedule_type) - if checker is None: - return False - return checker(action, now) - - -def _calc_next_daily(action: dict, now: datetime) -> Optional[str]: - """Calculate next run for a daily action.""" - target_time = action.get("time", "00:00") - try: - target_h, target_m = map(int, target_time.split(":")) - except (ValueError, AttributeError) as e: - logger.info("[actions_registry] calc_next_run daily time parse failed: %s", e) - return None - next_dt = now.replace(hour=target_h, minute=target_m, second=0, microsecond=0) - if next_dt <= now: - next_dt += timedelta(days=1) - return next_dt.isoformat() - - -def _calc_next_hourly(action: dict, now: datetime) -> Optional[str]: - """Calculate next run for an hourly action.""" - target_m_str = action.get("time", "0") - try: - target_m = int(target_m_str) - except (ValueError, TypeError) as e: - logger.info("[actions_registry] calc_next_run hourly time parse failed: %s", e) - return None - next_dt = now.replace(minute=target_m, second=0, microsecond=0) - if next_dt <= now: - next_dt += timedelta(hours=1) - return next_dt.isoformat() - - -def _calc_next_interval(action: dict, now: datetime) -> Optional[str]: - """Calculate next run for an interval action.""" - interval = action.get("interval_minutes", 60) - last_run = action.get("last_run") - if not last_run: - return now.isoformat() - try: - last_dt = datetime.fromisoformat(last_run) - return (last_dt + timedelta(minutes=interval)).isoformat() - except (ValueError, TypeError) as e: - logger.info("[actions_registry] calc_next_run interval last_run parse failed: %s", e) - return now.isoformat() - - -def calc_next_run(action: dict) -> Optional[str]: - """Calculate the next run time for an action. Returns ISO string or None.""" - now = datetime.now() - schedule_type = action.get("schedule_type", "") - - if schedule_type == "daily": - return _calc_next_daily(action, now) - if schedule_type == "hourly": - return _calc_next_hourly(action, now) - if schedule_type == "interval": - return _calc_next_interval(action, now) - if schedule_type == "once": - due_date = action.get("due_date") - if due_date and not action.get("completed"): - return due_date - return None - return None - - -def _next_due_interval(action: dict) -> str: - """Human-readable next due string for interval actions.""" - interval = action.get("interval_minutes", 60) - last_run = action.get("last_run") - if not last_run: - return "now" - try: - last_dt = datetime.fromisoformat(last_run) - next_dt = last_dt + timedelta(minutes=interval) - if next_dt <= datetime.now(): - return "now" - return next_dt.strftime("%H:%M") - except (ValueError, TypeError) as e: - logger.info("[actions_registry] next_due_str interval last_run parse failed: %s", e) - return "now" - - -def next_due_str(action: dict) -> str: - """Human-readable next due string for display.""" - schedule_type = action.get("schedule_type", "") - - if schedule_type == "daily": - return f"daily @ {action.get('time', '00:00')}" - if schedule_type == "hourly": - m = action.get("time", "0") - return f"hourly @ :{int(m):02d}" - if schedule_type == "interval": - return _next_due_interval(action) - if schedule_type == "once": - return action.get("due_date", "unknown") - return "unknown" - - -# ============================================= -# PLUGIN MIGRATION -# ============================================= - - -def migrate_plugins() -> int: - """ - Scan plugins/ directory and auto-register any plugins not yet in the registry. - - Maps PLUGIN_CONFIG fields to action fields. Preserves last_run timestamps - from .last_run.json. - - Returns: - Number of newly migrated plugins. - """ - registry = load_registry() - existing_plugins = {a["plugin_file"] for a in registry["actions"] if a.get("plugin_file")} - - # Load last_run data for timestamp preservation - last_run_file = PLUGINS_DIR / ".last_run.json" - last_run_map = {} - if last_run_file.exists(): - try: - last_run_map = json.loads(last_run_file.read_text(encoding="utf-8")) - except (json.JSONDecodeError, OSError) as e: - logger.warning("[actions_registry] Failed to load last_run.json: %s", e) - - # Discover plugins - migrated = 0 - for plugin_path in sorted(PLUGINS_DIR.glob("*.py")): - if plugin_path.name.startswith("_"): - continue - - plugin_name = plugin_path.stem - if plugin_name in existing_plugins: - continue - - # Import plugin to read PLUGIN_CONFIG - try: - import importlib - - # Use absolute package path for plugin import - spec_name = f"aipass.daemon.apps.plugins.{plugin_name}" - module = importlib.import_module(spec_name) - - if not hasattr(module, "PLUGIN_CONFIG"): - continue - - config = module.PLUGIN_CONFIG - except Exception as e: - logger.warning("[actions_registry] Failed to import plugin %s: %s", plugin_name, e) - continue - - # Map PLUGIN_CONFIG to action fields - action_id = _get_next_id(registry) - action = { - "id": action_id, - "name": config.get("name", plugin_name), - "type": "plugin", - "schedule_type": config.get("schedule", "interval"), - "time": config.get("time"), - "interval_minutes": config.get("interval_minutes"), - "due_date": None, - "target_branch": config.get("branch", ""), - "prompt": config.get("prompt", ""), - "fresh": config.get("fresh", True), - "max_turns": config.get("max_turns", 50), - "enabled": config.get("enabled", False), - "self_dispatch": config.get("self_dispatch", False), - "plugin_file": plugin_name, - "last_run": last_run_map.get(config.get("name", plugin_name)), - "next_run": None, - "created": datetime.now().isoformat(), - "completed": None, - } - - # Calculate next_run from last_run - action["next_run"] = calc_next_run(action) - - registry["actions"].append(action) - migrated += 1 - logger.info("[actions_registry] Migrated plugin: %s -> action %s", plugin_name, action_id) - - if migrated > 0: - save_registry(registry) - logger.info("[actions_registry] Migration complete: %d plugin(s) migrated", migrated) - - return migrated diff --git a/src/aipass/daemon/apps/handlers/schedule/__init__.py b/src/aipass/daemon/apps/handlers/schedule/__init__.py index 4a7d5be6..137cbc2f 100644 --- a/src/aipass/daemon/apps/handlers/schedule/__init__.py +++ b/src/aipass/daemon/apps/handlers/schedule/__init__.py @@ -2,10 +2,12 @@ # META DATA HEADER # Name: __init__.py - Schedule Handlers Package # Date: 2026-02-04 -# Version: 1.0.0 +# Version: 2.0.0 # Category: daemon/handlers/schedule # # CHANGELOG (Max 5 entries): +# - v2.0.0 (2026-06-25): task_registry archived (TDPLAN-0008); package +# now exposes runstate + discovery for the live .daemon/ scheduler. # - v1.0.0 (2026-02-04): Initial package setup # # CODE STANDARDS: @@ -14,25 +16,26 @@ # ============================================= """ -Schedule handlers for daemon's scheduled follow-ups system. +Schedule handlers for daemon's decentralized .daemon/ scheduler. + +task_registry (fire-and-forget follow-ups) has been archived — superseded +by the per-branch .daemon/schedule.json model (DPLAN-0204). """ -from aipass.daemon.apps.handlers.schedule.task_registry import ( - load_tasks, - save_tasks, - create_task, - delete_task, - get_due_tasks, - mark_completed, - parse_due_date, +from aipass.daemon.apps.handlers.schedule.runstate import ( + load_runstate, + save_runstate, + update_job_runstate, + is_job_due, + job_key, ) +from aipass.daemon.apps.handlers.schedule.discovery import discover_jobs __all__ = [ - "load_tasks", - "save_tasks", - "create_task", - "delete_task", - "get_due_tasks", - "mark_completed", - "parse_due_date", + "load_runstate", + "save_runstate", + "update_job_runstate", + "is_job_due", + "job_key", + "discover_jobs", ] diff --git a/src/aipass/daemon/apps/handlers/schedule/runstate.py b/src/aipass/daemon/apps/handlers/schedule/runstate.py index bc5be838..85353f3a 100644 --- a/src/aipass/daemon/apps/handlers/schedule/runstate.py +++ b/src/aipass/daemon/apps/handlers/schedule/runstate.py @@ -250,7 +250,7 @@ def update_job_runstate( schedule: dict, timestamp: Optional[str] = None, ) -> None: - """Update last_run and next_run for a job after firing.""" + """Update runstate for a job after successful firing.""" if timestamp is None: timestamp = datetime.now().isoformat() @@ -258,6 +258,9 @@ def update_job_runstate( entry = runstate.setdefault("jobs", {}).setdefault(key, {}) entry["last_run"] = timestamp entry["next_run"] = _calc_next_run(schedule, timestamp) + entry["last_status"] = "success" + entry["last_success_at"] = timestamp + entry["last_error"] = None if schedule.get("type") == "once": entry["completed"] = timestamp @@ -265,6 +268,28 @@ def update_job_runstate( json_handler.log_operation("update_job_runstate", {"key": key}) +def record_job_failure( + runstate: dict, + owner: str, + job_id: str, + error_msg: str, + status: str = "failed", + timestamp: Optional[str] = None, +) -> None: + """Record a failed job firing in runstate.""" + if timestamp is None: + timestamp = datetime.now().isoformat() + + key = job_key(owner, job_id) + entry = runstate.setdefault("jobs", {}).setdefault(key, {}) + entry["last_run"] = timestamp + entry["last_status"] = status + entry["last_failure_at"] = timestamp + entry["last_error"] = error_msg[:500] + + json_handler.log_operation("record_job_failure", {"key": key, "status": status}) + + def prune_orphans(runstate: dict, active_keys: set) -> int: """Remove runstate entries for jobs that no longer exist. Returns count pruned.""" jobs = runstate.get("jobs", {}) diff --git a/src/aipass/daemon/apps/handlers/schedule/task_registry.py b/src/aipass/daemon/apps/handlers/schedule/task_registry.py deleted file mode 100644 index 6dd70815..00000000 --- a/src/aipass/daemon/apps/handlers/schedule/task_registry.py +++ /dev/null @@ -1,588 +0,0 @@ -# =================== AIPass ==================== -# Name: task_registry.py -# Description: DAEMON Scheduled Tasks Registry -# Version: 1.0.0 -# Created: 2026-02-04 -# Modified: 2026-02-04 -# ============================================= - -""" -Handler for scheduled task storage and operations. - -Fire-and-forget follow-up system for DAEMON. -Tasks are stored in daemon_json/schedule.json and processed -when their due date arrives. -""" - -import json -import uuid -from pathlib import Path -from datetime import datetime, timedelta -from typing import Dict, List, Any, Optional -import re - -from aipass.prax import logger -from aipass.daemon.apps.handlers.json import json_handler - -# ============================================= -# CONSTANTS -# ============================================= - -_DAEMON_ROOT = Path(__file__).resolve().parents[3] # src/aipass/daemon/ -SCHEDULE_JSON_PATH = _DAEMON_ROOT / "daemon_json" / "schedule.json" - -DEFAULT_SCHEDULE_DATA: Dict[str, Any] = {"tasks": []} - -# ============================================= -# JSON FILE OPERATIONS -# ============================================= - - -def _ensure_json_exists() -> None: - """Ensure schedule.json exists, create with defaults if missing.""" - SCHEDULE_JSON_PATH.parent.mkdir(parents=True, exist_ok=True) - - if not SCHEDULE_JSON_PATH.exists(): - with open(SCHEDULE_JSON_PATH, "w", encoding="utf-8") as f: - json.dump(DEFAULT_SCHEDULE_DATA, f, indent=2, ensure_ascii=False) - - -def ensure_lock_dir() -> Dict[str, Any]: - """Ensure the daemon_json directory exists for lock files. - - Returns: - Dict with 'path' (str) of the lock file directory. - """ - lock_dir = SCHEDULE_JSON_PATH.parent - lock_dir.mkdir(parents=True, exist_ok=True) - return {"path": str(lock_dir)} - - -def load_tasks() -> List[Dict[str, Any]]: - """ - Load all tasks from schedule.json. - - Returns: - List of task dictionaries - """ - _ensure_json_exists() - - try: - with open(SCHEDULE_JSON_PATH, "r", encoding="utf-8") as f: - data = json.load(f) - return data.get("tasks", []) - except (json.JSONDecodeError, IOError) as e: - logger.error("[task_registry] Failed to load schedule.json: %s", e) - return [] - - -def save_tasks(tasks: List[Dict[str, Any]]) -> bool: - """ - Save tasks to schedule.json. - - Args: - tasks: List of task dictionaries to save - - Returns: - True if successful, False otherwise - """ - _ensure_json_exists() - - try: - data = {"tasks": tasks} - with open(SCHEDULE_JSON_PATH, "w", encoding="utf-8") as f: - json.dump(data, f, indent=2, ensure_ascii=False) - return True - except IOError as e: - logger.error("[task_registry] Failed to save schedule.json: %s", e) - return False - - -# ============================================= -# DATE PARSING -# ============================================= - - -def parse_due_date(date_str: str) -> str: - """ - Parse various date formats to ISO 8601 date string. - - Supports: - - "7d" -> 7 days from now - - "1w" -> 1 week from now - - "2w" -> 2 weeks from now - - "2026-02-11" -> exact date (ISO 8601) - - Args: - date_str: Date string in supported format - - Returns: - ISO 8601 date string (YYYY-MM-DD) - - Raises: - ValueError: If date format is invalid - """ - date_str = date_str.strip() - today = datetime.now().date() - - # Check for relative day format: "7d", "14d", etc. - day_match = re.match(r"^(\d+)d$", date_str, re.IGNORECASE) - if day_match: - days = int(day_match.group(1)) - future_date = today + timedelta(days=days) - return future_date.isoformat() - - # Check for relative week format: "1w", "2w", etc. - week_match = re.match(r"^(\d+)w$", date_str, re.IGNORECASE) - if week_match: - weeks = int(week_match.group(1)) - future_date = today + timedelta(weeks=weeks) - return future_date.isoformat() - - # Check for ISO 8601 date format: "2026-02-11" - iso_match = re.match(r"^(\d{4})-(\d{2})-(\d{2})$", date_str) - if iso_match: - try: - # Validate it's a real date - year = int(iso_match.group(1)) - month = int(iso_match.group(2)) - day = int(iso_match.group(3)) - parsed_date = datetime(year, month, day).date() - return parsed_date.isoformat() - except ValueError as e: - raise ValueError(f"Invalid date: {date_str}") from e - - raise ValueError(f"Invalid date format: '{date_str}'. Use '7d' (days), '1w' (weeks), or 'YYYY-MM-DD' (ISO date)") - - -# ============================================= -# TASK OPERATIONS -# ============================================= - - -def _generate_task_id() -> str: - """Generate 16-character UUID for task ID.""" - return uuid.uuid4().hex[:16] - - -def create_task(task: str, due_date: str, recipient: str, message: str) -> Dict[str, Any]: - """ - Create a new scheduled task. - - Args: - task: Brief description of the task/follow-up - due_date: When to trigger (supports "7d", "1w", "YYYY-MM-DD") - recipient: Target branch (e.g., "@devpulse") - message: Message to deliver when due - - Returns: - Created task dictionary - - Raises: - ValueError: If due_date format is invalid - """ - json_handler.log_operation("task_created") - parsed_due = parse_due_date(due_date) - - new_task: Dict[str, Any] = { - "id": _generate_task_id(), - "created": datetime.now().date().isoformat(), - "due_date": parsed_due, - "task": task, - "recipient": recipient, - "message": message, - "status": "pending", - } - - tasks = load_tasks() - tasks.append(new_task) - save_tasks(tasks) - - return new_task - - -def delete_task(task_id: str) -> bool: - """ - Delete a task by ID. - - Args: - task_id: 8-character task ID - - Returns: - True if task was found and deleted, False otherwise - """ - tasks = load_tasks() - original_count = len(tasks) - - tasks = [t for t in tasks if t.get("id") != task_id] - - if len(tasks) < original_count: - save_tasks(tasks) - return True - - return False - - -def get_due_tasks() -> List[Dict[str, Any]]: - """ - Get all tasks that are due (due_date <= today). - - Only returns tasks with status 'pending' - excludes 'dispatching' and 'completed'. - - Returns: - List of tasks that are due for processing - """ - tasks = load_tasks() - today = datetime.now().date().isoformat() - - due_tasks = [t for t in tasks if t.get("status") == "pending" and t.get("due_date", "") <= today] - - return due_tasks - - -def mark_dispatching(task_id: str) -> bool: - """ - Mark a task as currently being dispatched. - - Prevents re-dispatch while email is being sent. - - Args: - task_id: 8-character task ID - - Returns: - True if task was found and marked, False otherwise - """ - tasks = load_tasks() - - for task in tasks: - if task.get("id") == task_id: - task["status"] = "dispatching" - task["dispatch_started"] = datetime.now().isoformat() - save_tasks(tasks) - return True - - return False - - -def mark_pending(task_id: str) -> bool: - """ - Reset a task to pending status (for retry after failed dispatch). - - Args: - task_id: 8-character task ID - - Returns: - True if task was found and reset, False otherwise - """ - tasks = load_tasks() - - for task in tasks: - if task.get("id") == task_id: - task["status"] = "pending" - task.pop("dispatch_started", None) - save_tasks(tasks) - return True - - return False - - -def _is_stale_dispatch(started: str, cutoff: datetime) -> bool: - """Check if a dispatch_started timestamp is older than the cutoff.""" - try: - start_time = datetime.fromisoformat(started) - return start_time < cutoff - except ValueError as e: - logger.warning("[task_registry] Invalid dispatch_started timestamp, resetting task: %s", e) - return True - - -def recover_stale_dispatches(max_age_minutes: int = 5) -> int: - """ - Reset tasks stuck in 'dispatching' status for too long. - - Called before processing to recover from crashed dispatches. - - Args: - max_age_minutes: Maximum time a task can be in dispatching status - - Returns: - Number of tasks recovered - """ - tasks = load_tasks() - recovered = 0 - cutoff = datetime.now() - timedelta(minutes=max_age_minutes) - - for task in tasks: - if task.get("status") != "dispatching": - continue - started = task.get("dispatch_started") - if not started: - continue - if _is_stale_dispatch(started, cutoff): - task["status"] = "pending" - task.pop("dispatch_started", None) - recovered += 1 - - if recovered: - save_tasks(tasks) - - return recovered - - -def mark_completed(task_id: str) -> bool: - """ - Mark a task as completed. - - Args: - task_id: 8-character task ID - - Returns: - True if task was found and marked, False otherwise - """ - tasks = load_tasks() - - for task in tasks: - if task.get("id") == task_id: - task["status"] = "completed" - task["completed_date"] = datetime.now().date().isoformat() - save_tasks(tasks) - return True - - return False - - -def get_task_by_id(task_id: str) -> Optional[Dict[str, Any]]: - """ - Get a single task by ID. - - Args: - task_id: 8-character task ID - - Returns: - Task dictionary if found, None otherwise - """ - tasks = load_tasks() - - for task in tasks: - if task.get("id") == task_id: - return task - - return None - - -def get_pending_tasks() -> List[Dict[str, Any]]: - """ - Get all pending tasks (not yet due or completed). - - Returns: - List of pending tasks - """ - tasks = load_tasks() - return [t for t in tasks if t.get("status") == "pending"] - - -# ============================================= -# BATCH PROCESSING -# ============================================= - - -def _safe_mark_pending(task_id: str) -> None: - """Best-effort reset a task to pending, logging on failure.""" - try: - mark_pending(task_id) - except Exception as pending_err: - logger.error("[task_registry] Failed to reset task %s to pending: %s", task_id[:8], pending_err) - - -def process_due_tasks_batch( - send_email_fn=None, - stale_max_age: int = 5, -) -> Dict[str, Any]: - """ - Process all due tasks: recover stale, dispatch emails, track results. - - This is the implementation logic for batch task processing. - The module layer handles display; this handler returns raw data. - - Args: - send_email_fn: Callable to send email (to_branch, subject, message, ...). - If None, email dispatch is skipped. - stale_max_age: Maximum minutes before a dispatching task is considered stale. - - Returns: - Dict with keys: due, success, failed, recovered, errors (list of str), - processed_tasks (list of dicts with id, recipient, task, status). - """ - import time - - results: Dict[str, Any] = { - "due": 0, - "success": 0, - "failed": 0, - "recovered": 0, - "errors": [], - "processed_tasks": [], - } - - # Recover any stale dispatches - try: - recovered = recover_stale_dispatches(max_age_minutes=stale_max_age) - results["recovered"] = recovered - except Exception as e: - logger.warning("[task_registry] Stale dispatch recovery failed: %s", e) - results["errors"].append(f"Stale recovery: {e}") - - # Get due tasks - try: - due_tasks = get_due_tasks() - except Exception as e: - logger.error("[task_registry] Failed to load due tasks: %s", e) - results["errors"].append(f"Load tasks: {e}") - return results - - results["due"] = len(due_tasks) - - if not due_tasks: - return results - - for task in due_tasks: - task_id = task.get("id", "") - recipient = task.get("recipient", "") - task_desc = task.get("task", "") - message = task.get("message", "") - - task_result = { - "id": task_id, - "recipient": recipient, - "task": task_desc, - "status": "pending", - } - - # Mark as dispatching (prevents re-dispatch) - try: - mark_dispatching(task_id) - except Exception as e: - logger.error("[task_registry] Failed to mark task %s as dispatching: %s", task_id[:8], e) - results["errors"].append(f"Mark dispatching {task_id[:8]}: {e}") - results["failed"] += 1 - task_result["status"] = "error" - task_result["error"] = str(e) - results["processed_tasks"].append(task_result) - continue - - # Build email body - email_body = f"{task_desc}" - if message: - email_body += f"\n\nDetails:\n{message}" - - # Send the email - if send_email_fn is None: - mark_pending(task_id) - results["failed"] += 1 - task_result["status"] = "skipped" - task_result["error"] = "email function not available" - results["errors"].append(f"Email unavailable for {task_id[:8]}") - results["processed_tasks"].append(task_result) - continue - - try: - email_sent = send_email_fn( - to_branch=recipient, - subject=f"[SCHEDULED] {task_desc}", - message=email_body, - from_branch="@daemon", - auto_execute=True, - reply_to="@devpulse", - ) - - if email_sent: - mark_completed(task_id) - results["success"] += 1 - task_result["status"] = "sent" - else: - mark_pending(task_id) - results["failed"] += 1 - task_result["status"] = "failed" - task_result["error"] = "email send returned False" - results["errors"].append(f"Email failed: {task_id[:8]} -> {recipient}") - - except Exception as e: - logger.error("[task_registry] Email dispatch error for task %s: %s", task_id[:8], e) - _safe_mark_pending(task_id) - results["failed"] += 1 - task_result["status"] = "error" - task_result["error"] = str(e) - results["errors"].append(f"Email error {task_id[:8]}: {e}") - - results["processed_tasks"].append(task_result) - - # Small delay between dispatches (prevents thundering herd) - time.sleep(1.0) - - return results - - -# ============================================= -# MAIN - Testing -# ============================================= - -if __name__ == "__main__": - from rich.console import Console - from rich.panel import Panel - from rich.table import Table - - console = Console() - - console.print() - console.print(Panel.fit("[bold cyan]TASK REGISTRY - Handler Test[/bold cyan]", border_style="bright_blue")) - console.print() - - # Test date parsing - console.print("[yellow]Testing date parsing:[/yellow]") - test_dates = ["7d", "1w", "2w", "2026-03-15"] - for d in test_dates: - try: - result = parse_due_date(d) - console.print(f" {d} -> {result}") - except ValueError as e: - logger.warning("Date parse test failed for %s: %s", d, e) - console.print(f" {d} -> [red]ERROR: {e}[/red]") - - # Test invalid date - try: - parse_due_date("invalid") - except ValueError as e: - logger.info("Expected parse failure for 'invalid': %s", e) - console.print(f" invalid -> [green]Correctly raised: {e}[/green]") - - console.print() - console.print("[yellow]Testing task creation:[/yellow]") - - # Create a test task - test_task = create_task( - task="Test backup health check", - due_date="7d", - recipient="@devpulse", - message="Please verify backup systems are healthy", - ) - console.print(f" Created task: {test_task['id']}") - console.print(f" Due: {test_task['due_date']}") - - # Show all tasks - console.print() - console.print("[yellow]Current tasks:[/yellow]") - all_tasks = load_tasks() - - table = Table(show_header=True) - table.add_column("ID", style="cyan") - table.add_column("Task", style="white") - table.add_column("Due", style="yellow") - table.add_column("Status", style="green") - - for t in all_tasks: - table.add_row(t.get("id", "?"), t.get("task", "?")[:30], t.get("due_date", "?"), t.get("status", "?")) - - console.print(table) - console.print() - console.print(f"[dim]Schedule file: {SCHEDULE_JSON_PATH}[/dim]") - console.print() diff --git a/src/aipass/daemon/apps/handlers/schedule/telegram_notifier.py b/src/aipass/daemon/apps/handlers/schedule/telegram_notifier.py new file mode 100644 index 00000000..1b8485c4 --- /dev/null +++ b/src/aipass/daemon/apps/handlers/schedule/telegram_notifier.py @@ -0,0 +1,53 @@ +# =================== AIPass ==================== +# Name: telegram_notifier.py +# Description: Scheduler lifecycle notifications via Telegram +# Version: 3.0.0 +# Created: 2026-02-15 +# Modified: 2026-06-25 +# ============================================= + +""" +Scheduler lifecycle notifications — emit running/complete/failed pings +via the @skills telegram notifier (secret slug: telegram/scheduler). + +Fail-soft: if the secret is missing or the send fails, returns False +and never raises. The tick MUST keep firing jobs regardless. +""" + +from aipass.prax import logger +from aipass.daemon.apps.handlers.json import json_handler # noqa: F401 + + +def _send(message: str) -> bool: + """Send via skills notifier, fail-soft. + + Cross-branch import authorized by TDPLAN-0008 contract — + daemon emits lifecycle pings through the skills telegram notifier. + """ + try: + from aipass.skills.lib.telegram.apps.handlers.notifier import ( # noqa: E501 + send_telegram_notification, + ) + + return send_telegram_notification(message) + except Exception as e: + logger.info("[telegram_notifier] Send failed (non-fatal): %s", e) + return False + + +def notify_triggered(owner: str, job_id: str) -> bool: + """Emit 'running' ping when a job fires.""" + json_handler.log_operation("notify_triggered", {"owner": owner, "job_id": job_id}) + return _send(f"\U0001f535 {owner}/{job_id} running") + + +def notify_complete(owner: str, job_id: str, summary: str) -> bool: + """Emit 'complete' ping after successful dispatch.""" + json_handler.log_operation("notify_complete", {"owner": owner, "job_id": job_id}) + return _send(f"✅ {owner}/{job_id} dispatched\n{summary}") + + +def notify_error(owner: str, job_id: str, error: str) -> bool: + """Emit 'failed' ping on dispatch failure.""" + json_handler.log_operation("notify_error", {"owner": owner, "job_id": job_id}) + return _send(f"❌ {owner}/{job_id} FAILED\n{error}") diff --git a/src/aipass/daemon/apps/modules/actions.py b/src/aipass/daemon/apps/modules/actions.py index dcb13576..65cf2457 100644 --- a/src/aipass/daemon/apps/modules/actions.py +++ b/src/aipass/daemon/apps/modules/actions.py @@ -1,560 +1,51 @@ # =================== AIPass ==================== # Name: actions.py -# Description: Action Registry CLI Module -# Version: 1.0.0 +# Description: Action Registry CLI Module (RETIRED) +# Version: 2.0.0 # Created: 2026-03-02 -# Modified: 2026-03-02 +# Modified: 2026-06-25 # ============================================= """ -CLI interface for the numbered action registry. +RETIRED — the numbered action registry CLI. + +Superseded by the decentralized .daemon/schedule.json model (DPLAN-0204). +Author jobs directly in a branch's .daemon/schedule.json file. +See: drone @daemon run --help (full schema + schedule types) """ -# ============================================= -# IMPORTS -# ============================================= - -import sys from typing import List from aipass.prax import logger - -from aipass.cli.apps.modules import console, error as cli_error -from aipass.daemon.apps.handlers.actions.actions_registry import ( - list_actions, - get_action, - toggle_action, - delete_action, - create_action, - migrate_plugins, - next_due_str, -) +from aipass.cli.apps.modules import console from aipass.daemon.apps.handlers.json import json_handler -def _header(text): - console.print(f"\n[bold cyan]{'=' * 70}[/bold cyan]") - console.print(f"[bold cyan] {text}[/bold cyan]") - console.print(f"[bold cyan]{'=' * 70}[/bold cyan]") - - -def _success(text): - console.print(f"[green]OK:[/green] {text}") - - -def _error(text): - cli_error(text) - - -# ============================================= -# CONSTANTS -# ============================================= - -MODULE_NAME = "actions" - - -# ============================================= -# INTROSPECTION -# ============================================= - - def print_introspection(): - """Display module introspection info.""" + """Display retirement notice.""" console.print() - console.print("[bold cyan]actions Module[/bold cyan]") + console.print("[bold cyan]actions Module[/bold cyan] [yellow](RETIRED)[/yellow]") console.print() - console.print("[dim]CLI interface for the numbered action registry (DPLAN-043)[/dim]") + console.print("[dim]This CLI has been retired. Jobs now live in per-branch .daemon/schedule.json files.[/dim]") + console.print("[dim]Run [bold]drone @daemon run --help[/bold] for the schema and schedule types.[/dim]") + console.print("[dim]Run [bold]drone @daemon queue[/bold] to see the unified job queue.[/dim]") console.print() - console.print("[yellow]Connected Handlers:[/yellow]") - console.print(" handlers/actions/") - console.print( - " [cyan]*[/cyan] actions_registry.py" - " [dim](list_actions, get_action," - " toggle_action, delete_action, create_action," - " migrate_plugins, next_due_str — registry CRUD)[/dim]" - ) - console.print() - - -# ============================================= -# OUTPUT FORMATTING -# ============================================= - - -def _format_schedule(action: dict) -> str: - """Build schedule display string for an action.""" - schedule_type = action.get("schedule_type", "") - if schedule_type == "daily": - return f"daily @ {action.get('time', '??:??')}" - if schedule_type == "hourly": - m = action.get("time", "0") - return f"hourly @ :{int(m):02d}" - if schedule_type == "interval": - mins = action.get("interval_minutes", 0) - if mins >= 60: - return f"every {mins // 60}h" - return f"every {mins}m" - if schedule_type == "once": - return f"once: {action.get('due_date', '?')}" - return schedule_type - - -def _print_actions_table(actions: list) -> None: - """Display formatted action list as a table.""" - console.print() - _header("Action Registry") - console.print() - - if not actions: - console.print("[dim]No actions registered. Run 'actions migrate' to import plugins.[/dim]") - console.print() - return - - # Header row - console.print(f" {'ID':<6} {'ON':<4} {'NAME':<24} {'TYPE':<10} {'TARGET':<16} {'SCHEDULE':<20} {'NEXT DUE':<16}") - console.print(" " + "-" * 96) - - for action in actions: - action_id = action.get("id", "????") - enabled = "[green]ON[/green] " if action.get("enabled") else "[red]OFF[/red]" - name = action.get("name", "")[:22] - action_type = action.get("type", "")[:8] - target = action.get("target_branch", "")[:14] - - schedule_str = _format_schedule(action) - next_due = next_due_str(action) - - console.print( - f" {action_id:<6} {enabled:<4} {name:<24} {action_type:<10} " - f" {target:<16} {schedule_str:<20} {next_due:<16}" - ) - - console.print() - enabled_count = sum(1 for a in actions if a.get("enabled")) - console.print(f" [dim]Total: {len(actions)} actions ({enabled_count} enabled)[/dim]") - console.print() - - -def _print_action_detail(action: dict) -> None: - """Display detailed view of a single action.""" - console.print() - _header(f"Action {action['id']}: {action['name']}") - console.print() - - fields = [ - ("ID", action.get("id")), - ("Name", action.get("name")), - ("Type", action.get("type")), - ("Enabled", "[green]ON[/green]" if action.get("enabled") else "[red]OFF[/red]"), - ("Schedule", action.get("schedule_type")), - ("Time", action.get("time")), - ("Interval", f"{action.get('interval_minutes')}m" if action.get("interval_minutes") else None), - ("Due Date", action.get("due_date")), - ("Target", action.get("target_branch")), - ("Fresh", action.get("fresh")), - ("Max Turns", action.get("max_turns")), - ("Self Dispatch", action.get("self_dispatch")), - ("Plugin File", action.get("plugin_file")), - ("Last Run", action.get("last_run", "never")[:19] if action.get("last_run") else "never"), - ("Next Run", next_due_str(action)), - ("Created", action.get("created", "")[:19]), - ("Completed", action.get("completed")), - ] - - for label, value in fields: - if value is None: - continue - console.print(f" [cyan]{label:<16}[/cyan] {value}") - - # Show prompt (truncated for readability) - prompt = action.get("prompt", "") - if prompt: - console.print() - console.print(" [cyan]Prompt:[/cyan]") - # Show first 200 chars - display_prompt = prompt[:200] - if len(prompt) > 200: - display_prompt += "..." - for line in display_prompt.split("\n"): - console.print(f" [dim]{line}[/dim]") - - console.print() - - -def print_help() -> None: - """Display help using Rich formatted output.""" - console.print() - _header("Actions -- Numbered Action Registry") - console.print() - - console.print("[yellow]USAGE:[/yellow]") - console.print(" drone @daemon actions list") - console.print(" drone @daemon actions info") - console.print(" drone @daemon actions on") - console.print(" drone @daemon actions off") - console.print(' drone @daemon actions set reminder "message" [--to @branch]') - console.print(' drone @daemon actions set schedule @branch "prompt" [time]') - console.print(" drone @daemon actions migrate") - console.print(" drone @daemon actions delete ") - console.print() - - console.print("[yellow]COMMANDS:[/yellow]") - console.print(" list List all registered actions with status") - console.print(" info Show detailed view of a single action") - console.print(" on Enable an action") - console.print(" off Disable an action") - console.print(" set Create a new reminder or schedule") - console.print(" migrate Import existing plugins into registry") - console.print(" delete Remove an action from the registry") - console.print() - - console.print("[yellow]SET REMINDER:[/yellow]") - console.print(' set reminder 2026-03-11 "Check VERA progress"') - console.print(' set reminder 7d "Follow up on PR review" --to @flow') - console.print(" [dim]Date formats: YYYY-MM-DD, 1d, 7d, 1w, 2w[/dim]") - console.print() - - console.print("[yellow]SET SCHEDULE:[/yellow]") - console.print(' set schedule @seedgo "Run audit" daily 04:00') - console.print(' set schedule @daemon "Heartbeat" interval 240') - console.print(' set schedule @flow "Check plans" hourly 30') - console.print(" [dim]Types: daily HH:MM, hourly MM, interval MINUTES[/dim]") - console.print() - - console.print("[yellow]EXAMPLES:[/yellow]") - console.print(" actions list # See all actions") - console.print(" actions 0003 off # Disable action 3") - console.print(" actions 0003 on # Re-enable it") - console.print(' actions set reminder 2026-03-11 "check VERA" # One-shot reminder') - console.print() - - -# ============================================= -# SUBCOMMAND HANDLERS -# ============================================= - - -def _handle_list(_args: List[str]) -> bool: - """Handle 'actions list' subcommand.""" - actions = list_actions() - _print_actions_table(actions) - logger.info("[DAEMON] actions: Action list displayed") - return True - - -def _handle_toggle(action_id: str, enable: bool) -> bool: - """Handle 'actions on/off' subcommand.""" - action = get_action(action_id) - if action is None: - _error(f"Action not found: {action_id}") - return True # Error displayed - - toggle_action(action_id, enable) - state = "enabled" if enable else "disabled" - _success(f"Action {action_id} ({action['name']}) {state}") - logger.info("[DAEMON] actions: Action toggled") - return True - - -def _handle_info(action_id: str) -> bool: - """Handle 'actions info' subcommand.""" - action = get_action(action_id) - if action is None: - _error(f"Action not found: {action_id}") - return True # Error displayed - - _print_action_detail(action) - logger.info("[DAEMON] actions: Action info displayed") - return True - - -def _handle_set_reminder(args: List[str]) -> bool: - """Handle 'actions set reminder "message" [--to @branch]'.""" - if len(args) < 2: - _error('Usage: actions set reminder "message" [--to @branch]') - return True # Error displayed - - date_str = args[0] - message = args[1] - target_branch = "@devpulse" # Default reminder target - - # Parse --to flag - if "--to" in args: - to_idx = args.index("--to") - if to_idx + 1 < len(args): - target_branch = args[to_idx + 1] - - # Parse date - due_date = _parse_date(date_str) - if not due_date: - _error(f"Invalid date format: {date_str}") - console.print("[dim]Valid formats: YYYY-MM-DD, 1d, 7d, 1w, 2w[/dim]") - return True # Error displayed - - action = create_action( - name=message[:50], - action_type="reminder", - schedule_type="once", - target_branch=target_branch, - prompt=message, - due_date=due_date, - fresh=True, - max_turns=10, - enabled=True, - ) - - _success(f"Reminder created: {action['id']}") - console.print(f" [dim]Due:[/dim] {due_date}") - console.print(f" [dim]To:[/dim] {target_branch}") - console.print(f" [dim]Message:[/dim] {message[:60]}") - console.print() - logger.info("[DAEMON] actions: Reminder set") - return True - - -def _handle_set_schedule(args: List[str]) -> bool: - """Handle 'actions set schedule @branch "prompt" [time_spec]'.""" - if len(args) < 3: - _error('Usage: actions set schedule @branch "prompt" [time_spec]') - return True # Error displayed - - target_branch = args[0] - prompt = args[1] - schedule_type = args[2] - - time_val = None - interval_minutes = None - - if schedule_type not in ("daily", "hourly", "interval"): - _error(f"Unknown schedule type: {schedule_type}") - console.print("[dim]Valid types: daily, hourly, interval[/dim]") - return True # Error displayed - - if len(args) < 4: - _error(f"{schedule_type.title()} schedule requires a time/value argument") - return True # Error displayed - - if schedule_type in ("daily", "hourly"): - time_val = args[3] - else: - try: - interval_minutes = int(args[3]) - except ValueError: - logger.warning("Invalid interval minutes value: %s", args[3]) - _error(f"Invalid interval minutes: {args[3]}") - return True # Error displayed - - # Generate a name from the prompt - name = prompt[:50].replace(" ", "_").lower() - - action = create_action( - name=name, - action_type="schedule", - schedule_type=schedule_type, - target_branch=target_branch, - prompt=prompt, - time=time_val, - interval_minutes=interval_minutes, - fresh=True, - max_turns=50, - enabled=True, - ) - - _success(f"Schedule created: {action['id']}") - console.print(f" [dim]Name:[/dim] {action['name']}") - console.print(f" [dim]Target:[/dim] {target_branch}") - console.print(f" [dim]Type:[/dim] {schedule_type}") - if time_val: - console.print(f" [dim]Time:[/dim] {time_val}") - if interval_minutes: - console.print(f" [dim]Every:[/dim] {interval_minutes} minutes") - console.print() - logger.info("[DAEMON] actions: Schedule set") - return True - - -def _handle_migrate(_args: List[str]) -> bool: - """Handle 'actions migrate' -- import plugins into registry.""" - console.print() - console.print("[dim]Scanning plugins/ for unregistered plugins...[/dim]") - - count = migrate_plugins() - - if count > 0: - _success(f"Migrated {count} plugin(s) into the action registry") - else: - console.print("[dim]All plugins already registered (or none found).[/dim]") - - # Show the updated list - actions = list_actions() - _print_actions_table(actions) - logger.info("[DAEMON] actions: Plugin migration completed") - return True - - -def _handle_delete(args: List[str]) -> bool: - """Handle 'actions delete '.""" - if not args: - _error("Action ID required: actions delete ") - return True # Error displayed - - action_id = args[0] - action = get_action(action_id) - if action is None: - _error(f"Action not found: {action_id}") - return True # Error displayed - - delete_action(action_id) - _success(f"Deleted action {action_id}: {action['name']}") - logger.info("[DAEMON] actions: Action deleted") - return True - - -# ============================================= -# DATE PARSING -# ============================================= - - -def _parse_date(date_str: str) -> str: - """ - Parse a date string into ISO format. - - Supports: YYYY-MM-DD, 1d, 7d, 1w, 2w - - Returns: - ISO date string or empty string on failure. - """ - from datetime import datetime, timedelta - - date_str = date_str.strip() - - # Relative dates - if date_str.endswith("d"): - try: - days = int(date_str[:-1]) - return (datetime.now() + timedelta(days=days)).strftime("%Y-%m-%d") - except ValueError: - logger.warning("Invalid relative day format: %s", date_str) - return "" - elif date_str.endswith("w"): - try: - weeks = int(date_str[:-1]) - return (datetime.now() + timedelta(weeks=weeks)).strftime("%Y-%m-%d") - except ValueError: - logger.warning("Invalid relative week format: %s", date_str) - return "" - - # ISO date - try: - datetime.strptime(date_str, "%Y-%m-%d") - return date_str - except ValueError: - logger.warning("Invalid ISO date format: %s", date_str) - return "" - - -# ============================================= -# ORCHESTRATION -# ============================================= - - -def _route_set_subcommand(args: List[str]) -> bool: - """Route 'actions set reminder ...' / 'actions set schedule ...'.""" - if len(args) < 2: - _error("Usage: actions set ...") - return True # Error displayed - set_type = args[1] - if set_type == "reminder": - return _handle_set_reminder(args[2:]) - if set_type == "schedule": - return _handle_set_schedule(args[2:]) - _error(f"Unknown set type: {set_type}. Use 'reminder' or 'schedule'.") - return True # Error displayed - - -def _route_action_id(action_id: str, args: List[str]) -> bool: - """Route 'actions <4-digit-id> [on|off|info]'.""" - if len(args) < 2: - return _handle_info(action_id) - sub_action = args[1] - if sub_action == "on": - return _handle_toggle(action_id, True) - if sub_action == "off": - return _handle_toggle(action_id, False) - if sub_action == "info": - return _handle_info(action_id) - _error(f"Unknown action command: {sub_action}. Use 'on', 'off', or 'info'.") - return True # Error displayed def handle_command(command: str, args: List[str]) -> bool: - """ - Handle 'actions' command and route to subcommands. - - Args: - command: Command name (should be 'actions') - args: Command arguments - - Returns: - True if handled, False otherwise - """ + """Handle 'actions' command — retired, shows migration notice.""" if command != "actions": return False - try: - # No args -- introspection gate - if not args: - print_introspection() - return True + if not args: + print_introspection() + return True - # Help flag - if args[0] in ["--help", "-h", "help"]: - print_help() - return True + if args[0] in ("--help", "-h", "help"): + print_introspection() + return True - subcommand = args[0] - - json_handler.log_operation("actions_command", {"subcommand": args[0] if args else "introspection"}) - - # Named subcommands - if subcommand == "list": - return _handle_list(args[1:]) - if subcommand == "migrate": - return _handle_migrate(args[1:]) - if subcommand == "delete": - return _handle_delete(args[1:]) - if subcommand == "set": - return _route_set_subcommand(args) - - # Check if first arg is an action ID (4-digit numeric) - if subcommand.isdigit() and len(subcommand) == 4: - return _route_action_id(subcommand, args) - - _error(f"Unknown subcommand: {subcommand}") - console.print("[dim]Run 'actions --help' for available commands[/dim]") - return True # Command was handled (error displayed) - - except Exception as e: - logger.error("[actions] Error in actions command: %s", e, exc_info=True) - _error(f"Error: {e}") - return True # Error displayed - - -# ============================================= -# MAIN ENTRY -# ============================================= - - -def main() -> None: - """Main entry point for direct execution.""" - args = sys.argv[1:] - - if not args or args[0] in ["--help", "-h", "help"]: - print_help() - return - - handle_command("actions", args) - - -if __name__ == "__main__": - main() + json_handler.log_operation("actions_command_retired", {"args": args[:2]}) + logger.info("[actions] Retired CLI invoked") + print_introspection() + return True diff --git a/src/aipass/daemon/apps/modules/queue.py b/src/aipass/daemon/apps/modules/queue.py new file mode 100644 index 00000000..e0b672b7 --- /dev/null +++ b/src/aipass/daemon/apps/modules/queue.py @@ -0,0 +1,191 @@ +# =================== AIPass ==================== +# Name: queue.py +# Description: Unified job queue view (drone @daemon queue) +# Version: 1.0.0 +# Created: 2026-06-25 +# Modified: 2026-06-25 +# ============================================= + +""" +Unified queue view — aggregates .daemon/schedule.json jobs joined to runstate. + +Human-readable Rich table (default) or --json matching the frozen contract +consumed by @skills' scheduler bot. +""" + +import json +from datetime import datetime, timezone +from typing import List, Optional + +from aipass.prax import logger +from aipass.cli.apps.modules import console +from aipass.daemon.apps.handlers.json import json_handler +from aipass.daemon.apps.handlers.schedule.discovery import discover_jobs +from aipass.daemon.apps.handlers.schedule.runstate import ( + load_runstate, + get_job_state, +) + +HANDLED_COMMANDS = {"queue"} + + +def print_introspection(): + """Display module introspection info.""" + console.print() + console.print("[bold cyan]queue Module[/bold cyan]") + console.print() + console.print("[dim]Unified job queue view — .daemon/ jobs joined to runstate[/dim]") + console.print() + console.print("[yellow]Reads:[/yellow]") + console.print(" [cyan]*[/cyan] src/aipass/*/.daemon/*.json [dim](per-branch schedule files)[/dim]") + console.print(" [cyan]*[/cyan] daemon_json/daemon_runstate.json [dim](last_run/status state)[/dim]") + console.print() + + +def print_help(): + """Display usage information.""" + console.print("\n[bold cyan]queue — Unified Job Queue View[/bold cyan]") + console.print("\n[yellow]USAGE:[/yellow]") + console.print(" drone @daemon queue Show job queue (Rich table)") + console.print(" drone @daemon queue --json Machine-readable JSON (frozen schema)") + console.print(" drone @daemon queue --help Show this help message") + console.print() + + +def _schedule_human(job: dict) -> str: + """Build human-readable schedule string.""" + sched = job.get("schedule", {}) + sched_type = sched.get("type", "") + if sched_type == "once": + return sched.get("due_date", "?") + if sched_type == "daily": + return f"daily @ {sched.get('time', '??:??')}" + if sched_type == "hourly": + m = sched.get("time", "0") + return f"hourly @ :{int(m):02d}" + if sched_type == "interval": + mins = sched.get("interval_minutes", 0) + if mins >= 60: + return f"every {mins // 60}h" + return f"every {mins}m" + return sched_type + + +def _compute_next_run(job: dict, state: dict) -> Optional[str]: + """Determine next_run from runstate or schedule.""" + if state.get("completed"): + return None + next_run = state.get("next_run") + if next_run: + return next_run + sched = job.get("schedule", {}) + if sched.get("type") == "once": + due = sched.get("due_date") + if due and "T" not in due: + return f"{due}T09:00:00" + return due + return None + + +def _build_queue(jobs: list, runstate: dict) -> list: + """Build unified queue entries from discovered jobs + runstate.""" + entries = [] + for job in jobs: + state = get_job_state(runstate, job["owner"], job["id"]) + if state.get("completed"): + continue + + owner = job["owner"] + if owner.startswith("@"): + pass + elif "@" in owner: + owner = f"@{owner.split('@')[0]}" + + prompt = job.get("prompt", "") + preview = prompt[:80] + "..." if len(prompt) > 80 else prompt + + entries.append( + { + "owner": owner, + "id": job["id"], + "enabled": job.get("enabled", True), + "type": job["schedule"].get("type", ""), + "schedule_human": _schedule_human(job), + "next_run": _compute_next_run(job, state), + "last_run": state.get("last_run"), + "last_status": state.get("last_status"), + "last_error": state.get("last_error"), + "prompt_preview": preview, + "wake": job.get("wake", {}), + } + ) + return entries + + +def _print_rich_table(entries: list) -> None: + """Print queue as a Rich table.""" + console.print() + console.print("[bold cyan]Job Queue[/bold cyan]") + console.print() + + if not entries: + console.print("[dim]No jobs in queue.[/dim]") + console.print() + return + + console.print( + f" {'OWNER':<14} {'ID':<20} {'ON':<4} {'TYPE':<9} {'SCHEDULE':<18} {'LAST STATUS':<12} {'NEXT RUN':<20}" + ) + console.print(" " + "-" * 97) + + for e in entries: + enabled = "[green]ON[/green] " if e["enabled"] else "[red]OFF[/red]" + last_status = e.get("last_status") or "-" + next_run = (e.get("next_run") or "-")[:19] + console.print( + f" {e['owner']:<14} {e['id']:<20} {enabled:<4} {e['type']:<9} " + f"{e['schedule_human']:<18} {last_status:<12} {next_run:<20}" + ) + + console.print() + enabled_count = sum(1 for e in entries if e["enabled"]) + console.print(f" [dim]Total: {len(entries)} job(s) ({enabled_count} enabled)[/dim]") + console.print() + + +def _build_json_output(entries: list) -> dict: + """Build frozen-schema JSON output for @skills consumption.""" + return { + "generated_at": datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"), + "count": len(entries), + "jobs": entries, + } + + +def handle_command(command: str, args: List[str]) -> bool: + """Handle 'queue' command from daemon CLI router.""" + if command not in HANDLED_COMMANDS: + return False + + if not args: + print_introspection() + return True + + if args[0] in ("--help", "-h"): + print_help() + return True + + json_handler.log_operation("queue_command", {"json": "--json" in args}) + + jobs = discover_jobs() + runstate = load_runstate() + entries = _build_queue(jobs, runstate) + + if "--json" in args: + output = _build_json_output(entries) + console.print(json.dumps(output, indent=2)) + else: + _print_rich_table(entries) + + logger.info("[queue] Queue displayed (%d jobs)", len(entries)) + return True diff --git a/src/aipass/daemon/apps/modules/run.py b/src/aipass/daemon/apps/modules/run.py index f6160891..1101fd26 100644 --- a/src/aipass/daemon/apps/modules/run.py +++ b/src/aipass/daemon/apps/modules/run.py @@ -27,6 +27,7 @@ from aipass.daemon.apps.handlers.schedule.runstate import ( save_runstate, is_job_due, update_job_runstate, + record_job_failure, job_key, prune_orphans, ) @@ -112,18 +113,36 @@ def _log(message: str) -> None: console.print(f"[{timestamp}] {message}") -def _fire_job(job: dict) -> bool: - """Fire a single job via direct wake_branch import (DPLAN-0204 path A).""" +def _should_notify(job: dict) -> bool: + """Check if this job should emit telegram notifications.""" + return job.get("notify", True) + + +def _fire_job(job: dict) -> tuple: + """Fire a single job via direct wake_branch import (DPLAN-0204 path A). + + Returns (ok: bool, error_msg: str). + """ # Cross-branch handler import authorized by DPLAN-0204 §2.8 from aipass.ai_mail.apps.handlers.dispatch.wake import wake_branch # noqa: E402 + from aipass.daemon.apps.handlers.schedule.telegram_notifier import ( + notify_triggered, + notify_complete, + notify_error, + ) owner = job["owner"] + job_id = job["id"] prompt = job["prompt"] wake = job.get("wake", {}) fresh = wake.get("fresh", True) model = wake.get("model") + notify = _should_notify(job) - _log(f"FIRE: {owner}/{job['id']} -> wake_branch({owner}, fresh={fresh}, model={model})") + _log(f"FIRE: {owner}/{job_id} -> wake_branch({owner}, fresh={fresh}, model={model})") + + if notify: + notify_triggered(owner, job_id) try: status, ok = wake_branch( @@ -135,16 +154,24 @@ def _fire_job(job: dict) -> bool: model=model, ) if ok: - _log(f"OK: {owner}/{job['id']} — {status.summary}") - logger.info("[run] Fired %s/%s successfully", owner, job["id"]) + _log(f"OK: {owner}/{job_id} — {status.summary}") + logger.info("[run] Fired %s/%s successfully", owner, job_id) + if notify: + notify_complete(owner, job_id, status.summary) + return True, "" else: - _log(f"FAIL: {owner}/{job['id']} — {status.summary}") - logger.warning("[run] Failed to fire %s/%s: %s", owner, job["id"], status.summary) - return ok + msg = status.summary + _log(f"FAIL: {owner}/{job_id} — {msg}") + logger.warning("[run] Failed to fire %s/%s: %s", owner, job_id, msg) + if notify: + notify_error(owner, job_id, msg) + return False, msg except Exception as e: - logger.error("[run] Exception firing %s/%s: %s", owner, job["id"], e) - _log(f"ERROR: {owner}/{job['id']} — {e}") - return False + logger.error("[run] Exception firing %s/%s: %s", owner, job_id, e) + _log(f"ERROR: {owner}/{job_id} — {e}") + if notify: + notify_error(owner, job_id, str(e)) + return False, str(e) def run_tick(dry_run: bool = False) -> dict: @@ -209,13 +236,14 @@ def run_tick(dry_run: bool = False) -> dict: # Step 4: Fire due jobs for job in due_jobs: - ok = _fire_job(job) + ok, error_msg = _fire_job(job) if ok: results["fired"] += 1 update_job_runstate(runstate, job["owner"], job["id"], job["schedule"]) - save_runstate(runstate) else: results["failed"] += 1 + record_job_failure(runstate, job["owner"], job["id"], error_msg) + save_runstate(runstate) if job != due_jobs[-1]: time.sleep(1.0) diff --git a/src/aipass/daemon/apps/modules/schedule.py b/src/aipass/daemon/apps/modules/schedule.py index c3020771..924d6d3c 100644 --- a/src/aipass/daemon/apps/modules/schedule.py +++ b/src/aipass/daemon/apps/modules/schedule.py @@ -1,436 +1,51 @@ # =================== AIPass ==================== # Name: schedule.py -# Description: DAEMON Scheduled Follow-ups Module -# Version: 1.0.0 +# Description: DAEMON Scheduled Follow-ups Module (RETIRED) +# Version: 2.0.0 # Created: 2026-02-04 -# Modified: 2026-02-04 +# Modified: 2026-06-25 # ============================================= """ -CLI interface for fire-and-forget scheduled follow-ups. +RETIRED — the old fire-and-forget follow-up CLI. + +Superseded by the decentralized .daemon/schedule.json model (DPLAN-0204). +Author jobs directly in a branch's .daemon/schedule.json file. +See: drone @daemon run --help (full schema + schedule types) """ -# ============================================= -# IMPORTS -# ============================================= - -import sys -import argparse -import subprocess -from pathlib import Path from typing import List from aipass.prax import logger - -from aipass.cli.apps.modules import console, error as cli_error +from aipass.cli.apps.modules import console from aipass.daemon.apps.handlers.json import json_handler -from aipass.daemon.apps.handlers.schedule.task_registry import ( - load_tasks, - create_task, - delete_task, - parse_due_date, - process_due_tasks_batch, - ensure_lock_dir, -) - -# File lock for single-instance execution -try: - from filelock import FileLock, Timeout - - FILELOCK_AVAILABLE = True -except ImportError: - FILELOCK_AVAILABLE = False - FileLock = None # type: ignore[assignment,misc] - Timeout = None # type: ignore[assignment,misc] - logger.info("Optional: filelock not available") - - -def _header(text): - console.print(f"\n[bold cyan]{'=' * 70}[/bold cyan]") - console.print(f"[bold cyan] {text}[/bold cyan]") - console.print(f"[bold cyan]{'=' * 70}[/bold cyan]") - - -def _success(text): - console.print(f"[green]OK:[/green] {text}") - - -def _error(text): - cli_error(text) - - -def _send_email_via_drone( - to_branch, subject, message, from_branch="@daemon", auto_execute=True, reply_to=None, **kwargs -): - """Send email via drone @ai_mail send subprocess.""" - cmd = ["drone", "@ai_mail", "send", to_branch, subject, message] - if auto_execute: - cmd.append("--dispatch") - try: - result = subprocess.run(cmd, capture_output=True, text=True, timeout=DRONE_SUBPROCESS_TIMEOUT) - return result.returncode == 0 - except (subprocess.SubprocessError, OSError) as e: - logger.warning("Drone email subprocess failed: %s", e) - return False - - -AI_MAIL_AVAILABLE = True -send_email_direct = _send_email_via_drone - -# ============================================= -# CONSTANTS -# ============================================= - -MODULE_NAME = "schedule" - -# Constants -DRONE_SUBPROCESS_TIMEOUT = 15 # seconds -STALE_DISPATCH_MAX_AGE = 5 # minutes -LOCK_ACQUIRE_TIMEOUT = 0 # seconds (non-blocking) - - -# ============================================= -# INTROSPECTION -# ============================================= def print_introspection(): - """Display module introspection info.""" + """Display retirement notice.""" console.print() - console.print("[bold cyan]schedule Module[/bold cyan]") + console.print("[bold cyan]schedule Module[/bold cyan] [yellow](RETIRED)[/yellow]") console.print() - console.print("[dim]CLI interface for fire-and-forget scheduled follow-ups[/dim]") + console.print("[dim]This CLI has been retired. Jobs now live in per-branch .daemon/schedule.json files.[/dim]") + console.print("[dim]Run [bold]drone @daemon run --help[/bold] for the schema and schedule types.[/dim]") + console.print("[dim]Run [bold]drone @daemon queue[/bold] to see the unified job queue.[/dim]") console.print() - console.print("[yellow]Connected Handlers:[/yellow]") - console.print(" handlers/schedule/") - console.print( - " [cyan]*[/cyan] task_registry.py" - " [dim](load_tasks, create_task, delete_task," - " get_due_tasks, mark_completed, parse_due_date," - " mark_dispatching, mark_pending, recover_stale_dispatches," - " process_due_tasks_batch, ensure_lock_dir" - " — task CRUD and processing)[/dim]" - ) - console.print() - - -_DAEMON_ROOT = Path(__file__).resolve().parents[3] # src/aipass/daemon/ -JSON_DIR = _DAEMON_ROOT / "daemon_json" - -# ============================================= -# OUTPUT FORMATTING -# ============================================= - - -def _print_task_list(tasks: List[dict]) -> None: - """Print formatted task list to console.""" - console.print() - _header("Scheduled Tasks") - console.print() - - pending_tasks = [t for t in tasks if t.get("status") == "pending"] - completed_tasks = [t for t in tasks if t.get("status") == "completed"] - - if not pending_tasks: - console.print("[dim]No pending scheduled tasks.[/dim]") - else: - console.print("[bold cyan]PENDING TASKS[/bold cyan]") - console.print(f"{'ID':<10} {'DUE':<20} {'TO':<15} {'TASK':<40}") - console.print("-" * 85) - - for task in pending_tasks: - task_id = task.get("id", "")[:8] - due = task.get("due_date", "") - recipient = task.get("recipient", "") - task_text = task.get("task", "")[:38] - console.print(f"{task_id:<10} {due:<20} {recipient:<15} {task_text:<40}") - - console.print() - console.print(f"[dim]Total: {len(pending_tasks)} pending, {len(completed_tasks)} completed[/dim]") - console.print() - - -def _print_help() -> None: - """Display help using Rich formatted output.""" - console.print() - _header("Schedule Module - Fire-and-Forget Follow-ups") - console.print() - - console.print("[yellow]USAGE:[/yellow]") - console.print(' drone @daemon schedule create "task" --due 7d --to @branch --message "details"') - console.print(" drone @daemon schedule list") - console.print(" drone @daemon schedule delete ") - console.print(" drone @daemon schedule run-due") - console.print() - - console.print("[yellow]COMMANDS:[/yellow]") - console.print(" create Create a new scheduled task") - console.print(" list List all pending scheduled tasks") - console.print(" delete Delete a scheduled task by ID") - console.print(" run-due Execute all due tasks (sends emails, marks complete)") - console.print() - - console.print("[yellow]CREATE OPTIONS:[/yellow]") - console.print(" --due (Required) Due date: 1d, 7d, 2w, 1m, or ISO date (2026-02-15)") - console.print(" --to (Required) Recipient branch (e.g., @flow, @seedgo)") - console.print(" --message (Optional) Additional details for the follow-up") - console.print() - - console.print("[yellow]EXAMPLES:[/yellow]") - console.print(" # Remind Flow to check on a plan in 7 days") - console.print(' schedule create "Check FPLAN-0290 status" --due 7d --to @flow') - console.print() - console.print(" # Follow up with Seedgo about code review in 2 weeks") - console.print(' schedule create "Code review follow-up" --due 2w --to @seedgo --message "Review PR #45"') - console.print() - console.print(" # Check all due tasks and send reminder emails") - console.print(" schedule run-due") - console.print() - - -# ============================================= -# SUBCOMMAND HANDLERS -# ============================================= - - -def _handle_create(args: List[str]) -> bool: - """Handle schedule create subcommand.""" - parser = argparse.ArgumentParser(prog="schedule create", add_help=False) - parser.add_argument("task", nargs="?", help="Task description") - parser.add_argument("--due", required=True, help="Due date (1d, 7d, 2w, 1m, or ISO date)") - parser.add_argument("--to", required=True, dest="recipient", help="Recipient branch") - parser.add_argument("--message", default="", help="Additional message details") - - try: - parsed = parser.parse_args(args) - except SystemExit: - logger.warning("Invalid arguments for schedule create") - _error('Usage: schedule create "task" --due --to @branch [--message "details"]') - return False - - if not parsed.task: - _error("Task description is required") - console.print('[dim]Usage: schedule create "task" --due --to @branch[/dim]') - return False - - # Parse and validate due date - due_date = parse_due_date(parsed.due) - if not due_date: - _error(f"Invalid due date format: {parsed.due}") - console.print("[dim]Valid formats: 1d, 7d, 2w, 1m, or ISO date (2026-02-15)[/dim]") - return False - - # Create the task - try: - new_task = create_task(task=parsed.task, due_date=due_date, recipient=parsed.recipient, message=parsed.message) - task_id = new_task.get("id", "") - - _success(f"Scheduled task created: {task_id[:8]}") - console.print(f" [dim]Task:[/dim] {parsed.task}") - console.print(f" [dim]Due:[/dim] {due_date}") - console.print(f" [dim]To:[/dim] {parsed.recipient}") - if parsed.message: - console.print(f" [dim]Msg:[/dim] {parsed.message[:50]}...") - console.print() - - logger.info(f"[DAEMON] Scheduled task created: {task_id[:8]} -> {parsed.recipient}") - return True - - except Exception as e: - _error(f"Failed to create task: {e}") - logger.error(f"[DAEMON] Failed to create scheduled task: {e}", exc_info=True) - return False - - -def _handle_list(_args: List[str]) -> bool: - """Handle schedule list subcommand.""" - try: - tasks = load_tasks() - _print_task_list(tasks) - logger.info("[DAEMON] schedule: Task list displayed") - return True - - except Exception as e: - _error(f"Failed to load tasks: {e}") - logger.error(f"[DAEMON] Failed to load scheduled tasks: {e}", exc_info=True) - return False - - -def _handle_delete(args: List[str]) -> bool: - """Handle schedule delete subcommand.""" - if not args: - _error("Task ID is required") - console.print("[dim]Usage: schedule delete [/dim]") - return False - - task_id = args[0] - - try: - deleted = delete_task(task_id) - if deleted: - _success(f"Task deleted: {task_id[:8]}") - logger.info(f"[DAEMON] Scheduled task deleted: {task_id[:8]}") - return True - else: - _error(f"Task not found: {task_id[:8]}") - return False - - except Exception as e: - _error(f"Failed to delete task: {e}") - logger.error(f"[DAEMON] Failed to delete scheduled task: {e}", exc_info=True) - return False - - -def _handle_run_due(_args: List[str]) -> bool: - """Handle schedule run-due subcommand with single-instance lock.""" - if not FILELOCK_AVAILABLE: - console.print("[dim]filelock not available, running without lock.[/dim]") - return _process_due_tasks() - - lock_file = JSON_DIR / "schedule.lock" - ensure_lock_dir() - - # Try to acquire lock (non-blocking) - # FILELOCK_AVAILABLE guard above ensures these are not None - lock = FileLock(lock_file, timeout=LOCK_ACQUIRE_TIMEOUT) # type: ignore[misc] - try: - with lock.acquire(timeout=LOCK_ACQUIRE_TIMEOUT): - return _process_due_tasks() - except Timeout: # type: ignore[misc] - logger.warning("Schedule run-due already in progress, skipping") - console.print("[dim]Schedule run-due already in progress, skipping.[/dim]") - return True - - -def _display_task_result(task_result: dict) -> None: - """Display a single processed task result.""" - task_id = task_result.get("id", "")[:8] - recipient = task_result.get("recipient", "") - task_desc = task_result.get("task", "")[:40] - status = task_result.get("status", "") - - if status == "sent": - _success(f"Sent to {recipient}: {task_desc}") - logger.info(f"[DAEMON] Scheduled email sent: {task_id} -> {recipient}") - elif status == "skipped": - _error(f"ai_mail not available, cannot send to {recipient}") - elif status == "failed": - _error(f"Failed to send to {recipient}: {task_desc}") - logger.error(f"[DAEMON] Scheduled email failed: {task_id} -> {recipient}") - elif status == "error": - _error(f"Error sending to {recipient}: {task_result.get('error', '')}") - logger.error(f"[DAEMON] Scheduled email error: {task_id} -> {recipient}: {task_result.get('error', '')}") - - -def _process_due_tasks() -> bool: - """Process due tasks -- delegates to handler, formats output.""" - try: - # Delegate to handler for all implementation logic - email_fn = send_email_direct if AI_MAIL_AVAILABLE else None - results = process_due_tasks_batch(send_email_fn=email_fn, stale_max_age=STALE_DISPATCH_MAX_AGE) - - # Display results (module responsibility) - if results["recovered"]: - console.print(f"[dim]Recovered {results['recovered']} stale dispatch(es)[/dim]") - - if results["due"] == 0: - console.print("[dim]No tasks due at this time.[/dim]") - return True - - console.print() - _header(f"Running {results['due']} Due Task(s)") - console.print() - - for task_result in results.get("processed_tasks", []): - _display_task_result(task_result) - - console.print() - console.print(f"[bold]Results:[/bold] {results['success']} sent, {results['failed']} failed") - console.print() - - if results["failed"] > 0: - logger.warning("[DAEMON] %d scheduled task(s) failed to send", results["failed"]) - else: - logger.info("[DAEMON] schedule: Processed due tasks") - return True # Command was handled (failures are logged, not routing errors) - - except Exception as e: - _error(f"Failed to run due tasks: {e}") - logger.error(f"[DAEMON] Failed to run due tasks: {e}", exc_info=True) - return False - - -# ============================================= -# ORCHESTRATION -# ============================================= def handle_command(command: str, args: List[str]) -> bool: - """ - Handle 'schedule' command. - - Args: - command: Command name (should be 'schedule') - args: Command arguments (subcommand + subcommand args) - - Returns: - True if handled, False otherwise - """ + """Handle 'schedule' command — retired, shows migration notice.""" if command != "schedule": return False - try: - # No args -- introspection gate - if not args: - print_introspection() - return True + if not args: + print_introspection() + return True - # Handle help flag - if args[0] in ["--help", "-h", "help"]: - _print_help() - return True + if args[0] in ("--help", "-h", "help"): + print_introspection() + return True - subcommand = args[0] - subargs = args[1:] - - json_handler.log_operation("schedule_command", {"subcommand": args[0] if args else "list"}) - - # Route to subcommand handlers - if subcommand == "create": - return _handle_create(subargs) - if subcommand == "list": - return _handle_list(subargs) - if subcommand == "delete": - return _handle_delete(subargs) - if subcommand == "run-due": - return _handle_run_due(subargs) - - _error(f"Unknown subcommand: {subcommand}") - console.print("[dim]Run 'schedule --help' for available commands[/dim]") - return False - - except Exception as e: - logger.error(f"[DAEMON] Error in schedule command: {e}", exc_info=True) - _error(f"Error: {e}") - return False - - -# ============================================= -# MAIN ENTRY -# ============================================= - - -def main() -> None: - """Main entry point for direct execution.""" - args = sys.argv[1:] - - if len(args) == 0 or args[0] in ["--help", "-h", "help"]: - _print_help() - return - - # First arg is subcommand when called directly - handle_command("schedule", args) - - -if __name__ == "__main__": - main() + json_handler.log_operation("schedule_command_retired", {"args": args[:2]}) + logger.info("[schedule] Retired CLI invoked") + print_introspection() + return True diff --git a/src/aipass/daemon/tests/test_actions_module.py b/src/aipass/daemon/tests/test_actions_module.py deleted file mode 100644 index 7451f5c3..00000000 --- a/src/aipass/daemon/tests/test_actions_module.py +++ /dev/null @@ -1,456 +0,0 @@ -# =================== AIPass ==================== -# Name: test_actions_module.py -# Description: Tests for the actions CLI module -# Version: 1.0.0 -# Created: 2026-04-02 -# Modified: 2026-04-02 -# ============================================= - -"""Tests for the actions CLI module (apps/modules/actions.py).""" - -from datetime import datetime, timedelta -from unittest.mock import patch - -MODULE = "aipass.daemon.apps.modules.actions" - - -# ============================================= -# FIXTURES -# ============================================= - - -def _make_action( - action_id: str = "0001", - name: str = "test_action", - enabled: bool = True, - schedule_type: str = "daily", - time: str = "08:00", - action_type: str = "schedule", - target_branch: str = "@seedgo", - interval_minutes: int | None = None, - due_date: str | None = None, - prompt: str = "Run tests", -) -> dict: - """Build a sample action dict for tests.""" - action: dict = { - "id": action_id, - "name": name, - "enabled": enabled, - "schedule_type": schedule_type, - "time": time, - "type": action_type, - "target_branch": target_branch, - "prompt": prompt, - "created": "2026-03-01T00:00:00", - "last_run": None, - } - if interval_minutes is not None: - action["interval_minutes"] = interval_minutes - if due_date is not None: - action["due_date"] = due_date - return action - - -# ============================================= -# handle_command — routing -# ============================================= - - -@patch(f"{MODULE}.json_handler") -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -class TestHandleCommand: - """Tests for handle_command routing.""" - - def test_wrong_command_returns_false(self, _err, _con, _jh): - from aipass.daemon.apps.modules.actions import handle_command - - assert handle_command("not_actions", []) is False - - def test_no_args_shows_introspection(self, _err, mock_console, _jh): - from aipass.daemon.apps.modules.actions import handle_command - - result = handle_command("actions", []) - assert result is True - # introspection prints "actions Module" - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("actions Module" in c for c in calls) - - def test_help_flag(self, _err, mock_console, _jh): - from aipass.daemon.apps.modules.actions import handle_command - - assert handle_command("actions", ["--help"]) is True - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("USAGE" in c for c in calls) - - def test_help_word(self, _err, mock_console, _jh): - from aipass.daemon.apps.modules.actions import handle_command - - assert handle_command("actions", ["help"]) is True - - @patch(f"{MODULE}.list_actions", return_value=[]) - @patch(f"{MODULE}.next_due_str", return_value="--") - def test_list_subcommand(self, _nds, _la, _err, mock_console, mock_jh): - from aipass.daemon.apps.modules.actions import handle_command - - assert handle_command("actions", ["list"]) is True - mock_jh.log_operation.assert_called_once() - - @patch(f"{MODULE}.migrate_plugins", return_value=2) - @patch(f"{MODULE}.list_actions", return_value=[]) - @patch(f"{MODULE}.next_due_str", return_value="--") - def test_migrate_subcommand(self, _nds, _la, mock_migrate, _err, _con, mock_jh): - from aipass.daemon.apps.modules.actions import handle_command - - assert handle_command("actions", ["migrate"]) is True - mock_migrate.assert_called_once() - - @patch(f"{MODULE}.get_action") - @patch(f"{MODULE}.delete_action") - def test_delete_with_valid_id(self, mock_del, mock_get, _err, _con, _jh): - from aipass.daemon.apps.modules.actions import handle_command - - mock_get.return_value = _make_action() - assert handle_command("actions", ["delete", "0001"]) is True - mock_del.assert_called_once_with("0001") - - def test_delete_missing_id(self, mock_err, _con, _jh): - from aipass.daemon.apps.modules.actions import handle_command - - assert handle_command("actions", ["delete"]) is True - mock_err.assert_called() - - @patch(f"{MODULE}.create_action") - @patch(f"{MODULE}._parse_date", return_value="2026-04-09") - def test_set_reminder_valid(self, _pd, mock_create, _err, _con, _jh): - from aipass.daemon.apps.modules.actions import handle_command - - mock_create.return_value = _make_action(action_id="0099") - assert handle_command("actions", ["set", "reminder", "7d", "Check PR"]) is True - mock_create.assert_called_once() - - @patch(f"{MODULE}.create_action") - @patch(f"{MODULE}._parse_date", return_value="2026-04-09") - def test_set_schedule_valid(self, _pd, mock_create, _err, _con, _jh): - from aipass.daemon.apps.modules.actions import handle_command - - mock_create.return_value = _make_action(action_id="0088") - assert handle_command("actions", ["set", "schedule", "@seedgo", "Run audit", "daily", "04:00"]) is True - mock_create.assert_called_once() - - @patch(f"{MODULE}.get_action") - @patch(f"{MODULE}.next_due_str", return_value="--") - def test_action_id_routes(self, _nds, mock_get, _err, _con, _jh): - from aipass.daemon.apps.modules.actions import handle_command - - mock_get.return_value = _make_action(action_id="0003") - assert handle_command("actions", ["0003", "info"]) is True - mock_get.assert_called_with("0003") - - def test_unknown_subcommand(self, mock_err, mock_console, _jh): - from aipass.daemon.apps.modules.actions import handle_command - - assert handle_command("actions", ["foobar"]) is True - mock_err.assert_called() - - -# ============================================= -# _handle_toggle -# ============================================= - - -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -class TestHandleToggle: - @patch(f"{MODULE}.toggle_action") - @patch(f"{MODULE}.get_action") - def test_enable_success(self, mock_get, mock_toggle, _err, _con): - from aipass.daemon.apps.modules.actions import _handle_toggle - - mock_get.return_value = _make_action() - assert _handle_toggle("0001", True) is True - mock_toggle.assert_called_once_with("0001", True) - - @patch(f"{MODULE}.toggle_action") - @patch(f"{MODULE}.get_action") - def test_disable_success(self, mock_get, mock_toggle, _err, _con): - from aipass.daemon.apps.modules.actions import _handle_toggle - - mock_get.return_value = _make_action() - assert _handle_toggle("0001", False) is True - mock_toggle.assert_called_once_with("0001", False) - - @patch(f"{MODULE}.get_action", return_value=None) - def test_not_found(self, _get, mock_err, _con): - from aipass.daemon.apps.modules.actions import _handle_toggle - - assert _handle_toggle("9999", True) is True - mock_err.assert_called() - - -# ============================================= -# _handle_info -# ============================================= - - -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -class TestHandleInfo: - @patch(f"{MODULE}.next_due_str", return_value="--") - @patch(f"{MODULE}.get_action") - def test_info_success(self, mock_get, _nds, _err, mock_console): - from aipass.daemon.apps.modules.actions import _handle_info - - mock_get.return_value = _make_action() - assert _handle_info("0001") is True - # Should print detail header containing the action name - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("test_action" in c for c in calls) - - @patch(f"{MODULE}.get_action", return_value=None) - def test_info_not_found(self, _get, mock_err, _con): - from aipass.daemon.apps.modules.actions import _handle_info - - assert _handle_info("9999") is True - mock_err.assert_called() - - -# ============================================= -# _handle_set_reminder -# ============================================= - - -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -class TestHandleSetReminder: - def test_missing_args(self, mock_err, _con): - from aipass.daemon.apps.modules.actions import _handle_set_reminder - - assert _handle_set_reminder(["7d"]) is True - mock_err.assert_called() - - @patch(f"{MODULE}.create_action") - @patch(f"{MODULE}._parse_date", return_value="2026-04-09") - def test_with_to_flag(self, _pd, mock_create, _err, _con): - from aipass.daemon.apps.modules.actions import _handle_set_reminder - - mock_create.return_value = _make_action(action_id="0050") - assert _handle_set_reminder(["7d", "Follow up", "--to", "@flow"]) is True - call_kwargs = mock_create.call_args[1] - assert call_kwargs["target_branch"] == "@flow" - - @patch(f"{MODULE}._parse_date", return_value="") - def test_invalid_date(self, _pd, mock_err, _con): - from aipass.daemon.apps.modules.actions import _handle_set_reminder - - assert _handle_set_reminder(["xyz", "Some msg"]) is True - mock_err.assert_called() - - -# ============================================= -# _handle_set_schedule -# ============================================= - - -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -@patch(f"{MODULE}.logger") -class TestHandleSetSchedule: - def test_invalid_type(self, _log, mock_err, _con): - from aipass.daemon.apps.modules.actions import _handle_set_schedule - - assert _handle_set_schedule(["@branch", "prompt", "weekly"]) is True - mock_err.assert_called() - - def test_missing_time_arg(self, _log, mock_err, _con): - from aipass.daemon.apps.modules.actions import _handle_set_schedule - - assert _handle_set_schedule(["@branch", "prompt", "daily"]) is True - mock_err.assert_called() - - def test_interval_non_numeric(self, _log, mock_err, _con): - from aipass.daemon.apps.modules.actions import _handle_set_schedule - - assert _handle_set_schedule(["@b", "prompt", "interval", "abc"]) is True - mock_err.assert_called() - - @patch(f"{MODULE}.create_action") - def test_daily_success(self, mock_create, _log, _err, _con): - from aipass.daemon.apps.modules.actions import _handle_set_schedule - - mock_create.return_value = _make_action(action_id="0070") - assert _handle_set_schedule(["@seedgo", "Run audit", "daily", "04:00"]) is True - kw = mock_create.call_args[1] - assert kw["schedule_type"] == "daily" - assert kw["time"] == "04:00" - - @patch(f"{MODULE}.create_action") - def test_hourly_success(self, mock_create, _log, _err, _con): - from aipass.daemon.apps.modules.actions import _handle_set_schedule - - mock_create.return_value = _make_action(action_id="0071") - assert _handle_set_schedule(["@flow", "Check plans", "hourly", "30"]) is True - kw = mock_create.call_args[1] - assert kw["schedule_type"] == "hourly" - - @patch(f"{MODULE}.create_action") - def test_interval_success(self, mock_create, _log, _err, _con): - from aipass.daemon.apps.modules.actions import _handle_set_schedule - - mock_create.return_value = _make_action(action_id="0072") - assert _handle_set_schedule(["@daemon", "Heartbeat", "interval", "240"]) is True - kw = mock_create.call_args[1] - assert kw["interval_minutes"] == 240 - - -# ============================================= -# _handle_delete -# ============================================= - - -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -class TestHandleDelete: - def test_no_args(self, mock_err, _con): - from aipass.daemon.apps.modules.actions import _handle_delete - - assert _handle_delete([]) is True - mock_err.assert_called() - - @patch(f"{MODULE}.get_action", return_value=None) - def test_not_found(self, _get, mock_err, _con): - from aipass.daemon.apps.modules.actions import _handle_delete - - assert _handle_delete(["9999"]) is True - mock_err.assert_called() - - @patch(f"{MODULE}.delete_action") - @patch(f"{MODULE}.get_action") - def test_success(self, mock_get, mock_del, _err, _con): - from aipass.daemon.apps.modules.actions import _handle_delete - - mock_get.return_value = _make_action(action_id="0005") - assert _handle_delete(["0005"]) is True - mock_del.assert_called_once_with("0005") - - -# ============================================= -# _parse_date -# ============================================= - - -@patch(f"{MODULE}.logger") -class TestParseDate: - def test_relative_days(self, _log): - from aipass.daemon.apps.modules.actions import _parse_date - - result = _parse_date("7d") - expected = (datetime.now() + timedelta(days=7)).strftime("%Y-%m-%d") - assert result == expected - - def test_relative_weeks(self, _log): - from aipass.daemon.apps.modules.actions import _parse_date - - result = _parse_date("2w") - expected = (datetime.now() + timedelta(weeks=2)).strftime("%Y-%m-%d") - assert result == expected - - def test_iso_format(self, _log): - from aipass.daemon.apps.modules.actions import _parse_date - - assert _parse_date("2026-04-15") == "2026-04-15" - - def test_invalid_format(self, _log): - from aipass.daemon.apps.modules.actions import _parse_date - - assert _parse_date("not-a-date") == "" - - def test_invalid_relative_day(self, _log): - from aipass.daemon.apps.modules.actions import _parse_date - - assert _parse_date("xd") == "" - - def test_invalid_relative_week(self, _log): - from aipass.daemon.apps.modules.actions import _parse_date - - assert _parse_date("xw") == "" - - -# ============================================= -# _format_schedule -# ============================================= - - -class TestFormatSchedule: - def test_daily(self): - from aipass.daemon.apps.modules.actions import _format_schedule - - assert _format_schedule({"schedule_type": "daily", "time": "08:00"}) == "daily @ 08:00" - - def test_hourly(self): - from aipass.daemon.apps.modules.actions import _format_schedule - - assert _format_schedule({"schedule_type": "hourly", "time": "30"}) == "hourly @ :30" - - def test_interval_minutes(self): - from aipass.daemon.apps.modules.actions import _format_schedule - - assert _format_schedule({"schedule_type": "interval", "interval_minutes": 45}) == "every 45m" - - def test_interval_hours(self): - from aipass.daemon.apps.modules.actions import _format_schedule - - assert _format_schedule({"schedule_type": "interval", "interval_minutes": 120}) == "every 2h" - - def test_once(self): - from aipass.daemon.apps.modules.actions import _format_schedule - - assert _format_schedule({"schedule_type": "once", "due_date": "2026-04-10"}) == "once: 2026-04-10" - - def test_unknown_type(self): - from aipass.daemon.apps.modules.actions import _format_schedule - - assert _format_schedule({"schedule_type": "custom"}) == "custom" - - -# ============================================= -# _route_set_subcommand / _route_action_id -# ============================================= - - -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -class TestRouting: - def test_route_set_too_few_args(self, mock_err, _con): - from aipass.daemon.apps.modules.actions import _route_set_subcommand - - assert _route_set_subcommand(["set"]) is True - mock_err.assert_called() - - def test_route_set_unknown_type(self, mock_err, _con): - from aipass.daemon.apps.modules.actions import _route_set_subcommand - - assert _route_set_subcommand(["set", "bogus"]) is True - mock_err.assert_called() - - @patch(f"{MODULE}.get_action", return_value=None) - def test_route_action_id_no_sub_defaults_to_info(self, mock_get, mock_err, _con): - from aipass.daemon.apps.modules.actions import _route_action_id - - assert _route_action_id("0001", ["0001"]) is True - mock_get.assert_called_with("0001") - - @patch(f"{MODULE}.get_action") - @patch(f"{MODULE}.toggle_action") - def test_route_action_id_on(self, mock_toggle, mock_get, _err, _con): - from aipass.daemon.apps.modules.actions import _route_action_id - - mock_get.return_value = _make_action() - assert _route_action_id("0001", ["0001", "on"]) is True - mock_toggle.assert_called_once_with("0001", True) - - def test_route_action_id_unknown_sub(self, mock_err, _con): - from aipass.daemon.apps.modules.actions import _route_action_id - - assert _route_action_id("0001", ["0001", "banana"]) is True - mock_err.assert_called() diff --git a/src/aipass/daemon/tests/test_actions_registry.py b/src/aipass/daemon/tests/test_actions_registry.py deleted file mode 100644 index 221ff6d1..00000000 --- a/src/aipass/daemon/tests/test_actions_registry.py +++ /dev/null @@ -1,360 +0,0 @@ -# ===================AIPASS==================== -# META DATA HEADER -# Name: test_actions_registry.py - Action Registry Tests -# Date: 2026-03-02 -# Version: 1.1.0 -# Category: daemon/tests -# -# CHANGELOG (Max 5 entries): -# - v1.1.0 (2026-03-07): Adapted for AIPass public repo -# * Removed sys.path manipulation, uses package imports -# - v1.0.0 (2026-03-02): Initial creation - DPLAN-043 tests -# -# CODE STANDARDS: -# - Pytest conventions -# - Temp dir isolation (no writes to real registry) -# ============================================= - -"""Tests for the action registry handler.""" - -import json -from datetime import datetime, timedelta - -import pytest - -from aipass.daemon.apps.handlers.actions import actions_registry as _reg_mod - -create_action = _reg_mod.create_action -get_action = _reg_mod.get_action -list_actions = _reg_mod.list_actions -toggle_action = _reg_mod.toggle_action -delete_action = _reg_mod.delete_action -update_last_run = _reg_mod.update_last_run -mark_reminder_completed = _reg_mod.mark_reminder_completed -is_action_due = _reg_mod.is_action_due -calc_next_run = _reg_mod.calc_next_run -next_due_str = _reg_mod.next_due_str - - -@pytest.fixture(autouse=True) -def clean_registry(tmp_path): - """Isolate REGISTRY_FILE to a temp dir for every test.""" - test_registry = tmp_path / "actions_registry.json" - original = _reg_mod.REGISTRY_FILE - _reg_mod.REGISTRY_FILE = test_registry - yield test_registry - _reg_mod.REGISTRY_FILE = original - - -# ============================================= -# CRUD TESTS -# ============================================= - - -class TestCreate: - def test_create_action_basic(self, clean_registry): - """Create a simple schedule action and verify fields.""" - action = create_action( - name="test_audit", - action_type="schedule", - schedule_type="daily", - target_branch="@seedgo", - prompt="Run audit", - time="04:00", - fresh=True, - max_turns=20, - ) - assert action["id"] == "0001" - assert action["name"] == "test_audit" - assert action["type"] == "schedule" - assert action["schedule_type"] == "daily" - assert action["time"] == "04:00" - assert action["target_branch"] == "@seedgo" - assert action["enabled"] is True - assert action["last_run"] is None - assert action["completed"] is None - - def test_create_sequential_ids(self, clean_registry): - """IDs should be sequential: 0001, 0002, 0003...""" - a1 = create_action(name="first", action_type="schedule", schedule_type="daily") - a2 = create_action(name="second", action_type="schedule", schedule_type="daily") - a3 = create_action(name="third", action_type="reminder", schedule_type="once") - assert a1["id"] == "0001" - assert a2["id"] == "0002" - assert a3["id"] == "0003" - - def test_create_reminder(self, clean_registry): - """Create a one-shot reminder action.""" - action = create_action( - name="Check VERA progress", - action_type="reminder", - schedule_type="once", - target_branch="@devpulse", - prompt="Check VERA progress", - due_date="2026-03-11", - ) - assert action["type"] == "reminder" - assert action["schedule_type"] == "once" - assert action["due_date"] == "2026-03-11" - - def test_create_persists_to_json(self, clean_registry): - """Action should be persisted to the JSON file.""" - create_action(name="persisted", action_type="schedule", schedule_type="daily") - data = json.loads(clean_registry.read_text()) - assert len(data["actions"]) == 1 - assert data["actions"][0]["name"] == "persisted" - assert data["next_id"] == 2 - - -class TestGet: - def test_get_existing(self, clean_registry): - """Get an action by ID.""" - create_action(name="findme", action_type="schedule", schedule_type="daily") - action = get_action("0001") - assert action is not None - assert action["name"] == "findme" - - def test_get_missing(self, clean_registry): - """Get returns None for nonexistent ID.""" - assert get_action("9999") is None - - -class TestList: - def test_list_all(self, clean_registry): - """List returns all non-completed actions.""" - create_action(name="a", action_type="schedule", schedule_type="daily") - create_action(name="b", action_type="schedule", schedule_type="hourly") - actions = list_actions() - assert len(actions) == 2 - - def test_list_excludes_completed(self, clean_registry): - """Completed reminders should be excluded by default.""" - create_action(name="done", action_type="reminder", schedule_type="once", due_date="2026-01-01") - mark_reminder_completed("0001") - assert len(list_actions()) == 0 - assert len(list_actions(include_completed=True)) == 1 - - -class TestToggle: - def test_toggle_off(self, clean_registry): - """Toggle an action off.""" - create_action(name="toggleme", action_type="schedule", schedule_type="daily") - assert toggle_action("0001", False) is True - action = get_action("0001") - assert action is not None - assert action["enabled"] is False - - def test_toggle_on(self, clean_registry): - """Toggle an action back on.""" - create_action(name="toggleme", action_type="schedule", schedule_type="daily", enabled=False) - assert toggle_action("0001", True) is True - action = get_action("0001") - assert action is not None - assert action["enabled"] is True - - def test_toggle_missing(self, clean_registry): - """Toggle returns False for nonexistent ID.""" - assert toggle_action("9999", True) is False - - -class TestDelete: - def test_delete_existing(self, clean_registry): - """Delete an action by ID.""" - create_action(name="deleteme", action_type="schedule", schedule_type="daily") - assert delete_action("0001") is True - assert get_action("0001") is None - - def test_delete_missing(self, clean_registry): - """Delete returns False for nonexistent ID.""" - assert delete_action("9999") is False - - -# ============================================= -# DUE CHECKING TESTS -# ============================================= - - -class TestIsDue: - def test_daily_due_at_correct_time(self, clean_registry): - """Daily action is due when current time matches.""" - now = datetime.now() - action = { - "enabled": True, - "completed": None, - "schedule_type": "daily", - "time": f"{now.hour:02d}:{now.minute:02d}", - "last_run": None, - } - assert is_action_due(action) is True - - def test_daily_not_due_wrong_time(self, clean_registry): - """Daily action is not due at wrong time (12 hours away from now).""" - from datetime import datetime - - now = datetime.now() - # Pick a time 12 hours away — always outside the 15-min fuzzy window - far_hour = (now.hour + 12) % 24 - action = { - "enabled": True, - "completed": None, - "schedule_type": "daily", - "time": f"{far_hour:02d}:00", - "last_run": None, - } - assert is_action_due(action) is False - - def test_daily_not_due_already_ran_today(self, clean_registry): - """Daily action not due if already ran today.""" - now = datetime.now() - action = { - "enabled": True, - "completed": None, - "schedule_type": "daily", - "time": f"{now.hour:02d}:{now.minute:02d}", - "last_run": now.isoformat(), - } - assert is_action_due(action) is False - - def test_interval_due_never_run(self, clean_registry): - """Interval action is due if never run before.""" - action = { - "enabled": True, - "completed": None, - "schedule_type": "interval", - "interval_minutes": 60, - "last_run": None, - } - assert is_action_due(action) is True - - def test_interval_due_enough_time_elapsed(self, clean_registry): - """Interval action is due when enough time has passed.""" - past = (datetime.now() - timedelta(minutes=120)).isoformat() - action = { - "enabled": True, - "completed": None, - "schedule_type": "interval", - "interval_minutes": 60, - "last_run": past, - } - assert is_action_due(action) is True - - def test_interval_not_due_too_soon(self, clean_registry): - """Interval action is not due when too little time has passed.""" - recent = (datetime.now() - timedelta(minutes=5)).isoformat() - action = { - "enabled": True, - "completed": None, - "schedule_type": "interval", - "interval_minutes": 60, - "last_run": recent, - } - assert is_action_due(action) is False - - def test_once_due_past_date(self, clean_registry): - """Reminder is due when due_date is in the past.""" - action = { - "enabled": True, - "completed": None, - "schedule_type": "once", - "due_date": "2026-01-01", - } - assert is_action_due(action) is True - - def test_once_not_due_future_date(self, clean_registry): - """Reminder is not due when due_date is in the future.""" - action = { - "enabled": True, - "completed": None, - "schedule_type": "once", - "due_date": "2099-12-31", - } - assert is_action_due(action) is False - - def test_disabled_never_due(self, clean_registry): - """Disabled action is never due.""" - action = { - "enabled": False, - "completed": None, - "schedule_type": "interval", - "interval_minutes": 1, - "last_run": None, - } - assert is_action_due(action) is False - - def test_completed_never_due(self, clean_registry): - """Completed action is never due.""" - action = { - "enabled": True, - "completed": "2026-03-01T12:00:00", - "schedule_type": "once", - "due_date": "2026-01-01", - } - assert is_action_due(action) is False - - -# ============================================= -# NEXT RUN TESTS -# ============================================= - - -class TestCalcNextRun: - def test_daily_next_run(self, clean_registry): - """Daily action calculates next run correctly.""" - action = {"schedule_type": "daily", "time": "04:00", "last_run": None} - result = calc_next_run(action) - assert result is not None - assert "04:00:00" in result - - def test_interval_next_run(self, clean_registry): - """Interval action calculates next run from last_run + interval.""" - last = datetime.now().isoformat() - action = {"schedule_type": "interval", "interval_minutes": 60, "last_run": last} - result = calc_next_run(action) - assert result is not None - - def test_once_next_run(self, clean_registry): - """Reminder returns due_date as next run.""" - action = {"schedule_type": "once", "due_date": "2026-03-11", "completed": None} - assert calc_next_run(action) == "2026-03-11" - - -class TestNextDueStr: - def test_daily_str(self, clean_registry): - action = {"schedule_type": "daily", "time": "04:00"} - assert next_due_str(action) == "daily @ 04:00" - - def test_hourly_str(self, clean_registry): - action = {"schedule_type": "hourly", "time": "30"} - assert next_due_str(action) == "hourly @ :30" - - def test_once_str(self, clean_registry): - action = {"schedule_type": "once", "due_date": "2026-03-11"} - assert next_due_str(action) == "2026-03-11" - - -# ============================================= -# UPDATE TESTS -# ============================================= - - -class TestUpdateLastRun: - def test_update_last_run(self, clean_registry): - """Update last_run sets timestamp and recalculates next_run.""" - create_action(name="test", action_type="schedule", schedule_type="interval", interval_minutes=60) - ts = "2026-03-02T12:00:00" - assert update_last_run("0001", ts) is True - action = get_action("0001") - assert action is not None - assert action["last_run"] == ts - assert action["next_run"] is not None - - -class TestMarkCompleted: - def test_mark_reminder_completed(self, clean_registry): - """Marking a reminder completed sets completed timestamp and disables it.""" - create_action(name="reminder", action_type="reminder", schedule_type="once", due_date="2026-03-01") - assert mark_reminder_completed("0001") is True - action = get_action("0001") - assert action is not None - assert action["completed"] is not None - assert action["enabled"] is False diff --git a/src/aipass/daemon/tests/test_run_module.py b/src/aipass/daemon/tests/test_run_module.py index 99e65acb..c223d7de 100644 --- a/src/aipass/daemon/tests/test_run_module.py +++ b/src/aipass/daemon/tests/test_run_module.py @@ -1,3 +1,11 @@ +# =================== AIPass ==================== +# Name: test_run_module.py +# Description: Tests for the drone @daemon run module +# Version: 1.1.0 +# Created: 2026-06-15 +# Modified: 2026-06-25 +# ============================================= + """Tests for the drone @daemon run module (decentralized scheduler tick).""" from unittest.mock import patch @@ -63,7 +71,7 @@ class TestRunTick: assert results["due"] == 0 @patch("aipass.daemon.apps.modules.run.save_runstate") - @patch("aipass.daemon.apps.modules.run._fire_job", return_value=True) + @patch("aipass.daemon.apps.modules.run._fire_job", return_value=(True, "")) @patch("aipass.daemon.apps.modules.run.discover_jobs") @patch("aipass.daemon.apps.modules.run.load_runstate", return_value={"jobs": {}}) def test_fires_due_job(self, mock_rs, mock_discover, mock_fire, mock_save): @@ -84,7 +92,7 @@ class TestRunTick: mock_save.assert_called() @patch("aipass.daemon.apps.modules.run.save_runstate") - @patch("aipass.daemon.apps.modules.run._fire_job", return_value=False) + @patch("aipass.daemon.apps.modules.run._fire_job", return_value=(False, "wake failed")) @patch("aipass.daemon.apps.modules.run.discover_jobs") @patch("aipass.daemon.apps.modules.run.load_runstate", return_value={"jobs": {}}) def test_failed_fire_counted(self, mock_rs, mock_discover, mock_fire, mock_save): diff --git a/src/aipass/daemon/tests/test_schedule_module.py b/src/aipass/daemon/tests/test_schedule_module.py deleted file mode 100644 index 3f2bf336..00000000 --- a/src/aipass/daemon/tests/test_schedule_module.py +++ /dev/null @@ -1,352 +0,0 @@ -# =================== AIPass ==================== -# Name: test_schedule_module.py -# Description: Tests for the schedule CLI module -# Version: 1.0.0 -# Created: 2026-04-03 -# Modified: 2026-04-03 -# ============================================= - -"""Tests for the schedule CLI module (apps/modules/schedule.py).""" - -from unittest.mock import patch, MagicMock - -MODULE = "aipass.daemon.apps.modules.schedule" - - -# ============================================= -# handle_command -- routing basics -# ============================================= - - -@patch(f"{MODULE}.json_handler") -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -@patch(f"{MODULE}.logger") -class TestHandleCommandRouting: - """Tests for handle_command routing.""" - - def test_wrong_command_returns_false(self, _log, _err, _con, _jh): - from aipass.daemon.apps.modules.schedule import handle_command - - assert handle_command("not_schedule", []) is False - - def test_no_args_shows_introspection(self, _log, _err, mock_con, _jh): - from aipass.daemon.apps.modules.schedule import handle_command - - result = handle_command("schedule", []) - assert result is True - calls = [str(c) for c in mock_con.print.call_args_list] - assert any("schedule Module" in c for c in calls) - - def test_help_flag(self, _log, _err, mock_con, _jh): - from aipass.daemon.apps.modules.schedule import handle_command - - result = handle_command("schedule", ["--help"]) - assert result is True - calls = [str(c) for c in mock_con.print.call_args_list] - assert any("USAGE" in c for c in calls) - - def test_unknown_subcommand(self, _log, mock_err, _con, _jh): - from aipass.daemon.apps.modules.schedule import handle_command - - result = handle_command("schedule", ["foobar"]) - assert result is False - mock_err.assert_called() - - -# ============================================= -# handle_command -- list subcommand -# ============================================= - - -@patch(f"{MODULE}.json_handler") -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -@patch(f"{MODULE}.logger") -class TestListSubcommand: - """Tests for 'schedule list' subcommand.""" - - @patch(f"{MODULE}.load_tasks", return_value=[]) - def test_list_success(self, mock_load, _log, _err, mock_con, mock_jh): - from aipass.daemon.apps.modules.schedule import handle_command - - result = handle_command("schedule", ["list"]) - assert result is True - mock_load.assert_called_once() - - @patch(f"{MODULE}.load_tasks", side_effect=RuntimeError("disk error")) - def test_list_exception(self, _load, _log, mock_err, _con, _jh): - from aipass.daemon.apps.modules.schedule import handle_command - - result = handle_command("schedule", ["list"]) - assert result is False - mock_err.assert_called() - - -# ============================================= -# handle_command -- delete subcommand -# ============================================= - - -@patch(f"{MODULE}.json_handler") -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -@patch(f"{MODULE}.logger") -class TestDeleteSubcommand: - """Tests for 'schedule delete' subcommand.""" - - def test_delete_no_args_shows_error(self, _log, mock_err, _con, _jh): - from aipass.daemon.apps.modules.schedule import handle_command - - result = handle_command("schedule", ["delete"]) - assert result is False - mock_err.assert_called() - - @patch(f"{MODULE}.delete_task", return_value=True) - def test_delete_success(self, mock_del, _log, _err, _con, _jh): - from aipass.daemon.apps.modules.schedule import handle_command - - result = handle_command("schedule", ["delete", "abc123"]) - assert result is True - mock_del.assert_called_once_with("abc123") - - @patch(f"{MODULE}.delete_task", return_value=False) - def test_delete_not_found(self, mock_del, _log, mock_err, _con, _jh): - from aipass.daemon.apps.modules.schedule import handle_command - - result = handle_command("schedule", ["delete", "abc123"]) - assert result is False - mock_err.assert_called() - - -# ============================================= -# handle_command -- run-due subcommand -# ============================================= - - -@patch(f"{MODULE}.json_handler") -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -@patch(f"{MODULE}.logger") -class TestRunDueSubcommand: - """Tests for 'schedule run-due' subcommand.""" - - @patch( - f"{MODULE}.process_due_tasks_batch", - return_value={ - "recovered": 0, - "due": 0, - "success": 0, - "failed": 0, - "processed_tasks": [], - }, - ) - @patch(f"{MODULE}.FILELOCK_AVAILABLE", False) - def test_run_due_without_lock(self, mock_batch, _log, _err, mock_con, _jh): - from aipass.daemon.apps.modules.schedule import handle_command - - result = handle_command("schedule", ["run-due"]) - assert result is True - mock_batch.assert_called_once() - - @patch( - f"{MODULE}.process_due_tasks_batch", - return_value={ - "recovered": 1, - "due": 2, - "success": 1, - "failed": 1, - "processed_tasks": [ - {"id": "a1", "recipient": "@flow", "task": "Check plan", "status": "sent"}, - {"id": "a2", "recipient": "@seedgo", "task": "Audit", "status": "failed"}, - ], - }, - ) - @patch(f"{MODULE}.FILELOCK_AVAILABLE", False) - def test_run_due_processes_tasks(self, mock_batch, _log, _err, mock_con, _jh): - from aipass.daemon.apps.modules.schedule import handle_command - - result = handle_command("schedule", ["run-due"]) - assert result is True - mock_batch.assert_called_once() - calls = " ".join(str(c) for c in mock_con.print.call_args_list) - assert "1 sent" in calls - assert "1 failed" in calls - - -# ============================================= -# _handle_create -# ============================================= - - -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -@patch(f"{MODULE}.logger") -class TestHandleCreate: - """Tests for _handle_create.""" - - @patch(f"{MODULE}.create_task", return_value={"id": "task-001"}) - @patch(f"{MODULE}.parse_due_date", return_value="2026-04-10") - def test_create_valid(self, _due, mock_create, _log, _err, _con): - from aipass.daemon.apps.modules.schedule import _handle_create - - result = _handle_create(["Follow up", "--due", "7d", "--to", "@flow"]) - assert result is True - mock_create.assert_called_once() - - def test_create_missing_task(self, _log, mock_err, _con): - from aipass.daemon.apps.modules.schedule import _handle_create - - result = _handle_create(["--due", "7d", "--to", "@flow"]) - assert result is False - mock_err.assert_called() - - @patch(f"{MODULE}.parse_due_date", return_value=None) - def test_create_invalid_due(self, _due, _log, mock_err, _con): - from aipass.daemon.apps.modules.schedule import _handle_create - - result = _handle_create(["Task text", "--due", "xyz", "--to", "@flow"]) - assert result is False - mock_err.assert_called() - - -# ============================================= -# _process_due_tasks -# ============================================= - - -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -@patch(f"{MODULE}.logger") -class TestProcessDueTasks: - """Tests for _process_due_tasks.""" - - @patch( - f"{MODULE}.process_due_tasks_batch", - return_value={ - "recovered": 0, - "due": 0, - "success": 0, - "failed": 0, - "processed_tasks": [], - }, - ) - def test_no_tasks_due(self, mock_batch, _log, _err, mock_con): - from aipass.daemon.apps.modules.schedule import _process_due_tasks - - result = _process_due_tasks() - assert result is True - calls = " ".join(str(c) for c in mock_con.print.call_args_list) - assert "No tasks due" in calls - - @patch( - f"{MODULE}.process_due_tasks_batch", - return_value={ - "recovered": 0, - "due": 2, - "success": 1, - "failed": 1, - "processed_tasks": [ - {"id": "t1", "recipient": "@flow", "task": "Check", "status": "sent"}, - {"id": "t2", "recipient": "@seedgo", "task": "Audit", "status": "failed"}, - ], - }, - ) - def test_mix_sent_failed(self, mock_batch, _log, _err, mock_con): - from aipass.daemon.apps.modules.schedule import _process_due_tasks - - result = _process_due_tasks() - assert result is True - calls = " ".join(str(c) for c in mock_con.print.call_args_list) - assert "1 sent" in calls - assert "1 failed" in calls - - -# ============================================= -# _display_task_result -# ============================================= - - -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.cli_error") -@patch(f"{MODULE}.logger") -class TestDisplayTaskResult: - """Tests for _display_task_result per-status output.""" - - def test_status_sent(self, mock_log, _err, mock_con): - from aipass.daemon.apps.modules.schedule import _display_task_result - - _display_task_result( - { - "id": "t1", - "recipient": "@flow", - "task": "Check plan", - "status": "sent", - } - ) - calls = " ".join(str(c) for c in mock_con.print.call_args_list) - assert "OK" in calls or "Sent" in calls or "@flow" in calls - - def test_status_skipped(self, _log, mock_err, _con): - from aipass.daemon.apps.modules.schedule import _display_task_result - - _display_task_result( - { - "id": "t2", - "recipient": "@seedgo", - "task": "Audit", - "status": "skipped", - } - ) - mock_err.assert_called() - - def test_status_failed(self, _log, mock_err, _con): - from aipass.daemon.apps.modules.schedule import _display_task_result - - _display_task_result( - { - "id": "t3", - "recipient": "@daemon", - "task": "Heartbeat", - "status": "failed", - } - ) - mock_err.assert_called() - - def test_status_error(self, _log, mock_err, _con): - from aipass.daemon.apps.modules.schedule import _display_task_result - - _display_task_result( - { - "id": "t4", - "recipient": "@drone", - "task": "Ping", - "status": "error", - "error": "timeout", - } - ) - mock_err.assert_called() - - -# ============================================= -# _send_email_via_drone -# ============================================= - - -@patch(f"{MODULE}.logger") -class TestSendEmailViaDrone: - """Tests for _send_email_via_drone subprocess wrapper.""" - - @patch("subprocess.run") - def test_success(self, mock_run, _log): - from aipass.daemon.apps.modules.schedule import _send_email_via_drone - - mock_run.return_value = MagicMock(returncode=0) - assert _send_email_via_drone("@flow", "subj", "body") is True - mock_run.assert_called_once() - - @patch("subprocess.run", side_effect=OSError("no drone")) - def test_failure(self, _run, _log): - from aipass.daemon.apps.modules.schedule import _send_email_via_drone - - assert _send_email_via_drone("@flow", "subj", "body") is False diff --git a/src/aipass/daemon/tests/test_scheduler_bot.py b/src/aipass/daemon/tests/test_scheduler_bot.py new file mode 100644 index 00000000..2af59e22 --- /dev/null +++ b/src/aipass/daemon/tests/test_scheduler_bot.py @@ -0,0 +1,356 @@ +# =================== AIPass ==================== +# Name: test_scheduler_bot.py +# Description: Tests for TDPLAN-0008 Phase 1 — scheduler bot daemon layer +# Version: 1.0.0 +# Created: 2026-06-25 +# Modified: 2026-06-25 +# ============================================= + +"""Tests for TDPLAN-0008 Phase 1: status capture, queue view, lifecycle notifications, archive.""" + +from datetime import datetime +from pathlib import Path +from unittest.mock import patch + +import pytest + +from aipass.daemon.apps.handlers.schedule.runstate import ( + update_job_runstate, + record_job_failure, +) +from aipass.daemon.apps.handlers.schedule.telegram_notifier import ( + notify_triggered, + notify_complete, + notify_error, +) +from aipass.daemon.apps.modules.queue import ( + _build_queue, + _build_json_output, + _schedule_human, + handle_command, +) + + +# ── Fixtures ────────────────────────────────────────── + + +@pytest.fixture +def once_job(): + """A once-type job that fires on due_date.""" + return { + "owner": "@api", + "id": "data-check", + "enabled": True, + "schedule": {"type": "once", "due_date": datetime.now().strftime("%Y-%m-%d")}, + "wake": {"fresh": True, "model": "haiku"}, + "prompt": "Check the live data and report back.", + } + + +@pytest.fixture +def interval_job(): + """An interval-type job.""" + return { + "owner": "@commons", + "id": "rotation", + "enabled": True, + "schedule": {"type": "interval", "interval_minutes": 120}, + "wake": {"fresh": True}, + "prompt": "Rotate community content.", + } + + +# ── Test 1: once job fires on due_date, then marks completed ─── + + +class TestOnceJobLifecycle: + """Verify once jobs fire on/after due_date and mark completed.""" + + def test_fires_on_due_date(self, once_job): + """Once job with today's due_date gets last_status=success + completed.""" + runstate = {"jobs": {}} + update_job_runstate(runstate, "@api", "data-check", once_job["schedule"]) + entry = runstate["jobs"]["@api/data-check"] + assert entry["last_status"] == "success" + assert "completed" in entry + assert entry["last_success_at"] is not None + assert entry["last_error"] is None + + def test_completed_once_excluded_from_queue(self, once_job): + """Completed once jobs are filtered out of the queue view.""" + runstate = { + "jobs": { + "@api/data-check": { + "completed": "2026-06-25T10:00:00", + "last_run": "2026-06-25T10:00:00", + } + } + } + entries = _build_queue([once_job], runstate) + assert len(entries) == 0 + + +# ── Test 2: fire emits notify_triggered then notify_complete; failed emits notify_error ── + + +class TestLifecycleNotifications: + """Verify telegram notifications fire on real job events.""" + + @patch("aipass.daemon.apps.handlers.schedule.telegram_notifier._send") + def test_triggered_sends_running(self, mock_send): + """notify_triggered sends running ping.""" + mock_send.return_value = True + result = notify_triggered("@api", "data-check") + assert result is True + call_msg = mock_send.call_args[0][0] + assert "@api/data-check" in call_msg + assert "running" in call_msg + + @patch("aipass.daemon.apps.handlers.schedule.telegram_notifier._send") + def test_complete_sends_dispatched(self, mock_send): + """notify_complete sends dispatched ping with summary.""" + mock_send.return_value = True + result = notify_complete("@api", "data-check", "Agent spawned OK") + assert result is True + call_msg = mock_send.call_args[0][0] + assert "dispatched" in call_msg + assert "Agent spawned OK" in call_msg + + @patch("aipass.daemon.apps.handlers.schedule.telegram_notifier._send") + def test_error_sends_failed(self, mock_send): + """notify_error sends failed ping with error detail.""" + mock_send.return_value = True + result = notify_error("@api", "data-check", "timeout") + assert result is True + call_msg = mock_send.call_args[0][0] + assert "FAILED" in call_msg + assert "timeout" in call_msg + + +# ── Test 3: runstate records last_status + last_error on both paths ── + + +class TestRunstateStatusCapture: + """Verify runstate tracks status on success AND failure.""" + + def test_success_records_status(self): + """Successful fire sets last_status=success, last_success_at, clears error.""" + runstate = {"jobs": {}} + schedule = {"type": "interval", "interval_minutes": 60} + update_job_runstate(runstate, "@commons", "test", schedule) + entry = runstate["jobs"]["@commons/test"] + assert entry["last_status"] == "success" + assert entry["last_success_at"] is not None + assert entry["last_error"] is None + + def test_failure_records_status(self): + """Failed fire sets last_status=failed, last_failure_at, last_error.""" + runstate = {"jobs": {}} + record_job_failure(runstate, "@commons", "test", "branch locked") + entry = runstate["jobs"]["@commons/test"] + assert entry["last_status"] == "failed" + assert entry["last_failure_at"] is not None + assert entry["last_error"] == "branch locked" + + def test_failure_preserves_last_run(self): + """Failure path still records last_run timestamp.""" + runstate = {"jobs": {}} + record_job_failure(runstate, "@x", "y", "err") + entry = runstate["jobs"]["@x/y"] + assert "last_run" in entry + + def test_error_truncated(self): + """Long error messages are truncated to 500 chars.""" + runstate = {"jobs": {}} + record_job_failure(runstate, "@x", "y", "x" * 1000) + entry = runstate["jobs"]["@x/y"] + assert len(entry["last_error"]) == 500 + + +# ── Test 4: queue --json returns frozen schema ── + + +class TestQueueJsonSchema: + """Verify queue --json output matches the frozen contract.""" + + def test_schema_structure(self, interval_job): + """JSON output has generated_at, count, jobs array.""" + runstate = {"jobs": {}} + entries = _build_queue([interval_job], runstate) + output = _build_json_output(entries) + assert "generated_at" in output + assert "count" in output + assert isinstance(output["jobs"], list) + assert output["count"] == len(output["jobs"]) + + def test_job_fields(self, interval_job): + """Each job in output has all frozen-schema fields.""" + runstate = {"jobs": {}} + entries = _build_queue([interval_job], runstate) + job_out = entries[0] + required_fields = [ + "owner", + "id", + "enabled", + "type", + "schedule_human", + "next_run", + "last_run", + "last_status", + "last_error", + "prompt_preview", + "wake", + ] + for field in required_fields: + assert field in job_out, f"Missing field: {field}" + + def test_type_values(self, once_job, interval_job): + """Type field matches schedule type.""" + runstate = {"jobs": {}} + once_entries = _build_queue([once_job], runstate) + assert once_entries[0]["type"] == "once" + interval_entries = _build_queue([interval_job], runstate) + assert interval_entries[0]["type"] == "interval" + + def test_schedule_human_formats(self): + """schedule_human renders each type correctly.""" + assert _schedule_human({"schedule": {"type": "once", "due_date": "2026-07-02"}}) == "2026-07-02" + assert _schedule_human({"schedule": {"type": "daily", "time": "04:00"}}) == "daily @ 04:00" + assert _schedule_human({"schedule": {"type": "hourly", "time": "30"}}) == "hourly @ :30" + assert _schedule_human({"schedule": {"type": "interval", "interval_minutes": 120}}) == "every 2h" + assert _schedule_human({"schedule": {"type": "interval", "interval_minutes": 30}}) == "every 30m" + + +# ── Test 5: empty tick emits ZERO telegram calls ── + + +class TestEmptyTickNoNotify: + """Verify no telegram calls on empty ticks (no due jobs).""" + + @patch("aipass.daemon.apps.modules.run.discover_jobs", return_value=[]) + @patch("aipass.daemon.apps.handlers.schedule.telegram_notifier._send") + def test_no_jobs_no_send(self, mock_send, mock_discover): + """Empty tick with no discovered jobs makes zero telegram calls.""" + from aipass.daemon.apps.modules.run import run_tick + + run_tick(dry_run=True) + mock_send.assert_not_called() + + @patch("aipass.daemon.apps.modules.run.discover_jobs") + @patch("aipass.daemon.apps.modules.run.load_runstate") + @patch("aipass.daemon.apps.handlers.schedule.telegram_notifier._send") + def test_no_due_jobs_no_send(self, mock_send, mock_rs, mock_discover): + """Tick with jobs but none due makes zero telegram calls.""" + from aipass.daemon.apps.modules.run import run_tick + + mock_discover.return_value = [ + { + "owner": "@commons", + "id": "test", + "enabled": True, + "schedule": {"type": "interval", "interval_minutes": 9999}, + "wake": {}, + "prompt": "test", + } + ] + mock_rs.return_value = {"jobs": {"@commons/test": {"last_run": datetime.now().isoformat()}}} + run_tick() + mock_send.assert_not_called() + + +# ── Test 6: send is fail-soft ── + + +class TestFailSoft: + """Verify telegram send failures don't block job firing.""" + + @patch("aipass.daemon.apps.handlers.schedule.telegram_notifier._send") + def test_send_returns_false_on_failure(self, mock_send): + """When _send fails, notify functions return False gracefully.""" + mock_send.return_value = False + assert notify_triggered("@x", "y") is False + assert notify_complete("@x", "y", "s") is False + assert notify_error("@x", "y", "e") is False + + @patch( + "aipass.skills.lib.telegram.apps.handlers.notifier.send_telegram_notification", + side_effect=Exception("connection refused"), + ) + def test_exception_caught(self, mock_notifier): + """Exception in send_telegram_notification is caught, returns False.""" + from aipass.daemon.apps.handlers.schedule.telegram_notifier import _send + + result = _send("test message") + assert result is False + + @patch("aipass.daemon.apps.modules.run.save_runstate") + @patch("aipass.daemon.apps.modules.run.record_job_failure") + @patch("aipass.daemon.apps.modules.run.update_job_runstate") + @patch("aipass.daemon.apps.modules.run.discover_jobs") + @patch("aipass.daemon.apps.modules.run.load_runstate", return_value={"jobs": {}}) + def test_fire_continues_when_notify_fails(self, mock_rs, mock_discover, mock_update, mock_fail, mock_save): + """Job fires and records status even when telegram is down.""" + from aipass.daemon.apps.modules.run import run_tick + + mock_discover.return_value = [ + { + "owner": "@commons", + "id": "test", + "enabled": True, + "schedule": {"type": "interval", "interval_minutes": 1}, + "wake": {"fresh": True}, + "prompt": "test", + "notify": True, + } + ] + + with patch("aipass.daemon.apps.modules.run._fire_job", return_value=(True, "")): + results = run_tick() + assert results["fired"] == 1 + + +# ── Test 7: dormant registries archived ── + + +class TestDormantArchived: + """Verify dormant registries are archived and no live code references them.""" + + def test_data_files_archived(self): + """schedule.json and actions_registry.json moved to .archive.""" + daemon_root = Path(__file__).resolve().parents[1] + archive_dir = daemon_root / "daemon_json" / ".archive" + assert (archive_dir / "schedule.json").exists() + assert (archive_dir / "actions_registry.json").exists() + + def test_original_data_files_gone(self): + """Original data files no longer at daemon_json/ root.""" + daemon_root = Path(__file__).resolve().parents[1] + assert not (daemon_root / "daemon_json" / "schedule.json").exists() + assert not (daemon_root / "daemon_json" / "actions_registry.json").exists() + + def test_handler_files_archived(self): + """task_registry.py and actions_registry.py moved to .archive.""" + daemon_root = Path(__file__).resolve().parents[1] + sched_archive = daemon_root / "apps" / "handlers" / "schedule" / ".archive" + actions_archive = daemon_root / "apps" / "handlers" / "actions" / ".archive" + assert (sched_archive / "task_registry.py").exists() + assert (actions_archive / "actions_registry.py").exists() + + def test_schedule_module_retired(self): + """Schedule module handle_command shows retirement notice, not old CRUD.""" + from aipass.daemon.apps.modules.schedule import handle_command as sched_cmd + + assert sched_cmd("schedule", []) is True + assert sched_cmd("schedule", ["create", "test"]) is True + + def test_actions_module_retired(self): + """Actions module handle_command shows retirement notice, not old CRUD.""" + from aipass.daemon.apps.modules.actions import handle_command as act_cmd + + assert act_cmd("actions", []) is True + assert act_cmd("actions", ["list"]) is True + + def test_queue_command_wired(self): + """drone @daemon queue is routable.""" + assert handle_command("queue", []) is True + assert handle_command("notqueue", []) is False diff --git a/src/aipass/daemon/tests/test_scheduler_ops.py b/src/aipass/daemon/tests/test_scheduler_ops.py deleted file mode 100644 index 9be64627..00000000 --- a/src/aipass/daemon/tests/test_scheduler_ops.py +++ /dev/null @@ -1,117 +0,0 @@ -# =================== AIPass ==================== -# Name: test_scheduler_ops.py -# Description: Tests for the scheduler_ops facade module -# Version: 1.0.0 -# Created: 2026-04-03 -# Modified: 2026-04-03 -# ============================================= - -"""Tests for the scheduler_ops facade module (apps/modules/scheduler_ops.py).""" - -from unittest.mock import patch - -MODULE = "aipass.daemon.apps.modules.scheduler_ops" - - -# ============================================= -# handle_command — routing -# ============================================= - - -@patch(f"{MODULE}.json_handler") -@patch(f"{MODULE}.console") -@patch(f"{MODULE}.logger") -class TestHandleCommand: - """Tests for handle_command routing.""" - - def test_wrong_command_returns_false(self, _log, _con, _jh): - from aipass.daemon.apps.modules.scheduler_ops import handle_command - - assert handle_command("not-scheduler-ops", []) is False - - def test_no_args_shows_introspection(self, _log, mock_console, _jh): - from aipass.daemon.apps.modules.scheduler_ops import handle_command - - result = handle_command("scheduler-ops", []) - assert result is True - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("scheduler_ops Module" in c for c in calls) - - def test_help_flag_shows_introspection(self, _log, mock_console, _jh): - from aipass.daemon.apps.modules.scheduler_ops import handle_command - - assert handle_command("scheduler-ops", ["--help"]) is True - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("scheduler_ops Module" in c for c in calls) - - def test_h_flag_shows_introspection(self, _log, mock_console, _jh): - from aipass.daemon.apps.modules.scheduler_ops import handle_command - - assert handle_command("scheduler-ops", ["-h"]) is True - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("scheduler_ops Module" in c for c in calls) - - def test_help_word_shows_introspection(self, _log, mock_console, _jh): - from aipass.daemon.apps.modules.scheduler_ops import handle_command - - assert handle_command("scheduler-ops", ["help"]) is True - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("scheduler_ops Module" in c for c in calls) - - def test_status_arg_shows_registry_info(self, _log, mock_console, mock_jh): - from aipass.daemon.apps.modules.scheduler_ops import handle_command - - assert handle_command("scheduler-ops", ["status"]) is True - mock_jh.log_operation.assert_called_once_with("scheduler_ops_status") - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("Scheduler Ops" in c for c in calls) - - def test_status_prints_task_registry_availability(self, _log, mock_console, mock_jh): - from aipass.daemon.apps.modules.scheduler_ops import handle_command - - handle_command("scheduler-ops", ["status"]) - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("Task registry" in c for c in calls) - - def test_status_prints_action_registry_availability(self, _log, mock_console, mock_jh): - from aipass.daemon.apps.modules.scheduler_ops import handle_command - - handle_command("scheduler-ops", ["status"]) - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("Action registry" in c for c in calls) - - -# ============================================= -# Module-level availability flags -# ============================================= - - -class TestRegistryAvailability: - """Verify that registry imports succeed in the test environment.""" - - def test_task_registry_available(self): - from aipass.daemon.apps.modules.scheduler_ops import TASK_REGISTRY_AVAILABLE - - assert TASK_REGISTRY_AVAILABLE is True - - def test_action_registry_available(self): - from aipass.daemon.apps.modules.scheduler_ops import ACTION_REGISTRY_AVAILABLE - - assert ACTION_REGISTRY_AVAILABLE is True - - -# ============================================= -# print_introspection -# ============================================= - - -@patch(f"{MODULE}.console") -class TestPrintIntrospection: - """Tests for print_introspection output.""" - - def test_prints_module_header(self, mock_console): - from aipass.daemon.apps.modules.scheduler_ops import print_introspection - - print_introspection() - calls = [str(c) for c in mock_console.print.call_args_list] - assert any("scheduler_ops Module" in c for c in calls) diff --git a/src/aipass/daemon/tests/test_task_registry.py b/src/aipass/daemon/tests/test_task_registry.py deleted file mode 100644 index d27696d5..00000000 --- a/src/aipass/daemon/tests/test_task_registry.py +++ /dev/null @@ -1,588 +0,0 @@ -# ===================AIPASS==================== -# META DATA HEADER -# Name: test_task_registry.py - Task Registry Tests -# Date: 2026-03-24 -# Version: 1.0.0 -# Category: daemon/tests -# -# CHANGELOG (Max 5 entries): -# - v1.0.0 (2026-03-24): Initial creation - task_registry handler tests -# -# CODE STANDARDS: -# - Pytest conventions -# - Temp dir isolation (no writes to real registry) -# ============================================= - -"""Tests for the scheduled task registry handler.""" - -import json -from datetime import datetime, timedelta -from unittest.mock import patch - -import pytest - -from aipass.daemon.apps.handlers.schedule import task_registry as _mod - -parse_due_date = _mod.parse_due_date -create_task = _mod.create_task -load_tasks = _mod.load_tasks -save_tasks = _mod.save_tasks -get_due_tasks = _mod.get_due_tasks -mark_dispatching = _mod.mark_dispatching -mark_completed = _mod.mark_completed -mark_pending = _mod.mark_pending -recover_stale_dispatches = _mod.recover_stale_dispatches -delete_task = _mod.delete_task -get_task_by_id = _mod.get_task_by_id -get_pending_tasks = _mod.get_pending_tasks -ensure_lock_dir = _mod.ensure_lock_dir - - -@pytest.fixture(autouse=True) -def isolate_registry(tmp_path): - """Redirect SCHEDULE_JSON_PATH to a temp dir for every test.""" - test_file = tmp_path / "schedule.json" - original = _mod.SCHEDULE_JSON_PATH - _mod.SCHEDULE_JSON_PATH = test_file - yield test_file - _mod.SCHEDULE_JSON_PATH = original - - -# ============================================= -# DATE PARSING TESTS -# ============================================= - - -class TestParseDueDate: - def test_days_format(self): - """'7d' should resolve to 7 days from today.""" - result = parse_due_date("7d") - expected = (datetime.now().date() + timedelta(days=7)).isoformat() - assert result == expected - - def test_days_format_single_digit(self): - """'1d' should resolve to tomorrow.""" - result = parse_due_date("1d") - expected = (datetime.now().date() + timedelta(days=1)).isoformat() - assert result == expected - - def test_weeks_format(self): - """'1w' should resolve to 1 week from today.""" - result = parse_due_date("1w") - expected = (datetime.now().date() + timedelta(weeks=1)).isoformat() - assert result == expected - - def test_weeks_format_multiple(self): - """'2w' should resolve to 2 weeks from today.""" - result = parse_due_date("2w") - expected = (datetime.now().date() + timedelta(weeks=2)).isoformat() - assert result == expected - - def test_iso_date_format(self): - """'2026-06-15' should pass through as-is.""" - result = parse_due_date("2026-06-15") - assert result == "2026-06-15" - - def test_whitespace_stripped(self): - """Leading/trailing whitespace should be stripped.""" - result = parse_due_date(" 7d ") - expected = (datetime.now().date() + timedelta(days=7)).isoformat() - assert result == expected - - def test_case_insensitive_days(self): - """'7D' should work the same as '7d'.""" - result = parse_due_date("7D") - expected = (datetime.now().date() + timedelta(days=7)).isoformat() - assert result == expected - - def test_case_insensitive_weeks(self): - """'2W' should work the same as '2w'.""" - result = parse_due_date("2W") - expected = (datetime.now().date() + timedelta(weeks=2)).isoformat() - assert result == expected - - def test_invalid_format_raises(self): - """Unsupported format should raise ValueError.""" - with pytest.raises(ValueError, match="Invalid date format"): - parse_due_date("next tuesday") - - def test_invalid_iso_date_raises(self): - """Invalid calendar date in ISO format should raise ValueError.""" - with pytest.raises(ValueError, match="Invalid date"): - parse_due_date("2026-02-30") - - def test_empty_string_raises(self): - """Empty string should raise ValueError.""" - with pytest.raises(ValueError, match="Invalid date format"): - parse_due_date("") - - def test_zero_days(self): - """'0d' should resolve to today.""" - result = parse_due_date("0d") - expected = datetime.now().date().isoformat() - assert result == expected - - -# ============================================= -# LOAD / SAVE TESTS -# ============================================= - - -class TestLoadSave: - def test_load_creates_file_if_missing(self, isolate_registry): - """load_tasks should create schedule.json if it does not exist.""" - assert not isolate_registry.exists() - tasks = load_tasks() - assert tasks == [] - assert isolate_registry.exists() - - def test_load_returns_empty_on_fresh_file(self): - """Fresh schedule.json should have no tasks.""" - tasks = load_tasks() - assert tasks == [] - - def test_save_and_load_roundtrip(self, isolate_registry): - """save_tasks then load_tasks should return the same data.""" - sample = [{"id": "abc123", "task": "test", "status": "pending"}] - assert save_tasks(sample) is True - loaded = load_tasks() - assert len(loaded) == 1 - assert loaded[0]["id"] == "abc123" - - def test_save_overwrites_existing(self, isolate_registry): - """Saving new tasks should fully replace existing data.""" - save_tasks([{"id": "first", "status": "pending"}]) - save_tasks([{"id": "second", "status": "pending"}]) - loaded = load_tasks() - assert len(loaded) == 1 - assert loaded[0]["id"] == "second" - - def test_load_handles_corrupt_json(self, isolate_registry): - """Corrupt JSON should return empty list, not crash.""" - isolate_registry.parent.mkdir(parents=True, exist_ok=True) - isolate_registry.write_text("{invalid json", encoding="utf-8") - tasks = load_tasks() - assert tasks == [] - - -# ============================================= -# CREATE TASK TESTS -# ============================================= - - -class TestCreateTask: - @patch.object(_mod.json_handler, "log_operation") - def test_create_basic(self, mock_log): - """Create a task and verify all fields.""" - task = create_task( - task="Check backup health", - due_date="7d", - recipient="@devpulse", - message="Verify backup systems", - ) - assert task["task"] == "Check backup health" - assert task["recipient"] == "@devpulse" - assert task["message"] == "Verify backup systems" - assert task["status"] == "pending" - assert len(task["id"]) == 16 - assert task["id"].isalnum() - assert task["created"] == datetime.now().date().isoformat() - mock_log.assert_called_once_with("task_created") - - @patch.object(_mod.json_handler, "log_operation") - def test_create_persists_to_json(self, mock_log, isolate_registry): - """Created task should be saved to the JSON file.""" - create_task( - task="persisted task", - due_date="1d", - recipient="@seedgo", - message="msg", - ) - raw = json.loads(isolate_registry.read_text(encoding="utf-8")) - assert len(raw["tasks"]) == 1 - assert raw["tasks"][0]["task"] == "persisted task" - - @patch.object(_mod.json_handler, "log_operation") - def test_create_multiple_tasks(self, mock_log): - """Multiple tasks should accumulate in the registry.""" - create_task(task="t1", due_date="1d", recipient="@a", message="m1") - create_task(task="t2", due_date="2d", recipient="@b", message="m2") - tasks = load_tasks() - assert len(tasks) == 2 - assert tasks[0]["task"] == "t1" - assert tasks[1]["task"] == "t2" - - def test_create_invalid_date_raises(self): - """create_task should propagate ValueError from bad due_date.""" - with pytest.raises(ValueError): - create_task(task="bad", due_date="xyz", recipient="@a", message="m") - - -# ============================================= -# DUE TASKS TESTS -# ============================================= - - -class TestDueTasks: - def test_overdue_task_returned(self, isolate_registry): - """A pending task with a past due_date should be returned.""" - yesterday = (datetime.now().date() - timedelta(days=1)).isoformat() - save_tasks( - [ - { - "id": "past01", - "due_date": yesterday, - "status": "pending", - "task": "overdue", - } - ] - ) - due = get_due_tasks() - assert len(due) == 1 - assert due[0]["id"] == "past01" - - def test_today_task_returned(self, isolate_registry): - """A pending task due today should be returned.""" - today = datetime.now().date().isoformat() - save_tasks( - [ - { - "id": "today01", - "due_date": today, - "status": "pending", - "task": "due today", - } - ] - ) - due = get_due_tasks() - assert len(due) == 1 - assert due[0]["id"] == "today01" - - def test_future_task_not_returned(self, isolate_registry): - """A pending task with a future due_date should not be returned.""" - future = (datetime.now().date() + timedelta(days=30)).isoformat() - save_tasks( - [ - { - "id": "future01", - "due_date": future, - "status": "pending", - "task": "future task", - } - ] - ) - due = get_due_tasks() - assert len(due) == 0 - - def test_dispatching_task_excluded(self, isolate_registry): - """Tasks with status 'dispatching' should not be returned.""" - yesterday = (datetime.now().date() - timedelta(days=1)).isoformat() - save_tasks( - [ - { - "id": "disp01", - "due_date": yesterday, - "status": "dispatching", - "task": "already dispatching", - } - ] - ) - due = get_due_tasks() - assert len(due) == 0 - - def test_completed_task_excluded(self, isolate_registry): - """Tasks with status 'completed' should not be returned.""" - yesterday = (datetime.now().date() - timedelta(days=1)).isoformat() - save_tasks( - [ - { - "id": "done01", - "due_date": yesterday, - "status": "completed", - "task": "done", - } - ] - ) - due = get_due_tasks() - assert len(due) == 0 - - def test_empty_registry_returns_empty(self): - """Empty registry should return empty list.""" - due = get_due_tasks() - assert due == [] - - -# ============================================= -# STATUS TRANSITION TESTS -# ============================================= - - -class TestStatusTransitions: - def _seed_task(self, task_id: str = "abc12345abcd1234", status: str = "pending"): - """Helper to seed a single task.""" - save_tasks( - [ - { - "id": task_id, - "task": "test", - "status": status, - "due_date": "2026-01-01", - } - ] - ) - return task_id - - def test_mark_dispatching_success(self): - """mark_dispatching should set status and dispatch_started.""" - tid = self._seed_task() - assert mark_dispatching(tid) is True - task = get_task_by_id(tid) - assert task is not None - assert task["status"] == "dispatching" - assert "dispatch_started" in task - - def test_mark_dispatching_missing(self): - """mark_dispatching returns False for nonexistent ID.""" - assert mark_dispatching("nonexistent_id__") is False - - def test_mark_completed_success(self): - """mark_completed should set status and completed_date.""" - tid = self._seed_task() - assert mark_completed(tid) is True - task = get_task_by_id(tid) - assert task is not None - assert task["status"] == "completed" - assert task["completed_date"] == datetime.now().date().isoformat() - - def test_mark_completed_missing(self): - """mark_completed returns False for nonexistent ID.""" - assert mark_completed("nonexistent_id__") is False - - def test_mark_pending_success(self): - """mark_pending should reset status and remove dispatch_started.""" - tid = self._seed_task(status="dispatching") - # Add dispatch_started to simulate real scenario - tasks = load_tasks() - tasks[0]["dispatch_started"] = datetime.now().isoformat() - save_tasks(tasks) - - assert mark_pending(tid) is True - task = get_task_by_id(tid) - assert task is not None - assert task["status"] == "pending" - assert "dispatch_started" not in task - - def test_mark_pending_missing(self): - """mark_pending returns False for nonexistent ID.""" - assert mark_pending("nonexistent_id__") is False - - def test_full_lifecycle(self): - """pending -> dispatching -> completed lifecycle.""" - tid = self._seed_task() - task = get_task_by_id(tid) - assert task is not None - assert task["status"] == "pending" - - mark_dispatching(tid) - task = get_task_by_id(tid) - assert task is not None - assert task["status"] == "dispatching" - - mark_completed(tid) - task = get_task_by_id(tid) - assert task is not None - assert task["status"] == "completed" - - -# ============================================= -# RECOVER STALE DISPATCHES TESTS -# ============================================= - - -class TestRecoverStale: - def test_recovers_stale_task(self, isolate_registry): - """Task stuck in dispatching beyond max_age should be reset.""" - stale_time = (datetime.now() - timedelta(minutes=10)).isoformat() - save_tasks( - [ - { - "id": "stale01", - "task": "stale dispatch", - "status": "dispatching", - "dispatch_started": stale_time, - "due_date": "2026-01-01", - } - ] - ) - recovered = recover_stale_dispatches(max_age_minutes=5) - assert recovered == 1 - task = get_task_by_id("stale01") - assert task is not None - assert task["status"] == "pending" - assert "dispatch_started" not in task - - def test_does_not_recover_recent_dispatch(self, isolate_registry): - """Task dispatching within max_age should not be recovered.""" - recent_time = (datetime.now() - timedelta(minutes=1)).isoformat() - save_tasks( - [ - { - "id": "recent01", - "task": "recent dispatch", - "status": "dispatching", - "dispatch_started": recent_time, - "due_date": "2026-01-01", - } - ] - ) - recovered = recover_stale_dispatches(max_age_minutes=5) - assert recovered == 0 - task = get_task_by_id("recent01") - assert task is not None - assert task["status"] == "dispatching" - - def test_recovers_invalid_timestamp(self, isolate_registry): - """Task with unparseable dispatch_started should be recovered.""" - save_tasks( - [ - { - "id": "bad_ts01", - "task": "bad timestamp", - "status": "dispatching", - "dispatch_started": "not-a-date", - "due_date": "2026-01-01", - } - ] - ) - recovered = recover_stale_dispatches(max_age_minutes=5) - assert recovered == 1 - task = get_task_by_id("bad_ts01") - assert task is not None - assert task["status"] == "pending" - - def test_pending_tasks_untouched(self, isolate_registry): - """Pending tasks should not be affected by recovery.""" - save_tasks( - [ - { - "id": "ok01", - "task": "normal pending", - "status": "pending", - "due_date": "2026-01-01", - } - ] - ) - recovered = recover_stale_dispatches(max_age_minutes=5) - assert recovered == 0 - task = get_task_by_id("ok01") - assert task is not None - assert task["status"] == "pending" - - def test_empty_registry_returns_zero(self): - """Recovery on empty registry should return 0.""" - assert recover_stale_dispatches() == 0 - - -# ============================================= -# DELETE TASK TESTS -# ============================================= - - -class TestDeleteTask: - def test_delete_existing(self, isolate_registry): - """Deleting an existing task returns True and removes it.""" - save_tasks([{"id": "del01", "task": "to delete", "status": "pending"}]) - assert delete_task("del01") is True - assert get_task_by_id("del01") is None - assert load_tasks() == [] - - def test_delete_missing(self): - """Deleting a nonexistent task returns False.""" - assert delete_task("nonexistent_id__") is False - - def test_delete_preserves_other_tasks(self, isolate_registry): - """Deleting one task should leave others intact.""" - save_tasks( - [ - {"id": "keep01", "task": "keep this", "status": "pending"}, - {"id": "del02", "task": "delete this", "status": "pending"}, - ] - ) - delete_task("del02") - remaining = load_tasks() - assert len(remaining) == 1 - assert remaining[0]["id"] == "keep01" - - def test_delete_from_empty_registry(self): - """Delete on empty registry should return False without error.""" - assert delete_task("anything") is False - - -# ============================================= -# GET PENDING TASKS TESTS -# ============================================= - - -class TestGetPendingTasks: - """Tests for get_pending_tasks().""" - - def test_returns_only_pending(self, isolate_registry): - """Only tasks with status 'pending' are returned.""" - save_tasks( - [ - {"id": "pend01", "task": "pending one", "status": "pending"}, - {"id": "pend02", "task": "pending two", "status": "pending"}, - {"id": "done01", "task": "done", "status": "completed"}, - ] - ) - result = get_pending_tasks() - assert len(result) == 2 - assert all(t["status"] == "pending" for t in result) - - def test_excludes_dispatching_and_completed(self, isolate_registry): - """Tasks with dispatching or completed status are excluded.""" - save_tasks( - [ - {"id": "disp01", "task": "dispatching", "status": "dispatching"}, - {"id": "done01", "task": "completed", "status": "completed"}, - {"id": "pend01", "task": "pending", "status": "pending"}, - ] - ) - result = get_pending_tasks() - assert len(result) == 1 - assert result[0]["id"] == "pend01" - - def test_empty_registry_returns_empty(self): - """Empty registry returns empty list.""" - result = get_pending_tasks() - assert result == [] - - -# ============================================= -# ENSURE LOCK DIR TESTS -# ============================================= - - -class TestEnsureLockDir: - """Tests for ensure_lock_dir().""" - - def test_creates_directory_if_missing(self, isolate_registry): - """Creates the lock directory when it does not exist.""" - lock_dir = isolate_registry.parent - if lock_dir.exists(): - import shutil - - shutil.rmtree(lock_dir) - assert not lock_dir.exists() - - result = ensure_lock_dir() - assert lock_dir.exists() - assert lock_dir.is_dir() - assert result["path"] == str(lock_dir) - - def test_returns_dict_with_path_key(self, isolate_registry): - """Return value is a dict containing the 'path' key.""" - result = ensure_lock_dir() - assert isinstance(result, dict) - assert "path" in result - assert isinstance(result["path"], str)