@spawn: owner + registry_id written into the SEALED registry entries (authority lives in registry, not the self-editable passport). ensure_project_has_owner() now keys off citizen_class=manager (was earliest-created, which mislabeled @aipass) and writes the registry entry. get_owner()/is_owner() resolvers added. 315 tests, seedgo 100%. @hooks: new registry_gate PreToolUse handler seals *_REGISTRY.json — blocks raw writes/tee/sed/rm + Edit/Write/MultiEdit, redirects to drone @spawn; per-clause bypass defeats compound-command smuggling; reads allowed. 82 tests, seedgo 100%. @ai_mail: wake-back reslope — SKIP_SENDERS blocklist replaced by an is_owner allowlist. Only the project owner is woken when their dispatched agent completes; all other guards intact (depth cap, lock, occupancy, honest messaging, dispatch_wake.log). seedgo 100% on the changed file. devpulse cross-part verify (REAL unmocked resolver): is_owner resolves devpulse-only; gate 13/13 incl compound-smuggle blocked + reads/drone-@spawn allowed; wake-back wakes owner / skips non-owner / respects depth-cap; 195 new-suite tests green together. Note: AIPASS_REGISTRY.json is gitignored — this ships the CODE; owner data regenerates per-install via ensure_project_has_owner. Still open (PART 4): gate watchdog+feedback on is_owner; portability of owner-only privileges across projects.
407 lines
14 KiB
Python
407 lines
14 KiB
Python
# =================== AIPass ====================
|
|
# Name: dispatch.py
|
|
# Description: Dispatch Module
|
|
# Version: 3.0.0
|
|
# Created: 2026-02-02
|
|
# Modified: 2026-02-02
|
|
# =============================================
|
|
|
|
"""
|
|
Dispatch Module
|
|
|
|
Orchestrates dispatch commands: status tracking and daemon management.
|
|
Delegates all business logic to handlers.
|
|
"""
|
|
|
|
import os
|
|
import sys
|
|
from pathlib import Path
|
|
from typing import List
|
|
|
|
from aipass.prax.apps.modules.logger import system_logger as logger
|
|
from aipass.cli.apps.modules import console, error
|
|
from aipass.ai_mail.apps.handlers.json import json_handler
|
|
from aipass.ai_mail.apps.handlers.dispatch.status import load_dispatch_log, check_pid_status, calculate_age
|
|
|
|
if sys.platform == "win32":
|
|
os.environ.setdefault("PYTHONUTF8", "1")
|
|
for _stream in (sys.stdout, sys.stderr):
|
|
_reconfigure = getattr(_stream, "reconfigure", None)
|
|
if _reconfigure is not None:
|
|
_reconfigure(encoding="utf-8", errors="replace")
|
|
|
|
|
|
def print_help() -> None:
|
|
"""Print help for dispatch commands."""
|
|
help_text = """
|
|
Dispatch Module - Agent dispatch management
|
|
|
|
COMMANDS:
|
|
dispatch @target "Subject" "Body" - Send dispatch email + wake target
|
|
dispatch status - Show last 5 dispatch spawns with current status
|
|
dispatch daemon - Start the continuous dispatch daemon
|
|
dispatch wake @branch - Wake only (no email sent)
|
|
|
|
DISPATCH (send + wake):
|
|
drone @ai_mail dispatch @branch "Subject" "Body" # Send + wake
|
|
drone @ai_mail dispatch @branch "Subject" "Body" --fresh # Send + fresh wake
|
|
drone @ai_mail dispatch @branch "Subject" "Body" --model opus # Send + wake with Opus
|
|
drone @ai_mail dispatch @branch "Subject" "Body" --no-memory-save
|
|
|
|
WAKE ONLY:
|
|
drone @ai_mail dispatch wake @branch # Wake with default inbox check
|
|
drone @ai_mail dispatch wake @branch "custom" # Wake with custom prompt
|
|
drone @ai_mail dispatch wake @branch --model opus # Wake with Opus model
|
|
drone wake @branch # Shortcut via drone
|
|
|
|
MODEL OPTIONS:
|
|
--model opus Claude Opus 4.6 (default — full reasoning for all tasks)
|
|
--model sonnet Claude Sonnet 4.6 (cost-effective for lighter tasks)
|
|
--model haiku Claude Haiku 4.5 (fastest, simplest tasks)
|
|
|
|
DAEMON:
|
|
The daemon polls branch inboxes for --dispatch emails and spawns agents.
|
|
Kill switch: touch <repo_root>/.aipass/autonomous_pause
|
|
Config: safety_config.json
|
|
"""
|
|
console.print(help_text)
|
|
|
|
|
|
def handle_command(command: str, args: List[str]) -> bool:
|
|
"""
|
|
Handle dispatch commands.
|
|
|
|
Args:
|
|
command: Command name
|
|
args: Command arguments
|
|
|
|
Returns:
|
|
True if command handled, False otherwise
|
|
"""
|
|
if command != "dispatch":
|
|
return False
|
|
|
|
if args and args[0] in ["--help", "-h", "help"]:
|
|
print_help()
|
|
return True
|
|
|
|
if not args:
|
|
print_introspection()
|
|
return True
|
|
|
|
subcommand = args[0]
|
|
|
|
json_handler.log_operation("dispatch_command", {"subcommand": subcommand})
|
|
|
|
if subcommand == "status":
|
|
return _orchestrate_status()
|
|
elif subcommand == "daemon":
|
|
return _orchestrate_daemon()
|
|
elif subcommand == "wake":
|
|
return _orchestrate_wake(args[1:])
|
|
elif subcommand.startswith("@") or subcommand.startswith("/"):
|
|
return _orchestrate_dispatch_send(args)
|
|
else:
|
|
error(f"Unknown dispatch subcommand: {subcommand}")
|
|
print_help()
|
|
return False
|
|
|
|
|
|
def _orchestrate_status() -> bool:
|
|
"""Orchestrate dispatch status display."""
|
|
logger.info("[dispatch] Showing dispatch status")
|
|
|
|
dispatches = load_dispatch_log()
|
|
|
|
if not dispatches:
|
|
console.print("\n[dim]No dispatches recorded yet.[/dim]")
|
|
return True
|
|
|
|
recent = dispatches[-5:][::-1]
|
|
|
|
console.print("\n[bold]DISPATCH STATUS[/bold]")
|
|
console.print("─" * 50)
|
|
|
|
active_count = 0
|
|
for entry in recent:
|
|
branch = entry.get("branch", "unknown")
|
|
pid = entry.get("pid")
|
|
timestamp = entry.get("timestamp", "")
|
|
spawn_status = entry.get("status", "unknown")
|
|
|
|
if spawn_status == "spawned" and pid:
|
|
current_status = check_pid_status(pid)
|
|
elif spawn_status == "failed":
|
|
current_status = "FAILED"
|
|
else:
|
|
current_status = "UNKNOWN"
|
|
|
|
if current_status == "RUNNING":
|
|
active_count += 1
|
|
|
|
age_str = calculate_age(timestamp)
|
|
|
|
if current_status == "RUNNING":
|
|
status_display = "[green]RUNNING[/green]"
|
|
elif current_status == "COMPLETED":
|
|
status_display = "[dim]COMPLETED[/dim]"
|
|
elif current_status == "FAILED":
|
|
status_display = "[red]FAILED[/red]"
|
|
else:
|
|
status_display = f"[yellow]{current_status}[/yellow]"
|
|
|
|
pid_display = f"PID {pid}" if pid else "NO PID"
|
|
console.print(f" {branch:<12} {pid_display:<12} {status_display:<18} {age_str}")
|
|
|
|
console.print("─" * 50)
|
|
console.print(f"[dim]Active: {active_count} | Total: {len(recent)}[/dim]\n")
|
|
|
|
return True
|
|
|
|
|
|
def _orchestrate_wake(args: List[str]) -> bool:
|
|
"""Orchestrate manual branch wake."""
|
|
if not args or args[0] in ["--help", "-h", "help"]:
|
|
console.print("\n[bold]Wake - Manual branch spawn[/bold]")
|
|
console.print(' Usage: dispatch wake @branch ["custom message"]')
|
|
console.print(' Or: drone wake @branch ["custom message"]\n')
|
|
return True
|
|
|
|
# Parse --fresh, --sender, --model flags
|
|
use_fresh = "--fresh" in args
|
|
use_sender = "@devpulse"
|
|
use_model = None
|
|
filtered = []
|
|
i = 0
|
|
while i < len(args):
|
|
if args[i] == "--fresh":
|
|
i += 1
|
|
continue
|
|
if args[i] == "--sender" and i + 1 < len(args):
|
|
use_sender = args[i + 1]
|
|
i += 2
|
|
continue
|
|
if args[i] == "--model" and i + 1 < len(args):
|
|
use_model = args[i + 1]
|
|
i += 2
|
|
continue
|
|
filtered.append(args[i])
|
|
i += 1
|
|
|
|
if not filtered:
|
|
error("Missing branch argument")
|
|
return False
|
|
|
|
branch_email = filtered[0]
|
|
custom_message = filtered[1] if len(filtered) > 1 else None
|
|
|
|
from aipass.ai_mail.apps.handlers.dispatch.wake import is_wake_blocked
|
|
|
|
if is_wake_blocked(branch_email):
|
|
error(
|
|
f"target {branch_email} is protected from manual wake. "
|
|
f'Use \'drone @ai_mail dispatch {branch_email} "Subject" "Body"\' to send work instead.'
|
|
)
|
|
return True
|
|
|
|
logger.info(f"[dispatch] Manual wake requested for {branch_email}")
|
|
console.print(f"\n⏳ Waking {branch_email}...")
|
|
|
|
from aipass.ai_mail.apps.handlers.dispatch.wake import wake_branch
|
|
|
|
dispatch_status, success = wake_branch(
|
|
branch_email, custom_message, fresh=use_fresh, sender=use_sender, model=use_model
|
|
)
|
|
|
|
# Print step-by-step status
|
|
console.print(dispatch_status.format())
|
|
|
|
return success
|
|
|
|
|
|
def _orchestrate_dispatch_send(args: List[str]) -> bool:
|
|
"""Orchestrate combined dispatch: send email with --dispatch flag + wake branch."""
|
|
# Parse flags
|
|
use_fresh = False
|
|
no_memory_save = False
|
|
from_branch = None
|
|
use_model = None
|
|
filtered = []
|
|
i = 0
|
|
while i < len(args):
|
|
if args[i] == "--fresh":
|
|
use_fresh = True
|
|
i += 1
|
|
continue
|
|
if args[i] == "--no-memory-save":
|
|
no_memory_save = True
|
|
i += 1
|
|
continue
|
|
if args[i] == "--no-watchdog":
|
|
i += 1
|
|
continue
|
|
if args[i] == "--from" and i + 1 < len(args):
|
|
from_branch = args[i + 1]
|
|
i += 2
|
|
continue
|
|
if args[i] == "--model" and i + 1 < len(args):
|
|
use_model = args[i + 1]
|
|
i += 2
|
|
continue
|
|
filtered.append(args[i])
|
|
i += 1
|
|
|
|
if len(filtered) < 3:
|
|
error('Usage: dispatch @target "Subject" "Body" [--fresh] [--no-memory-save]')
|
|
return True
|
|
|
|
target = filtered[0]
|
|
subject = filtered[1]
|
|
body = filtered[2]
|
|
|
|
logger.info(f"[dispatch] Combined dispatch: send + wake for {target}")
|
|
json_handler.log_operation("dispatch_send_and_wake", {"target": target, "subject": subject, "fresh": use_fresh})
|
|
|
|
# --- Step 1: Send dispatch email ---
|
|
console.print(f"\nSending dispatch email to {target}...")
|
|
|
|
from aipass.ai_mail.apps.handlers.email.send import resolve_sender_info, send_to_single
|
|
from aipass.ai_mail.apps.handlers.email.create import create_email_file, load_email_file
|
|
from aipass.ai_mail.apps.handlers.email.delivery import deliver_email_to_branch
|
|
from aipass.ai_mail.apps.handlers.email.header import prepend_dispatch_header
|
|
from aipass.ai_mail.apps.handlers.email.error_dispatch import dispatch_send_error, on_email_delivered
|
|
from aipass.ai_mail.apps.handlers.users.user import get_current_user
|
|
from aipass.ai_mail.apps.handlers.registry.read import get_branch_by_email
|
|
from aipass.ai_mail.apps.handlers.paths import find_repo_root
|
|
|
|
try:
|
|
from aipass.ai_mail.apps.handlers.central_writer import update_central
|
|
except ImportError as e:
|
|
logger.warning("[dispatch] central_writer import unavailable: %s", e)
|
|
update_central = None
|
|
|
|
_ai_mail_dir = Path(__file__).resolve().parents[2]
|
|
_repo_root = find_repo_root()
|
|
|
|
def _delivery_callback(branch_path, new_count, opened_count, total):
|
|
on_email_delivered(
|
|
branch_path,
|
|
new_count,
|
|
opened_count,
|
|
total,
|
|
update_central_fn=update_central,
|
|
)
|
|
|
|
try:
|
|
user_info = resolve_sender_info(from_branch, _repo_root, _ai_mail_dir, get_branch_by_email, get_current_user)
|
|
message = prepend_dispatch_header(body, no_memory_save=no_memory_save)
|
|
|
|
send_ok, send_error = send_to_single(
|
|
target,
|
|
subject,
|
|
message,
|
|
user_info,
|
|
True,
|
|
no_memory_save,
|
|
None,
|
|
target,
|
|
create_email_file,
|
|
load_email_file,
|
|
deliver_email_to_branch,
|
|
_delivery_callback,
|
|
json_handler.log_operation,
|
|
update_central,
|
|
)
|
|
|
|
if not send_ok:
|
|
error(f"Send failed: {send_error}")
|
|
dispatch_send_error(target, subject, send_error or "", deliver_email_to_branch)
|
|
return False
|
|
|
|
console.print(f"[green]Email sent to {target}[/green]")
|
|
|
|
try:
|
|
from aipass.trigger.apps.modules.core import trigger
|
|
|
|
trigger.fire("email_dispatched", to=target, subject=subject)
|
|
except Exception as e:
|
|
logger.warning("[dispatch] trigger fire failed: %s", e)
|
|
|
|
except Exception as e:
|
|
logger.error(f"[dispatch] Send phase failed: {e}")
|
|
error(f"Send failed: {e}")
|
|
return False
|
|
|
|
# --- Step 2: Wake the branch ---
|
|
console.print(f"\nWaking {target}...")
|
|
|
|
from aipass.ai_mail.apps.handlers.dispatch.wake import wake_branch
|
|
|
|
dispatch_status, wake_ok = wake_branch(
|
|
target, fresh=use_fresh, sender=user_info.get("email_address", "@ai_mail"), model=use_model
|
|
)
|
|
console.print(dispatch_status.format())
|
|
|
|
if not wake_ok:
|
|
logger.warning("[dispatch] Wake failed for %s — email was sent", target)
|
|
error(f"Email sent but wake failed — retry: drone @ai_mail dispatch wake {target}")
|
|
else:
|
|
console.print(f"[dim]Wake-back enabled — sender will be woken when {target} completes (if available)[/dim]")
|
|
|
|
return True
|
|
|
|
|
|
def _orchestrate_daemon() -> bool:
|
|
"""Orchestrate daemon startup."""
|
|
logger.info("[dispatch] Starting dispatch daemon")
|
|
console.print("\n[bold]Starting dispatch daemon...[/bold]")
|
|
|
|
from aipass.ai_mail.apps.handlers.dispatch.daemon import run_daemon
|
|
|
|
run_daemon()
|
|
return True
|
|
|
|
|
|
def print_introspection():
|
|
"""Display module introspection info."""
|
|
console.print()
|
|
console.print("[bold cyan]dispatch Module[/bold cyan]")
|
|
console.print(
|
|
"[dim]Orchestrates dispatch commands: combined send+wake,"
|
|
" status tracking, daemon management, and manual wake.[/dim]"
|
|
)
|
|
console.print()
|
|
console.print("[yellow]Connected Handlers:[/yellow]")
|
|
console.print(" [cyan]handlers/dispatch/[/cyan]")
|
|
console.print(" - [cyan]status.py[/cyan] [dim](load_dispatch_log — load dispatch log entries)[/dim]")
|
|
console.print(
|
|
" - [cyan]status.py[/cyan] [dim](check_pid_status — check if a spawned process is still running)[/dim]"
|
|
)
|
|
console.print(" - [cyan]status.py[/cyan] [dim](calculate_age — calculate age string from timestamp)[/dim]")
|
|
console.print(" - [cyan]wake.py[/cyan] [dim](wake_branch — manually wake a branch by spawning an agent)[/dim]")
|
|
console.print(" - [cyan]daemon.py[/cyan] [dim](run_daemon — start the continuous dispatch daemon)[/dim]")
|
|
console.print(" [cyan]handlers/email/[/cyan] [dim](used by combined dispatch)[/dim]")
|
|
console.print(" - [cyan]send.py[/cyan] [dim](resolve_sender_info, send_to_single — send email pipeline)[/dim]")
|
|
console.print(" - [cyan]create.py[/cyan] [dim](create_email_file, load_email_file — email file creation)[/dim]")
|
|
console.print(" - [cyan]delivery.py[/cyan] [dim](deliver_email_to_branch — inbox delivery)[/dim]")
|
|
console.print(" - [cyan]header.py[/cyan] [dim](prepend_dispatch_header — dispatch header injection)[/dim]")
|
|
console.print()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
if len(sys.argv) == 1:
|
|
print_help()
|
|
sys.exit(0)
|
|
|
|
if sys.argv[1] in ["--help", "-h", "help"]:
|
|
print_help()
|
|
sys.exit(0)
|
|
|
|
command = sys.argv[1]
|
|
remaining_args = sys.argv[2:] if len(sys.argv) > 2 else []
|
|
|
|
if handle_command(command, remaining_args):
|
|
sys.exit(0)
|
|
else:
|
|
sys.exit(1)
|