Skip to main content

qualia_client_core/
chat_relay.rs

1//! Group chat relay — publish to daemon inbox and pull from peer relay endpoints.
2
3use std::collections::HashMap;
4use std::path::{Path, PathBuf};
5use std::sync::atomic::{AtomicBool, Ordering};
6use std::sync::OnceLock;
7use std::time::Duration;
8
9use serde::{Deserialize, Serialize};
10
11use crate::chat_session::{ChatError, Role};
12
13static RELAY_POLLER_STARTED: AtomicBool = AtomicBool::new(false);
14
15#[derive(Debug, Clone, Serialize, Deserialize)]
16pub struct RelayEnvelope {
17    pub session_id: String,
18    pub lamport: u64,
19    pub role: String,
20    pub content: String,
21    pub author_did: String,
22    pub author_name: Option<String>,
23    pub reply_to_fragment: Option<String>,
24    pub timestamp: u64,
25    pub signature_hex: String,
26    #[serde(default, skip_serializing_if = "Option::is_none")]
27    pub sub_agent_of: Option<String>,
28    #[serde(default, skip_serializing_if = "Option::is_none")]
29    pub agent_did: Option<String>,
30    #[serde(default, skip_serializing_if = "Option::is_none")]
31    pub model_id: Option<String>,
32    #[serde(default, skip_serializing_if = "Option::is_none")]
33    pub agent_backend: Option<String>,
34    #[serde(default, skip_serializing_if = "Option::is_none")]
35    pub outcome_sharing: Option<crate::chat_agents::OutcomeSharingPolicy>,
36}
37
38#[derive(Debug, Clone, Serialize, Deserialize)]
39struct RelayCursor {
40    per_session: HashMap<String, u64>,
41}
42
43fn cursor_path() -> PathBuf {
44    crate::state::app_meta_dir().join("relay_cursors.json")
45}
46
47fn load_cursors() -> RelayCursor {
48    fs_read_json(cursor_path()).unwrap_or(RelayCursor {
49        per_session: HashMap::new(),
50    })
51}
52
53fn save_cursors(cursors: &RelayCursor) -> Result<(), String> {
54    let path = cursor_path();
55    if let Some(parent) = path.parent() {
56        std::fs::create_dir_all(parent).map_err(|e| e.to_string())?;
57    }
58    let text = serde_json::to_string_pretty(cursors).map_err(|e| e.to_string())?;
59    std::fs::write(path, text).map_err(|e| e.to_string())
60}
61
62fn fs_read_json<T: for<'de> Deserialize<'de>>(path: PathBuf) -> Option<T> {
63    let text = std::fs::read_to_string(path).ok()?;
64    serde_json::from_str(&text).ok()
65}
66
67pub fn local_relay_base_url() -> String {
68    let port = crate::api::get_active_daemon_port();
69    if let Some(profile) = profile_relay_base() {
70        if !profile.is_empty() {
71            return profile;
72        }
73    }
74    format!("http://127.0.0.1:{port}")
75}
76
77fn profile_relay_base() -> Option<String> {
78    let profile = crate::user_profile::load_profile();
79    profile.relay_base_url.clone()
80}
81
82fn sign_envelope(envelope: &mut RelayEnvelope) {
83    let payload = serde_json::json!({
84        "session_id": envelope.session_id,
85        "lamport": envelope.lamport,
86        "role": envelope.role,
87        "content": envelope.content,
88        "author_did": envelope.author_did,
89        "author_name": envelope.author_name,
90        "reply_to_fragment": envelope.reply_to_fragment,
91        "timestamp": envelope.timestamp,
92        "sub_agent_of": envelope.sub_agent_of,
93        "agent_did": envelope.agent_did,
94        "model_id": envelope.model_id,
95        "agent_backend": envelope.agent_backend,
96        "outcome_sharing": envelope.outcome_sharing,
97    });
98    let payload_str = serde_json::to_string(&payload).unwrap_or_default();
99    let state = match crate::state::APP_STATE.get() {
100        Some(s) => s,
101        None => return,
102    };
103    let vault = state.key_vault.lock().unwrap();
104    let key = vault.derive_key(&format!("relay:{}", envelope.author_did));
105    let sig = vault.sign_payload(&key, payload_str.as_bytes());
106    envelope.signature_hex = hex::encode(sig.to_bytes());
107}
108
109pub fn message_to_envelope(
110    session_id: &str,
111    msg: &crate::chat_session::ChatMessage,
112) -> RelayEnvelope {
113    let profile = crate::user_profile::load_profile();
114    let mut envelope = RelayEnvelope {
115        session_id: session_id.to_string(),
116        lamport: msg.lamport,
117        role: msg.role.as_str().to_string(),
118        content: msg.content.clone(),
119        author_did: msg
120            .author_did
121            .clone()
122            .unwrap_or_else(|| profile.public_did.clone()),
123        author_name: msg.author_name.clone(),
124        reply_to_fragment: msg.reply_to_fragment.clone(),
125        timestamp: msg.timestamp,
126        signature_hex: String::new(),
127        sub_agent_of: msg.sub_agent_of.clone(),
128        agent_did: msg.agent_did.clone(),
129        model_id: msg.model_id.clone(),
130        agent_backend: msg.agent_backend.clone(),
131        outcome_sharing: msg.outcome_sharing.clone(),
132    };
133    sign_envelope(&mut envelope);
134    envelope
135}
136
137pub fn publish_session_message(
138    storage_root: &Path,
139    session_id: &str,
140    lamport: u64,
141) -> Result<(), String> {
142    let session =
143        crate::chat_session::load_session(storage_root, session_id).map_err(|e| e.to_string())?;
144    let msg = session
145        .messages
146        .iter()
147        .find(|m| m.lamport == lamport)
148        .ok_or_else(|| format!("message lamport {lamport} not found"))?;
149
150    let envelope = message_to_envelope(session_id, msg);
151    publish_envelope(&session.meta.participants, &envelope)
152}
153
154pub fn publish_envelope(
155    participants: &[crate::chat_session::ChatParticipant],
156    envelope: &RelayEnvelope,
157) -> Result<(), String> {
158    let local = local_relay_base_url();
159    let mut endpoints: Vec<String> = vec![local.clone()];
160
161    for p in participants {
162        if let Some(contact) = crate::social_connect::find_contact_by_did(&p.did) {
163            if let Some(url) = contact.relay_endpoint {
164                if !url.is_empty() && !endpoints.contains(&url) {
165                    endpoints.push(url);
166                }
167            }
168        }
169    }
170
171    let client = reqwest::blocking::Client::builder()
172        .timeout(Duration::from_secs(8))
173        .build()
174        .map_err(|e| e.to_string())?;
175
176    let mut last_err = None;
177    for base in endpoints {
178        let url = format!("{}/chat/publish", base.trim_end_matches('/'));
179        match client.post(&url).json(envelope).send() {
180            Ok(resp) if resp.status().is_success() => {}
181            Ok(resp) => {
182                last_err = Some(format!("relay {url} returned {}", resp.status()));
183            }
184            Err(e) => {
185                last_err = Some(format!("relay {url} failed: {e}"));
186            }
187        }
188    }
189
190    if let Some(err) = last_err {
191        eprintln!("[chat_relay] {err}");
192    }
193    Ok(())
194}
195
196#[derive(Debug, Deserialize)]
197pub struct PullResponse {
198    messages: Vec<RelayEnvelope>,
199    latest_lamport: u64,
200}
201
202pub fn pull_from_relay(
203    base_url: &str,
204    session_id: &str,
205    since_lamport: u64,
206) -> Result<PullResponse, String> {
207    let url = format!(
208        "{}/chat/pull?session_id={}&since_lamport={}",
209        base_url.trim_end_matches('/'),
210        urlencoding_path(session_id),
211        since_lamport
212    );
213    let client = reqwest::blocking::Client::builder()
214        .timeout(Duration::from_secs(8))
215        .build()
216        .map_err(|e| e.to_string())?;
217    let resp = client.get(&url).send().map_err(|e| e.to_string())?;
218    if !resp.status().is_success() {
219        return Err(format!("pull failed: {}", resp.status()));
220    }
221    resp.json::<PullResponse>().map_err(|e| e.to_string())
222}
223
224fn urlencoding_path(s: &str) -> String {
225    s.chars()
226        .map(|c| match c {
227            'A'..='Z' | 'a'..='z' | '0'..='9' | '-' | '_' | '.' | '~' => c.to_string(),
228            _ => format!("%{:02X}", c as u8),
229        })
230        .collect()
231}
232
233fn ingest_remote_message(
234    storage_root: &Path,
235    session_id: &str,
236    envelope: &RelayEnvelope,
237) -> Result<bool, ChatError> {
238    let profile = crate::user_profile::load_profile();
239    if envelope.author_did == profile.public_did {
240        return Ok(false);
241    }
242
243    let session = crate::chat_session::load_session(storage_root, session_id)?;
244    if session.messages.iter().any(|m| {
245        m.lamport == envelope.lamport
246            && m.author_did.as_deref() == Some(envelope.author_did.as_str())
247    }) {
248        return Ok(false);
249    }
250
251    let role = Role::from_str(&envelope.role)?;
252
253    if role == Role::Agent {
254        let preview = crate::chat_session::ChatMessage {
255            lamport: envelope.lamport,
256            role,
257            content: envelope.content.clone(),
258            timestamp: envelope.timestamp,
259            content_hash: 0,
260            author_did: Some(envelope.author_did.clone()),
261            author_name: envelope.author_name.clone(),
262            reply_to_fragment: envelope.reply_to_fragment.clone(),
263            source: Some(format!("relay:{}", envelope.author_did)),
264            sub_agent_of: envelope.sub_agent_of.clone(),
265            agent_did: envelope.agent_did.clone(),
266            model_id: envelope.model_id.clone(),
267            agent_backend: envelope.agent_backend.clone(),
268            outcome_sharing: envelope.outcome_sharing.clone(),
269        };
270        crate::chat_agents::validate_ingested_agent_message(&preview, &session.meta.participants)
271            .map_err(|e| ChatError::InvalidSession(e))?;
272        if !crate::chat_agents::can_view_agent_outcome(
273            &preview,
274            &profile.public_did,
275            &session.meta.participants,
276        ) {
277            return Ok(false);
278        }
279    }
280
281    let lamport = crate::chat_session::append_relay_message_with_agent_meta(
282        storage_root,
283        session_id,
284        role,
285        &envelope.content,
286        envelope.reply_to_fragment.clone(),
287        Some(format!("relay:{}", envelope.author_did)),
288        Some(envelope.author_did.clone()),
289        envelope.author_name.clone(),
290        envelope.sub_agent_of.clone(),
291        envelope.agent_did.clone(),
292        envelope.model_id.clone(),
293        envelope.agent_backend.clone(),
294        envelope.outcome_sharing.clone(),
295    )?;
296    let _ = lamport;
297    notify_session_updated(session_id);
298
299    Ok(true)
300}
301
302/// Apply an inbound [`RelayEnvelope`] to a session, regardless of the transport it arrived on.
303///
304/// This is the transport-agnostic ingest path shared by the HTTP relay and the SocialWebNet mesh
305/// (`chat_mesh`): it deduplicates by `(lamport, author_did)`, validates agent messages, and appends
306/// the message. Returns `true` if the message was newly applied, `false` if it was a duplicate or
307/// our own echo. Envelope signature verification is the caller's/session policy's concern, exactly as
308/// for the relay pull path.
309pub fn apply_incoming_envelope(
310    storage_root: &Path,
311    session_id: &str,
312    envelope: &RelayEnvelope,
313) -> Result<bool, ChatError> {
314    ingest_remote_message(storage_root, session_id, envelope)
315}
316
317pub fn sync_session_relay(storage_root: &Path, session_id: &str) -> Result<usize, String> {
318    let session =
319        crate::chat_session::load_session(storage_root, session_id).map_err(|e| e.to_string())?;
320
321    let mut cursors = load_cursors();
322    let since = *cursors.per_session.get(session_id).unwrap_or(&0);
323    let mut ingested = 0usize;
324    let mut latest = since;
325
326    let mut endpoints = vec![local_relay_base_url()];
327    for p in &session.meta.participants {
328        if let Some(contact) = crate::social_connect::find_contact_by_did(&p.did) {
329            if let Some(url) = contact.relay_endpoint {
330                if !url.is_empty() && !endpoints.contains(&url) {
331                    endpoints.push(url);
332                }
333            }
334        }
335    }
336
337    for base in endpoints {
338        let pull = match pull_from_relay(&base, session_id, since) {
339            Ok(p) => p,
340            Err(e) => {
341                eprintln!("[chat_relay] pull {base}: {e}");
342                continue;
343            }
344        };
345        latest = latest.max(pull.latest_lamport);
346        for env in pull.messages {
347            if ingest_remote_message(storage_root, session_id, &env).unwrap_or(false) {
348                ingested += 1;
349            }
350        }
351    }
352
353    cursors.per_session.insert(session_id.to_string(), latest);
354    save_cursors(&cursors)?;
355    Ok(ingested)
356}
357
358pub fn sync_all_group_sessions() -> Result<usize, String> {
359    let state = crate::state::APP_STATE
360        .get()
361        .ok_or("APP_STATE not initialized")?;
362    let storage = state.config.lock().unwrap().storage_path.clone();
363    let storage_path = Path::new(&storage);
364
365    let sessions = crate::chat_session::list_sessions(storage_path).map_err(|e| e.to_string())?;
366    let mut total = 0usize;
367    for summary in sessions {
368        if summary.session_kind != crate::chat_session::SessionKind::Group {
369            continue;
370        }
371        total += sync_session_relay(storage_path, &summary.id)?;
372    }
373    Ok(total)
374}
375
376pub fn start_relay_poller() {
377    if RELAY_POLLER_STARTED.swap(true, Ordering::SeqCst) {
378        return;
379    }
380
381    std::thread::spawn(|| loop {
382        if let Err(e) = sync_all_group_sessions() {
383            eprintln!("[chat_relay] poller: {e}");
384        }
385        std::thread::sleep(Duration::from_secs(4));
386    });
387}
388
389static RELAY_NOTIFY: OnceLock<std::sync::Mutex<Option<tokio::sync::broadcast::Sender<String>>>> =
390    OnceLock::new();
391
392pub fn relay_notify_channel() -> tokio::sync::broadcast::Sender<String> {
393    RELAY_NOTIFY
394        .get_or_init(|| {
395            let (tx, _) = tokio::sync::broadcast::channel(64);
396            std::sync::Mutex::new(Some(tx))
397        })
398        .lock()
399        .unwrap()
400        .clone()
401        .expect("relay notify tx")
402}
403
404pub fn notify_session_updated(session_id: &str) {
405    let _ = relay_notify_channel().send(session_id.to_string());
406}