Skip to main content

qualia_client_core/
chat_graph.rs

1//! Chat graph — selectable fragments and reply edges forming a DAG off the linear chat.
2
3use std::fs::{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
10use crate::chat_session::{ChatError, Role};
11
12const OBJECT_HASH_MASK: u64 = 0x0FFF_FFFF_FFFF_FFFF;
13const LAMPORT_SHIFT: u32 = 32;
14const LAMPORT_MASK: u64 = 0x1FFF_FFFF;
15
16#[derive(Debug, Clone, Serialize, Deserialize)]
17pub struct ChatFragment {
18    pub fragment_id: String,
19    pub message_lamport: u64,
20    pub anchor_start: u32,
21    pub anchor_end: u32,
22    pub anchor_text: String,
23    pub author_did: Option<String>,
24    pub author_name: Option<String>,
25    pub created_at: u64,
26}
27
28#[derive(Debug, Clone, Serialize, Deserialize)]
29pub struct ChatGraphEdge {
30    pub child_fragment_id: String,
31    pub parent_fragment_id: String,
32    pub reply_message_lamport: u64,
33    pub created_at: u64,
34    #[serde(default)]
35    pub branch_type_id: Option<String>,
36    #[serde(default)]
37    pub branch_label: Option<String>,
38    #[serde(default)]
39    pub branch_emoji: Option<String>,
40    #[serde(default)]
41    pub wordnet_grounding_hash: Option<String>,
42}
43
44#[derive(Debug, Clone, Serialize, Deserialize)]
45pub struct ChatGraphSnapshot {
46    pub fragments: Vec<ChatFragment>,
47    pub edges: Vec<ChatGraphEdge>,
48}
49
50fn fragments_path(storage_root: &Path, session_id: &str) -> PathBuf {
51    storage_root
52        .join("Chats")
53        .join(session_id)
54        .join("fragments.jsonl")
55}
56
57fn edges_path(storage_root: &Path, session_id: &str) -> PathBuf {
58    storage_root
59        .join("Chats")
60        .join(session_id)
61        .join("graph_edges.jsonl")
62}
63
64fn unix_now() -> u64 {
65    std::time::SystemTime::now()
66        .duration_since(std::time::UNIX_EPOCH)
67        .unwrap_or_default()
68        .as_secs()
69}
70
71pub fn fragment_id_for_span(session_id: &str, lamport: u64, start: u32, end: u32) -> String {
72    let raw = q_hash(&format!("frag:{session_id}:{lamport}:{start}:{end}"));
73    format!("{raw:016x}")
74}
75
76fn session_subject_hash(session_id: &str) -> u64 {
77    q_hash(&format!("chat:session:{session_id}"))
78}
79
80fn build_fragment_quin(session_id: &str, fragment_id_hex: &str, lamport: u64) -> NQuin {
81    let subject = session_subject_hash(session_id);
82    let predicate = q_hash("chat:hasFragment");
83    let object = u64::from_str_radix(fragment_id_hex, 16).unwrap_or(0) & OBJECT_HASH_MASK;
84    let context = q_hash(&format!("msg:{lamport}"));
85    let metadata = (lamport & LAMPORT_MASK) << LAMPORT_SHIFT;
86    let parity = subject ^ predicate ^ object ^ context ^ metadata;
87    NQuin {
88        subject,
89        predicate,
90        object,
91        context,
92        metadata,
93        parity,
94    }
95}
96
97fn build_reply_edge_quin(
98    session_id: &str,
99    child_fragment_id_hex: &str,
100    parent_fragment_id_hex: &str,
101    reply_lamport: u64,
102) -> NQuin {
103    let subject = u64::from_str_radix(child_fragment_id_hex, 16).unwrap_or(0) & OBJECT_HASH_MASK;
104    let predicate = q_hash("chat:repliesTo");
105    let object = u64::from_str_radix(parent_fragment_id_hex, 16).unwrap_or(0) & OBJECT_HASH_MASK;
106    let context = session_subject_hash(session_id);
107    let metadata = (reply_lamport & LAMPORT_MASK) << LAMPORT_SHIFT;
108    let parity = subject ^ predicate ^ object ^ context ^ metadata;
109    NQuin {
110        subject,
111        predicate,
112        object,
113        context,
114        metadata,
115        parity,
116    }
117}
118
119fn append_jsonl<T: Serialize>(path: &Path, row: &T) -> Result<(), ChatError> {
120    if let Some(parent) = path.parent() {
121        std::fs::create_dir_all(parent)?;
122    }
123    let mut file = OpenOptions::new().create(true).append(true).open(path)?;
124    writeln!(file, "{}", serde_json::to_string(row)?)?;
125    Ok(())
126}
127
128fn read_fragments(path: &Path) -> Result<Vec<ChatFragment>, ChatError> {
129    if !path.is_file() {
130        return Ok(Vec::new());
131    }
132    let file = File::open(path)?;
133    let reader = BufReader::new(file);
134    let mut out = Vec::new();
135    for line in reader.lines() {
136        let line = line?;
137        if line.trim().is_empty() {
138            continue;
139        }
140        out.push(serde_json::from_str(&line)?);
141    }
142    Ok(out)
143}
144
145fn read_edges(path: &Path) -> Result<Vec<ChatGraphEdge>, ChatError> {
146    if !path.is_file() {
147        return Ok(Vec::new());
148    }
149    let file = File::open(path)?;
150    let reader = BufReader::new(file);
151    let mut out = Vec::new();
152    for line in reader.lines() {
153        let line = line?;
154        if line.trim().is_empty() {
155            continue;
156        }
157        out.push(serde_json::from_str(&line)?);
158    }
159    Ok(out)
160}
161
162pub fn load_graph(storage_root: &Path, session_id: &str) -> Result<ChatGraphSnapshot, ChatError> {
163    Ok(ChatGraphSnapshot {
164        fragments: read_fragments(&fragments_path(storage_root, session_id))?,
165        edges: read_edges(&edges_path(storage_root, session_id))?,
166    })
167}
168
169pub fn create_fragment_from_selection(
170    storage_root: &Path,
171    session_id: &str,
172    message_lamport: u64,
173    message_content: &str,
174    anchor_start: u32,
175    anchor_end: u32,
176) -> Result<ChatFragment, ChatError> {
177    let start = anchor_start.min(message_content.len() as u32);
178    let end = anchor_end.max(start).min(message_content.len() as u32);
179    let anchor_text = message_content[start as usize..end as usize]
180        .trim()
181        .to_string();
182    if anchor_text.is_empty() {
183        return Err(ChatError::InvalidSession(
184            "Selected fragment is empty.".to_string(),
185        ));
186    }
187
188    let fragment_id = fragment_id_for_span(session_id, message_lamport, start, end);
189    let profile = crate::user_profile::load_profile();
190    let fragment = ChatFragment {
191        fragment_id: fragment_id.clone(),
192        message_lamport,
193        anchor_start: start,
194        anchor_end: end,
195        anchor_text,
196        author_did: profile
197            .sharing
198            .share_public_did
199            .then(|| profile.public_did.clone()),
200        author_name: profile
201            .sharing
202            .share_display_name
203            .then(|| profile.display_name.clone()),
204        created_at: unix_now(),
205    };
206
207    let frag_path = fragments_path(storage_root, session_id);
208    let existing = read_fragments(&frag_path)?;
209    if !existing.iter().any(|f| f.fragment_id == fragment_id) {
210        append_jsonl(&frag_path, &fragment)?;
211        let wal_path = storage_root.join("Chats").join(session_id).join("chat.wal");
212        let quin = build_fragment_quin(session_id, &fragment_id, message_lamport);
213        let mut wal = WriteAheadLog::open(&wal_path)
214            .map_err(|e| ChatError::Wal(format!("Cannot open chat.wal: {e}")))?;
215        wal.append_mutation(&quin)
216            .map_err(|e| ChatError::Wal(format!("fragment quin append failed: {e}")))?;
217    }
218
219    Ok(fragment)
220}
221
222pub fn link_reply_to_fragment(
223    storage_root: &Path,
224    session_id: &str,
225    parent_fragment_id: &str,
226    reply_message_lamport: u64,
227    child_fragment_id: Option<&str>,
228    anchor_text: Option<&str>,
229    reply_text: Option<&str>,
230    branch_type_override: Option<&str>,
231) -> Result<ChatGraphEdge, ChatError> {
232    let child_id = child_fragment_id
233        .map(|s| s.to_string())
234        .unwrap_or_else(|| fragment_id_for_span(session_id, reply_message_lamport, 0, 0));
235
236    let classification = if let Some(override_id) = branch_type_override {
237        crate::chat_ontology::list_branch_types(storage_root)
238            .into_iter()
239            .find(|t| t.id == override_id)
240            .map(|t| crate::chat_ontology::BranchClassification {
241                branch_type_id: t.id,
242                label: t.label,
243                emoji: t.emoji,
244                confidence: 1.0,
245                wordnet_grounding_hash: t.wordnet_grounding_hash,
246            })
247    } else {
248        match (anchor_text, reply_text) {
249            (Some(a), Some(r)) => Some(crate::chat_ontology::classify_branch(storage_root, a, r)),
250            _ => None,
251        }
252    };
253
254    let edge = ChatGraphEdge {
255        child_fragment_id: child_id.clone(),
256        parent_fragment_id: parent_fragment_id.to_string(),
257        reply_message_lamport,
258        created_at: unix_now(),
259        branch_type_id: classification.as_ref().map(|c| c.branch_type_id.clone()),
260        branch_label: classification.as_ref().map(|c| c.label.clone()),
261        branch_emoji: classification.as_ref().map(|c| c.emoji.clone()),
262        wordnet_grounding_hash: classification.and_then(|c| c.wordnet_grounding_hash),
263    };
264
265    append_jsonl(&edges_path(storage_root, session_id), &edge)?;
266
267    let wal_path = storage_root.join("Chats").join(session_id).join("chat.wal");
268    let quin = build_reply_edge_quin(
269        session_id,
270        &child_id,
271        parent_fragment_id,
272        reply_message_lamport,
273    );
274    let mut wal = WriteAheadLog::open(&wal_path)
275        .map_err(|e| ChatError::Wal(format!("Cannot open chat.wal: {e}")))?;
276    wal.append_mutation(&quin)
277        .map_err(|e| ChatError::Wal(format!("reply edge quin append failed: {e}")))?;
278
279    Ok(edge)
280}
281
282pub fn build_thread_context_block(
283    storage_root: &Path,
284    session_id: &str,
285    target_fragment_id: &str,
286    max_depth: usize,
287) -> Result<String, ChatError> {
288    let graph = load_graph(storage_root, session_id)?;
289    let session = crate::chat_session::load_session(storage_root, session_id)?;
290
291    let mut lines = vec!["[Chat graph thread context]".to_string()];
292    let mut current = target_fragment_id.to_string();
293    let mut depth = 0;
294
295    while depth < max_depth {
296        let fragment = graph.fragments.iter().find(|f| f.fragment_id == current);
297        if let Some(f) = fragment {
298            lines.push(format!(
299                "fragment {} (msg #{}) anchor=\"{}\"",
300                f.fragment_id, f.message_lamport, f.anchor_text
301            ));
302            if let Some(msg) = session
303                .messages
304                .iter()
305                .find(|m| m.lamport == f.message_lamport)
306            {
307                let author = msg.author_name.as_deref().unwrap_or(match msg.role {
308                    Role::User => "user",
309                    Role::Agent => "agent",
310                });
311                lines.push(format!("  full_message[{author}]: {}", msg.content));
312            }
313        }
314
315        let parent_edge = graph.edges.iter().find(|e| e.child_fragment_id == current);
316        match parent_edge {
317            Some(e) => {
318                current = e.parent_fragment_id.clone();
319                depth += 1;
320            }
321            None => break,
322        }
323    }
324
325    let child_edges: Vec<_> = graph
326        .edges
327        .iter()
328        .filter(|e| e.parent_fragment_id == target_fragment_id)
329        .collect();
330    if !child_edges.is_empty() {
331        lines.push("direct_replies:".to_string());
332        for e in child_edges {
333            if let Some(reply_msg) = session
334                .messages
335                .iter()
336                .find(|m| m.lamport == e.reply_message_lamport)
337            {
338                let branch = e.branch_emoji.as_deref().unwrap_or("💬");
339                let label = e.branch_label.as_deref().unwrap_or("Comment");
340                lines.push(format!(
341                    "  → {branch} {label} msg #{}: {}",
342                    e.reply_message_lamport, reply_msg.content
343                ));
344            }
345        }
346    }
347
348    Ok(lines.join("\n"))
349}
350
351pub fn append_message_with_reply(
352    storage_root: &Path,
353    session_id: &str,
354    role: Role,
355    content: &str,
356    reply_to_fragment_id: Option<&str>,
357) -> Result<u64, ChatError> {
358    crate::chat_session::append_message_with_options(
359        storage_root,
360        session_id,
361        role,
362        content,
363        reply_to_fragment_id.map(|s| s.to_string()),
364        None,
365    )
366}