Skip to main content

kcode_k1_web_code_workspace/
lib.rs

1#![forbid(unsafe_code)]
2
3use getrandom::fill;
4use kcode_k1_peering::K1Peering;
5pub use kcode_k1_transaction_id::TxId;
6use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem, SubsystemId, TxId as KtoTxId};
7pub use kcode_k1_web_code_document::CodeDocument;
8pub use kcode_k1_web_code_workspace_format::Workspace;
9use kcode_k1_web_code_workspace_format::{Action, OperationId, Projection};
10pub use kcode_k1_web_package::WebFamily;
11use std::{
12    collections::{HashMap, hash_map::Entry},
13    sync::{Arc, Mutex, MutexGuard},
14    time::{Duration, Instant},
15};
16
17type Evidence = (KtoTxId, Result<Workspace, String>);
18
19struct Pending {
20    action: Action,
21    evidence: Option<Evidence>,
22}
23
24struct State {
25    projection: Projection,
26    pending: HashMap<OperationId, Pending>,
27    fault: Option<String>,
28}
29
30impl State {
31    fn ready(&self) -> Result<(), String> {
32        self.fault.clone().map_or(Ok(()), Err)
33    }
34
35    fn poison(&mut self, message: impl Into<String>) -> String {
36        self.fault.get_or_insert_with(|| message.into()).clone()
37    }
38
39    fn record(&mut self, action: &Action, evidence: Evidence) -> Result<(), String> {
40        let operation = action.operation_id();
41        let problem = match self.pending.get_mut(&operation) {
42            None => None,
43            Some(pending) if pending.action != *action => {
44                Some("workspace callback action mismatch")
45            }
46            Some(pending) if pending.evidence.is_some() => {
47                Some("duplicate workspace callback evidence")
48            }
49            Some(pending) => {
50                pending.evidence = Some(evidence);
51                None
52            }
53        };
54        problem.map_or(Ok(()), |message| Err(self.poison(message)))
55    }
56}
57
58struct Core(Mutex<State>);
59
60impl Core {
61    fn lock(&self) -> Result<MutexGuard<'_, State>, String> {
62        self.0.lock().map_err(|_| "workspace lock poisoned".into())
63    }
64
65    fn apply(&self, transaction: KtoTxId, payload: &[u8]) -> Result<(), String> {
66        let action = Action::decode(payload)?;
67        let workspace_transaction = TxId::from_bytes(*transaction.as_bytes());
68        let mut state = self.lock()?;
69        state.ready()?;
70        let outcome = match state
71            .projection
72            .apply(workspace_transaction, action.clone())
73        {
74            Ok(value) => value,
75            Err(error) => {
76                let error = state.poison(format!("workspace projection failed: {error}"));
77                return Err(error);
78            }
79        };
80        let evidence = (
81            transaction,
82            outcome.result().cloned().map_err(str::to_owned),
83        );
84        state.record(&action, evidence)
85    }
86}
87
88impl Subsystem for Core {
89    fn submit_txn(&self, id: KtoTxId, payload: &[u8]) -> Result<(), String> {
90        self.apply(id, payload)
91    }
92
93    fn reorg(&self) -> Result<(), String> {
94        let mut state = self.lock()?;
95        state.pending.clear();
96        Err(state.poison("Web code workspaces unavailable after reorganization"))
97    }
98}
99
100#[derive(Clone)]
101pub struct K1WebCodeWorkspace {
102    core: Arc<Core>,
103    peering: Arc<K1Peering>,
104    subsystem: SubsystemId,
105}
106
107impl K1WebCodeWorkspace {
108    pub fn open(ordering: Arc<K1TxnOrdering>, peering: Arc<K1Peering>) -> Result<Self, String> {
109        let began = Instant::now();
110        let core = Arc::new(Core(Mutex::new(State {
111            projection: Projection::new(),
112            pending: HashMap::new(),
113            fault: None,
114        })));
115        let subsystem = SubsystemId::from_str("k1-web-ws")?;
116        ordering.register_subsystem(subsystem, None, core.clone())?;
117        core.lock()?.ready()?;
118        if began.elapsed() > Duration::from_millis(100) {
119            eprintln!("{{\"level\":\"warning\",\"event\":\"k1_web_code_workspace_open_slow\"}}");
120        }
121        Ok(Self {
122            core,
123            peering,
124            subsystem,
125        })
126    }
127
128    pub fn create(&self, document: CodeDocument) -> Result<Workspace, String> {
129        let action = self.reserve(|operation| Action::create(operation, document.clone()))?;
130        self.submit(action)
131    }
132
133    pub fn branch(&self, document: CodeDocument) -> Result<Workspace, String> {
134        let action = self.reserve(|operation| Action::branch(operation, document.clone()))?;
135        self.submit(action)
136    }
137
138    pub fn overwrite(
139        &self,
140        workspace: TxId,
141        expected: TxId,
142        document: CodeDocument,
143    ) -> Result<Workspace, String> {
144        let action = self.reserve(|operation| {
145            Action::overwrite(operation, workspace, expected, document.clone())
146        })?;
147        self.submit(action)
148    }
149
150    pub fn get(&self, workspace: TxId) -> Result<Option<Workspace>, String> {
151        let state = self.core.lock()?;
152        state.ready()?;
153        Ok(state.projection.get(workspace))
154    }
155
156    pub fn latest(&self, family: &WebFamily) -> Result<Option<Workspace>, String> {
157        let state = self.core.lock()?;
158        state.ready()?;
159        Ok(state.projection.latest(family))
160    }
161
162    fn reserve(&self, mut build: impl FnMut(OperationId) -> Action) -> Result<Action, String> {
163        loop {
164            let mut operation = [0; 16];
165            fill(&mut operation).map_err(|error| error.to_string())?;
166            let action = build(operation);
167            let mut state = self.core.lock()?;
168            state.ready()?;
169            if let Entry::Vacant(slot) = state.pending.entry(operation) {
170                slot.insert(Pending {
171                    action: action.clone(),
172                    evidence: None,
173                });
174                return Ok(action);
175            }
176        }
177    }
178
179    fn submit(&self, action: Action) -> Result<Workspace, String> {
180        let operation = action.operation_id();
181        let submitted = self.peering.submit_txn(self.subsystem, &action.encode());
182        let mut state = self.core.lock()?;
183        let Some(pending) = state.pending.remove(&operation) else {
184            return Err(state.poison("missing workspace pending operation"));
185        };
186        let result = correlate(submitted, pending.evidence).map_err(|error| state.poison(error))?;
187        state.ready()?;
188        result
189    }
190}
191
192fn correlate(
193    submitted: Result<KtoTxId, String>,
194    evidence: Option<Evidence>,
195) -> Result<Result<Workspace, String>, &'static str> {
196    match (submitted, evidence) {
197        (Ok(id), Some((seen, outcome))) if id == seen => Ok(outcome),
198        (Ok(_), Some(_)) => Err("workspace callback transaction mismatch"),
199        (Ok(_), None) => Err("missing workspace callback evidence"),
200        (Err(_), Some((_, outcome))) => Ok(outcome),
201        (Err(error), None) => Ok(Err(error)),
202    }
203}
204
205#[cfg(test)]
206mod tests {
207    use super::*;
208    use kcode_k1_web_code_workspace_testkit::{Candidate, run_tests};
209
210    struct Subject;
211
212    impl Candidate for Subject {
213        type Driver = K1WebCodeWorkspace;
214        fn open(
215            ordering: Arc<K1TxnOrdering>,
216            peering: Arc<K1Peering>,
217        ) -> Result<Self::Driver, String> {
218            K1WebCodeWorkspace::open(ordering, peering)
219        }
220        fn create(driver: &Self::Driver, document: CodeDocument) -> Result<Workspace, String> {
221            driver.create(document)
222        }
223        fn branch(driver: &Self::Driver, document: CodeDocument) -> Result<Workspace, String> {
224            driver.branch(document)
225        }
226        fn overwrite(
227            driver: &Self::Driver,
228            workspace: TxId,
229            expected: TxId,
230            document: CodeDocument,
231        ) -> Result<Workspace, String> {
232            driver.overwrite(workspace, expected, document)
233        }
234        fn get(driver: &Self::Driver, workspace: TxId) -> Result<Option<Workspace>, String> {
235            driver.get(workspace)
236        }
237        fn latest(driver: &Self::Driver, family: &WebFamily) -> Result<Option<Workspace>, String> {
238            driver.latest(family)
239        }
240    }
241
242    #[test]
243    fn conformance() {
244        run_tests::<Subject>();
245    }
246}