Skip to main content

qualia_core_db/services/
chat_relay_daemon.rs

1//! HTTP relay inbox for group chat messages (daemon `/chat/publish` + `/chat/pull`).
2
3#![cfg(not(target_arch = "wasm32"))]
4
5use std::fs::{File, OpenOptions};
6use std::io::{BufRead, BufReader, Write};
7use std::path::PathBuf;
8use std::sync::{Arc, Mutex};
9
10use axum::{
11    extract::{Query, State},
12    http::StatusCode,
13    response::IntoResponse,
14    routing::{get, post},
15    Json, Router,
16};
17use serde::{Deserialize, Serialize};
18
19#[derive(Debug, Clone, Serialize, Deserialize)]
20pub struct RelayEnvelope {
21    pub session_id: String,
22    pub lamport: u64,
23    pub role: String,
24    pub content: String,
25    pub author_did: String,
26    pub author_name: Option<String>,
27    pub reply_to_fragment: Option<String>,
28    pub timestamp: u64,
29    pub signature_hex: String,
30}
31
32#[derive(Debug, Clone, Serialize, Deserialize)]
33pub struct RelayPullResponse {
34    pub messages: Vec<RelayEnvelope>,
35    pub latest_lamport: u64,
36}
37
38#[derive(Clone)]
39pub struct ChatState {
40    pub storage_path: String,
41    pub vault: Arc<Mutex<crate::key_vault::KeyVault>>,
42}
43
44fn relay_root(storage_path: &str) -> PathBuf {
45    PathBuf::from(storage_path).join("ChatRelay")
46}
47
48fn inbox_path(storage_path: &str, session_id: &str) -> PathBuf {
49    relay_root(storage_path)
50        .join(session_id)
51        .join("inbox.jsonl")
52}
53
54fn append_inbox(storage_path: &str, envelope: &RelayEnvelope) -> Result<(), String> {
55    let path = inbox_path(storage_path, &envelope.session_id);
56    if let Some(parent) = path.parent() {
57        std::fs::create_dir_all(parent).map_err(|e| e.to_string())?;
58    }
59    let mut file = OpenOptions::new()
60        .create(true)
61        .append(true)
62        .open(&path)
63        .map_err(|e| e.to_string())?;
64    let line = serde_json::to_string(envelope).map_err(|e| e.to_string())?;
65    writeln!(file, "{line}").map_err(|e| e.to_string())
66}
67
68fn read_inbox(
69    storage_path: &str,
70    session_id: &str,
71    since_lamport: u64,
72) -> Result<RelayPullResponse, String> {
73    let path = inbox_path(storage_path, session_id);
74    if !path.is_file() {
75        return Ok(RelayPullResponse {
76            messages: Vec::new(),
77            latest_lamport: since_lamport,
78        });
79    }
80
81    let file = File::open(&path).map_err(|e| e.to_string())?;
82    let reader = BufReader::new(file);
83    let mut messages = Vec::new();
84    let mut latest = since_lamport;
85
86    for line in reader.lines() {
87        let line = line.map_err(|e| e.to_string())?;
88        if line.trim().is_empty() {
89            continue;
90        }
91        let env: RelayEnvelope = serde_json::from_str(&line).map_err(|e| e.to_string())?;
92        if env.lamport > since_lamport {
93            messages.push(env.clone());
94        }
95        latest = latest.max(env.lamport);
96    }
97
98    messages.sort_by_key(|m| m.lamport);
99    Ok(RelayPullResponse {
100        messages,
101        latest_lamport: latest,
102    })
103}
104
105async fn publish_handler(
106    State(state): State<ChatState>,
107    Json(envelope): Json<RelayEnvelope>,
108) -> impl IntoResponse {
109    if envelope.content.is_empty() || envelope.session_id.is_empty() {
110        return (
111            StatusCode::BAD_REQUEST,
112            Json(serde_json::json!({"error": "invalid envelope"})),
113        );
114    }
115
116    if !envelope.signature_hex.is_empty() {
117        if let Ok(sig_bytes) = hex::decode(&envelope.signature_hex) {
118            if sig_bytes.len() == 64 {
119                let payload = serde_json::json!({
120                    "session_id": envelope.session_id,
121                    "lamport": envelope.lamport,
122                    "role": envelope.role,
123                    "content": envelope.content,
124                    "author_did": envelope.author_did,
125                    "author_name": envelope.author_name,
126                    "reply_to_fragment": envelope.reply_to_fragment,
127                    "timestamp": envelope.timestamp,
128                });
129                if let Ok(payload_str) = serde_json::to_string(&payload) {
130                    let vault = state.vault.lock().unwrap();
131                    let key = vault.derive_key(&format!("relay:{}", envelope.author_did));
132                    let pk = ed25519_dalek::VerifyingKey::from(&key);
133                    let mut sig_arr = [0u8; 64];
134                    sig_arr.copy_from_slice(&sig_bytes);
135                    if crate::key_vault::KeyVault::verify_signature(
136                        pk.as_bytes(),
137                        payload_str.as_bytes(),
138                        &sig_arr,
139                    )
140                    .is_err()
141                    {
142                        return (
143                            StatusCode::UNAUTHORIZED,
144                            Json(serde_json::json!({"error": "signature invalid"})),
145                        );
146                    }
147                }
148            }
149        }
150    }
151
152    match append_inbox(&state.storage_path, &envelope) {
153        Ok(()) => (
154            StatusCode::OK,
155            Json(serde_json::json!({"ok": true, "lamport": envelope.lamport})),
156        ),
157        Err(e) => (
158            StatusCode::INTERNAL_SERVER_ERROR,
159            Json(serde_json::json!({"error": e})),
160        ),
161    }
162}
163
164async fn pull_handler(
165    State(state): State<ChatState>,
166    Query(params): Query<std::collections::HashMap<String, String>>,
167) -> impl IntoResponse {
168    let session_id = params.get("session_id").cloned().unwrap_or_default();
169    let since = params
170        .get("since_lamport")
171        .and_then(|s| s.parse::<u64>().ok())
172        .unwrap_or(0);
173
174    if session_id.is_empty() {
175        return (
176            StatusCode::BAD_REQUEST,
177            Json(serde_json::json!({"error": "session_id required"})).into_response(),
178        );
179    }
180
181    match read_inbox(&state.storage_path, &session_id, since) {
182        Ok(resp) => (
183            StatusCode::OK,
184            Json(serde_json::json!(resp)).into_response(),
185        ),
186        Err(e) => (
187            StatusCode::INTERNAL_SERVER_ERROR,
188            Json(serde_json::json!({"error": e})).into_response(),
189        ),
190    }
191}
192
193pub fn chat_relay_routes(
194    storage_path: String,
195    vault: Arc<Mutex<crate::key_vault::KeyVault>>,
196) -> Router {
197    let state = ChatState {
198        storage_path,
199        vault,
200    };
201    let relay = Router::new()
202        .route("/publish", post(publish_handler))
203        .route("/pull", get(pull_handler))
204        .with_state(state.clone());
205    Router::new()
206        .nest("/chat", relay)
207        // Legacy root paths (pre-I2 clients)
208        .route("/publish", post(publish_handler))
209        .route("/pull", get(pull_handler))
210        .with_state(state)
211}