- Reset master to upstream/main (16,697 commits) - Overlay 2,271 local-only files (skills, tools, workspace, configs, apps) - Restore IDENTITY.md and USER.md templates - Build verified, gateway running, Discord working Co-Authored-By: Claude Opus 4.6 <[email protected]>
234 lines
7.4 KiB
Python
234 lines
7.4 KiB
Python
#!/usr/bin/env python3
|
|
import argparse
|
|
import json
|
|
import os
|
|
import time
|
|
import uuid
|
|
from dataclasses import dataclass
|
|
from pathlib import Path
|
|
from typing import Iterable, List
|
|
|
|
from james.lib.log import get_logger
|
|
from james.lib.schema import validate_event
|
|
|
|
log = get_logger("mail")
|
|
|
|
|
|
@dataclass
|
|
class EmailItem:
|
|
message_id: str
|
|
sender: str
|
|
subject: str
|
|
snippet: str
|
|
received_at: int
|
|
labels: List[str]
|
|
|
|
|
|
DEFAULT_FIXTURE = [
|
|
{
|
|
"id": "msg-001",
|
|
"from": "[email protected]",
|
|
"subject": "Invoice for February",
|
|
"snippet": "Your invoice is attached. Please remit payment by Friday.",
|
|
"received_at": 1738706400,
|
|
"labels": ["INBOX", "INVOICE"],
|
|
},
|
|
{
|
|
"id": "msg-002",
|
|
"from": "[email protected]",
|
|
"subject": "New login detected",
|
|
"snippet": "We detected a new login from a new device.",
|
|
"received_at": 1738708200,
|
|
"labels": ["INBOX", "SECURITY"],
|
|
},
|
|
{
|
|
"id": "msg-003",
|
|
"from": "[email protected]",
|
|
"subject": "Action required: status update",
|
|
"snippet": "Please reply with the latest project status by EOD.",
|
|
"received_at": 1738710000,
|
|
"labels": ["INBOX"],
|
|
},
|
|
{
|
|
"id": "msg-004",
|
|
"from": "[email protected]",
|
|
"subject": "Weekly newsletter",
|
|
"snippet": "This week in infrastructure: five takeaways.",
|
|
"received_at": 1738711800,
|
|
"labels": ["INBOX", "NEWSLETTER"],
|
|
},
|
|
]
|
|
|
|
|
|
def redis_url() -> str:
|
|
url = os.environ.get("REDIS_URL")
|
|
if url:
|
|
return url
|
|
host = os.environ.get("REDIS_HOST", "127.0.0.1")
|
|
port = os.environ.get("REDIS_PORT", "6379")
|
|
password = os.environ.get("REDIS_PASSWORD")
|
|
auth = f":{password}@" if password else ""
|
|
return f"redis://{auth}{host}:{port}/0"
|
|
|
|
|
|
def load_fixture(path: str | None) -> List[EmailItem]:
|
|
raw = DEFAULT_FIXTURE
|
|
if path:
|
|
data = json.loads(Path(path).read_text())
|
|
if isinstance(data, dict) and "messages" in data:
|
|
raw = data["messages"]
|
|
elif isinstance(data, list):
|
|
raw = data
|
|
else:
|
|
raise ValueError("fixture must be a list or {messages: [...]} object")
|
|
|
|
items: List[EmailItem] = []
|
|
for entry in raw:
|
|
items.append(
|
|
EmailItem(
|
|
message_id=str(entry.get("id", "")),
|
|
sender=str(entry.get("from", "")),
|
|
subject=str(entry.get("subject", "")),
|
|
snippet=str(entry.get("snippet", "")),
|
|
received_at=int(entry.get("received_at", int(time.time()))),
|
|
labels=list(entry.get("labels", [])),
|
|
)
|
|
)
|
|
return items
|
|
|
|
|
|
def classify_email(email: EmailItem) -> str:
|
|
text = f"{email.subject} {email.snippet}".lower()
|
|
if any(token in text for token in ["invoice", "receipt", "billing", "payment", "subscription"]):
|
|
return "invoice"
|
|
if any(token in text for token in ["security", "login", "password", "2fa", "suspicious", "verification"]):
|
|
return "security_alert"
|
|
if any(token in text for token in ["action required", "please", "reply", "respond", "urgent", "request"]):
|
|
return "action_required"
|
|
if any(token in text for token in ["newsletter", "unsubscribe", "promo", "sale", "digest"]):
|
|
return "newsletter"
|
|
return "fyi"
|
|
|
|
|
|
def severity_for_classification(classification: str) -> str:
|
|
mapping = {
|
|
"security_alert": "high",
|
|
"action_required": "medium",
|
|
"invoice": "medium",
|
|
"newsletter": "low",
|
|
"fyi": "low",
|
|
}
|
|
return mapping.get(classification, "low")
|
|
|
|
|
|
def draft_reply_for(classification: str, email: EmailItem) -> str | None:
|
|
if classification == "security_alert":
|
|
return "Please confirm whether this login was expected."
|
|
if classification == "invoice":
|
|
return "Thanks, invoice received. We will process and confirm shortly."
|
|
if classification == "action_required":
|
|
return "Acknowledged. I will follow up with the requested status update."
|
|
return None
|
|
|
|
|
|
def base_event(event_type: str, severity: str, payload: dict) -> dict:
|
|
return {
|
|
"id": f"evt-{time.strftime('%Y%m%d')}-{uuid.uuid4().hex[:8]}",
|
|
"source": os.environ.get("JAMES_NODE", "mail-agent"),
|
|
"type": event_type,
|
|
"timestamp": int(time.time()),
|
|
"severity": severity,
|
|
"payload": payload,
|
|
}
|
|
|
|
|
|
def build_event(email: EmailItem, dry_run: bool) -> dict:
|
|
classification = classify_email(email)
|
|
severity = severity_for_classification(classification)
|
|
event_type = f"mail.{classification}"
|
|
payload = {
|
|
"email_id": email.message_id,
|
|
"from": email.sender,
|
|
"subject": email.subject,
|
|
"snippet": email.snippet,
|
|
"received_at": email.received_at,
|
|
"labels": email.labels,
|
|
"classification": classification,
|
|
"summary": f"{classification.replace('_', ' ')} email from {email.sender}",
|
|
"draft_reply": draft_reply_for(classification, email),
|
|
"dry_run": dry_run,
|
|
}
|
|
return base_event(event_type, severity, payload)
|
|
|
|
|
|
def publish_events(events: Iterable[dict], channel: str, stdout_only: bool) -> int:
|
|
if stdout_only:
|
|
for event in events:
|
|
print(json.dumps(event))
|
|
return 0
|
|
|
|
try:
|
|
import redis
|
|
except Exception: # pragma: no cover - runtime guard
|
|
log.error("redis is required for publish", exc_info=True)
|
|
return 2
|
|
|
|
try:
|
|
r = redis.Redis.from_url(redis_url())
|
|
r.ping()
|
|
except Exception:
|
|
log.error("redis connection failed", exc_info=True)
|
|
return 1
|
|
|
|
for event in events:
|
|
r.publish(channel, json.dumps(event))
|
|
log.info("published", event_id=event.get("id"), channel=channel, event_type=event.get("type"))
|
|
return 0
|
|
|
|
|
|
def main() -> int:
|
|
parser = argparse.ArgumentParser(description="James Mail Agent (dry-run)")
|
|
parser.add_argument("--fixture", help="Path to JSON fixture with sample messages")
|
|
parser.add_argument("--max", type=int, default=0, help="Max messages to process (0 = all)")
|
|
parser.add_argument("--stdout", action="store_true", help="Print events to stdout instead of Redis")
|
|
parser.add_argument("--dry-run", action="store_true", help="Run without any external side effects")
|
|
args = parser.parse_args()
|
|
|
|
dry_run = args.dry_run or os.environ.get("JAMES_MAIL_DRY_RUN", "1") == "1"
|
|
if not dry_run:
|
|
log.error("live mail fetch not implemented; use --dry-run")
|
|
return 2
|
|
|
|
channel = os.environ.get("JAMES_MAIL_CHANNEL", "james.events.mail")
|
|
stdout_only = args.stdout or os.environ.get("JAMES_MAIL_STDOUT", "0") == "1"
|
|
|
|
try:
|
|
emails = load_fixture(args.fixture or os.environ.get("JAMES_MAIL_FIXTURE"))
|
|
except Exception as exc:
|
|
log.error("failed to load fixture", error=str(exc))
|
|
return 1
|
|
|
|
if args.max > 0:
|
|
emails = emails[: args.max]
|
|
|
|
events: List[dict] = []
|
|
for email in emails:
|
|
event = build_event(email, dry_run)
|
|
try:
|
|
validate_event(event)
|
|
except Exception:
|
|
log.error("invalid event", exc_info=True)
|
|
continue
|
|
events.append(event)
|
|
|
|
if not events:
|
|
log.info("no events generated")
|
|
return 0
|
|
|
|
log.info("prepared events", count=len(events), channel=channel, dry_run=dry_run, stdout_only=stdout_only)
|
|
return publish_events(events, channel, stdout_only)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|