Compare commits

...

3 Commits

Author SHA1 Message Date
root
7706a4047c auto-repair: commit 7 uncommitted file(s) — 2026-05-23 2026-05-23 06:00:00 +00:00
Claude Opus 4.6
cd553883ed docs: scope Telepathy inbound path migration (SMTP + Muse Inbox)
Analyzed current architecture across /data/muse/smtp_handler.py and
/data/telepathy/app.py. Inbound path splits SMTP receiving (port 25)
and email dispatch (port 8114), both of which must be consolidated
into Elixir supervision tree.

Summary:
- Muse Inbox (aiosmtpd + aiohttp) ~360 lines
- Telepathy API (FastAPI) ~143 + 204 lines
- Proposed Elixir: ~650-870 lines (SMTPHandler, InboxStore, Pipeline)
- Effort: 2-3 focused sessions
- Recommendation: Schedule for next session (no urgent blockers identified)

Scope document saved to INBOUND_MIGRATION_SCOPE.md with architecture,
complexity breakdown, blockers, risks, and next-step checklist.

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
2026-05-21 06:05:59 +00:00
Claude Opus 4.6
cd3a136bee config: guard runtime env overrides from clobbering test config
SYMBIONT_PORT=8111 is set in the live service shell, causing runtime.exs
to override test.exs's port: 0 even under MIX_ENV=test. Wrap all env-based
overrides in a config_env() != :test guard so tests can run against the
live host without port collision.

40 tests, 0 failures.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-20 06:07:08 +00:00
8 changed files with 782 additions and 16 deletions

136
INBOUND_MIGRATION_SCOPE.md Normal file
View File

@ -0,0 +1,136 @@
# Inbound Email Path Migration Scope
**Date:** 2026-05-21
**Status:** Scoping complete; ready for implementation planning
**Successor to:** Telepathy JMAP outbound (shipped in fe5896b)
---
## 1. Current State
The inbound email path is split across **two Python services**:
### Muse Inbox (`/data/muse/smtp_handler.py` — 360 lines)
- **Purpose:** SMTP server on port 25 accepting inbound mail for `muse@hydrascale.net` and `muse@cortex.hydrascale.net`
- **Architecture:** aiosmtpd + aiohttp async event loop
- **Key features:**
- Parses MIME emails
- Logs raw emails to `inbox.jsonl` (similar to Telepathy's `messages.jsonl`)
- Dispatches to Symbiont API `/task` endpoint (port 8111) with task prompt + context
- Implements exponential backoff retry (3 attempts) on transient Symbiont failures
- Sends reply via Telepathy `/email` (JMAP)
- Marks inbox entries as `processed=True` after success
### Telepathy (`/data/telepathy/app.py` — 143 lines; `/data/telepathy/mailer.py` — 204 lines)
- **Purpose:** HTTP API (FastAPI on port 8114) for message persistence and JMAP email sending
- **Endpoints:** `POST /messages`, `GET /messages`, `GET /messages/unread`, `POST /messages/{id}/read`, `POST /email`, `GET /health`
- **Dependencies:** requests, FastAPI, Pydantic
### Data Flow (Merged View)
```
Inbound email (SMTP:25, Muse)
→ Parse MIME + log to inbox.jsonl
→ POST /task to Symbiont (8111)
→ Claude CLI produces reply
→ POST /email to Telepathy (8114)
→ JMAP send via Fastmail
→ Mark processed in inbox.jsonl
→ (Optionally) Mark Telepathy message as read
```
---
## 2. Proposed Elixir Architecture
**Goal:** Consolidate both services into Symbiont.Telepathy supervision tree on port 8111 (single app).
### New Modules
| Module | Purpose | Est. Lines |
|--------|---------|-----------|
| `Symbiont.Telepathy.SMTPHandler` | gen_smtp_server_session callbacks | 100-150 |
| `Symbiont.Telepathy.InboxStore` | JSONL-backed inbox (like MessageStore) | 80-120 |
| `Symbiont.Telepathy.Pipeline` | Orchestrates receive → Symbiont dispatch → reply → mark read | 150-200 |
| HTTP endpoints (extend `Symbiont.API.ex`) | `/messages`, `/email`, `/health` | 80-120 |
| Tests | SMTP integration, dispatch retry, pipeline | 150-200 |
### Supervision Tree (Updated)
```
Symbiont.Telepathy.Supervisor (one_for_one)
+-- MessageStore (JSONL: Fastmail-fetched emails)
+-- InboxStore (JSONL: SMTP-received emails)
+-- JMAP (GenServer: Fastmail session + sender)
+-- SMTPHandler (wraps :gen_smtp listener on port 25)
+-- TaskSupervisor (spawns Pipeline.process tasks)
```
### Key Design Decisions
1. **Inbox vs Messages:** Separate stores because they track different email directions (received vs polled from Fastmail). Both feed the UI and dispatch logic.
2. **Retry logic:** Implement as a simple loop in Pipeline.process (not a separate backoff supervisor) since Muse inbox already handles transient retries.
3. **API consolidation:** POST `/email` becomes a direct call to `JMAP.send_email/3` instead of an HTTP hop; remove port 8114 entirely.
---
## 3. Complexity Estimate
**Total LOC: ~650-870 across modules + tests**
| Category | Effort | Notes |
|----------|--------|-------|
| SMTP handler (gen_smtp) | Medium (100-150) | Straightforward protocol impl; risk is MIME edge cases |
| Pipeline + dispatch | Medium (150-200) | Needs retry logic, timeout handling, error reporting |
| Inbox store | Small (80-120) | Clone MessageStore pattern; simple JSONL ops |
| API endpoints | Small (80-120) | Thin wrappers around store/sender; mostly copy-paste |
| Tests + integration | Medium (150-200) | Must verify SMTP flow, Symbiont dispatch, reply send |
| **Refactor/Config** | **Small** | Update Application.ex, mix.exs (add :gen_smtp), runtime.exs |
**Blockers:**
- None identified. Both Muse Inbox and Telepathy run standalone with no inter-service locks. Can be developed/tested in parallel with Python services running.
**Risks:**
1. **Email loss on crashes:** If Elixir app restarts mid-Pipeline.process, inbox entry is not marked processed. Mitigation: restore from unprocessed entries on startup.
2. **gen_smtp limitations:** Erlang SMTP impl is less polished than Python aiosmtpd. May hit MIME parsing edge cases. Mitigation: test with real emails from inbox.jsonl.
3. **Symbiont dispatch timeout:** Task may take >600s. Current SMTP client will time out waiting for reply. Mitigation: return 250 OK immediately (like Muse does now), dispatch async.
---
## 4. Recommendation
**Start next session (after Michael confirms).**
**Rationale:**
- Work is well-scoped with a clear blueprint and existing code to port
- JMAP client already shipped; this is a natural successor
- System currently stable (both services running)
- Estimated 2-3 focused sessions including tests and cutover
- **No explicit urgency signal from Michael** (reflection logs don't mention inbound as a blocker)
- Higher-priority items remain blocked (blog project, permissions)
**Decision factors to verify with Michael:**
1. Is the "19 unread" bug urgent enough to fast-track this? (Elixir port fixes it by design)
2. Are there other blockers (blog, permissions) that should be prioritized first?
3. Can SMTP listen on a non-privileged port during development (e.g., 10025) for testing?
---
## 5. Next Steps (If Approved)
1. Confirm with Michael: urgency, priority relative to blog/permissions, port 25 availability
2. Add `:gen_smtp ~> 1.2` to `mix.exs`
3. Implement SMTPHandler as gen_smtp_server_session behaviour
4. Implement InboxStore (clone MessageStore, s/messages/inbox/)
5. Implement Pipeline with retry loop and Symbiont dispatch
6. Add HTTP endpoints to API.ex
7. Integration test: send email to muse@, verify dispatch, verify reply
8. Cutover: stop Python services, verify Elixir handles full flow
9. Decommission `/data/muse/` and `/data/telepathy/` services
---
**Appendix: Python Source Files for Reference**
- `/data/muse/smtp_handler.py` (360 lines) — aiosmtpd handler + Symbiont dispatch logic
- `/data/telepathy/app.py` (143 lines) — FastAPI message/email endpoints
- `/data/telepathy/mailer.py` (204 lines) — JMAP sender (already ported to Elixir)
- `/data/muse/inbox.jsonl` — inbox log (schema: timestamp, from, subject, message_id, body, processed)
- `/data/telepathy/messages.jsonl` — messages log (schema: id, timestamp, source, subject, content, read)

View File

@ -10,7 +10,8 @@ config :symbiont,
config :symbiont, :telepathy,
enabled: true,
poll_interval_ms: 60_000
poll_interval_ms: 60_000,
smtp_port: 25
config :logger,
level: :info

View File

@ -1,5 +1,8 @@
import Config
# Skip env-based overrides in test mode so that test.exs values (port: 0,
# telepathy: disabled) are not stomped by shell env vars from the live service.
if config_env() != :test do
if port = System.get_env("SYMBIONT_PORT") do
config :symbiont, port: String.to_integer(port)
end
@ -11,3 +14,8 @@ end
if System.get_env("TELEPATHY_ENABLED") == "false" do
config :symbiont, :telepathy, enabled: false
end
if smtp_port = System.get_env("SMTP_PORT") do
config :symbiont, :telepathy, smtp_port: String.to_integer(smtp_port)
end
end

View File

@ -11,7 +11,8 @@ defmodule Symbiont.Application do
├── Symbiont.Heartbeat — periodic health checks + queue processing
├── Symbiont.Telepathy.Supervisor — email communication layer (if enabled)
│ ├── Symbiont.Telepathy.MessageStore — JSONL-backed message persistence
│ └── Symbiont.Telepathy.JMAP — Fastmail JMAP client + polling
│ ├── Symbiont.Telepathy.JMAP — Fastmail JMAP client + polling
│ └── Symbiont.Telepathy.SMTP — gen_smtp listener on port 25
└── Bandit (Symbiont.API) — HTTP API
In test mode, the supervisor starts empty — tests manage their own processes.
@ -65,10 +66,11 @@ defmodule Symbiont.Application do
telepathy_cfg = Application.get_env(:symbiont, :telepathy, [])
telepathy_enabled = Keyword.get(telepathy_cfg, :enabled, true)
poll_ms = Keyword.get(telepathy_cfg, :poll_interval_ms, 60_000)
smtp_port = Keyword.get(telepathy_cfg, :smtp_port, 25)
telepathy_children =
if telepathy_enabled do
[{Symbiont.Telepathy.Supervisor, poll_interval_ms: poll_ms}]
[{Symbiont.Telepathy.Supervisor, poll_interval_ms: poll_ms, smtp_port: smtp_port}]
else
[]
end

View File

@ -0,0 +1,309 @@
defmodule Symbiont.Telepathy.Pipeline do
@moduledoc """
Inbound email processing pipeline.
Called from `Symbiont.Telepathy.SMTP.handle_DATA/4` via a
`Task.Supervisor` task. Runs entirely asynchronously — the SMTP "250 OK"
has already been sent before this module is invoked.
## Pipeline steps
1. **Store** — Persist the inbound email in `MessageStore` (messages.jsonl).
2. **Dispatch** — POST to the Symbiont `/task` endpoint with the full email
context and a prompt asking Muse to respond. Retries on transient errors
using an exponential backoff schedule of 10s → 30s → 90s (4 total attempts).
3. **Reply** — Send the response via `Symbiont.Telepathy.JMAP.send_email/3`.
4. **Mark read** — Flip `processed: true` in `inbox.jsonl` and `read: true`
in `MessageStore`.
## Error handling
Transient failures (connection errors, HTTP 5xx) are retried. Permanent
failures (HTTP 4xx, empty replies) stop retrying immediately. After all
attempts are exhausted, the `inbox.jsonl` entry is left with `processed: false`
so it can be manually re-dispatched.
## Prompt fidelity
The prompt sent to `/task` is identical to the one used by the Python
`smtp_handler.py`, including the optional `context.md` prefix. This ensures
consistent Muse behaviour during the Python-to-Elixir migration period.
"""
require Logger
@symbiont_url "http://127.0.0.1:8111"
@inbox_log "/data/muse/inbox.jsonl"
@context_path "/data/muse/context.md"
# Retry schedule in milliseconds: attempt 1 → immediate, 2 → 10s, 3 → 30s, 4 → 90s
@backoff_schedule_ms [10_000, 30_000, 90_000]
# ---------------------------------------------------------------------------
# Public API
# ---------------------------------------------------------------------------
@doc """
Run the full inbound email pipeline.
Called by `Symbiont.Telepathy.SMTP` under `Symbiont.TaskSupervisor`.
Returns `:ok` on success, `:error` on terminal failure.
"""
def process(from_addr, subject, body, message_id) do
Logger.info("Pipeline: processing email from #{from_addr} (#{subject})")
# Step 1: store in MessageStore
telepathy_id = store_message(from_addr, subject, body)
# Step 2: dispatch to Symbiont with retry
case dispatch_with_retry(from_addr, subject, body, message_id) do
{:ok, reply_text} ->
send_and_mark(from_addr, subject, reply_text, message_id, telepathy_id)
{:error, reason} ->
Logger.error(
"Pipeline: dispatch exhausted for #{from_addr} " <>
"(subject=#{inspect(subject)}): #{inspect(reason)}"
)
:error
end
end
# ---------------------------------------------------------------------------
# Step 1: MessageStore
# ---------------------------------------------------------------------------
defp store_message(from_addr, subject, body) do
msg = %{
"content" => "From: #{from_addr}\nSubject: #{subject}\n\n#{String.slice(body, 0, 2000)}",
"source" => "email-in",
"subject" => "[Inbound] #{subject}"
}
case Symbiont.Telepathy.MessageStore.store(msg) do
{:ok, id} -> id
_ -> nil
end
end
# ---------------------------------------------------------------------------
# Step 2: Symbiont dispatch with backoff retry
# ---------------------------------------------------------------------------
defp dispatch_with_retry(from_addr, subject, body, message_id) do
do_dispatch(from_addr, subject, body, message_id, @backoff_schedule_ms, 1)
end
# No more backoff slots — last attempt.
defp do_dispatch(from_addr, subject, body, message_id, [], attempt) do
case dispatch_to_symbiont(from_addr, subject, body) do
{:ok, reply} ->
{:ok, reply}
{_, detail} ->
Logger.error(
"Pipeline: giving up after #{attempt} attempt(s) for #{from_addr} " <>
"(message_id=#{inspect(message_id)}): #{detail}"
)
{:error, detail}
end
end
defp do_dispatch(from_addr, subject, body, message_id, [delay | rest], attempt) do
case dispatch_to_symbiont(from_addr, subject, body) do
{:ok, reply} ->
{:ok, reply}
{:permanent, detail} ->
Logger.error(
"Pipeline: permanent failure for #{from_addr} " <>
"(message_id=#{inspect(message_id)}): #{detail}"
)
{:error, :permanent}
{:transient, detail} ->
Logger.warning(
"Pipeline: transient failure (attempt #{attempt}) for #{from_addr}: " <>
"#{detail}. Retrying in #{div(delay, 1000)}s..."
)
Process.sleep(delay)
do_dispatch(from_addr, subject, body, message_id, rest, attempt + 1)
end
end
defp dispatch_to_symbiont(from_addr, subject, body) do
context = load_context()
prompt = build_prompt(context, from_addr, subject, body)
payload = Jason.encode!(%{"task" => prompt, "force_tier" => "sonnet"})
url = ~c"#{@symbiont_url}/task"
headers = [{~c"Content-Type", ~c"application/json"}]
case :httpc.request(
:post,
{url, headers, ~c"application/json", payload},
[{:timeout, 600_000}, {:connect_timeout, 10_000}],
[]
) do
{:ok, {{_, 200, _}, _, resp_body}} ->
parse_symbiont_response(resp_body)
{:ok, {{_, status, _}, _, resp_body}} when status >= 500 ->
detail = resp_body |> IO.iodata_to_binary() |> String.slice(0, 200)
{:transient, "HTTP #{status}: #{detail}"}
{:ok, {{_, status, _}, _, resp_body}} ->
detail = resp_body |> IO.iodata_to_binary() |> String.slice(0, 200)
{:permanent, "HTTP #{status}: #{detail}"}
{:error, {:failed_connect, _}} ->
{:transient, "connection refused to #{@symbiont_url}"}
{:error, reason} ->
{:transient, "httpc error: #{inspect(reason)}"}
end
end
defp parse_symbiont_response(resp_body) do
case Jason.decode(IO.iodata_to_binary(resp_body)) do
{:ok, %{"result" => reply}} when is_binary(reply) and reply != "" ->
{:ok, String.trim(reply)}
{:ok, _} ->
{:permanent, "empty or missing 'result' in Symbiont response"}
{:error, _} ->
{:permanent, "invalid JSON from Symbiont"}
end
end
# ---------------------------------------------------------------------------
# Step 3 + 4: Reply and mark read
# ---------------------------------------------------------------------------
defp send_and_mark(from_addr, subject, reply_text, message_id, telepathy_id) do
re_subject = if String.starts_with?(subject, "Re:"), do: subject, else: "Re: #{subject}"
reply_to = extract_email_addr(from_addr)
case Symbiont.Telepathy.JMAP.send_email(reply_to, re_subject, reply_text) do
{:ok, _email_id} ->
Logger.info("Pipeline: reply sent to #{reply_to} (#{re_subject})")
mark_inbox_processed(message_id)
if telepathy_id, do: Symbiont.Telepathy.MessageStore.mark_read(telepathy_id)
:ok
{:error, reason} ->
Logger.error("Pipeline: JMAP send failed for #{from_addr}: #{inspect(reason)}")
:error
end
end
# Extract bare email from "Display Name <addr@domain>" or plain "addr@domain".
defp extract_email_addr(addr) do
case Regex.run(~r/<([^>]+)>/, addr) do
[_, email] -> String.trim(email)
_ -> String.trim(addr)
end
end
# ---------------------------------------------------------------------------
# inbox.jsonl: flip processed flag after successful dispatch
# ---------------------------------------------------------------------------
defp mark_inbox_processed(message_id) when is_binary(message_id) and message_id != "" do
case File.read(@inbox_log) do
{:ok, content} ->
lines =
content
|> String.split("\n", trim: true)
|> Enum.map(&update_processed_flag(&1, message_id))
tmp = @inbox_log <> ".tmp"
File.write!(tmp, Enum.join(lines, "\n") <> "\n")
File.rename!(tmp, @inbox_log)
Logger.info("Pipeline: marked inbox entry processed: #{message_id}")
{:error, _} ->
:ok
end
rescue
e -> Logger.warning("Pipeline: mark_inbox_processed failed: #{Exception.message(e)}")
end
defp mark_inbox_processed(_), do: :ok
defp update_processed_flag(line, message_id) do
case Jason.decode(line) do
{:ok, entry} when is_map(entry) ->
if entry["message_id"] == message_id and entry["processed"] == false do
entry
|> Map.put("processed", true)
|> Map.put("processed_at", DateTime.utc_now() |> DateTime.to_iso8601())
|> Jason.encode!()
else
line
end
_ ->
line
end
end
# ---------------------------------------------------------------------------
# Prompt construction
# ---------------------------------------------------------------------------
defp load_context do
case File.read(@context_path) do
{:ok, content} -> content
{:error, _} -> ""
end
end
defp build_prompt(context, from_addr, subject, body) do
"""
#{context}
---
You received an email from Michael.
From: #{from_addr}
Subject: #{subject}
#{body}
---
You are Muse. You have full system access on cortex via your tools (Bash, Read,
Edit, Write, etc.). The dispatcher runs you under --dangerously-skip-permissions,
so you do NOT need to ask for permission — just act.
Decide whether this email is:
(A) A REQUEST TO DO WORK (imperative): "please commit X", "fix the Y bug",
"check Z and tell me", "write a script that…", "investigate…", etc.
→ DO THE WORK NOW using your tools. Then write a brief reply summarizing
what you actually did (file paths touched, commits made, what you found).
Be specific: include real command output, real diffs, real numbers.
Do NOT promise to do it later — the only "later" that exists is the next
time he emails you. Either it gets done in this session or it doesn't.
(B) A CONVERSATIONAL MESSAGE or QUESTION: just write a reply, grounded in
real system state (run quick checks with your tools if useful).
For both, the reply should:
- Be concise (2-4 short paragraphs is usually right; longer if results demand it)
- Use real data from your tools, never invented numbers or fake commit hashes
- Sign off as "Muse"
- Contain NO preamble, NO markdown fences, NO meta-commentary about what you're
about to do — just the email body itself
Return ONLY the reply text.
"""
end
end

View File

@ -0,0 +1,303 @@
defmodule Symbiont.Telepathy.SMTP do
@moduledoc """
Muse Inbox SMTP receiver — gen_smtp_server_session behaviour.
Listens on port 25 (configurable via `:smtp_port`), accepts mail addressed
to `muse@hydrascale.net` or `muse@cortex.hydrascale.net`, and hands off
each accepted message to `Symbiont.Telepathy.Pipeline` for async processing.
## Data flow
SMTP client
→ handle_RCPT/2 — validate recipient
→ handle_DATA/4 — parse MIME via :mimemail, log to inbox.jsonl,
spawn Pipeline task via Task.Supervisor
→ "250 OK" — response returned to sender immediately
Processing (dispatch to Claude, send reply, mark read) happens asynchronously
under `Symbiont.TaskSupervisor` so the SMTP response is never delayed by
upstream latency. See `Symbiont.Telepathy.Pipeline` for the pipeline logic.
## Supervision
`Symbiont.Telepathy.Supervisor` starts this module as a `:supervisor`-type
child (gen_smtp_server is itself a supervisor internally).
## Inbox log
Every accepted message is appended to `/data/muse/inbox.jsonl` in the same
schema used by the Python `smtp_handler.py`, ensuring backward compatibility
with any tooling that reads that file. The `processed` flag is flipped to
`true` by `Pipeline` after successful dispatch.
"""
@behaviour :gen_smtp_server_session
require Logger
@accepted_recipients ["muse@hydrascale.net", "muse@cortex.hydrascale.net"]
@inbox_log "/data/muse/inbox.jsonl"
@domain "cortex.hydrascale.net"
defstruct from: nil, rcpt_to: []
# ---------------------------------------------------------------------------
# Supervision interface
# ---------------------------------------------------------------------------
@doc """
Returns a child spec suitable for `Supervisor.start_link/2`.
Passes `:smtp_port` through to `start_link/1`.
"""
def child_spec(opts) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
type: :supervisor,
restart: :permanent
}
end
@doc """
Starts the gen_smtp_server listener process.
`opts` must include `:smtp_port` (integer). Falls back to 25 if absent.
"""
def start_link(opts) do
port = Keyword.get(opts, :smtp_port, 25)
Logger.info("Telepathy SMTP: starting listener on port #{port}")
:gen_smtp_server.start_link(__MODULE__, [
[port: port, domain: @domain, address: {0, 0, 0, 0}]
])
end
# ---------------------------------------------------------------------------
# gen_smtp_server_session callbacks
# ---------------------------------------------------------------------------
@impl :gen_smtp_server_session
def init(_hostname, _session_count, _peer_info, _options) do
banner = "#{@domain} ESMTP Muse Inbox ready"
{:ok, banner, %__MODULE__{}}
end
@impl :gen_smtp_server_session
def handle_HELO(_hostname, state) do
{:ok, state}
end
@impl :gen_smtp_server_session
def handle_EHLO(_hostname, extensions, state) do
# Return the default extension list unchanged. SIZE and 8BITMIME are
# included by gen_smtp by default; we don't advertise AUTH.
{:ok, extensions, state}
end
@impl :gen_smtp_server_session
def handle_MAIL(from, state) do
{:ok, %{state | from: from}}
end
@impl :gen_smtp_server_session
def handle_MAIL_extension(_extension, state) do
{:ok, state}
end
@impl :gen_smtp_server_session
def handle_RCPT(to, state) do
addr =
to
|> to_string()
|> String.trim()
|> String.trim_leading("<")
|> String.trim_trailing(">")
|> String.downcase()
if addr in @accepted_recipients do
{:ok, %{state | rcpt_to: [to | state.rcpt_to]}}
else
Logger.info("Telepathy SMTP: rejected recipient #{addr}")
{:error, "550 Recipient #{addr} not accepted", state}
end
end
@impl :gen_smtp_server_session
def handle_RCPT_extension(_extension, state) do
{:ok, state}
end
@doc """
Core DATA handler. Called once the entire message has been received.
1. Parses the raw SMTP payload via `:mimemail.decode/1`.
2. Appends a raw entry (processed: false) to `inbox.jsonl`.
3. Fires a `Symbiont.Telepathy.Pipeline.process/4` task under
`Symbiont.TaskSupervisor` — returns 250 immediately without waiting.
"""
@impl :gen_smtp_server_session
def handle_DATA(from, to, data, state) do
case parse_email(data) do
{:ok, parsed} ->
entry = %{
"timestamp" => DateTime.utc_now() |> DateTime.to_iso8601(),
"from" => parsed.from,
"to" => Enum.join(to, ", "),
"subject" => parsed.subject,
"message_id" => parsed.message_id,
"body" => String.slice(parsed.body, 0, 5000),
"processed" => false
}
append_inbox_log(entry)
Task.Supervisor.start_child(Symbiont.TaskSupervisor, fn ->
Symbiont.Telepathy.Pipeline.process(
parsed.from,
parsed.subject,
parsed.body,
parsed.message_id
)
end)
ref = :crypto.strong_rand_bytes(4) |> Base.encode16(case: :lower)
Logger.info("Telepathy SMTP: accepted #{ref} from #{from} (#{parsed.subject})")
{:ok, ref, state}
{:error, reason} ->
Logger.error("Telepathy SMTP: parse error from #{from}: #{inspect(reason)}")
{:error, "451 Error processing message", state}
end
end
@impl :gen_smtp_server_session
def handle_RSET(state) do
{:ok, %{state | from: nil, rcpt_to: []}}
end
@impl :gen_smtp_server_session
def handle_VRFY(_address, state) do
{:error, "252 Cannot VRFY user, but will accept message and attempt delivery", state}
end
@impl :gen_smtp_server_session
def handle_other(_verb, _args, state) do
{"500 Command unrecognized", state}
end
@impl :gen_smtp_server_session
def handle_AUTH(_type, _username, _password, state) do
{:error, state}
end
@impl :gen_smtp_server_session
def handle_STARTTLS(state) do
{:ok, state}
end
@impl :gen_smtp_server_session
def code_change(_old_vsn, state, _extra) do
{:ok, state}
end
@impl :gen_smtp_server_session
def terminate(reason, state) do
if reason not in [:normal, :shutdown] do
Logger.debug("Telepathy SMTP: session terminated (#{inspect(reason)})")
end
{:ok, state}
end
# ---------------------------------------------------------------------------
# Private: MIME parsing via :mimemail
# ---------------------------------------------------------------------------
defp parse_email(data) do
try do
{type, subtype, headers, _params, body} = :mimemail.decode(data)
from = find_header(headers, "From") || "unknown"
subject = find_header(headers, "Subject") || "(no subject)"
message_id = find_header(headers, "Message-ID") || ""
text_body = extract_body(type, subtype, headers, body)
{:ok, %{from: from, subject: subject, message_id: message_id, body: text_body}}
rescue
e -> {:error, Exception.message(e)}
catch
kind, value -> {:error, {kind, value}}
end
end
# Plain text — use directly.
defp extract_body("text", "plain", _headers, body) when is_binary(body) do
String.trim(body)
end
# HTML with no plain-text alternative — return a placeholder matching the
# Python smtp_handler.py behaviour.
defp extract_body("text", "html", headers, _body) do
from = find_header(headers, "From") || "sender"
"[HTML content from #{from}]"
end
# Multipart — walk parts, prefer text/plain, fall back to HTML placeholder.
defp extract_body("multipart", _subtype, _headers, parts) when is_list(parts) do
plain = find_part(parts, "text", "plain")
html = find_part(parts, "text", "html")
cond do
plain != nil ->
{_, _, _, _, body} = plain
String.trim(body)
html != nil ->
{_, _, part_headers, _, _} = html
from = find_header(part_headers, "From") || "sender"
"[HTML content from #{from}]"
true ->
""
end
end
defp extract_body(_type, _subtype, _headers, body) when is_binary(body) do
String.trim(body)
end
defp extract_body(_type, _subtype, _headers, _body), do: ""
defp find_part(parts, type, subtype) do
Enum.find(parts, fn
{^type, ^subtype, _, _, _} -> true
_ -> false
end)
end
defp find_header(headers, name) do
name_lower = String.downcase(name)
Enum.find_value(headers, nil, fn
{n, v} ->
n_str = to_string(n)
if String.downcase(n_str) == name_lower, do: to_string(v), else: nil
_ ->
nil
end)
end
# ---------------------------------------------------------------------------
# Private: inbox.jsonl logging
# ---------------------------------------------------------------------------
defp append_inbox_log(entry) do
File.mkdir_p!(Path.dirname(@inbox_log))
File.write!(@inbox_log, Jason.encode!(entry) <> "\n", [:append])
rescue
e -> Logger.warning("Telepathy SMTP: inbox log write failed: #{Exception.message(e)}")
end
end

View File

@ -2,9 +2,13 @@ defmodule Symbiont.Telepathy.Supervisor do
@moduledoc """
Supervision subtree for Telepathy (email communication layer).
Children:
- MessageStore — JSONL-backed message persistence
Children (start order matters — rest_for_one strategy):
- MessageStore — JSONL-backed message persistence (must start first)
- JMAP — Fastmail JMAP client with periodic inbox polling
- SMTP — gen_smtp listener on port 25 (inbound email receiver)
The `rest_for_one` strategy ensures that if MessageStore crashes, JMAP and
SMTP are also restarted (they depend on MessageStore being available).
"""
use Supervisor
@ -16,10 +20,12 @@ defmodule Symbiont.Telepathy.Supervisor do
def init(opts) do
messages_path = Keyword.get(opts, :messages_path, "/data/telepathy/messages.jsonl")
poll_interval_ms = Keyword.get(opts, :poll_interval_ms, 60_000)
smtp_port = Keyword.get(opts, :smtp_port, 25)
children = [
{Symbiont.Telepathy.MessageStore, path: messages_path},
{Symbiont.Telepathy.JMAP, poll_interval_ms: poll_interval_ms}
{Symbiont.Telepathy.JMAP, poll_interval_ms: poll_interval_ms},
{Symbiont.Telepathy.SMTP, smtp_port: smtp_port}
]
Supervisor.init(children, strategy: :rest_for_one)

View File

@ -25,7 +25,8 @@ defmodule Symbiont.MixProject do
{:bandit, "~> 1.0"},
{:plug, "~> 1.15"},
{:jason, "~> 1.4"},
{:exqlite, "~> 0.27"}
{:exqlite, "~> 0.27"},
{:gen_smtp, "~> 1.2"}
]
end