Skip to main content

kcode_k1_rust_projection/
lib.rs

1use kcode_k1_peering::K1Peering;
2use kcode_k1_rust_package::{AuthorityId, LibraryFamily, LibraryId, SourceFile, SourcePackage};
3use kcode_k1_rust_transaction::{decode, encode};
4use kcode_k1_transaction_id::TxId as AuthorityTxId;
5use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem, SubsystemId, TxId};
6use semver::{Version, VersionReq};
7use std::collections::BTreeMap;
8use std::fs;
9use std::path::{Path, PathBuf};
10use std::str::FromStr;
11use std::sync::{Arc, Mutex, MutexGuard};
12use std::time::{Duration, Instant};
13
14const SNAPSHOT: &str = "rust-projection.snapshot";
15const SUBSYSTEM: &str = "k1-rust-libs";
16
17#[derive(Clone, Copy, Debug, Eq, PartialEq)]
18pub enum PublishStatus {
19    Published,
20    Idempotent,
21    Conflict,
22}
23
24#[derive(Clone, Debug, Eq, PartialEq)]
25pub struct PublishOutcome {
26    status: PublishStatus,
27    winning: TxId,
28    submitted: Option<TxId>,
29}
30
31impl PublishOutcome {
32    pub const fn status(&self) -> PublishStatus {
33        self.status
34    }
35
36    pub const fn winning(&self) -> TxId {
37        self.winning
38    }
39
40    pub const fn submitted(&self) -> Option<TxId> {
41        self.submitted
42    }
43}
44
45pub struct K1RustProjection {
46    state: Arc<Mutex<State>>,
47    peering: Arc<K1Peering>,
48}
49
50struct State {
51    root: PathBuf,
52    control: PathBuf,
53    cursor: Option<TxId>,
54    winners: BTreeMap<LibraryId, Winner>,
55    available: bool,
56}
57
58struct Winner {
59    transaction: TxId,
60    package: SourcePackage,
61}
62
63struct Handler(Arc<Mutex<State>>);
64struct OpenTimer(Instant);
65type Saved = (Option<TxId>, BTreeMap<LibraryId, Winner>);
66
67impl Drop for OpenTimer {
68    fn drop(&mut self) {
69        if self.0.elapsed() > Duration::from_millis(100) {
70            eprintln!(r#"{{"level":"warning","event":"k1_rust_projection_open_slow"}}"#);
71        }
72    }
73}
74
75impl K1RustProjection {
76    pub fn open(
77        projection_root: impl AsRef<Path>,
78        control_root: impl AsRef<Path>,
79        ordering: Arc<K1TxnOrdering>,
80        peering: Arc<K1Peering>,
81    ) -> Result<Self, String> {
82        let _timer = OpenTimer(Instant::now());
83        let projection_root = projection_root.as_ref().to_path_buf();
84        let control = control_root.as_ref().to_path_buf();
85        directory(&control)?;
86        let saved = read_snapshot(&control, &projection_root, &ordering)?;
87        let rebuilding = saved.is_none();
88        let (root, cursor, winners) = match saved {
89            Some((cursor, winners)) => (projection_root.clone(), cursor, winners),
90            None => (staging_root(&projection_root)?, None, BTreeMap::new()),
91        };
92        directory(&root)?;
93        let state = Arc::new(Mutex::new(State {
94            root,
95            control,
96            cursor,
97            winners,
98            available: true,
99        }));
100        let handler = Arc::new(Handler(state.clone()));
101        let after = lock(&state)?.cursor;
102        if let Err(error) =
103            ordering.register_subsystem(SubsystemId::from_str(SUBSYSTEM)?, after, handler)
104        {
105            lock(&state)?.available = false;
106            return Err(error);
107        }
108        if rebuilding && let Err(error) = promote(&state, &projection_root) {
109            lock(&state)?.available = false;
110            return Err(error);
111        }
112        Ok(Self { state, peering })
113    }
114
115    pub fn load(&self, id: &LibraryId) -> Result<Option<SourcePackage>, String> {
116        let root = available_state(&self.state)?.root.clone();
117        load_from(&root, id)
118    }
119
120    pub fn resolve(
121        &self,
122        family: &LibraryFamily,
123        requirement: &VersionReq,
124    ) -> Result<Option<SourcePackage>, String> {
125        let state = available_state(&self.state)?;
126        Ok(state
127            .winners
128            .iter()
129            .filter(|(id, _)| id.family() == family && requirement.matches(id.version()))
130            .max_by(|(left, _), (right, _)| left.version().cmp(right.version()))
131            .map(|(_, winner)| winner.package.clone()))
132    }
133
134    pub fn publish(&self, package: &SourcePackage) -> Result<PublishOutcome, String> {
135        drop(available_state(&self.state)?);
136        let payload = encode(package).map_err(|error| error.to_string())?;
137        match self
138            .peering
139            .submit_txn(SubsystemId::from_str(SUBSYSTEM)?, &payload)
140        {
141            Ok(id) => self.reconcile(package, Some(id)),
142            Err(error) => self.reconcile(package, None).or(Err(error)),
143        }
144    }
145
146    pub fn cursor(&self) -> Result<Option<TxId>, String> {
147        Ok(available_state(&self.state)?.cursor)
148    }
149
150    pub fn is_available(&self) -> bool {
151        lock(&self.state).is_ok_and(|state| state.available)
152    }
153
154    fn reconcile(
155        &self,
156        package: &SourcePackage,
157        submitted: Option<TxId>,
158    ) -> Result<PublishOutcome, String> {
159        let state = available_state(&self.state)?;
160        let winner = state
161            .winners
162            .get(package.id())
163            .ok_or_else(|| "publication result has no canonical winner".to_owned())?;
164        let status = if winner.package != *package {
165            PublishStatus::Conflict
166        } else if submitted == Some(winner.transaction) {
167            PublishStatus::Published
168        } else {
169            PublishStatus::Idempotent
170        };
171        Ok(PublishOutcome {
172            status,
173            winning: winner.transaction,
174            submitted,
175        })
176    }
177}
178
179impl Subsystem for Handler {
180    fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
181        let mut state = lock(&self.0)?;
182        if !state.available {
183            return Ok(());
184        }
185        let result = apply(&mut state, id, payload);
186        state.available &= result.is_ok();
187        result
188    }
189
190    fn reorg(&self) -> Result<(), String> {
191        lock(&self.0)?.available = false;
192        Ok(())
193    }
194}
195
196fn lock(state: &Arc<Mutex<State>>) -> Result<MutexGuard<'_, State>, String> {
197    state.lock().map_err(|_| "projection mutex poisoned".into())
198}
199
200fn available_state(state: &Arc<Mutex<State>>) -> Result<MutexGuard<'_, State>, String> {
201    let state = lock(state)?;
202    state
203        .available
204        .then_some(state)
205        .ok_or("projection is unavailable".into())
206}
207
208fn apply(state: &mut State, id: TxId, payload: &[u8]) -> Result<(), String> {
209    if let Ok(package) = decode(payload) {
210        claim(state, id, package)?;
211    }
212    state.cursor = Some(id);
213    write_snapshot(state)
214}
215
216fn claim(state: &mut State, id: TxId, package: SourcePackage) -> Result<(), String> {
217    let key = package.id().clone();
218    if state.winners.contains_key(&key) {
219        return Ok(());
220    }
221    match load_from(&state.root, &key)? {
222        Some(existing) if existing != package => {
223            return Err("materialized package disagrees with replay".into());
224        }
225        Some(_) => {}
226        None => materialize(&state.root, &package)?,
227    }
228    state.winners.insert(
229        key,
230        Winner {
231            transaction: id,
232            package,
233        },
234    );
235    Ok(())
236}
237
238fn read_snapshot(
239    control: &Path,
240    root: &Path,
241    ordering: &K1TxnOrdering,
242) -> Result<Option<Saved>, String> {
243    let path = control.join(SNAPSHOT);
244    let Some(metadata) = entry_metadata(&path)? else {
245        return Ok(None);
246    };
247    if !metadata.file_type().is_file() {
248        return Err("snapshot is not an ordinary file".into());
249    }
250    let text = fs::read_to_string(path).map_err(|error| error.to_string())?;
251    let Some(text) = text.strip_suffix('\n') else {
252        return Ok(None);
253    };
254    let mut lines = text.split('\n');
255    if lines.next() != Some("K1RUSTPROJECTION1") {
256        return Ok(None);
257    }
258    let cursor = match lines.next() {
259        Some("-") => None,
260        Some(value) => match canonical_txid(value) {
261            Some(id) => Some(id),
262            None => return Ok(None),
263        },
264        None => return Ok(None),
265    };
266    if cursor.is_some_and(|id| !ordering.contains(id)) {
267        return Ok(None);
268    }
269    let mut winners = BTreeMap::new();
270    for line in lines {
271        let Some((transaction, id)) = snapshot_coordinate(line) else {
272            return Ok(None);
273        };
274        if !ordering.contains(transaction)
275            || winners
276                .last_key_value()
277                .is_some_and(|(previous, _)| previous >= &id)
278        {
279            return Ok(None);
280        }
281        let package = match load_from(root, &id) {
282            Ok(Some(package)) if package.id() == &id => package,
283            Ok(_) => return Ok(None),
284            Err(error) if error.starts_with("invalid materialized package: ") => return Ok(None),
285            Err(error) => return Err(error),
286        };
287        winners.insert(
288            id,
289            Winner {
290                transaction,
291                package,
292            },
293        );
294    }
295    if cursor.is_none() && !winners.is_empty() {
296        return Ok(None);
297    }
298    Ok(Some((cursor, winners)))
299}
300
301fn canonical_txid(text: &str) -> Option<TxId> {
302    let id = TxId::from_str(text).ok()?;
303    (id.to_string() == text).then_some(id)
304}
305
306fn snapshot_coordinate(line: &str) -> Option<(TxId, LibraryId)> {
307    let mut fields = line.split('|');
308    let transaction = canonical_txid(fields.next()?)?;
309    let authority_text = fields.next()?;
310    let name = fields.next()?;
311    let version_text = fields.next()?;
312    (fields.next().is_none() && authority_text.len() == 24).then_some(())?;
313    let authority = AuthorityTxId::from_str(authority_text).ok()?;
314    (authority.to_string() == authority_text).then_some(())?;
315    let family = LibraryFamily::new(AuthorityId::new(authority), name.to_owned()).ok()?;
316    let version = Version::from_str(version_text).ok()?;
317    (family.logical_name() == name && version.to_string() == version_text).then_some(())?;
318    Some((transaction, LibraryId::new(family, version).ok()?))
319}
320
321fn write_snapshot(state: &State) -> Result<(), String> {
322    let path = state.control.join(SNAPSHOT);
323    let temporary = state
324        .control
325        .join(format!("{SNAPSHOT}.{}.tmp", std::process::id()));
326    remove_stale(&temporary, false)?;
327    let mut text = format!(
328        "K1RUSTPROJECTION1\n{}\n",
329        state.cursor.map_or_else(|| "-".into(), |id| id.to_string())
330    );
331    for winner in state.winners.values() {
332        let id = winner.package.id();
333        text.push_str(&winner.transaction.to_string());
334        text.push('|');
335        text.push_str(&id.family().authority().to_string());
336        text.push('|');
337        text.push_str(id.family().logical_name());
338        text.push('|');
339        text.push_str(&id.version().to_string());
340        text.push('\n');
341    }
342    fs::write(&temporary, text).map_err(|error| error.to_string())?;
343    fs::rename(temporary, path).map_err(|error| error.to_string())
344}
345
346fn promote(state: &Arc<Mutex<State>>, projection_root: &Path) -> Result<(), String> {
347    let mut state = lock(state)?;
348    let backup = projection_root.with_extension(format!("previous.{}", std::process::id()));
349    remove_stale(&backup, true)?;
350    let replaced = match entry_metadata(projection_root)? {
351        Some(metadata) if metadata.file_type().is_dir() => true,
352        Some(_) => return Err("projection root is not a real directory".into()),
353        None => false,
354    };
355    if replaced {
356        fs::rename(projection_root, &backup).map_err(|error| error.to_string())?;
357    }
358    fs::rename(&state.root, projection_root).map_err(|error| error.to_string())?;
359    state.root = projection_root.to_path_buf();
360    write_snapshot(&state)?;
361    if replaced {
362        remove_stale(&backup, true)?;
363    }
364    Ok(())
365}
366
367fn staging_root(root: &Path) -> Result<PathBuf, String> {
368    let parent = root
369        .parent()
370        .ok_or_else(|| "projection root needs a parent".to_owned())?;
371    directory(parent)?;
372    let stage = root.with_extension(format!("staging.{}", std::process::id()));
373    remove_stale(&stage, true)?;
374    fs::create_dir(&stage).map_err(|error| error.to_string())?;
375    Ok(stage)
376}
377
378fn materialize(root: &Path, package: &SourcePackage) -> Result<(), String> {
379    let parent = coordinate_directory(root, package.id(), true)?
380        .ok_or_else(|| "package coordinate could not be created".to_owned())?;
381    let target = parent.join(package.id().version().to_string());
382    if entry_metadata(&target)?.is_some() {
383        return Err("package target already exists".into());
384    }
385    let stage = parent.join(format!(
386        ".{}.stage.{}",
387        package.id().version(),
388        std::process::id()
389    ));
390    remove_stale(&stage, true)?;
391    fs::create_dir(&stage).map_err(|error| error.to_string())?;
392    for source in package.files() {
393        let path = stage.join(source.path());
394        directory(
395            path.parent()
396                .ok_or_else(|| "source file has no parent".to_owned())?,
397        )?;
398        fs::write(path, source.bytes()).map_err(|error| error.to_string())?;
399    }
400    fs::rename(stage, target).map_err(|error| error.to_string())
401}
402
403fn load_from(root: &Path, id: &LibraryId) -> Result<Option<SourcePackage>, String> {
404    let Some(parent) = coordinate_directory(root, id, false)? else {
405        return Ok(None);
406    };
407    let target = parent.join(id.version().to_string());
408    match entry_metadata(&target)? {
409        None => return Ok(None),
410        Some(metadata) if !metadata.file_type().is_dir() => {
411            return Err("package is not a real directory".into());
412        }
413        Some(_) => {}
414    }
415    let mut files = Vec::new();
416    collect(&target, &target, &mut files)?;
417    SourcePackage::new(id.clone(), files)
418        .map(Some)
419        .map_err(|error| format!("invalid materialized package: {error}"))
420}
421
422fn coordinate_directory(
423    root: &Path,
424    id: &LibraryId,
425    create: bool,
426) -> Result<Option<PathBuf>, String> {
427    let family = id.family();
428    let authority = family.authority().to_string();
429    let mut path = root.to_path_buf();
430    for component in [authority.as_str(), family.logical_name()] {
431        path.push(component);
432        match entry_metadata(&path)? {
433            Some(metadata) if metadata.file_type().is_dir() => {}
434            Some(_) => return Err(format!("{} is not a real directory", path.display())),
435            None if create => fs::create_dir(&path).map_err(|error| error.to_string())?,
436            None => return Ok(None),
437        }
438    }
439    Ok(Some(path))
440}
441
442fn collect(root: &Path, directory_path: &Path, files: &mut Vec<SourceFile>) -> Result<(), String> {
443    for entry in fs::read_dir(directory_path).map_err(|error| error.to_string())? {
444        let path = entry.map_err(|error| error.to_string())?.path();
445        let metadata = fs::symlink_metadata(&path).map_err(|error| error.to_string())?;
446        if metadata.file_type().is_dir() {
447            collect(root, &path, files)?;
448        } else if metadata.file_type().is_file() {
449            let relative = path
450                .strip_prefix(root)
451                .map_err(|error| error.to_string())?
452                .to_str()
453                .ok_or("source path is not valid UTF-8")?
454                .replace('\\', "/");
455            let bytes = fs::read(path).map_err(|error| error.to_string())?;
456            files.push(SourceFile::new(relative, bytes));
457        } else {
458            return Err("source contains a symlink or non-file entry".into());
459        }
460    }
461    Ok(())
462}
463
464fn entry_metadata(path: &Path) -> Result<Option<fs::Metadata>, String> {
465    match fs::symlink_metadata(path) {
466        Ok(metadata) => Ok(Some(metadata)),
467        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
468        Err(error) => Err(error.to_string()),
469    }
470}
471
472fn remove_stale(path: &Path, directory: bool) -> Result<(), String> {
473    match (entry_metadata(path)?, directory) {
474        (None, _) => Ok(()),
475        (Some(metadata), true) if metadata.file_type().is_dir() => {
476            fs::remove_dir_all(path).map_err(|error| error.to_string())
477        }
478        (Some(metadata), false) if metadata.file_type().is_file() => {
479            fs::remove_file(path).map_err(|error| error.to_string())
480        }
481        _ => Err(format!(
482            "{} is not an expected temporary entry",
483            path.display()
484        )),
485    }
486}
487
488fn directory(path: &Path) -> Result<(), String> {
489    match entry_metadata(path)? {
490        Some(metadata) if metadata.file_type().is_dir() => Ok(()),
491        Some(_) => Err(format!("{} is not a real directory", path.display())),
492        None => fs::create_dir_all(path).map_err(|error| error.to_string()),
493    }
494}