Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
165 changes: 136 additions & 29 deletions src/lingtai/kernel/services/mail.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
import uuid
from abc import ABC, abstractmethod
from pathlib import Path
from typing import Callable
from typing import Callable, Iterator

from ..handshake import is_agent, is_alive, manifest, resolve_address

Expand Down Expand Up @@ -130,6 +130,13 @@ def __init__(
self._poll_thread: threading.Thread | None = None
self._poll_stop = threading.Event()
self._seen: set[str] = set()
# Own-inbox scans can be expensive on large/slow external-volume
# mailboxes. Keep iterator progress across poll ticks so Phase 2 can
# be sliced instead of restarting from the first historical entry.
self._own_inbox_iter: Iterator[Path] | None = None
self._own_inbox_slice_seconds = 0.05
self._own_inbox_slice_entries = 200
self._mail_poll_slow_seconds = 1.0

# ------------------------------------------------------------------
# address
Expand Down Expand Up @@ -247,40 +254,139 @@ def _poll_loop() -> None:
while not self._poll_stop.is_set():
# Phase 1 — subscribed pseudo-agent outboxes. Poll these
# before historical own-inbox scans so urgent human/TUI
# wake mail cannot starve behind old inbox entries. Isolated
# from Phase 2 so a persistent OSError here cannot skip the
# same tick's own-inbox scan.
try:
for pseudo_dir in self._pseudo_agent_dirs:
self._poll_pseudo_outbox(pseudo_dir, on_message)
except OSError:
pass
# wake mail cannot starve behind old inbox entries.
self._poll_pseudo_outboxes(on_message)

# Phase 2 — own inbox, sliced. A large historical inbox on a
# slow external volume must not monopolize the poll thread for
# minutes; after every bounded slice we immediately check
# pseudo outboxes again before sleeping.
self._poll_own_inbox_slice(on_message)
self._poll_pseudo_outboxes(on_message)

# Phase 2 — own inbox.
try:
if self._inbox_dir.is_dir():
for entry in self._inbox_dir.iterdir():
# _seen only holds handled directory names or
# pseudo-claim UUIDs, so skip before the stat.
if entry.name in self._seen:
continue
if not entry.is_dir():
continue
msg_file = entry / "message.json"
if msg_file.is_file():
try:
payload = json.loads(msg_file.read_text(encoding="utf-8"))
on_message(payload)
except (json.JSONDecodeError, OSError):
pass
self._seen.add(entry.name)
except OSError:
pass
self._poll_stop.wait(0.5)

self._poll_thread = threading.Thread(target=_poll_loop, daemon=True)
self._poll_thread.start()

def _poll_pseudo_outboxes(self, on_message: Callable[[dict], None]) -> None:
"""Poll subscribed pseudo-agent outboxes, isolating per-dir errors."""
start = time.monotonic()
dirs = 0
for pseudo_dir in self._pseudo_agent_dirs:
dirs += 1
try:
self._poll_pseudo_outbox(pseudo_dir, on_message)
except OSError:
logger.debug(
"pseudo-agent outbox poll failed for %s",
pseudo_dir,
exc_info=True,
)
self._log_slow_mail_phase(
"pseudo_outboxes",
start,
dirs=dirs,
)

def _poll_own_inbox_slice(
self,
on_message: Callable[[dict], None],
*,
budget_seconds: float | None = None,
max_entries: int | None = None,
) -> bool:
"""Poll a bounded slice of own inbox entries.

Returns True when there is more iterator work to continue in a later
poll tick. Keeping iterator progress turns the historical own-inbox
scan into small slices, letting pseudo-agent outboxes be checked
between chunks.
"""
budget = self._own_inbox_slice_seconds if budget_seconds is None else budget_seconds
limit = self._own_inbox_slice_entries if max_entries is None else max_entries
if budget <= 0:
budget = self._own_inbox_slice_seconds
if limit <= 0:
limit = self._own_inbox_slice_entries

start = time.monotonic()
deadline = start + budget
visited = 0
skipped_seen = 0
dispatched = 0

try:
if not self._inbox_dir.is_dir():
self._reset_own_inbox_iter()
return False
if self._own_inbox_iter is None:
self._own_inbox_iter = iter(self._inbox_dir.iterdir())

while not self._poll_stop.is_set() and visited < limit:
try:
entry = next(self._own_inbox_iter)
except StopIteration:
self._reset_own_inbox_iter()
return False

visited += 1
# _seen only holds handled directory names or pseudo-claim
# UUIDs, so skip before the stat.
if entry.name in self._seen:
skipped_seen += 1
continue
if not entry.is_dir():
continue
msg_file = entry / "message.json"
if msg_file.is_file():
try:
payload = json.loads(msg_file.read_text(encoding="utf-8"))
on_message(payload)
dispatched += 1
except (json.JSONDecodeError, OSError):
pass
self._seen.add(entry.name)

if time.monotonic() >= deadline:
break
except OSError:
self._reset_own_inbox_iter()
return False
finally:
self._log_slow_mail_phase(
"own_inbox_slice",
start,
visited=visited,
skipped_seen=skipped_seen,
dispatched=dispatched,
has_more=self._own_inbox_iter is not None,
)

return self._own_inbox_iter is not None

def _reset_own_inbox_iter(self) -> None:
iterator = self._own_inbox_iter
self._own_inbox_iter = None
close = getattr(iterator, "close", None)
if callable(close):
close()

def _log_slow_mail_phase(self, phase: str, start: float, **fields: object) -> None:
elapsed = time.monotonic() - start
if elapsed < self._mail_poll_slow_seconds:
return
redacted_keys = {"message", "body", "content", "subject", "attachments"}
safe_fields = {key: value for key, value in fields.items() if key not in redacted_keys}
details = " ".join(f"{key}={value}" for key, value in sorted(safe_fields.items()))
logger.warning(
"filesystem mail poll phase slow phase=%s elapsed_ms=%.1f agent=%s %s",
phase,
elapsed * 1000,
self.address,
details,
)

def _poll_pseudo_outbox(
self,
pseudo_dir: Path,
Expand Down Expand Up @@ -491,6 +597,7 @@ def stop(self) -> None:
if self._poll_thread is not None:
self._poll_thread.join(timeout=3.0)
self._poll_thread = None
self._reset_own_inbox_iter()


def _runtime_probe_payload(payload: dict) -> dict | None:
Expand Down
96 changes: 96 additions & 0 deletions tests/test_filesystem_mail.py
Original file line number Diff line number Diff line change
Expand Up @@ -796,3 +796,99 @@ def test_pseudo_agent_outbox_skips_non_matching_to(tmp_path):
)
assert received == []
assert received == [], f"on_message must not fire for non-matching To: got {received}"


def test_own_inbox_slice_rechecks_pseudo_outbox_between_chunks(tmp_path):
"""A long own-inbox pass must not hide new human pseudo-mail until exhaustion."""
from lingtai.kernel.services.mail import FilesystemMailService

agent_dir = _make_agent_dir(tmp_path, "agent01")
inbox = agent_dir / "mailbox" / "inbox"
inbox.mkdir(parents=True, exist_ok=True)

human_dir = tmp_path / "human"
human_dir.mkdir()
(human_dir / ".agent.json").write_text(json.dumps({
"agent_name": "human",
"admin": None,
}))
for folder in ["outbox", "sent"]:
(human_dir / "mailbox" / folder).mkdir(parents=True, exist_ok=True)

def write_entry(folder, entry_id, payload):
entry = folder / entry_id
entry.mkdir(parents=True, exist_ok=True)
(entry / "message.json").write_text(json.dumps(payload), encoding="utf-8")
return entry

own_1 = write_entry(inbox, "own-1", {"message": "own-1"})
own_2 = write_entry(inbox, "own-2", {"message": "own-2"})

svc = FilesystemMailService(
agent_dir,
mailbox_rel="mailbox",
pseudo_agent_subscriptions=[str(human_dir)],
)
svc._own_inbox_slice_entries = 1
svc._own_inbox_slice_seconds = 60.0
# Make the slice order deterministic and avoid relying on filesystem order.
svc._own_inbox_iter = iter([own_1, own_2])

received = []

def on_message(payload):
received.append(payload["message"])
if payload["message"] == "own-1":
write_entry(
human_dir / "mailbox" / "outbox",
"human-between-slices",
{
"from": "human",
"to": ["agent01"],
"message": "human-between-slices",
},
)

assert svc._poll_own_inbox_slice(on_message) is True
assert received == ["own-1"]

# This is what the listen loop does after each bounded slice. The pseudo
# mail created while own-inbox work remained should be delivered before the
# next own-inbox entry is processed.
svc._poll_pseudo_outboxes(on_message)
assert received == ["own-1", "human-between-slices"]
assert not (human_dir / "mailbox" / "outbox" / "human-between-slices").exists()
assert (human_dir / "mailbox" / "sent" / "human-between-slices" / "message.json").exists()
assert (inbox / "human-between-slices" / "message.json").exists()

assert svc._poll_own_inbox_slice(on_message) is True
assert received == ["own-1", "human-between-slices", "own-2"]
# One more slice observes iterator exhaustion and resets the cursor.
assert svc._poll_own_inbox_slice(on_message) is False
assert received == ["own-1", "human-between-slices", "own-2"]


def test_mail_poll_slow_phase_telemetry_is_body_free(tmp_path, caplog):
"""Slow mail-poll telemetry should identify phase/agent without body text."""
import logging
import time

from lingtai.kernel.services.mail import FilesystemMailService

agent_dir = _make_agent_dir(tmp_path, "agent01")
svc = FilesystemMailService(agent_dir, mailbox_rel="mailbox")
svc._mail_poll_slow_seconds = 0.0

caplog.set_level(logging.WARNING)
svc._log_slow_mail_phase(
"own_inbox_slice",
time.monotonic(),
visited=1,
dispatched=1,
message="do-not-log-body",
)

assert "filesystem mail poll phase slow" in caplog.text
assert "phase=own_inbox_slice" in caplog.text
assert "agent=agent01" in caplog.text
assert "do-not-log-body" not in caplog.text