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