Skip to main content

domi_server/http/
mod.rs

1//! HTTP layer for the `domi-server` binary (Phase 2c-γ).
2//!
3//! Top-level orchestration lives here; concrete pieces live in sibling modules.
4
5pub mod args;
6pub mod handlers;
7pub mod router;
8pub mod state;
9pub mod ws;
10
11use std::sync::Arc;
12
13use tracing_subscriber::EnvFilter;
14
15use crate::events::EventWriter;
16use crate::serve::watcher::{NotifyWatcher, WatchEventKind, Watcher};
17
18use self::args::Args;
19use self::state::AppState;
20
21pub async fn run(args: Args) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
22    // 1. tracing init.
23    let filter = EnvFilter::try_new(&args.log_level).unwrap_or_else(|_| EnvFilter::new("info"));
24    tracing_subscriber::fmt().with_env_filter(filter).init();
25
26    // 2. Ensure state dir exists; resolve events.jsonl path.
27    std::fs::create_dir_all(&args.state)?;
28    std::fs::create_dir_all(&args.root)?;
29    let events_path = args.state.join("events.jsonl");
30
31    // 3. Construct EventWriter (sync).
32    let writer = Arc::new(EventWriter::new(&events_path));
33
34    // 4. Construct AppState.
35    let state = Arc::new(AppState::new(
36        args.root.clone(),
37        args.state.clone(),
38        writer,
39        256,
40        args.library_root.clone(),
41    ));
42
43    // 5. Spawn watcher logger.
44    spawn_watcher_logger(&args.root);
45
46    // 6. Build router.
47    let router = router::build_router(state.clone());
48
49    // 7. Bind.
50    let addr = format!("{}:{}", args.host, args.port);
51    let listener = tokio::net::TcpListener::bind(&addr).await?;
52    let bound = listener.local_addr()?;
53    tracing::info!(bound_url = %format!("http://{}/", bound), server_id = %state.server_id, "domi-server listening");
54
55    // 8. Serve with graceful shutdown.
56    axum::serve(listener, router)
57        .with_graceful_shutdown(shutdown_signal())
58        .await?;
59
60    Ok(())
61}
62
63fn spawn_watcher_logger(root: &std::path::Path) {
64    let mut watcher = match NotifyWatcher::new(root, 50) {
65        Ok(w) => w,
66        Err(e) => {
67            tracing::warn!(error = %e, "watcher init failed; continuing without");
68            return;
69        }
70    };
71    tokio::spawn(async move {
72        loop {
73            match watcher.next_event(500) {
74                Ok(Some(ev)) => {
75                    let kind = match ev.kind {
76                        WatchEventKind::Created => "created",
77                        WatchEventKind::Modified => "modified",
78                        WatchEventKind::Removed => "removed",
79                        WatchEventKind::Any => "any",
80                    };
81                    for p in &ev.paths {
82                        tracing::debug!(kind, path = %p.display(), "watcher");
83                    }
84                }
85                Ok(None) => continue,
86                Err(e) => {
87                    tracing::warn!(error = %e, "watcher error; stopping");
88                    break;
89                }
90            }
91        }
92    });
93}
94
95async fn shutdown_signal() {
96    let ctrl_c = async {
97        let _ = tokio::signal::ctrl_c().await;
98    };
99    #[cfg(unix)]
100    let sigterm = async {
101        let mut s = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
102            .expect("install SIGTERM handler");
103        s.recv().await;
104    };
105    #[cfg(not(unix))]
106    let sigterm = std::future::pending::<()>();
107    tokio::select! {
108        _ = ctrl_c => {}
109        _ = sigterm => {}
110    }
111    tracing::info!("shutdown signal received");
112}