1use anyhow::{Context, Result, anyhow};
8use onlyne_config::layout::RoleWorkspace;
9use parking_lot::Mutex;
10use serde::{Deserialize, Serialize};
11use serde_json::Value;
12use std::collections::BTreeMap;
13use std::fs::{File, OpenOptions};
14use std::io::{Read, Seek, SeekFrom, Write};
15use std::path::{Path, PathBuf};
16use std::sync::Arc;
17
18#[derive(Clone, Debug, PartialEq)]
20pub struct ContentRecord {
21 pub seq: u64,
22 pub task_id: String,
23 pub session_id: Option<String>,
24 pub at: String,
25 pub record: Value,
27}
28
29pub trait ContentSink: Send + Sync {
35 fn publish(&self, record: ContentRecord);
36}
37
38#[derive(Clone, Default)]
39pub(crate) struct ContentWriter {
40 inner: Arc<Mutex<ContentWriterState>>,
41}
42
43#[derive(Default)]
44struct ContentWriterState {
45 sink: Option<Arc<dyn ContentSink>>,
46 heads: BTreeMap<PathBuf, u64>,
47}
48
49#[derive(Debug, Serialize, Deserialize)]
50#[serde(deny_unknown_fields)]
51struct ContentIndexEntry {
52 seq: u64,
53 task_id: String,
54 #[serde(default, skip_serializing_if = "Option::is_none")]
55 session_id: Option<String>,
56 at: String,
57 journal: String,
58 offset: u64,
59 len: u64,
60}
61
62impl ContentWriter {
63 pub(crate) fn set_sink(&self, sink: Arc<dyn ContentSink>) {
64 self.inner.lock().sink = Some(sink);
65 }
66
67 pub(crate) fn append(
70 &self,
71 workspace: &Path,
72 task_id: &str,
73 session_id: Option<&str>,
74 journal: &Path,
75 record: &Value,
76 at: &str,
77 ) -> std::io::Result<()> {
78 let mut state = self.inner.lock();
79 let head = match state.heads.get(workspace).copied() {
80 Some(head) => head,
81 None => indexed_head(workspace)?,
82 };
83 let seq = head
84 .checked_add(1)
85 .ok_or_else(|| std::io::Error::other("content sequence exhausted"))?;
86 state.heads.insert(workspace.to_path_buf(), seq);
90 let (offset, len) = append_record(journal, record)?;
91 let journal_name = journal
92 .file_name()
93 .and_then(|name| name.to_str())
94 .ok_or_else(|| {
95 std::io::Error::new(
96 std::io::ErrorKind::InvalidInput,
97 format!("journal has no UTF-8 file name: {}", journal.display()),
98 )
99 })?;
100 let entry = ContentIndexEntry {
101 seq,
102 task_id: task_id.to_string(),
103 session_id: session_id.map(str::to_string),
104 at: at.to_string(),
105 journal: journal_name.to_string(),
106 offset,
107 len,
108 };
109 append_index(workspace, &entry)?;
110 if let Some(sink) = state.sink.as_ref() {
111 sink.publish(ContentRecord {
112 seq,
113 task_id: task_id.to_string(),
114 session_id: session_id.map(str::to_string),
115 at: at.to_string(),
116 record: record.clone(),
117 });
118 }
119 Ok(())
120 }
121}
122
123fn indexed_head(workspace: &Path) -> std::io::Result<u64> {
124 let path = RoleWorkspace::resolve(workspace).content_index_path();
125 let text = match std::fs::read_to_string(&path) {
126 Ok(text) => text,
127 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(0),
128 Err(error) => return Err(error),
129 };
130 text.lines()
131 .filter(|line| !line.trim().is_empty())
132 .try_fold(0, |head, line| {
133 let entry: ContentIndexEntry = serde_json::from_str(line)
134 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?;
135 Ok(head.max(entry.seq))
136 })
137}
138
139fn append_record(path: &Path, record: &Value) -> std::io::Result<(u64, u64)> {
140 if let Some(dir) = path.parent() {
141 std::fs::create_dir_all(dir)?;
142 }
143 let body = serde_json::to_vec(record).map_err(std::io::Error::other)?;
144 let mut file = OpenOptions::new().create(true).append(true).open(path)?;
145 let offset = file.metadata()?.len();
146 file.write_all(&body)?;
147 file.write_all(b"\n")?;
148 file.flush()?;
149 Ok((offset, body.len() as u64))
150}
151
152fn append_index(workspace: &Path, entry: &ContentIndexEntry) -> std::io::Result<()> {
153 let path = RoleWorkspace::resolve(workspace).content_index_path();
154 if let Some(dir) = path.parent() {
155 std::fs::create_dir_all(dir)?;
156 }
157 let mut line = serde_json::to_vec(entry).map_err(std::io::Error::other)?;
158 line.push(b'\n');
159 let mut file = OpenOptions::new().create(true).append(true).open(path)?;
160 file.write_all(&line)?;
161 file.flush()
162}
163
164pub fn read_content_records(workspace: &Path) -> Result<Vec<ContentRecord>> {
170 let index_path = RoleWorkspace::resolve(workspace).content_index_path();
171 let text = match std::fs::read_to_string(&index_path) {
172 Ok(text) => text,
173 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
174 Err(error) => return Err(error).with_context(|| format!("read {}", index_path.display())),
175 };
176 let logs = index_path
177 .parent()
178 .ok_or_else(|| anyhow!("content index has no parent: {}", index_path.display()))?;
179 let mut records = Vec::new();
180 for (line_number, line) in text.lines().enumerate() {
181 if line.trim().is_empty() {
182 continue;
183 }
184 let entry: ContentIndexEntry = serde_json::from_str(line).with_context(|| {
185 format!(
186 "decode content index {} line {}",
187 index_path.display(),
188 line_number + 1
189 )
190 })?;
191 let path = safe_journal_path(logs, &entry.journal)?;
192 let record = read_indexed_record(&path, entry.offset, entry.len)
193 .with_context(|| format!("read content seq {} from {}", entry.seq, path.display()))?;
194 records.push(ContentRecord {
195 seq: entry.seq,
196 task_id: entry.task_id,
197 session_id: entry.session_id,
198 at: entry.at,
199 record,
200 });
201 }
202 records.sort_by_key(|record| record.seq);
203 for pair in records.windows(2) {
204 if pair[0].seq == pair[1].seq {
205 return Err(anyhow!("duplicate content sequence {}", pair[0].seq));
206 }
207 }
208 Ok(records)
209}
210
211fn safe_journal_path(logs: &Path, journal: &str) -> Result<PathBuf> {
212 let name = Path::new(journal);
213 if name.components().count() != 1
214 || name.file_name().and_then(|part| part.to_str()) != Some(journal)
215 {
216 return Err(anyhow!("invalid content journal name {journal:?}"));
217 }
218 Ok(logs.join(name))
219}
220
221fn read_indexed_record(path: &Path, offset: u64, len: u64) -> Result<Value> {
222 let len: usize = len
223 .try_into()
224 .map_err(|_| anyhow!("indexed content length {len} does not fit memory"))?;
225 let mut file = File::open(path)?;
226 file.seek(SeekFrom::Start(offset))?;
227 let mut bytes = vec![0; len];
228 file.read_exact(&mut bytes)?;
229 let mut newline = [0_u8; 1];
230 file.read_exact(&mut newline)?;
231 if newline[0] != b'\n' {
232 return Err(anyhow!("indexed content is not one complete JSONL record"));
233 }
234 serde_json::from_slice(&bytes).context("decode indexed journal record")
235}