1use std::fs::{self, File, OpenOptions};
4use std::io::{BufRead, BufReader, Write};
5use std::path::{Path, PathBuf};
6
7use qualia_core_db::{q_hash, wal::WriteAheadLog, NQuin};
8use serde::{Deserialize, Serialize};
9
10const OBJECT_HASH_MASK: u64 = 0x0FFF_FFFF_FFFF_FFFF;
11
12pub fn content_hash_u64(content: &str) -> u64 {
15 q_hash(content) & OBJECT_HASH_MASK
16}
17const LAMPORT_SHIFT: u32 = 32;
18const LAMPORT_MASK: u64 = 0x1FFF_FFFF;
19
20#[derive(Debug)]
21pub enum ChatError {
22 NotFound(String),
23 InvalidSession(String),
24 Wal(String),
25 Io(std::io::Error),
26 Json(serde_json::Error),
27 Compact(String),
28}
29
30impl std::fmt::Display for ChatError {
31 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
32 match self {
33 ChatError::NotFound(id) => write!(f, "Chat session not found: {id}"),
34 ChatError::InvalidSession(msg) => write!(f, "Invalid session: {msg}"),
35 ChatError::Wal(msg) => write!(f, "WAL error: {msg}"),
36 ChatError::Io(e) => write!(f, "IO error: {e}"),
37 ChatError::Json(e) => write!(f, "JSON error: {e}"),
38 ChatError::Compact(msg) => write!(f, "Compaction error: {msg}"),
39 }
40 }
41}
42
43impl From<std::io::Error> for ChatError {
44 fn from(e: std::io::Error) -> Self {
45 ChatError::Io(e)
46 }
47}
48
49impl From<serde_json::Error> for ChatError {
50 fn from(e: serde_json::Error) -> Self {
51 ChatError::Json(e)
52 }
53}
54
55#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
56#[serde(rename_all = "lowercase")]
57pub enum SessionKind {
58 #[default]
59 Solo,
60 Group,
61}
62
63#[derive(Debug, Clone, Serialize, Deserialize)]
64pub struct ChatParticipant {
65 pub did: String,
66 pub display_name: String,
67 pub actor_id: String,
68 pub role: String,
69 pub joined_at: u64,
70}
71
72#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
73#[serde(rename_all = "lowercase")]
74pub enum Role {
75 User,
76 Agent,
77}
78
79impl Role {
80 pub fn as_str(self) -> &'static str {
81 match self {
82 Role::User => "user",
83 Role::Agent => "agent",
84 }
85 }
86
87 pub fn from_str(s: &str) -> Result<Self, ChatError> {
88 match s.to_lowercase().as_str() {
89 "user" => Ok(Role::User),
90 "agent" => Ok(Role::Agent),
91 _ => Err(ChatError::InvalidSession(format!("unknown role: {s}"))),
92 }
93 }
94}
95
96#[derive(Debug, Clone, Serialize, Deserialize)]
98pub struct OntologyScopeSummary {
99 pub id: String,
100 pub name: String,
101 pub quin_count: u64,
102 pub q42_path: String,
103 #[serde(default)]
104 pub domain: Option<String>,
105 #[serde(default)]
106 pub tags: Option<Vec<String>>,
107 #[serde(default)]
108 pub source: Option<String>,
109}
110
111#[derive(Debug, Clone, Serialize, Deserialize)]
113pub struct ChatEnvironment {
114 pub session_id: String,
115 pub active_model_profile_id: u64,
116 pub ontology_ids: Vec<String>,
117 pub prior_session_ids: Vec<String>,
118 pub graph_scope_hashes: Vec<u64>,
119 #[serde(default)]
120 pub lexicon_prefixes: Vec<u64>,
121 #[serde(default)]
122 pub capability_briefing: String,
123 #[serde(default)]
124 pub model_id: Option<String>,
125 #[serde(default)]
126 pub model_modality: String,
127 #[serde(default)]
128 pub context_window: u32,
129 #[serde(default)]
130 pub engine_capabilities: Vec<String>,
131 #[serde(default)]
132 pub installed_qapps: Vec<String>,
133 #[serde(default)]
134 pub ontology_summaries: Vec<OntologyScopeSummary>,
135 #[serde(default)]
136 pub daemon_reachable: bool,
137 #[serde(default)]
138 pub session_kind: SessionKind,
139 #[serde(default)]
140 pub participants: Vec<ChatParticipant>,
141 #[serde(default)]
143 pub graph_mutation: bool,
144 #[serde(default)]
145 pub axiom_bounds: crate::context_binding::AxiomBounds,
146}
147
148fn active_model_profile_id() -> u64 {
149 let path = crate::state::app_meta_dir().join("active_model.json");
150 if let Ok(text) = fs::read_to_string(path) {
151 if let Ok(record) = serde_json::from_str::<crate::model_lifecycle::ActiveModelRecord>(&text)
152 {
153 return record.profile_id;
154 }
155 }
156 0
157}
158
159impl ChatEnvironment {
160 pub fn default_for_session(session_id: &str, storage_root: &Path) -> Self {
161 let catalog = qualia_core_db::resource_catalog::load_default()
162 .unwrap_or_else(|_| qualia_core_db::resource_catalog::ResourceCatalog::empty());
163 crate::context_binding::compile_chat_environment(
164 storage_root,
165 &catalog,
166 &crate::context_binding::ChatEnvironmentConfig {
167 session_id: session_id.to_string(),
168 ontology_ids: Vec::new(),
169 prior_session_ids: Vec::new(),
170 session_kind: SessionKind::Solo,
171 participants: Vec::new(),
172 graph_mutation: false,
173 axiom_bounds: crate::context_binding::AxiomBounds::default(),
174 },
175 )
176 .unwrap_or_else(|_| {
177 let profile_id = active_model_profile_id();
178 let session_scope = q_hash(&format!("chat:session:{session_id}"));
179 Self {
180 session_id: session_id.to_string(),
181 active_model_profile_id: profile_id,
182 ontology_ids: Vec::new(),
183 prior_session_ids: Vec::new(),
184 graph_scope_hashes: vec![session_scope],
185 lexicon_prefixes: Vec::new(),
186 capability_briefing: String::new(),
187 model_id: None,
188 model_modality: "text".to_string(),
189 context_window: 4096,
190 engine_capabilities: Vec::new(),
191 installed_qapps: Vec::new(),
192 ontology_summaries: Vec::new(),
193 daemon_reachable: false,
194 session_kind: SessionKind::Solo,
195 participants: Vec::new(),
196 graph_mutation: false,
197 axiom_bounds: crate::context_binding::AxiomBounds::default(),
198 }
199 })
200 }
201
202 pub fn save_to_session_dir(&self, storage_root: &Path) -> Result<(), ChatError> {
203 let path = session_dir(storage_root, &self.session_id).join("environment.json");
204 fs::write(path, serde_json::to_string_pretty(self)?)?;
205 Ok(())
206 }
207}
208
209#[derive(Debug, Clone, Serialize, Deserialize)]
210pub struct SessionMeta {
211 pub id: String,
212 pub title: String,
213 pub created_at: u64,
214 pub updated_at: u64,
215 pub message_count: u64,
216 pub next_lamport: u64,
217 pub environment_ref: String,
218 #[serde(default)]
219 pub session_kind: SessionKind,
220 #[serde(default)]
221 pub participants: Vec<ChatParticipant>,
222 #[serde(default)]
223 pub owner_did: String,
224 #[serde(default)]
226 pub session_did: String,
227}
228
229#[derive(Debug, Clone, Serialize, Deserialize)]
231pub struct ChatSessionShareTarget {
232 pub session_id: String,
233 pub session_did: String,
234 pub title: String,
235 pub session_kind: SessionKind,
236 pub participant_count: u64,
237}
238
239pub fn compile_session_did(session_id: &str, kind: SessionKind) -> String {
240 let scope = match kind {
241 SessionKind::Solo => "solo",
242 SessionKind::Group => "group",
243 };
244 let digest = q_hash(&format!("qualia:chat:{scope}:{session_id}"));
245 format!("did:qualia:chat:{scope}:{digest:016x}")
246}
247
248fn ensure_session_did(meta: &mut SessionMeta) -> String {
249 if meta.session_did.is_empty() {
250 meta.session_did = compile_session_did(&meta.id, meta.session_kind);
251 }
252 meta.session_did.clone()
253}
254
255fn persist_session_did_if_needed(
256 storage_root: &Path,
257 meta: &mut SessionMeta,
258) -> Result<(), ChatError> {
259 let before = meta.session_did.clone();
260 let _ = ensure_session_did(meta);
261 if meta.session_did != before {
262 let path = session_meta_path(storage_root, &meta.id);
263 fs::write(path, serde_json::to_string_pretty(meta)?)?;
264 }
265 Ok(())
266}
267
268#[derive(Debug, Clone, Serialize, Deserialize)]
269pub struct ChatMessage {
270 pub lamport: u64,
271 pub role: Role,
272 pub content: String,
273 pub timestamp: u64,
274 pub content_hash: u64,
275 #[serde(default, skip_serializing_if = "Option::is_none")]
276 pub author_did: Option<String>,
277 #[serde(default, skip_serializing_if = "Option::is_none")]
278 pub author_name: Option<String>,
279 #[serde(default, skip_serializing_if = "Option::is_none")]
280 pub reply_to_fragment: Option<String>,
281 #[serde(default, skip_serializing_if = "Option::is_none")]
282 pub source: Option<String>,
283 #[serde(default, skip_serializing_if = "Option::is_none")]
285 pub sub_agent_of: Option<String>,
286 #[serde(default, skip_serializing_if = "Option::is_none")]
287 pub agent_did: Option<String>,
288 #[serde(default, skip_serializing_if = "Option::is_none")]
289 pub model_id: Option<String>,
290 #[serde(default, skip_serializing_if = "Option::is_none")]
291 pub agent_backend: Option<String>,
292 #[serde(default, skip_serializing_if = "Option::is_none")]
293 pub outcome_sharing: Option<crate::chat_agents::OutcomeSharingPolicy>,
294}
295
296#[derive(Debug, Clone, Serialize, Deserialize)]
297pub struct ChatSessionSummary {
298 pub id: String,
299 pub title: String,
300 pub created_at: u64,
301 pub updated_at: u64,
302 pub message_count: u64,
303 #[serde(default)]
304 pub session_kind: SessionKind,
305 #[serde(default)]
306 pub participant_count: u64,
307 #[serde(default)]
308 pub session_did: String,
309}
310
311#[derive(Debug, Clone, Serialize, Deserialize)]
312pub struct ChatSession {
313 pub meta: SessionMeta,
314 pub environment: ChatEnvironment,
315 pub messages: Vec<ChatMessage>,
316}
317
318#[derive(Debug, Clone, Serialize, Deserialize)]
319struct LastSessionPrefs {
320 last_session_id: Option<String>,
321}
322
323pub fn chats_dir(storage_root: &Path) -> PathBuf {
324 storage_root.join("Chats")
325}
326
327fn session_dir(storage_root: &Path, id: &str) -> PathBuf {
328 chats_dir(storage_root).join(id)
329}
330
331fn session_meta_path(storage_root: &Path, id: &str) -> PathBuf {
332 session_dir(storage_root, id).join("session.json")
333}
334
335fn environment_path(storage_root: &Path, id: &str) -> PathBuf {
336 session_dir(storage_root, id).join("environment.json")
337}
338
339fn messages_path(storage_root: &Path, id: &str) -> PathBuf {
340 session_dir(storage_root, id).join("messages.jsonl")
341}
342
343fn wal_path(storage_root: &Path, id: &str) -> PathBuf {
344 session_dir(storage_root, id).join("chat.wal")
345}
346
347fn q42_path(storage_root: &Path, id: &str) -> PathBuf {
348 session_dir(storage_root, id).join("chat.q42")
349}
350
351fn last_session_prefs_path() -> PathBuf {
352 crate::state::app_meta_dir().join("chat_prefs.json")
353}
354
355fn unix_now() -> u64 {
356 std::time::SystemTime::now()
357 .duration_since(std::time::UNIX_EPOCH)
358 .unwrap_or_default()
359 .as_secs()
360}
361
362fn new_session_id() -> String {
363 let ts = unix_now();
364 let rnd: u64 = rand::random();
365 format!("{ts}-{rnd:016x}")
366}
367
368fn session_subject_hash(session_id: &str) -> u64 {
369 q_hash(&format!("chat:session:{session_id}"))
370}
371
372fn message_object_hash(lamport: u64) -> u64 {
373 q_hash(&format!("msg:{lamport}")) & OBJECT_HASH_MASK
374}
375
376fn role_context_hash(role: Role) -> u64 {
377 q_hash(&format!("chat:role:{}", role.as_str()))
378}
379
380fn build_message_quin(session_id: &str, role: Role, lamport: u64, content_hash: u64) -> NQuin {
381 let subject = session_subject_hash(session_id);
382 let predicate = q_hash("chat:hasMessage");
383 let object = message_object_hash(lamport);
384 let context = role_context_hash(role);
385 let metadata = (lamport & LAMPORT_MASK) << LAMPORT_SHIFT | (content_hash & 0xFFFF_FFFF);
386 let parity = subject ^ predicate ^ object ^ context ^ metadata;
387 NQuin {
388 subject,
389 predicate,
390 object,
391 context,
392 metadata,
393 parity,
394 }
395}
396
397fn write_quins_to_q42(quins: &[NQuin], out_path: &Path) -> Result<u64, ChatError> {
398 let count = qualia_core_db::q42_volume::write_sorted_quins_volume(out_path, quins)?;
399 Ok(count as u64)
400}
401
402pub fn get_last_session_id() -> Option<String> {
403 let path = last_session_prefs_path();
404 let text = fs::read_to_string(path).ok()?;
405 let prefs: LastSessionPrefs = serde_json::from_str(&text).ok()?;
406 prefs.last_session_id
407}
408
409pub fn set_last_session_id(id: &str) -> Result<(), ChatError> {
410 let path = last_session_prefs_path();
411 if let Some(parent) = path.parent() {
412 fs::create_dir_all(parent)?;
413 }
414 let prefs = LastSessionPrefs {
415 last_session_id: Some(id.to_string()),
416 };
417 fs::write(path, serde_json::to_string_pretty(&prefs)?)?;
418 Ok(())
419}
420
421pub fn create_session(
422 storage_root: &Path,
423 title: Option<String>,
424 env: Option<ChatEnvironment>,
425) -> Result<String, ChatError> {
426 fs::create_dir_all(chats_dir(storage_root))?;
427 let id = new_session_id();
428 let dir = session_dir(storage_root, &id);
429 fs::create_dir_all(&dir)?;
430
431 let now = unix_now();
432 let title = title.unwrap_or_else(|| "New chat".to_string());
433 let environment =
434 env.unwrap_or_else(|| ChatEnvironment::default_for_session(&id, storage_root));
435
436 let profile = crate::user_profile::load_profile();
437 let owner_did = profile.public_did.clone();
438
439 let session_did = compile_session_did(&id, SessionKind::Solo);
440 let meta = SessionMeta {
441 id: id.clone(),
442 title: title.clone(),
443 created_at: now,
444 updated_at: now,
445 message_count: 0,
446 next_lamport: 1,
447 environment_ref: "environment.json".to_string(),
448 session_kind: SessionKind::Solo,
449 participants: Vec::new(),
450 owner_did,
451 session_did,
452 };
453
454 fs::write(
455 session_meta_path(storage_root, &id),
456 serde_json::to_string_pretty(&meta)?,
457 )?;
458 fs::write(
459 environment_path(storage_root, &id),
460 serde_json::to_string_pretty(&environment)?,
461 )?;
462 fs::File::create(messages_path(storage_root, &id))?;
463 let _ = WriteAheadLog::open(wal_path(storage_root, &id))
464 .map_err(|e| ChatError::Wal(format!("Cannot create chat.wal: {e}")))?;
465
466 set_last_session_id(&id)?;
467 let _ = crate::chat_agents::ensure_local_agent_config(storage_root, &id);
468 Ok(id)
469}
470
471pub fn list_sessions(storage_root: &Path) -> Result<Vec<ChatSessionSummary>, ChatError> {
472 let root = chats_dir(storage_root);
473 if !root.is_dir() {
474 return Ok(Vec::new());
475 }
476
477 let mut out = Vec::new();
478 for entry in fs::read_dir(&root)?.filter_map(Result::ok) {
479 if !entry.file_type().map(|t| t.is_dir()).unwrap_or(false) {
480 continue;
481 }
482 let id = entry.file_name().to_string_lossy().into_owned();
483 let meta_path = session_meta_path(storage_root, &id);
484 if !meta_path.is_file() {
485 continue;
486 }
487 let mut meta: SessionMeta = serde_json::from_str(&fs::read_to_string(&meta_path)?)?;
488 persist_session_did_if_needed(storage_root, &mut meta)?;
489 out.push(ChatSessionSummary {
490 id: meta.id,
491 title: meta.title,
492 created_at: meta.created_at,
493 updated_at: meta.updated_at,
494 message_count: meta.message_count,
495 session_kind: meta.session_kind,
496 participant_count: meta.participants.len() as u64,
497 session_did: meta.session_did,
498 });
499 }
500
501 out.sort_by(|a, b| b.updated_at.cmp(&a.updated_at));
502 Ok(out)
503}
504
505fn read_messages_jsonl(path: &Path) -> Result<Vec<ChatMessage>, ChatError> {
506 if !path.is_file() {
507 return Ok(Vec::new());
508 }
509 let file = File::open(path)?;
510 let reader = BufReader::new(file);
511 let mut messages = Vec::new();
512 for line in reader.lines() {
513 let line = line?;
514 if line.trim().is_empty() {
515 continue;
516 }
517 messages.push(serde_json::from_str(&line)?);
518 }
519 Ok(messages)
520}
521
522fn load_environment_or_default(
523 storage_root: &Path,
524 id: &str,
525) -> Result<ChatEnvironment, ChatError> {
526 let env_path = environment_path(storage_root, id);
527 if !env_path.is_file() {
528 return Ok(ChatEnvironment::default_for_session(id, storage_root));
529 }
530
531 let text = fs::read_to_string(&env_path)?;
532 if text.trim().is_empty() {
533 log::warn!(
534 "Chat session {} had an empty environment.json; regenerating default environment",
535 id
536 );
537 let env = ChatEnvironment::default_for_session(id, storage_root);
538 let _ = env.save_to_session_dir(storage_root);
539 return Ok(env);
540 }
541
542 match serde_json::from_str(&text) {
543 Ok(env) => Ok(env),
544 Err(err) => {
545 log::warn!(
546 "Chat session {} had an invalid environment.json ({}); regenerating default environment",
547 id,
548 err
549 );
550 let env = ChatEnvironment::default_for_session(id, storage_root);
551 let _ = env.save_to_session_dir(storage_root);
552 Ok(env)
553 }
554 }
555}
556
557pub fn load_session(storage_root: &Path, id: &str) -> Result<ChatSession, ChatError> {
558 let meta_path = session_meta_path(storage_root, id);
559 if !meta_path.is_file() {
560 return Err(ChatError::NotFound(id.to_string()));
561 }
562 let mut meta: SessionMeta = serde_json::from_str(&fs::read_to_string(&meta_path)?)?;
563 persist_session_did_if_needed(storage_root, &mut meta)?;
564 let environment = load_environment_or_default(storage_root, id)?;
565 let messages = read_messages_jsonl(&messages_path(storage_root, id))?;
566 Ok(ChatSession {
567 meta,
568 environment,
569 messages,
570 })
571}
572
573pub fn append_message(
574 storage_root: &Path,
575 id: &str,
576 role: Role,
577 content: &str,
578) -> Result<u64, ChatError> {
579 append_message_with_options(storage_root, id, role, content, None, None)
580}
581
582pub fn append_message_with_options(
583 storage_root: &Path,
584 id: &str,
585 role: Role,
586 content: &str,
587 reply_to_fragment: Option<String>,
588 source: Option<String>,
589) -> Result<u64, ChatError> {
590 append_message_with_author(
591 storage_root,
592 id,
593 role,
594 content,
595 reply_to_fragment,
596 source,
597 None,
598 None,
599 None,
600 )
601}
602
603pub fn append_message_with_author(
604 storage_root: &Path,
605 id: &str,
606 role: Role,
607 content: &str,
608 reply_to_fragment: Option<String>,
609 source: Option<String>,
610 override_author_did: Option<String>,
611 override_author_name: Option<String>,
612 branch_type_override: Option<String>,
613) -> Result<u64, ChatError> {
614 let meta_path = session_meta_path(storage_root, id);
615 if !meta_path.is_file() {
616 return Err(ChatError::NotFound(id.to_string()));
617 }
618
619 let mut meta: SessionMeta = serde_json::from_str(&fs::read_to_string(&meta_path)?)?;
620 let lamport = meta.next_lamport;
621 meta.next_lamport += 1;
622 meta.message_count += 1;
623 meta.updated_at = unix_now();
624
625 let content_hash = content_hash_u64(content);
626 let quin = build_message_quin(id, role, lamport, content_hash);
627
628 let mut wal = WriteAheadLog::open(wal_path(storage_root, id))
629 .map_err(|e| ChatError::Wal(format!("Cannot open chat.wal: {e}")))?;
630 wal.append_mutation(&quin)
631 .map_err(|e| ChatError::Wal(format!("append failed: {e}")))?;
632
633 let (author_did, author_name) =
634 if override_author_did.is_some() || override_author_name.is_some() {
635 (override_author_did, override_author_name)
636 } else if meta.session_kind == SessionKind::Group && role == Role::User {
637 let profile = crate::user_profile::load_profile();
638 (
639 Some(profile.public_did),
640 profile
641 .sharing
642 .share_display_name
643 .then(|| profile.display_name),
644 )
645 } else {
646 (None, None)
647 };
648
649 let relay_ingest = source.as_ref().is_some_and(|s| s.starts_with("relay:"));
650
651 let mut msg = ChatMessage {
652 lamport,
653 role,
654 content: content.to_string(),
655 timestamp: meta.updated_at,
656 content_hash,
657 author_did,
658 author_name,
659 reply_to_fragment: reply_to_fragment.clone(),
660 source,
661 sub_agent_of: None,
662 agent_did: None,
663 model_id: None,
664 agent_backend: None,
665 outcome_sharing: None,
666 };
667
668 if role == Role::Agent && !relay_ingest {
669 let _ = crate::chat_agents::decorate_local_agent_message(storage_root, id, &mut msg);
670 }
671
672 let mut jsonl = OpenOptions::new()
673 .create(true)
674 .append(true)
675 .open(messages_path(storage_root, id))?;
676 writeln!(jsonl, "{}", serde_json::to_string(&msg)?)?;
677
678 fs::write(meta_path, serde_json::to_string_pretty(&meta)?)?;
679
680 if let Some(parent_id) = reply_to_fragment.as_deref() {
681 let graph = crate::chat_graph::load_graph(storage_root, id).ok();
682 let anchor = graph
683 .as_ref()
684 .and_then(|g| g.fragments.iter().find(|f| f.fragment_id == parent_id))
685 .map(|f| f.anchor_text.as_str());
686 let _ = crate::chat_graph::link_reply_to_fragment(
687 storage_root,
688 id,
689 parent_id,
690 lamport,
691 None,
692 anchor,
693 Some(content),
694 branch_type_override.as_deref(),
695 );
696 }
697
698 if meta.session_kind == SessionKind::Group && !relay_ingest {
699 if role != Role::Agent
700 || crate::chat_agents::can_relay_agent_outcome(&msg, &meta.participants)
701 {
702 let _ = crate::chat_relay::publish_session_message(storage_root, id, lamport);
703 }
704 }
705
706 set_last_session_id(id)?;
707 Ok(lamport)
708}
709
710pub fn append_relay_message_with_agent_meta(
712 storage_root: &Path,
713 id: &str,
714 role: Role,
715 content: &str,
716 reply_to_fragment: Option<String>,
717 source: Option<String>,
718 author_did: Option<String>,
719 author_name: Option<String>,
720 sub_agent_of: Option<String>,
721 agent_did: Option<String>,
722 model_id: Option<String>,
723 agent_backend: Option<String>,
724 outcome_sharing: Option<crate::chat_agents::OutcomeSharingPolicy>,
725) -> Result<u64, ChatError> {
726 let meta_path = session_meta_path(storage_root, id);
727 if !meta_path.is_file() {
728 return Err(ChatError::NotFound(id.to_string()));
729 }
730
731 let mut meta: SessionMeta = serde_json::from_str(&fs::read_to_string(&meta_path)?)?;
732 let lamport = meta.next_lamport;
733 meta.next_lamport += 1;
734 meta.message_count += 1;
735 meta.updated_at = unix_now();
736
737 let content_hash = content_hash_u64(content);
738 let quin = build_message_quin(id, role, lamport, content_hash);
739
740 let mut wal = WriteAheadLog::open(wal_path(storage_root, id))
741 .map_err(|e| ChatError::Wal(format!("Cannot open chat.wal: {e}")))?;
742 wal.append_mutation(&quin)
743 .map_err(|e| ChatError::Wal(format!("append failed: {e}")))?;
744
745 let msg = ChatMessage {
746 lamport,
747 role,
748 content: content.to_string(),
749 timestamp: meta.updated_at,
750 content_hash,
751 author_did,
752 author_name,
753 reply_to_fragment,
754 source,
755 sub_agent_of,
756 agent_did,
757 model_id,
758 agent_backend,
759 outcome_sharing,
760 };
761
762 let mut jsonl = OpenOptions::new()
763 .create(true)
764 .append(true)
765 .open(messages_path(storage_root, id))?;
766 writeln!(jsonl, "{}", serde_json::to_string(&msg)?)?;
767 fs::write(meta_path, serde_json::to_string_pretty(&meta)?)?;
768 set_last_session_id(id)?;
769 Ok(lamport)
770}
771
772pub fn compact_session_to_q42(storage_root: &Path, id: &str) -> Result<PathBuf, ChatError> {
773 let dir = session_dir(storage_root, id);
774 if !dir.is_dir() {
775 return Err(ChatError::NotFound(id.to_string()));
776 }
777
778 let wal_file = wal_path(storage_root, id);
779 let mut wal = WriteAheadLog::open(&wal_file)
780 .map_err(|e| ChatError::Wal(format!("Cannot open chat.wal: {e}")))?;
781 let quins = wal
782 .recover()
783 .map_err(|e| ChatError::Compact(format!("WAL recover failed: {e}")))?;
784
785 if quins.is_empty() {
786 return Err(ChatError::Compact("No quins in session WAL".to_string()));
787 }
788
789 let out = q42_path(storage_root, id);
790 let count = write_quins_to_q42(&quins, &out)?;
791 if count == 0 {
792 return Err(ChatError::Compact("Wrote zero quins".to_string()));
793 }
794 Ok(out)
795}
796
797pub fn delete_session(storage_root: &Path, id: &str) -> Result<(), ChatError> {
798 let dir = session_dir(storage_root, id);
799 if !dir.is_dir() {
800 return Err(ChatError::NotFound(id.to_string()));
801 }
802 fs::remove_dir_all(&dir)?;
803
804 if get_last_session_id().as_deref() == Some(id) {
805 let path = last_session_prefs_path();
806 let prefs = LastSessionPrefs {
807 last_session_id: None,
808 };
809 let _ = fs::write(path, serde_json::to_string_pretty(&prefs)?);
810 }
811 Ok(())
812}
813
814fn owner_participant(profile: &crate::user_profile::UserProfile) -> ChatParticipant {
815 ChatParticipant {
816 did: profile.public_did.clone(),
817 display_name: profile.display_name.clone(),
818 actor_id: "self".to_string(),
819 role: "owner".to_string(),
820 joined_at: unix_now(),
821 }
822}
823
824fn contacts_for_dids(dids: &[String]) -> Vec<ChatParticipant> {
825 let contacts = crate::social_connect::list_chat_contacts();
826 let now = unix_now();
827 dids.iter()
828 .filter_map(|did| {
829 contacts
830 .iter()
831 .find(|c| c.did == *did)
832 .map(|c| ChatParticipant {
833 did: c.did.clone(),
834 display_name: c.display_name.clone(),
835 actor_id: c.actor_id.clone(),
836 role: "member".to_string(),
837 joined_at: now,
838 })
839 })
840 .collect()
841}
842
843fn sync_participants_to_environment(
844 storage_root: &Path,
845 id: &str,
846 meta: &SessionMeta,
847) -> Result<(), ChatError> {
848 let env_path = environment_path(storage_root, id);
849 let mut environment = load_environment_or_default(storage_root, id)?;
850 environment.session_kind = meta.session_kind;
851 environment.participants = meta.participants.clone();
852 fs::write(env_path, serde_json::to_string_pretty(&environment)?)?;
853 Ok(())
854}
855
856pub fn create_group_session(
857 storage_root: &Path,
858 title: Option<String>,
859 participant_dids: &[String],
860) -> Result<String, ChatError> {
861 let profile = crate::user_profile::load_profile();
862 if !profile.sharing.allow_group_chat_invites {
863 return Err(ChatError::InvalidSession(
864 "Group chat invites are disabled in your profile.".to_string(),
865 ));
866 }
867
868 let mut participants = vec![owner_participant(&profile)];
869 participants.extend(contacts_for_dids(participant_dids));
870
871 if participants.len() < 2 {
872 return Err(ChatError::InvalidSession(
873 "Select at least one friend to start a group chat.".to_string(),
874 ));
875 }
876
877 fs::create_dir_all(chats_dir(storage_root))?;
878 let id = new_session_id();
879 let dir = session_dir(storage_root, &id);
880 fs::create_dir_all(&dir)?;
881
882 let now = unix_now();
883 let title = title.unwrap_or_else(|| {
884 let names: Vec<_> = participants
885 .iter()
886 .filter(|p| p.role != "owner")
887 .map(|p| p.display_name.as_str())
888 .take(3)
889 .collect();
890 if names.is_empty() {
891 "Group chat".to_string()
892 } else {
893 format!("Group: {}", names.join(", "))
894 }
895 });
896
897 let mut environment = ChatEnvironment::default_for_session(&id, storage_root);
898 environment.session_kind = SessionKind::Group;
899 environment.participants = participants.clone();
900
901 let session_did = compile_session_did(&id, SessionKind::Group);
902 let meta = SessionMeta {
903 id: id.clone(),
904 title: title.clone(),
905 created_at: now,
906 updated_at: now,
907 message_count: 0,
908 next_lamport: 1,
909 environment_ref: "environment.json".to_string(),
910 session_kind: SessionKind::Group,
911 participants,
912 owner_did: profile.public_did,
913 session_did,
914 };
915
916 fs::write(
917 session_meta_path(storage_root, &id),
918 serde_json::to_string_pretty(&meta)?,
919 )?;
920 fs::write(
921 environment_path(storage_root, &id),
922 serde_json::to_string_pretty(&environment)?,
923 )?;
924 fs::File::create(messages_path(storage_root, &id))?;
925 let _ = WriteAheadLog::open(wal_path(storage_root, &id))
926 .map_err(|e| ChatError::Wal(format!("Cannot create chat.wal: {e}")))?;
927
928 set_last_session_id(&id)?;
929 let _ = crate::chat_agents::ensure_local_agent_config(storage_root, &id);
930 Ok(id)
931}
932
933pub fn add_participant(
934 storage_root: &Path,
935 id: &str,
936 participant_did: &str,
937) -> Result<Vec<ChatParticipant>, ChatError> {
938 let meta_path = session_meta_path(storage_root, id);
939 if !meta_path.is_file() {
940 return Err(ChatError::NotFound(id.to_string()));
941 }
942 let mut meta: SessionMeta = serde_json::from_str(&fs::read_to_string(&meta_path)?)?;
943 if meta.session_kind != SessionKind::Group {
944 meta.session_kind = SessionKind::Group;
945 }
946 if meta.participants.iter().any(|p| p.did == participant_did) {
947 return Ok(meta.participants.clone());
948 }
949
950 let added = contacts_for_dids(&[participant_did.to_string()]);
951 let Some(participant) = added.into_iter().next() else {
952 return Err(ChatError::InvalidSession(format!(
953 "No contact found for DID {participant_did}. Add them via Profile → Add Friend first."
954 )));
955 };
956
957 meta.participants.push(participant);
958 meta.updated_at = unix_now();
959 fs::write(&meta_path, serde_json::to_string_pretty(&meta)?)?;
960 sync_participants_to_environment(storage_root, id, &meta)?;
961 Ok(meta.participants.clone())
962}
963
964pub fn remove_participant(
965 storage_root: &Path,
966 id: &str,
967 participant_did: &str,
968) -> Result<Vec<ChatParticipant>, ChatError> {
969 let meta_path = session_meta_path(storage_root, id);
970 if !meta_path.is_file() {
971 return Err(ChatError::NotFound(id.to_string()));
972 }
973 let mut meta: SessionMeta = serde_json::from_str(&fs::read_to_string(&meta_path)?)?;
974 if meta.owner_did == participant_did {
975 return Err(ChatError::InvalidSession(
976 "Cannot remove the session owner.".to_string(),
977 ));
978 }
979 meta.participants.retain(|p| p.did != participant_did);
980 meta.updated_at = unix_now();
981 fs::write(&meta_path, serde_json::to_string_pretty(&meta)?)?;
982 sync_participants_to_environment(storage_root, id, &meta)?;
983 Ok(meta.participants.clone())
984}
985
986pub fn get_participants(storage_root: &Path, id: &str) -> Result<Vec<ChatParticipant>, ChatError> {
987 let meta_path = session_meta_path(storage_root, id);
988 if !meta_path.is_file() {
989 return Err(ChatError::NotFound(id.to_string()));
990 }
991 let meta: SessionMeta = serde_json::from_str(&fs::read_to_string(&meta_path)?)?;
992 Ok(meta.participants.clone())
993}
994
995pub fn list_session_share_targets(
997 storage_root: &Path,
998) -> Result<Vec<ChatSessionShareTarget>, ChatError> {
999 let summaries = list_sessions(storage_root)?;
1000 Ok(summaries
1001 .into_iter()
1002 .map(|s| ChatSessionShareTarget {
1003 session_id: s.id,
1004 session_did: s.session_did,
1005 title: s.title,
1006 session_kind: s.session_kind,
1007 participant_count: s.participant_count,
1008 })
1009 .collect())
1010}
1011
1012pub fn get_session_did(storage_root: &Path, session_id: &str) -> Result<String, ChatError> {
1013 let meta_path = session_meta_path(storage_root, session_id);
1014 if !meta_path.is_file() {
1015 return Err(ChatError::NotFound(session_id.to_string()));
1016 }
1017 let mut meta: SessionMeta = serde_json::from_str(&fs::read_to_string(&meta_path)?)?;
1018 persist_session_did_if_needed(storage_root, &mut meta)?;
1019 Ok(meta.session_did)
1020}
1021
1022pub fn rename_session(storage_root: &Path, id: &str, title: &str) -> Result<(), ChatError> {
1023 let meta_path = session_meta_path(storage_root, id);
1024 if !meta_path.is_file() {
1025 return Err(ChatError::NotFound(id.to_string()));
1026 }
1027 let mut meta: SessionMeta = serde_json::from_str(&fs::read_to_string(&meta_path)?)?;
1028 meta.title = title.trim().to_string();
1029 meta.updated_at = unix_now();
1030 fs::write(meta_path, serde_json::to_string_pretty(&meta)?)?;
1031 Ok(())
1032}
1033
1034#[cfg(test)]
1035mod tests {
1036 use super::*;
1037 use std::env;
1038
1039 fn temp_storage() -> PathBuf {
1040 let mut dir = env::temp_dir();
1041 dir.push(format!("qualia-chat-test-{}", rand::random::<u32>()));
1042 dir
1043 }
1044
1045 #[test]
1046 fn create_list_append_load_delete_roundtrip() {
1047 let storage = temp_storage();
1048 let id = create_session(&storage, Some("Test".into()), None).unwrap();
1049 assert_eq!(get_last_session_id().as_deref(), Some(id.as_str()));
1050 let sessions = list_sessions(&storage).unwrap();
1051 assert_eq!(sessions.len(), 1);
1052 assert_eq!(sessions[0].id, id);
1053
1054 let lamport = append_message(&storage, &id, Role::User, "hello").unwrap();
1055 assert_eq!(lamport, 1);
1056 append_message(&storage, &id, Role::Agent, "hi there").unwrap();
1057
1058 let session = load_session(&storage, &id).unwrap();
1059 assert_eq!(session.messages.len(), 2);
1060 assert_eq!(session.messages[0].content, "hello");
1061
1062 let q42 = compact_session_to_q42(&storage, &id).unwrap();
1063 assert!(q42.is_file());
1064 let bytes = fs::read(&q42).unwrap();
1065 assert!(
1066 bytes.starts_with(&qualia_core_db::q42_volume::Q42_MAGIC),
1067 "chat.q42 must be a unified v3 volume"
1068 );
1069 let volume = qualia_core_db::q42_volume::Q42Volume::open(&q42).unwrap();
1070 volume.verify_all_blocks().expect("chat.q42 ECC");
1071 assert_eq!(volume.read_all_quins().unwrap().len(), 2);
1072
1073 delete_session(&storage, &id).unwrap();
1074 assert!(load_session(&storage, &id).is_err());
1075 let _ = fs::remove_dir_all(&storage);
1076 }
1077}