kcode_k1_web_code_workspace/
lib.rs1#![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}