Skip to main content

onlyne_client/
content.rs

1//! Durable cursors for the session content a backend journals.
2//!
3//! The event JSON remains in the task journal exactly once.  The companion
4//! index stores only enough metadata to recover the role-wide sequence and read
5//! those original bytes back, whoever the eventual reader is.
6
7use 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/// One content record after the journal append that gave it a role-wide cursor.
19#[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    /// The exact JSON value appended to the task journal.
26    pub record: Value,
27}
28
29/// Client-owned edge of a backend's content stream.
30///
31/// The backend calls `publish` only after the journal object and its durable
32/// cursor index entry have both been appended. An implementation must do bounded
33/// bookkeeping only and must never block the journalling turn.
34pub 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    /// Append one unchanged JSON object, index its original bytes, then notify
68    /// the client. The lock is the role-wide ordering point across task turns.
69    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        // Reserve in memory before either append. A failed record is never
87        // published, and the next successful record must not reuse its number in
88        // this client run.
89        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
164/// Read the original journal objects named by the durable cursor index.
165///
166/// Entries are returned in role sequence order. A corrupt index or an indexed
167/// range that no longer names exactly one JSON object is an error rather than a
168/// silent stream gap.
169pub 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}