Skip to main content

qualia_client_core/
chat_session.rs

1//! Chat session persistence — WAL quins + JSON sidecar under `{storage}/Chats/`.
2
3use 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
12/// The content hash carried on a chat message: FNV-1a `q_hash` of the content, masked to the 60-bit
13/// object space. Public so interop layers (e.g. `solid_chat`) can reproduce it for imported messages.
14pub 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/// Per-ontology scope summary compiled into the chat environment.
97#[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/// Snapshot of chat inference scope and LLM-visible capabilities.
112#[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    /// When true, chat turns route through the neuro-symbolic sieve + WAL orchestrator path.
142    #[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    /// Stable DID for ontology / torrent sharing scoped to this chat session or group.
225    #[serde(default)]
226    pub session_did: String,
227}
228
229/// Target descriptor for ontology sharing UI (solo chats and group sessions).
230#[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    /// Human principal DID — agent messages are sub-agents of this participant.
284    #[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
710/// Ingest a relayed message with pre-set sub-agent metadata (skips local decoration / re-publish).
711pub 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
995/// Sessions and groups that can be selected as share targets (each has a stable `session_did`).
996pub 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}