mod auth;
mod bootstrap;
mod config;
mod dto;
mod error;
mod health;
mod identity;
#[cfg(feature = "nemo")]
mod registration;
mod request_id;
mod routes;
mod service;
mod shutdown;
mod state;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use axum::middleware;
use axum::Router;
#[cfg(feature = "nemo")]
use mcp_service4agent::register;
use mcp_service4agent::router::{DocRouter, RouteDoc};
use config::{Config, ConfigError};
use mail4agent_store_sqlite::{migrations, Db, DbConfig, DbError, MigrationRunner, SqliteMailStore};
use service::MailboxService;
use state::AppState;
#[derive(Debug, thiserror::Error)]
enum MainError {
#[error("config: {0}")]
Config(#[from] ConfigError),
#[error("database: {0}")]
Db(#[from] DbError),
#[error("build tokio runtime: {0}")]
Runtime(std::io::Error),
#[error("bootstrap: {0}")]
Bootstrap(#[from] bootstrap::BootstrapError),
#[error("bind {0}: {1}")]
Bind(SocketAddr, std::io::Error),
#[error("serve: {0}")]
Serve(std::io::Error),
}
fn main() -> Result<(), MainError> {
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")),
)
.init();
let config = Config::load()?;
let db_path = config.resolve_db_path()?;
tracing::info!(db = %db_path.display(), bind = %config.bind, "mail4agent starting");
let db_config = DbConfig::new(db_path);
let db = Db::open(&db_config)?;
db.run_migrations_blocking(MigrationRunner::new(migrations()))?;
let engine_store = SqliteMailStore::new(db.clone());
let reader_store = SqliteMailStore::new(db);
let service = Arc::new(MailboxService::new(engine_store, reader_store));
let runtime = tokio::runtime::Runtime::new().map_err(MainError::Runtime)?;
runtime.block_on(serve(service, config.bind))
}
async fn serve(service: Arc<MailboxService>, bind: SocketAddr) -> Result<(), MainError> {
bootstrap::ensure_bootstrap_operator(&service).await?;
let app = Arc::new(AppState { service: service.clone(), bind_addr: bind });
let health_state = Arc::new(health::HealthState::new("mail4agent", env!("CARGO_PKG_VERSION")));
let (health_router, health_endpoints) = DocRouter::<Arc<AppState>>::new()
.get(
"/health",
move || {
let health_state = health_state.clone();
async move { health::handle(health_state).await }
},
RouteDoc::new(
"Liveness, service name, version and uptime: {\"ok\":true,\"service\":\"mail4agent\",\"version\":...,\"started_at\":...,\"uptime_secs\":...,\"dependencies\":[],\"background_tasks\":[]}. The watchdog only reads the status code.",
)
.public(),
)
.into_parts();
let (mail_router, mail_endpoints) = DocRouter::<Arc<AppState>>::new()
.post(
"/mail/send",
routes::mail::send,
RouteDoc::bearer("Send a message to a participant or a room. The sender is derived from the credential and cannot be supplied. Returns the new message id and the sender's own address."),
)
.post(
"/mail/inbox",
routes::mail::inbox,
RouteDoc::bearer("Read messages addressed to the caller directly, plus messages to rooms the caller currently belongs to. Takes since_unix_ms and limit, plus wait_secs: long-polls (clamped to 60s, never refused for asking longer) when the inbox would otherwise answer empty, returning an empty page rather than an error on expiry."),
)
.post("/mail/ack", routes::mail::ack, RouteDoc::bearer("Acknowledge one message the caller may read. Idempotent per (message, reader)."))
.post("/mail/get", routes::mail::get, RouteDoc::bearer("Read one message by id, if the caller may read it."))
.post(
"/mail/unread",
routes::mail::unread,
RouteDoc::bearer("Unread count for the caller, or for another participant when the caller is an operator."),
)
.post(
"/mail/whoami",
routes::mail::whoami,
RouteDoc::bearer("The caller's own SESSION address (resolved from the connection, never declared), its account's label, room memberships, and its own card -- attested (kernel), corroborated (its own command line) and declared (its own claims) kept apart."),
)
.post(
"/mail/status",
routes::mail::status,
RouteDoc::bearer("The session declares what it is working on, its role, and which session spawned it. The only writer of that group; recorded as the session's own claim, never verified."),
)
.post(
"/mail/directory",
routes::mail::directory,
RouteDoc::bearer("Every registered account (id, label) with its live sessions nested under it (card, whether it is live), and every room the mailbox tracks (id, whether the caller is a member). Never returns a secret digest."),
)
.into_parts();
let mcp_server = routes::mcp::build();
let mcp_endpoints = mcp_server.route_docs();
let mcp_router: Router<Arc<AppState>> = mcp_server
.into_router()
.route_layer(middleware::from_fn_with_state(app.clone(), routes::mcp::resolve_caller_middleware));
let authenticated_guard =
auth::TierGuard { service: service.clone(), required: auth::Tier::Authenticated };
let authenticated_router: Router<Arc<AppState>> = mail_router
.merge(mcp_router)
.route_layer(middleware::from_fn_with_state(authenticated_guard, auth::require_tier));
let (admin_router, admin_endpoints) = DocRouter::<Arc<AppState>>::new()
.post(
"/admin/participant",
routes::admin::register_participant,
RouteDoc::bearer("Register a participant and return its secret ONCE. Only the digest is kept. Operator only."),
)
.post(
"/admin/participant/rotate",
routes::admin::rotate_participant,
RouteDoc::bearer("Issue a new secret for a participant and invalidate the old one. Operator only."),
)
.post(
"/admin/participant/remove",
routes::admin::remove_participant,
RouteDoc::bearer("Deregister a participant. Its messages are kept; it can no longer authenticate. Operator only."),
)
.post("/admin/room", routes::admin::create_room, RouteDoc::bearer("Create a room. Operator only."))
.post(
"/admin/room/member/add",
routes::admin::add_room_member,
RouteDoc::bearer("Add a participant to a room, which is what grants it read access to that room. Operator only."),
)
.post(
"/admin/room/member/remove",
routes::admin::remove_room_member,
RouteDoc::bearer("Remove a participant from a room. It stops reading that room from then on. Operator only."),
)
.post(
"/admin/listener",
routes::admin::set_listener,
RouteDoc::bearer("Register (or replace) the URL the mailbox POSTs a delivery notification to when mail arrives for an account or any of its sessions. Loopback only. Never carries the subject or body. Best-effort. Operator only."),
)
.post(
"/admin/listener/remove",
routes::admin::remove_listener,
RouteDoc::bearer("Remove an account's registered delivery listener, if any. Idempotent. Operator only."),
)
.into_parts();
let admin_guard = auth::TierGuard { service: service.clone(), required: auth::Tier::Admin };
let admin_router: Router<Arc<AppState>> =
admin_router.route_layer(middleware::from_fn_with_state(admin_guard, auth::require_tier));
let mut endpoints = health_endpoints;
endpoints.extend(mail_endpoints);
endpoints.extend(mcp_endpoints);
endpoints.extend(admin_endpoints);
let router = health_router
.merge(authenticated_router)
.merge(admin_router)
.with_state(app.clone())
.layer(middleware::from_fn(request_id::request_id_layer));
let listener = tokio::net::TcpListener::bind(bind).await.map_err(|e| MainError::Bind(bind, e))?;
tracing::info!(addr = %bind, "mail4agent listening");
#[cfg(feature = "nemo")]
{
let _reassert =
register::spawn(registration::service_manifest(endpoints, bind.port()), register::DEFAULT_REASSERT_INTERVAL);
}
#[cfg(not(feature = "nemo"))]
{
let _ = endpoints;
}
let (shutdown_tx, mut shutdown_rx) = tokio::sync::watch::channel(false);
let serve_handle = tokio::spawn(async move {
axum::serve(listener, router.into_make_service_with_connect_info::<SocketAddr>())
.with_graceful_shutdown(async move {
let _ = shutdown_rx.changed().await;
})
.await
});
shutdown::graceful_shutdown_signal().await;
tracing::info!("shutdown signal received");
let _ = shutdown_tx.send(true);
match tokio::time::timeout(Duration::from_secs(30), serve_handle).await {
Ok(Ok(Ok(()))) => tracing::info!("mail4agent stopped"),
Ok(Ok(Err(err))) => return Err(MainError::Serve(err)),
Ok(Err(join_err)) => tracing::error!(error = %join_err, "serve task panicked"),
Err(_) => tracing::warn!("shutdown drain exceeded 30s -- exiting anyway"),
}
Ok(())
}