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}