1use 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
302pub 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}