feat(system): feat(system): DPLAN-0141 Phase 3 — ai_mail audit rescue (email.py split + nesting fix + dead code) + devpulse cleanup (feedback stderr + bypass + README + DPLAN-0141)
Co-Authored-By: @devpulse <devpulse@aipass>
This commit is contained in:
@@ -354,6 +354,11 @@
|
||||
"file": "tests/test_daemon.py",
|
||||
"standard": "encapsulation",
|
||||
"reason": "Unit tests must import daemon internals directly to verify behavior. Module entry-point-only rule does not apply to tests."
|
||||
},
|
||||
{
|
||||
"file": "tests/test_daemon.py",
|
||||
"standard": "documentation",
|
||||
"reason": "mock_spawn_agent is a nested helper inside test methods, not a public function. AST checker misidentifies it."
|
||||
}
|
||||
],
|
||||
"notes": {
|
||||
|
||||
@@ -106,6 +106,22 @@ def _get_branches_list(registry: dict) -> list:
|
||||
# =============================================
|
||||
|
||||
|
||||
def _synthesize_external_branch(caller_branch: str) -> Optional[Dict]:
|
||||
"""Build a synthetic branch info dict from env vars for an external project."""
|
||||
caller_cwd = os.environ.get("AIPASS_CALLER_CWD", "")
|
||||
if not caller_cwd:
|
||||
return None
|
||||
cwd_path = Path(caller_cwd)
|
||||
name_key = caller_branch.lstrip("@").lower()
|
||||
return {
|
||||
"name": name_key,
|
||||
"path": str(cwd_path),
|
||||
"email": f"@{name_key}",
|
||||
"status": "active",
|
||||
"type": "external",
|
||||
}
|
||||
|
||||
|
||||
def detect_branch_from_pwd() -> Optional[Dict]:
|
||||
"""
|
||||
Detect which branch is calling based on current working directory.
|
||||
@@ -114,66 +130,29 @@ def detect_branch_from_pwd() -> Optional[Dict]:
|
||||
Then looks up branch info in AIPASS_REGISTRY.json.
|
||||
|
||||
Returns:
|
||||
Dict with branch info if detected:
|
||||
{
|
||||
"name": "SEEDGO",
|
||||
"path": "src/aipass/seedgo",
|
||||
"email": "@seedgo",
|
||||
"display_name": "Seedgo (Standards Branch)",
|
||||
...
|
||||
}
|
||||
None if no branch detected
|
||||
Dict with branch info if detected, or None.
|
||||
"""
|
||||
json_handler.log_operation("detect_branch_from_pwd", {"cwd": str(Path.cwd())})
|
||||
|
||||
try:
|
||||
# Primary: use explicit branch name passed by drone (works in Docker + local)
|
||||
caller_branch = os.environ.get("AIPASS_CALLER_BRANCH")
|
||||
if caller_branch:
|
||||
# Try contacts first (fastest, works for external projects)
|
||||
contact = _get_contact_info(caller_branch)
|
||||
if contact:
|
||||
return contact
|
||||
# Fall back to registry lookup
|
||||
branch_info = _lookup_branch_by_name(caller_branch)
|
||||
if branch_info:
|
||||
return branch_info
|
||||
# Last resort: synthesize from env vars (external project, first contact)
|
||||
caller_cwd = os.environ.get("AIPASS_CALLER_CWD", "")
|
||||
if caller_cwd:
|
||||
cwd_path = Path(caller_cwd)
|
||||
name_key = caller_branch.lstrip("@").lower()
|
||||
mailbox = cwd_path / ".ai_mail.local"
|
||||
if not mailbox.exists():
|
||||
# Walk up to find project-level mailbox
|
||||
for parent in [cwd_path] + list(cwd_path.parents):
|
||||
candidate = parent / ".ai_mail.local"
|
||||
if candidate.exists():
|
||||
mailbox = candidate
|
||||
break
|
||||
return {
|
||||
"name": name_key,
|
||||
"path": str(cwd_path),
|
||||
"email": f"@{name_key}",
|
||||
"status": "active",
|
||||
"type": "external",
|
||||
}
|
||||
return _synthesize_external_branch(caller_branch)
|
||||
|
||||
# Fallback: use caller's CWD for path-based detection (local only)
|
||||
caller_cwd = os.environ.get("AIPASS_CALLER_CWD")
|
||||
cwd = Path(caller_cwd) if caller_cwd else Path.cwd()
|
||||
|
||||
# Find branch root
|
||||
branch_root = find_branch_root(cwd)
|
||||
if not branch_root:
|
||||
return None
|
||||
|
||||
# Get branch info from registry
|
||||
branch_info = get_branch_info_from_registry(branch_root)
|
||||
if not branch_info:
|
||||
return None
|
||||
|
||||
return branch_info
|
||||
return get_branch_info_from_registry(branch_root)
|
||||
|
||||
except Exception as e:
|
||||
logger.warning("[identity] detect_branch_from_pwd() failed: %s", e)
|
||||
|
||||
@@ -29,12 +29,10 @@ _REPO_ROOT = _AI_MAIL_DIR.parents[2]
|
||||
|
||||
from aipass.prax import logger
|
||||
from aipass.cli.apps.modules import console, error
|
||||
from aipass.trigger.apps.modules.core import trigger
|
||||
|
||||
# Handlers - business logic providers
|
||||
from aipass.ai_mail.apps.handlers.email.dashboard_sync import push_dashboard_update
|
||||
from aipass.ai_mail.apps.handlers.email.delivery import deliver_email_to_branch
|
||||
from aipass.ai_mail.apps.handlers.email.create import create_email_file, load_email_file
|
||||
from aipass.ai_mail.apps.handlers.email.create import load_email_file
|
||||
from aipass.ai_mail.apps.handlers.email.format import format_email_list_item, format_email_header
|
||||
from aipass.ai_mail.apps.handlers.email.inbox_ops import load_inbox
|
||||
from aipass.ai_mail.apps.handlers.email.inbox_cleanup import (
|
||||
@@ -43,20 +41,12 @@ from aipass.ai_mail.apps.handlers.email.inbox_cleanup import (
|
||||
mark_as_closed_and_archive,
|
||||
)
|
||||
from aipass.ai_mail.apps.handlers.email.reply import get_email_by_id, send_reply
|
||||
from aipass.ai_mail.apps.handlers.email.header import prepend_dispatch_header
|
||||
from aipass.ai_mail.apps.handlers.users.user import get_current_user
|
||||
from aipass.ai_mail.apps.handlers.registry.read import get_all_branches, get_branch_by_email
|
||||
from aipass.ai_mail.apps.handlers.json import json_handler
|
||||
from aipass.ai_mail.apps.handlers.email.send import (
|
||||
resolve_sender_info,
|
||||
send_to_broadcast,
|
||||
send_to_single,
|
||||
collect_interactive_input,
|
||||
)
|
||||
from aipass.ai_mail.apps.handlers.email.error_dispatch import dispatch_send_error, on_email_delivered
|
||||
from aipass.ai_mail.apps.handlers.email.send_args import parse_send_args, resolve_dispatch_target
|
||||
from aipass.ai_mail.apps.handlers.email.close_ops import batch_close, batch_close_post_ops
|
||||
from aipass.ai_mail.apps.handlers.email.inbox_resolve import resolve_inbox_target
|
||||
from aipass.ai_mail.apps.modules.email_send import handle_send
|
||||
|
||||
try:
|
||||
from aipass.ai_mail.apps.handlers.central_writer import update_central
|
||||
@@ -65,18 +55,6 @@ except ImportError as e:
|
||||
update_central = None
|
||||
|
||||
|
||||
def _delivery_callback(branch_path, new_count, opened_count, total):
|
||||
"""Post-delivery callback: delegates to error_dispatch handler."""
|
||||
on_email_delivered(
|
||||
branch_path,
|
||||
new_count,
|
||||
opened_count,
|
||||
total,
|
||||
push_dashboard_fn=push_dashboard_update,
|
||||
update_central_fn=update_central,
|
||||
)
|
||||
|
||||
|
||||
def _resolve_branch_path() -> Path:
|
||||
"""Resolve the branch path for inbox operations.
|
||||
|
||||
@@ -150,182 +128,6 @@ def handle_command(command: str, args: List[str]) -> bool:
|
||||
return dispatch[command](args)
|
||||
|
||||
|
||||
def handle_send(args: List[str]) -> bool:
|
||||
"""Orchestrate email sending workflow."""
|
||||
json_handler.log_operation("send_email_initiated", {"args_count": len(args)})
|
||||
parsed = parse_send_args(args)
|
||||
|
||||
if parsed["mode"] == "error":
|
||||
error(parsed["error"])
|
||||
console.print(' Multiple: send @branch1 @branch2 "Subject" "Message"')
|
||||
return False
|
||||
|
||||
if parsed["mode"] == "interactive":
|
||||
return _send_interactive()
|
||||
|
||||
# Direct send
|
||||
recipients = parsed["recipients"]
|
||||
from_branch = parsed.get("from_branch")
|
||||
if len(recipients) == 1:
|
||||
target = resolve_dispatch_target(recipients[0], parsed["auto_execute"], _get_branch_info_fn())
|
||||
return _send_direct(
|
||||
recipients[0],
|
||||
parsed["subject"],
|
||||
parsed["message"],
|
||||
parsed["auto_execute"],
|
||||
parsed["reply_to"],
|
||||
target,
|
||||
parsed["no_memory_save"],
|
||||
from_branch=from_branch,
|
||||
)
|
||||
|
||||
console.print(f"\n[bold]Group send to {len(recipients)} recipients...[/bold]")
|
||||
ok = 0
|
||||
for r in recipients:
|
||||
target = resolve_dispatch_target(r, parsed["auto_execute"], _get_branch_info_fn())
|
||||
if _send_direct(
|
||||
r,
|
||||
parsed["subject"],
|
||||
parsed["message"],
|
||||
parsed["auto_execute"],
|
||||
parsed["reply_to"],
|
||||
target,
|
||||
parsed["no_memory_save"],
|
||||
from_branch=from_branch,
|
||||
):
|
||||
ok += 1
|
||||
console.print(f"\nGroup send complete: {ok}/{len(recipients)} delivered")
|
||||
return ok > 0
|
||||
|
||||
|
||||
def _get_branch_info_fn():
|
||||
"""Return branch info lookup fn for dispatch target resolution, or None."""
|
||||
try:
|
||||
from aipass.ai_mail.apps.handlers.users.branch_detection import get_branch_info_from_registry
|
||||
|
||||
return get_branch_info_from_registry
|
||||
except ImportError as e:
|
||||
logger.warning("[email] branch_detection import unavailable: %s", e)
|
||||
return None
|
||||
|
||||
|
||||
def _send_interactive() -> bool:
|
||||
"""Interactive email sending with prompts."""
|
||||
branches = get_all_branches()
|
||||
console.print("\nAI_Mail - Send Email\n" + "=" * 50)
|
||||
console.print("\nSelect recipient:")
|
||||
for i, b in enumerate(branches, 1):
|
||||
console.print(f" {i}. {b['name']} ({b['email']})")
|
||||
console.print(f" {len(branches) + 1}. ALL BRANCHES (broadcast)")
|
||||
console.print("Message (press Ctrl+D when done, Ctrl+C to cancel):")
|
||||
|
||||
result = collect_interactive_input(branches)
|
||||
if result is None:
|
||||
console.print("\nCancelled")
|
||||
return False
|
||||
|
||||
console.print("\n" + "=" * 50)
|
||||
console.print(f"To: {result['to']}\nSubject: {result['subject']}\nMessage:\n{result['message']}")
|
||||
console.print("=" * 50)
|
||||
return _send_direct(result["to"], result["subject"], result["message"])
|
||||
|
||||
|
||||
def _send_direct(
|
||||
to_branch,
|
||||
subject,
|
||||
message,
|
||||
auto_execute=False,
|
||||
reply_to=None,
|
||||
dispatched_to=None,
|
||||
no_memory_save=False,
|
||||
from_branch=None,
|
||||
) -> bool:
|
||||
"""Direct email send - thin wrapper over send handlers."""
|
||||
try:
|
||||
user_info = resolve_sender_info(from_branch, _REPO_ROOT, _AI_MAIL_DIR, get_branch_by_email, get_current_user)
|
||||
if auto_execute:
|
||||
message = prepend_dispatch_header(message, no_memory_save=no_memory_save)
|
||||
|
||||
if to_branch.lower() in ["all", "@all"]:
|
||||
return _send_broadcast(subject, message, user_info, auto_execute, no_memory_save, reply_to, dispatched_to)
|
||||
|
||||
success, error_msg = send_to_single(
|
||||
to_branch,
|
||||
subject,
|
||||
message,
|
||||
user_info,
|
||||
auto_execute,
|
||||
no_memory_save,
|
||||
reply_to,
|
||||
dispatched_to,
|
||||
create_email_file,
|
||||
load_email_file,
|
||||
deliver_email_to_branch,
|
||||
_delivery_callback,
|
||||
json_handler.log_operation,
|
||||
update_central,
|
||||
)
|
||||
|
||||
if success:
|
||||
label = "\\[dispatch: queued for daemon]" if auto_execute else ""
|
||||
console.print(f"[green]Email sent to {to_branch} {label}[/green]")
|
||||
if auto_execute:
|
||||
_fire_dispatch_trigger(to_branch, subject)
|
||||
return True
|
||||
else:
|
||||
error(f"Failed to deliver: {error_msg}")
|
||||
dispatch_send_error(to_branch, subject, error_msg or "", deliver_email_to_branch)
|
||||
return False
|
||||
except BrokenPipeError:
|
||||
logger.info("[email] Send: broken pipe (stdout closed early)")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"[email] Send failed: {e}")
|
||||
error(f"Error: {e}")
|
||||
dispatch_send_error(to_branch, subject, str(e), deliver_email_to_branch)
|
||||
return False
|
||||
|
||||
|
||||
def _fire_dispatch_trigger(to_branch: str, subject: str) -> None:
|
||||
"""Fire email_dispatched trigger event if auto_execute enabled."""
|
||||
try:
|
||||
trigger.fire("email_dispatched", to=to_branch, subject=subject)
|
||||
except Exception as e:
|
||||
logger.warning("[email] trigger fire for email_dispatched failed: %s", e)
|
||||
|
||||
|
||||
def _send_broadcast(subject, message, user_info, auto_execute, no_memory_save, reply_to, dispatched_to) -> bool:
|
||||
"""Broadcast send to all branches - display wrapper."""
|
||||
branches = get_all_branches()
|
||||
console.print(f"\nBroadcasting to {len(branches)} branches...")
|
||||
ok, success_count, total, results = send_to_broadcast(
|
||||
subject,
|
||||
message,
|
||||
user_info,
|
||||
auto_execute,
|
||||
no_memory_save,
|
||||
reply_to,
|
||||
dispatched_to,
|
||||
branches,
|
||||
create_email_file,
|
||||
load_email_file,
|
||||
deliver_email_to_branch,
|
||||
_delivery_callback,
|
||||
json_handler.log_operation,
|
||||
update_central,
|
||||
)
|
||||
if isinstance(results, str) or results is None:
|
||||
error("Failed to load email file for broadcast")
|
||||
return False
|
||||
for name, ok, err in results: # type: ignore[union-attr]
|
||||
if ok:
|
||||
console.print(f" [green]OK[/green] {name}")
|
||||
else:
|
||||
error(f"FAIL {name} ({err})")
|
||||
console.print(f"\nBroadcast complete: {success_count}/{total} delivered")
|
||||
return ok
|
||||
|
||||
|
||||
def handle_inbox(args: List[str]) -> bool:
|
||||
"""Orchestrate inbox viewing."""
|
||||
json_handler.log_operation("inbox_viewed")
|
||||
|
||||
@@ -0,0 +1,254 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: email_send.py
|
||||
# Description: Email Send Orchestration (extracted from email.py)
|
||||
# Version: 1.0.0
|
||||
# Created: 2026-04-22
|
||||
# Modified: 2026-04-22
|
||||
# =============================================
|
||||
|
||||
"""
|
||||
Email Send Orchestration
|
||||
|
||||
Handles the send/email command workflow: direct send, interactive send,
|
||||
broadcast, and dispatch trigger. Extracted from email.py to keep modules
|
||||
under the size threshold.
|
||||
"""
|
||||
|
||||
from pathlib import Path
|
||||
from typing import List
|
||||
|
||||
_AI_MAIL_DIR = Path(__file__).resolve().parents[2]
|
||||
_REPO_ROOT = _AI_MAIL_DIR.parents[2]
|
||||
|
||||
from aipass.prax import logger
|
||||
from aipass.cli.apps.modules import console, error
|
||||
from aipass.trigger.apps.modules.core import trigger
|
||||
|
||||
from aipass.ai_mail.apps.handlers.email.dashboard_sync import push_dashboard_update
|
||||
from aipass.ai_mail.apps.handlers.email.delivery import deliver_email_to_branch
|
||||
from aipass.ai_mail.apps.handlers.email.create import create_email_file, load_email_file
|
||||
from aipass.ai_mail.apps.handlers.email.header import prepend_dispatch_header
|
||||
from aipass.ai_mail.apps.handlers.json import json_handler
|
||||
from aipass.ai_mail.apps.handlers.registry.read import get_all_branches, get_branch_by_email
|
||||
from aipass.ai_mail.apps.handlers.users.user import get_current_user
|
||||
from aipass.ai_mail.apps.handlers.email.send import (
|
||||
resolve_sender_info,
|
||||
send_to_broadcast,
|
||||
send_to_single,
|
||||
collect_interactive_input,
|
||||
)
|
||||
from aipass.ai_mail.apps.handlers.email.error_dispatch import dispatch_send_error, on_email_delivered
|
||||
from aipass.ai_mail.apps.handlers.email.send_args import parse_send_args, resolve_dispatch_target
|
||||
|
||||
try:
|
||||
from aipass.ai_mail.apps.handlers.central_writer import update_central
|
||||
except ImportError as e:
|
||||
logger.warning("[email_send] central_writer import unavailable: %s", e)
|
||||
update_central = None
|
||||
|
||||
|
||||
def _delivery_callback(branch_path, new_count, opened_count, total):
|
||||
"""Post-delivery callback: delegates to error_dispatch handler."""
|
||||
on_email_delivered(
|
||||
branch_path,
|
||||
new_count,
|
||||
opened_count,
|
||||
total,
|
||||
push_dashboard_fn=push_dashboard_update,
|
||||
update_central_fn=update_central,
|
||||
)
|
||||
|
||||
|
||||
def _get_branch_info_fn():
|
||||
"""Return branch info lookup fn for dispatch target resolution, or None."""
|
||||
try:
|
||||
from aipass.ai_mail.apps.handlers.users.branch_detection import get_branch_info_from_registry
|
||||
|
||||
return get_branch_info_from_registry
|
||||
except ImportError as e:
|
||||
logger.warning("[email_send] branch_detection import unavailable: %s", e)
|
||||
return None
|
||||
|
||||
|
||||
def handle_send(args: List[str]) -> bool:
|
||||
"""Orchestrate email sending workflow."""
|
||||
json_handler.log_operation("send_email_initiated", {"args_count": len(args)})
|
||||
parsed = parse_send_args(args)
|
||||
|
||||
if parsed["mode"] == "error":
|
||||
error(parsed["error"])
|
||||
console.print(' Multiple: send @branch1 @branch2 "Subject" "Message"')
|
||||
return False
|
||||
|
||||
if parsed["mode"] == "interactive":
|
||||
return _send_interactive()
|
||||
|
||||
recipients = parsed["recipients"]
|
||||
from_branch = parsed.get("from_branch")
|
||||
if len(recipients) == 1:
|
||||
target = resolve_dispatch_target(recipients[0], parsed["auto_execute"], _get_branch_info_fn())
|
||||
return _send_direct(
|
||||
recipients[0],
|
||||
parsed["subject"],
|
||||
parsed["message"],
|
||||
parsed["auto_execute"],
|
||||
parsed["reply_to"],
|
||||
target,
|
||||
parsed["no_memory_save"],
|
||||
from_branch=from_branch,
|
||||
)
|
||||
|
||||
console.print(f"\n[bold]Group send to {len(recipients)} recipients...[/bold]")
|
||||
ok = 0
|
||||
for r in recipients:
|
||||
target = resolve_dispatch_target(r, parsed["auto_execute"], _get_branch_info_fn())
|
||||
if _send_direct(
|
||||
r,
|
||||
parsed["subject"],
|
||||
parsed["message"],
|
||||
parsed["auto_execute"],
|
||||
parsed["reply_to"],
|
||||
target,
|
||||
parsed["no_memory_save"],
|
||||
from_branch=from_branch,
|
||||
):
|
||||
ok += 1
|
||||
console.print(f"\nGroup send complete: {ok}/{len(recipients)} delivered")
|
||||
return ok > 0
|
||||
|
||||
|
||||
def _send_interactive() -> bool:
|
||||
"""Interactive email sending with prompts."""
|
||||
branches = get_all_branches()
|
||||
console.print("\nAI_Mail - Send Email\n" + "=" * 50)
|
||||
console.print("\nSelect recipient:")
|
||||
for i, b in enumerate(branches, 1):
|
||||
console.print(f" {i}. {b['name']} ({b['email']})")
|
||||
console.print(f" {len(branches) + 1}. ALL BRANCHES (broadcast)")
|
||||
console.print("Message (press Ctrl+D when done, Ctrl+C to cancel):")
|
||||
|
||||
result = collect_interactive_input(branches)
|
||||
if result is None:
|
||||
console.print("\nCancelled")
|
||||
return False
|
||||
|
||||
console.print("\n" + "=" * 50)
|
||||
console.print(f"To: {result['to']}\nSubject: {result['subject']}\nMessage:\n{result['message']}")
|
||||
console.print("=" * 50)
|
||||
return _send_direct(result["to"], result["subject"], result["message"])
|
||||
|
||||
|
||||
def _send_direct(
|
||||
to_branch,
|
||||
subject,
|
||||
message,
|
||||
auto_execute=False,
|
||||
reply_to=None,
|
||||
dispatched_to=None,
|
||||
no_memory_save=False,
|
||||
from_branch=None,
|
||||
) -> bool:
|
||||
"""Direct email send - thin wrapper over send handlers."""
|
||||
try:
|
||||
user_info = resolve_sender_info(from_branch, _REPO_ROOT, _AI_MAIL_DIR, get_branch_by_email, get_current_user)
|
||||
if auto_execute:
|
||||
message = prepend_dispatch_header(message, no_memory_save=no_memory_save)
|
||||
|
||||
if to_branch.lower() in ["all", "@all"]:
|
||||
return _send_broadcast(subject, message, user_info, auto_execute, no_memory_save, reply_to, dispatched_to)
|
||||
|
||||
success, error_msg = send_to_single(
|
||||
to_branch,
|
||||
subject,
|
||||
message,
|
||||
user_info,
|
||||
auto_execute,
|
||||
no_memory_save,
|
||||
reply_to,
|
||||
dispatched_to,
|
||||
create_email_file,
|
||||
load_email_file,
|
||||
deliver_email_to_branch,
|
||||
_delivery_callback,
|
||||
json_handler.log_operation,
|
||||
update_central,
|
||||
)
|
||||
|
||||
if success:
|
||||
label = "\\[dispatch: queued for daemon]" if auto_execute else ""
|
||||
console.print(f"[green]Email sent to {to_branch} {label}[/green]")
|
||||
if auto_execute:
|
||||
_fire_dispatch_trigger(to_branch, subject)
|
||||
return True
|
||||
else:
|
||||
error(f"Failed to deliver: {error_msg}")
|
||||
dispatch_send_error(to_branch, subject, error_msg or "", deliver_email_to_branch)
|
||||
return False
|
||||
except BrokenPipeError:
|
||||
logger.info("[email_send] Send: broken pipe (stdout closed early)")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"[email_send] Send failed: {e}")
|
||||
error(f"Error: {e}")
|
||||
dispatch_send_error(to_branch, subject, str(e), deliver_email_to_branch)
|
||||
return False
|
||||
|
||||
|
||||
def _fire_dispatch_trigger(to_branch: str, subject: str) -> None:
|
||||
"""Fire email_dispatched trigger event if auto_execute enabled."""
|
||||
try:
|
||||
trigger.fire("email_dispatched", to=to_branch, subject=subject)
|
||||
except Exception as e:
|
||||
logger.warning("[email_send] trigger fire for email_dispatched failed: %s", e)
|
||||
|
||||
|
||||
def _send_broadcast(subject, message, user_info, auto_execute, no_memory_save, reply_to, dispatched_to) -> bool:
|
||||
"""Broadcast send to all branches - display wrapper."""
|
||||
branches = get_all_branches()
|
||||
console.print(f"\nBroadcasting to {len(branches)} branches...")
|
||||
ok, success_count, total, results = send_to_broadcast(
|
||||
subject,
|
||||
message,
|
||||
user_info,
|
||||
auto_execute,
|
||||
no_memory_save,
|
||||
reply_to,
|
||||
dispatched_to,
|
||||
branches,
|
||||
create_email_file,
|
||||
load_email_file,
|
||||
deliver_email_to_branch,
|
||||
_delivery_callback,
|
||||
json_handler.log_operation,
|
||||
update_central,
|
||||
)
|
||||
if isinstance(results, str) or results is None:
|
||||
error("Failed to load email file for broadcast")
|
||||
return False
|
||||
for name, ok, err in results: # type: ignore[union-attr]
|
||||
if ok:
|
||||
console.print(f" [green]OK[/green] {name}")
|
||||
else:
|
||||
error(f"FAIL {name} ({err})")
|
||||
console.print(f"\nBroadcast complete: {success_count}/{total} delivered")
|
||||
return ok
|
||||
|
||||
|
||||
def print_introspection():
|
||||
"""Print module introspection for seedgo compliance."""
|
||||
console.print("\n" + "=" * 70)
|
||||
console.print("EMAIL SEND ORCHESTRATION")
|
||||
console.print("=" * 70)
|
||||
console.print("\nFunctions provided:")
|
||||
console.print(" - handle_send(args) -> bool")
|
||||
console.print(" - _send_direct(...) -> bool")
|
||||
console.print(" - _send_interactive() -> bool")
|
||||
console.print(" - _send_broadcast(...) -> bool")
|
||||
console.print(" - _fire_dispatch_trigger(to_branch, subject) -> None")
|
||||
console.print(" - _delivery_callback(branch_path, new_count, opened_count, total)")
|
||||
console.print()
|
||||
console.print("=" * 70 + "\n")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
print_introspection()
|
||||
@@ -511,9 +511,9 @@ def test_is_protected_branch_empty_string():
|
||||
|
||||
# ---- _has_test_token tests -----------------------------------
|
||||
|
||||
from aipass.ai_mail.apps.handlers.dispatch.daemon import ( # noqa: E402
|
||||
_has_test_token,
|
||||
_auto_ack_test_email,
|
||||
from aipass.ai_mail.apps.handlers.dispatch.test_token import ( # noqa: E402
|
||||
has_test_token as _has_test_token,
|
||||
auto_ack_test_email as _auto_ack_test_email,
|
||||
scan_and_ack_test_emails,
|
||||
TEST_TOKEN,
|
||||
)
|
||||
@@ -620,7 +620,7 @@ def test_scan_and_ack_test_emails_acks_matching(tmp_path):
|
||||
}
|
||||
(ai_mail_local / "inbox.json").write_text(json.dumps(inbox))
|
||||
|
||||
with patch("aipass.ai_mail.apps.handlers.dispatch.daemon._auto_ack_test_email", return_value=True) as mock_ack:
|
||||
with patch("aipass.ai_mail.apps.handlers.dispatch.test_token.auto_ack_test_email", return_value=True) as mock_ack:
|
||||
count = scan_and_ack_test_emails(branch_path, "@testbranch")
|
||||
|
||||
assert count == 1
|
||||
@@ -639,7 +639,7 @@ def test_scan_and_ack_test_emails_skips_closed(tmp_path):
|
||||
}
|
||||
(ai_mail_local / "inbox.json").write_text(json.dumps(inbox))
|
||||
|
||||
with patch("aipass.ai_mail.apps.handlers.dispatch.daemon._auto_ack_test_email") as mock_ack:
|
||||
with patch("aipass.ai_mail.apps.handlers.dispatch.test_token.auto_ack_test_email") as mock_ack:
|
||||
count = scan_and_ack_test_emails(branch_path, "@testbranch")
|
||||
|
||||
assert count == 0
|
||||
|
||||
@@ -1,24 +1,23 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: test_dispatch_watchdog.py
|
||||
# Description: Tests for watchdog auto-spawn in dispatch pipeline
|
||||
# Version: 1.0.0
|
||||
# Version: 1.1.0
|
||||
# Created: 2026-04-19
|
||||
# Modified: 2026-04-19
|
||||
# Modified: 2026-04-22
|
||||
# =============================================
|
||||
|
||||
"""Tests for _spawn_watchdog() — watchdog auto-spawn in dispatch pipeline."""
|
||||
|
||||
import json
|
||||
from unittest.mock import patch, MagicMock
|
||||
import pytest
|
||||
|
||||
import aipass.ai_mail.apps.modules.dispatch as dispatch_mod
|
||||
|
||||
|
||||
# Retrieve private function via module attribute access
|
||||
_spawn_watchdog = getattr(dispatch_mod, "_spawn_watchdog")
|
||||
|
||||
_POPEN_PATH = "aipass.ai_mail.apps.modules.dispatch.subprocess.Popen"
|
||||
_GET_BRANCH = "aipass.ai_mail.apps.handlers.registry.read.get_branch_by_email"
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
@@ -37,16 +36,14 @@ class TestSpawnWatchdog:
|
||||
"""_spawn_watchdog sets cwd=devpulse_path when spawning."""
|
||||
devpulse_path = tmp_path / "devpulse"
|
||||
devpulse_path.mkdir()
|
||||
registry_file = tmp_path / "AIPASS_REGISTRY.json"
|
||||
registry_file.write_text(
|
||||
json.dumps({"branches": [{"name": "DEVPULSE", "email": "@devpulse", "path": str(devpulse_path)}]}),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
fake_proc = MagicMock()
|
||||
fake_proc.pid = 42
|
||||
|
||||
with patch(_POPEN_PATH, return_value=fake_proc) as mock_popen:
|
||||
with (
|
||||
patch(_GET_BRANCH, return_value={"path": str(devpulse_path), "email": "@devpulse"}),
|
||||
patch(_POPEN_PATH, return_value=fake_proc) as mock_popen,
|
||||
):
|
||||
result = _spawn_watchdog("@drone", tmp_path)
|
||||
|
||||
assert result is True
|
||||
@@ -58,44 +55,25 @@ class TestSpawnWatchdog:
|
||||
|
||||
def test_returns_false_when_devpulse_not_in_registry(self, tmp_path):
|
||||
"""Returns False when devpulse not found in registry."""
|
||||
registry_file = tmp_path / "AIPASS_REGISTRY.json"
|
||||
registry_file.write_text(
|
||||
json.dumps({"branches": [{"name": "DRONE", "email": "@drone", "path": str(tmp_path / "drone")}]}),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
result = _spawn_watchdog("@drone", tmp_path)
|
||||
assert result is False
|
||||
|
||||
def test_returns_false_when_registry_missing(self, tmp_path):
|
||||
"""Returns False when AIPASS_REGISTRY.json does not exist."""
|
||||
result = _spawn_watchdog("@drone", tmp_path)
|
||||
with patch(_GET_BRANCH, return_value=None):
|
||||
result = _spawn_watchdog("@drone", tmp_path)
|
||||
assert result is False
|
||||
|
||||
def test_returns_false_when_devpulse_path_missing(self, tmp_path):
|
||||
"""Returns False when devpulse path from registry does not exist on disk."""
|
||||
registry_file = tmp_path / "AIPASS_REGISTRY.json"
|
||||
registry_file.write_text(
|
||||
json.dumps(
|
||||
{"branches": [{"name": "DEVPULSE", "email": "@devpulse", "path": str(tmp_path / "nonexistent")}]}
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
result = _spawn_watchdog("@drone", tmp_path)
|
||||
with patch(_GET_BRANCH, return_value={"path": str(tmp_path / "nonexistent"), "email": "@devpulse"}):
|
||||
result = _spawn_watchdog("@drone", tmp_path)
|
||||
assert result is False
|
||||
|
||||
def test_returns_false_when_drone_not_found(self, tmp_path):
|
||||
"""Returns False when 'drone' binary not on PATH (FileNotFoundError)."""
|
||||
devpulse_path = tmp_path / "devpulse"
|
||||
devpulse_path.mkdir()
|
||||
registry_file = tmp_path / "AIPASS_REGISTRY.json"
|
||||
registry_file.write_text(
|
||||
json.dumps({"branches": [{"name": "DEVPULSE", "email": "@devpulse", "path": str(devpulse_path)}]}),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
with patch(_POPEN_PATH, side_effect=FileNotFoundError("drone not found")):
|
||||
with (
|
||||
patch(_GET_BRANCH, return_value={"path": str(devpulse_path), "email": "@devpulse"}),
|
||||
patch(_POPEN_PATH, side_effect=FileNotFoundError("drone not found")),
|
||||
):
|
||||
result = _spawn_watchdog("@drone", tmp_path)
|
||||
|
||||
assert result is False
|
||||
@@ -104,16 +82,14 @@ class TestSpawnWatchdog:
|
||||
"""Resolves relative devpulse path relative to repo_root."""
|
||||
devpulse_path = tmp_path / "src" / "devpulse"
|
||||
devpulse_path.mkdir(parents=True)
|
||||
registry_file = tmp_path / "AIPASS_REGISTRY.json"
|
||||
registry_file.write_text(
|
||||
json.dumps({"branches": [{"name": "DEVPULSE", "email": "@devpulse", "path": "src/devpulse"}]}),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
fake_proc = MagicMock()
|
||||
fake_proc.pid = 99
|
||||
|
||||
with patch(_POPEN_PATH, return_value=fake_proc) as mock_popen:
|
||||
with (
|
||||
patch(_GET_BRANCH, return_value={"path": "src/devpulse", "email": "@devpulse"}),
|
||||
patch(_POPEN_PATH, return_value=fake_proc) as mock_popen,
|
||||
):
|
||||
result = _spawn_watchdog("@flow", tmp_path)
|
||||
|
||||
assert result is True
|
||||
|
||||
@@ -1,113 +0,0 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: test_identity.py
|
||||
# Description: Tests for the per-branch identity handler
|
||||
# Version: 1.0.0
|
||||
# Created: 2026-04-11
|
||||
# Modified: 2026-04-11
|
||||
# =============================================
|
||||
|
||||
"""Tests for per-branch identity handler (DPLAN-0121 Phase 5).
|
||||
|
||||
Bypass entries for architecture and encapsulation are in .seedgo/bypass.json.
|
||||
"""
|
||||
|
||||
import json
|
||||
import pytest
|
||||
from unittest.mock import patch
|
||||
|
||||
from aipass.ai_mail.apps.handlers.email.identity import (
|
||||
create_identity,
|
||||
read_identity,
|
||||
)
|
||||
|
||||
|
||||
# ---- Fixtures ------------------------------------------------
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _silence_json_handler():
|
||||
"""Prevent log_operation from writing real JSON files during tests."""
|
||||
with patch("aipass.ai_mail.apps.handlers.email.identity.json_handler") as mock_jh:
|
||||
mock_jh.log_operation.return_value = True
|
||||
yield mock_jh
|
||||
|
||||
|
||||
# ---- create_identity() tests --------------------------------
|
||||
|
||||
|
||||
def test_create_identity_writes_file(tmp_path):
|
||||
"""create_identity writes identity.json with correct fields."""
|
||||
ok = create_identity(tmp_path, "devpulse", "AIPass")
|
||||
assert ok is True
|
||||
|
||||
identity_file = tmp_path / ".ai_mail.local" / "identity.json"
|
||||
assert identity_file.exists()
|
||||
|
||||
with open(identity_file, "r", encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
|
||||
assert data["branch"] == "devpulse"
|
||||
assert data["project"] == "AIPass"
|
||||
assert data["inbox"] == str(tmp_path / ".ai_mail.local" / "inbox.json")
|
||||
|
||||
|
||||
def test_create_identity_lowercases_name(tmp_path):
|
||||
"""create_identity stores branch name in lowercase."""
|
||||
create_identity(tmp_path, "DEVPULSE", "AIPass")
|
||||
identity_file = tmp_path / ".ai_mail.local" / "identity.json"
|
||||
with open(identity_file, "r", encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
assert data["branch"] == "devpulse"
|
||||
|
||||
|
||||
def test_create_identity_creates_parent_dirs(tmp_path):
|
||||
"""create_identity creates .ai_mail.local/ if it does not exist."""
|
||||
branch_path = tmp_path / "nested" / "branch"
|
||||
# Directory does not exist yet
|
||||
assert not branch_path.exists()
|
||||
|
||||
ok = create_identity(branch_path, "newbranch", "AIPass")
|
||||
assert ok is True
|
||||
assert (branch_path / ".ai_mail.local" / "identity.json").exists()
|
||||
|
||||
|
||||
def test_create_identity_overwrites_existing(tmp_path):
|
||||
"""create_identity overwrites an existing identity.json."""
|
||||
create_identity(tmp_path, "alpha", "OldProject")
|
||||
create_identity(tmp_path, "alpha", "NewProject")
|
||||
|
||||
identity_file = tmp_path / ".ai_mail.local" / "identity.json"
|
||||
with open(identity_file, "r", encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
|
||||
assert data["project"] == "NewProject"
|
||||
|
||||
|
||||
# ---- read_identity() tests ----------------------------------
|
||||
|
||||
|
||||
def test_read_identity_reads_back(tmp_path):
|
||||
"""read_identity returns the dict written by create_identity."""
|
||||
create_identity(tmp_path, "devpulse", "AIPass")
|
||||
result = read_identity(tmp_path)
|
||||
|
||||
assert result is not None
|
||||
assert result["branch"] == "devpulse"
|
||||
assert result["project"] == "AIPass"
|
||||
assert "inbox" in result
|
||||
|
||||
|
||||
def test_read_identity_missing_file(tmp_path):
|
||||
"""read_identity returns None when identity.json does not exist."""
|
||||
result = read_identity(tmp_path)
|
||||
assert result is None
|
||||
|
||||
|
||||
def test_read_identity_corrupted_json(tmp_path):
|
||||
"""read_identity returns None when identity.json is invalid JSON."""
|
||||
mail_dir = tmp_path / ".ai_mail.local"
|
||||
mail_dir.mkdir(parents=True)
|
||||
(mail_dir / "identity.json").write_text("not json", encoding="utf-8")
|
||||
|
||||
result = read_identity(tmp_path)
|
||||
assert result is None
|
||||
Reference in New Issue
Block a user