use std::{
collections::HashMap,
fs::{File, OpenOptions},
io::{Read, Seek, SeekFrom, Write},
path::Path,
sync::{Arc, Mutex},
};
pub use kcode_k1_launch_node_codec::{
Authority, GroupId, NodeId, TargetId, TargetName, TxId, UserId,
};
use kcode_k1_launch_node_codec::{
PROJECTION_HEADER, ProjectionDecode, ProjectionRecord, SetAction, decode_projection_record,
};
use kcode_k1_peering::K1Peering;
use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem as OrderingSubsystem, SubsystemId};
struct Pending {
payload: Vec<u8>,
evidence: Option<TxId>,
}
struct Inner {
file: File,
bindings: HashMap<TargetId, NodeId>,
checkpoint: Option<TxId>,
pending: HashMap<[u8; 16], Pending>,
failed: bool,
}
struct Shared {
inner: Mutex<Inner>,
}
struct Driver {
shared: Arc<Shared>,
}
pub struct LaunchNodes {
shared: Arc<Shared>,
peering: Arc<K1Peering>,
}
impl LaunchNodes {
pub fn open(
path: &Path,
ordering: Arc<K1TxnOrdering>,
peering: Arc<K1Peering>,
) -> Result<Self, String> {
let mut projection = open_projection(path)?;
if let Some(checkpoint) = projection.checkpoint.to_owned() {
match ordering.get_txn(checkpoint) {
Ok(Some(_)) => {}
Ok(None) => {
reset_projection(&mut projection.file, path)?;
projection.bindings.clear();
projection.checkpoint = None;
}
Err(error) => return Err(format!("validate launch nodes checkpoint: {error}")),
}
}
let after = projection.checkpoint.to_owned();
let shared = Arc::new(Shared {
inner: Mutex::new(Inner {
file: projection.file,
bindings: projection.bindings,
checkpoint: projection.checkpoint,
pending: HashMap::new(),
failed: false,
}),
});
ordering.register_subsystem(
subsystem()?,
after,
Arc::new(Driver {
shared: shared.clone(),
}),
)?;
Ok(Self { shared, peering })
}
pub fn get(&self, target: &TargetId) -> Result<Option<NodeId>, String> {
let inner = self
.shared
.inner
.lock()
.map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
if inner.failed {
return Err("k1-launch-nodes facade is unavailable".into());
}
Ok(inner.bindings.get(target).cloned())
}
pub fn set(&self, target: TargetId, node: NodeId) -> Result<(), String> {
let (operation_id, payload) = loop {
let mut operation_id = [0_u8; 16];
getrandom::fill(&mut operation_id)
.map_err(|error| format!("generate launch nodes operation ID: {error}"))?;
let payload = SetAction::new(operation_id, target.clone(), node).encode();
let mut inner = self
.shared
.inner
.lock()
.map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
if inner.failed {
return Err("k1-launch-nodes facade is unavailable".into());
}
if inner.pending.contains_key(&operation_id) {
continue;
}
inner.pending.insert(
operation_id,
Pending {
payload: payload.clone(),
evidence: None,
},
);
break (operation_id, payload);
};
let submission = self.peering.submit_txn(subsystem()?, &payload);
let mut inner = self
.shared
.inner
.lock()
.map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
if inner.failed {
return Err("k1-launch-nodes facade is unavailable".into());
}
let Some(pending) = inner.pending.remove(&operation_id) else {
return Err(fault(
&mut inner,
"k1-launch-nodes callback evidence is missing",
));
};
match (submission, pending.evidence) {
(Ok(submitted), Some(callback)) if submitted == callback => Ok(()),
(Err(_), Some(_)) => Ok(()),
(Err(error), None) => Err(format!("submit k1-launch-nodes transaction: {error}")),
(Ok(_), None) => Err(fault(
&mut inner,
"k1-launch-nodes callback evidence is missing",
)),
(Ok(_), Some(_)) => Err(fault(
&mut inner,
"k1-launch-nodes callback transaction does not match submission",
)),
}
}
}
impl OrderingSubsystem for Driver {
fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
let action = match SetAction::decode(payload) {
Ok(action) => action,
Err(error) => {
return self
.shared
.fail(format!("decode k1-launch-nodes Set action: {error}"));
}
};
let operation_id = action.operation_id().to_owned();
let canonical_payload = action.encode();
let record =
ProjectionRecord::new(id, action.target().to_owned(), action.node().to_owned())
.encode();
let mut inner = self
.shared
.inner
.lock()
.map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
if inner.failed {
return Err("k1-launch-nodes facade is unavailable".into());
}
if inner.checkpoint == Some(id) {
return Err(fault(
&mut inner,
"k1-launch-nodes callback transaction is duplicated",
));
}
if let Some(pending) = inner.pending.get(&operation_id) {
if pending.payload != canonical_payload {
return Err(fault(
&mut inner,
"k1-launch-nodes callback action does not match reservation",
));
}
if pending.evidence.is_some() {
return Err(fault(
&mut inner,
"k1-launch-nodes callback evidence is duplicated",
));
}
}
let append = inner
.file
.seek(SeekFrom::End(0))
.and_then(|_| inner.file.write_all(&record))
.and_then(|()| inner.file.sync_data());
if let Err(error) = append {
return Err(fault(
&mut inner,
format!("append k1-launch-nodes projection record: {error}"),
));
}
inner
.bindings
.insert(action.target().to_owned(), action.node().to_owned());
inner.checkpoint = Some(id);
if let Some(pending) = inner.pending.get_mut(&operation_id) {
pending.evidence = Some(id);
}
Ok(())
}
fn reorg(&self) -> Result<(), String> {
let mut inner = self
.shared
.inner
.lock()
.map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
inner.failed = true;
inner.pending.clear();
Ok(())
}
}
impl Shared {
fn fail(&self, message: String) -> Result<(), String> {
let mut inner = self
.inner
.lock()
.map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
Err(fault(&mut inner, message))
}
}
struct Projection {
file: File,
bindings: HashMap<TargetId, NodeId>,
checkpoint: Option<TxId>,
}
fn open_projection(path: &Path) -> Result<Projection, String> {
let mut file = OpenOptions::new()
.create(true)
.truncate(false)
.read(true)
.write(true)
.open(path)
.map_err(|error| format!("open launch nodes projection: {error}"))?;
if file
.metadata()
.map_err(|error| format!("inspect launch nodes projection: {error}"))?
.len()
== 0
{
reset_projection(&mut file, path)?;
}
file.seek(SeekFrom::Start(0))
.map_err(|error| format!("seek launch nodes projection: {error}"))?;
let mut bytes = Vec::new();
file.read_to_end(&mut bytes)
.map_err(|error| format!("read launch nodes projection: {error}"))?;
if bytes.len() < PROJECTION_HEADER.len()
|| &bytes[..PROJECTION_HEADER.len()] != PROJECTION_HEADER
{
return Err("unsupported launch nodes projection format; no migration is available".into());
}
let mut bindings = HashMap::new();
let mut checkpoint = None;
let mut offset = PROJECTION_HEADER.len();
let mut corrupt = false;
while offset < bytes.len() {
match decode_projection_record(&bytes[offset..]) {
Ok(ProjectionDecode::Complete { record, consumed }) => {
bindings.insert(record.target().to_owned(), record.node().to_owned());
let callback: [u8; 12] = bytes[offset..offset + 12]
.try_into()
.map_err(|_| "launch nodes projection callback is unavailable".to_owned())?;
checkpoint = Some(TxId::from_bytes(callback));
offset += consumed;
}
Ok(ProjectionDecode::Incomplete) => {
repair_tail(&mut file, offset)?;
break;
}
Err(_) => {
corrupt = true;
break;
}
}
}
if corrupt {
reset_projection(&mut file, path)?;
bindings.clear();
checkpoint = None;
}
file.seek(SeekFrom::End(0))
.map_err(|error| format!("seek launch nodes append position: {error}"))?;
Ok(Projection {
file,
bindings,
checkpoint,
})
}
fn reset_projection(file: &mut File, path: &Path) -> Result<(), String> {
file.set_len(0)
.and_then(|()| file.seek(SeekFrom::Start(0)).map(|_| ()))
.and_then(|()| file.write_all(PROJECTION_HEADER))
.and_then(|()| file.sync_data())
.map_err(|error| format!("reset launch nodes projection: {error}"))?;
sync_parent(path)
}
fn repair_tail(file: &mut File, offset: usize) -> Result<(), String> {
file.set_len(offset as u64)
.and_then(|()| file.sync_data())
.map_err(|error| format!("repair launch nodes projection tail: {error}"))
}
fn sync_parent(path: &Path) -> Result<(), String> {
let parent = path
.parent()
.filter(|parent| !parent.as_os_str().is_empty())
.unwrap_or_else(|| Path::new("."));
File::open(parent)
.and_then(|directory| directory.sync_all())
.map_err(|error| format!("synchronize launch nodes projection directory: {error}"))
}
fn subsystem() -> Result<SubsystemId, String> {
SubsystemId::from_str("k1-launch-nodes")
}
fn fault(inner: &mut Inner, message: impl Into<String>) -> String {
inner.failed = true;
inner.pending.clear();
message.into()
}
#[cfg(test)]
mod tests {
use super::*;
use std::{
fs,
path::{Path, PathBuf},
sync::atomic::{AtomicU64, Ordering},
};
static NEXT: AtomicU64 = AtomicU64::new(0);
struct Temp {
root: PathBuf,
}
impl Temp {
fn new() -> Self {
let root = std::env::temp_dir().join(format!(
"k1-launch-nodes-{}-{}",
std::process::id(),
NEXT.fetch_add(1, Ordering::Relaxed)
));
fs::create_dir(&root).unwrap();
Self { root }
}
}
impl Drop for Temp {
fn drop(&mut self) {
fs::remove_dir_all(&self.root).unwrap();
}
}
fn services(root: &Path) -> (Arc<K1TxnOrdering>, Arc<K1Peering>) {
fs::create_dir_all(root).unwrap();
let ordering = Arc::new(K1TxnOrdering::open(&root.join("ordering")).unwrap());
let peering = Arc::new(K1Peering::open(&root.join("peering"), ordering.clone()).unwrap());
(ordering, peering)
}
fn target(authority: Authority, name: &str) -> TargetId {
TargetId::new(authority, TargetName::new(name.to_owned()).unwrap())
}
fn user(byte: u8) -> Authority {
Authority::User(UserId::from_tx_id(TxId::from_bytes([byte; 12])))
}
fn group(byte: u8) -> Authority {
Authority::Group(GroupId::new(TxId::from_bytes([byte; 12])))
}
fn record_count(path: &Path) -> usize {
let bytes = fs::read(path).unwrap();
let mut offset = PROJECTION_HEADER.len();
let mut count = 0;
while offset < bytes.len() {
match decode_projection_record(&bytes[offset..]).unwrap() {
ProjectionDecode::Complete { consumed, .. } => {
count += 1;
offset += consumed;
}
ProjectionDecode::Incomplete => panic!("unexpected incomplete record"),
}
}
count
}
#[test]
fn callbacks_preserve_authorities_and_last_canonical_write() {
let temp = Temp::new();
let (ordering, peering) = services(&temp.root.join("service"));
let store = LaunchNodes::open(&temp.root.join("projection"), ordering, peering).unwrap();
let user_target = target(user(1), "default/chat");
let group_target = target(group(1), "default/chat");
store.set(user_target.clone(), NodeId([1; 12])).unwrap();
store.set(group_target.clone(), NodeId([2; 12])).unwrap();
store.set(user_target.clone(), NodeId([3; 12])).unwrap();
assert_eq!(store.get(&user_target).unwrap(), Some(NodeId([3; 12])));
assert_eq!(store.get(&group_target).unwrap(), Some(NodeId([2; 12])));
}
#[test]
fn peering_error_without_callback_does_not_bind() {
let temp = Temp::new();
let ordering = Arc::new(K1TxnOrdering::open(&temp.root.join("ordering-a")).unwrap());
let other = Arc::new(K1TxnOrdering::open(&temp.root.join("ordering-b")).unwrap());
let peering = Arc::new(K1Peering::open(&temp.root.join("peering-b"), other).unwrap());
let store = LaunchNodes::open(&temp.root.join("projection"), ordering, peering).unwrap();
let key = target(user(2), "model/x");
assert!(store.set(key.clone(), NodeId([4; 12])).is_err());
assert_eq!(store.get(&key).unwrap(), None);
assert_eq!(record_count(&temp.root.join("projection")), 0);
}
#[test]
fn projection_repairs_rebuilds_and_replays_same_value_sets() {
let temp = Temp::new();
let service = temp.root.join("service");
let projection = temp.root.join("projection");
let key = target(group(3), "same/value");
let node = NodeId([5; 12]);
let (ordering, peering) = services(&service);
let store = LaunchNodes::open(&projection, ordering.clone(), peering.clone()).unwrap();
store.set(key.clone(), node).unwrap();
store.set(key.clone(), node).unwrap();
drop((store, peering, ordering));
let clean = fs::read(&projection).unwrap();
let mut incomplete = clean.clone();
incomplete.extend_from_slice(&[1, 2, 3]);
fs::write(&projection, incomplete).unwrap();
let (ordering, peering) = services(&service);
let store = LaunchNodes::open(&projection, ordering.clone(), peering.clone()).unwrap();
assert_eq!(store.get(&key).unwrap(), Some(node));
assert_eq!(fs::read(&projection).unwrap(), clean);
drop((store, peering, ordering));
let mut corrupt = fs::read(&projection).unwrap();
let last = corrupt.last_mut().unwrap();
*last ^= 0xff;
fs::write(&projection, corrupt).unwrap();
let (ordering, peering) = services(&service);
let store = LaunchNodes::open(&projection, ordering, peering).unwrap();
assert_eq!(store.get(&key).unwrap(), Some(node));
assert_eq!(record_count(&projection), 2);
}
#[test]
fn non_v3_projection_is_rejected_without_migration() {
let temp = Temp::new();
let projection = temp.root.join("projection");
fs::write(&projection, b"K1LNV2\0\0").unwrap();
let (ordering, peering) = services(&temp.root.join("service"));
let error = LaunchNodes::open(&projection, ordering, peering)
.err()
.unwrap();
assert_eq!(
error,
"unsupported launch nodes projection format; no migration is available"
);
assert_eq!(fs::read(projection).unwrap(), b"K1LNV2\0\0");
}
#[test]
fn malformed_callback_faults_facade() {
let temp = Temp::new();
let (ordering, peering) = services(&temp.root.join("malformed"));
let store =
LaunchNodes::open(&temp.root.join("malformed-projection"), ordering, peering).unwrap();
let driver = Driver {
shared: store.shared.clone(),
};
assert!(
driver
.submit_txn(TxId::from_bytes([8; 12]), b"malformed")
.is_err()
);
assert!(store.get(&target(user(8), "faulted")).is_err());
}
#[test]
fn complete_package_stays_below_the_managed_limit() {
let files = [
include_str!("../Cargo.toml"),
include_str!("../Documentation.md"),
include_str!("lib.rs"),
];
let count = files
.iter()
.flat_map(|file| file.lines())
.filter(|line| !line.trim().is_empty())
.count();
assert!(count < 500, "complete package has {count} nonblank lines");
}
}