1use kcode_k1_transaction_id::TxId;
2use kcode_k1_web_package::{DependencySelector, SourceFile, SourcePackage, WebFamily, WebId};
3use serde::{Deserialize, Serialize};
4use sha2::{Digest as _, Sha256};
5use std::fs;
6use std::path::{Path, PathBuf};
7use std::time::{Duration, Instant};
8
9mod storage;
10
11#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
12pub struct Digest(pub [u8; 32]);
13#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
14pub struct SourceIdentity {
15 pub generation: u64,
16 pub digest: Digest,
17}
18#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
19pub struct ResolutionIdentity {
20 pub generation: u64,
21 pub digest: Digest,
22}
23#[derive(Clone, Debug, Eq, PartialEq)]
24pub struct ResolutionEntry {
25 pub family: WebFamily,
26 pub selector: DependencySelector,
27 pub resolved: WebId,
28 pub winning: TxId,
29}
30#[derive(Clone, Debug, Eq, PartialEq)]
31pub struct ResolutionManifest {
32 candidate: WebId,
33 source: SourceIdentity,
34 projection_cursor: Option<TxId>,
35 entries: Vec<ResolutionEntry>,
36}
37
38impl ResolutionManifest {
39 pub fn new(
40 candidate: &SourcePackage,
41 source: SourceIdentity,
42 projection_cursor: Option<TxId>,
43 mut entries: Vec<ResolutionEntry>,
44 ) -> Result<Self, String> {
45 entries.sort_by(|a, b| (&a.family, &a.selector).cmp(&(&b.family, &b.selector)));
46 let value = Self {
47 candidate: candidate.id().clone(),
48 source,
49 projection_cursor,
50 entries,
51 };
52 value.validate(candidate)?;
53 Ok(value)
54 }
55 pub fn entries(&self) -> &[ResolutionEntry] {
56 &self.entries
57 }
58 pub fn route(
59 &self,
60 family: &WebFamily,
61 selector: &DependencySelector,
62 ) -> Option<&ResolutionEntry> {
63 self.entries
64 .iter()
65 .find(|entry| &entry.family == family && &entry.selector == selector)
66 }
67 fn validate(&self, package: &SourcePackage) -> Result<(), String> {
68 if self.source.generation == 0
69 || self.candidate != *package.id()
70 || self.source.digest != package_digest(package)
71 {
72 return Err("manifest source identity does not match candidate".into());
73 }
74 let mut expected = package
75 .dependencies()
76 .iter()
77 .map(|dependency| {
78 WebFamily::new(dependency.authority(), dependency.name().to_owned())
79 .map(|family| (family, dependency.selector().clone()))
80 .map_err(|e| e.to_string())
81 })
82 .collect::<Result<Vec<_>, _>>()?;
83 expected.sort();
84 let actual = self
85 .entries
86 .iter()
87 .map(|entry| (entry.family.clone(), entry.selector.clone()))
88 .collect::<Vec<_>>();
89 let invalid = self.entries.iter().any(|entry| {
90 entry.resolved.family() != &entry.family
91 || !entry.selector.matches(entry.resolved.version())
92 });
93 if invalid || actual != expected {
94 return Err("resolution entries are invalid".into());
95 }
96 Ok(())
97 }
98}
99
100#[derive(Clone, Debug, Eq, PartialEq)]
101pub struct CheckIdentity {
102 pub cache_epoch: u64,
103 pub source: SourceIdentity,
104 pub resolution: ResolutionIdentity,
105 pub graph: Digest,
106 pub projection_cursor: Option<TxId>,
107 pub checker_executable: Digest,
108 pub container_image: Digest,
109 pub chromium_version: String,
110 pub route_revision: String,
111 pub harness_revision: String,
112 pub check_policy_revision: String,
113 pub command_policy: String,
114}
115
116pub struct UserWebCache {
117 root: PathBuf,
118 epoch: u64,
119 identity: Vec<u8>,
120}
121#[derive(Deserialize, Serialize)]
122#[serde(deny_unknown_fields)]
123struct Saved {
124 epoch: u64,
125 source_counter: u64,
126 source: Option<(u64, [u8; 32], [u8; 32])>,
127 resolution_counter: u64,
128 resolution: Option<(u64, [u8; 32], Option<[u8; 12]>)>,
129}
130struct OpenTimer(Instant);
131impl Drop for OpenTimer {
132 fn drop(&mut self) {
133 if self.0.elapsed() > Duration::from_millis(100) {
134 eprintln!("{{\"level\":\"warning\",\"event\":\"k1_web_cache_open_slow\"}}");
135 }
136 }
137}
138
139impl UserWebCache {
140 pub fn open(
141 cache_root: impl AsRef<Path>,
142 user: TxId,
143 boot: &str,
144 schema: &str,
145 checker_policy: &str,
146 ) -> Result<Self, String> {
147 let _timer = OpenTimer(Instant::now());
148 nonempty(&[boot, schema, checker_policy])?;
149 storage::directory(cache_root.as_ref())?;
150 let root = cache_root.as_ref().join(user.to_string());
151 let identity = cache_identity(boot, schema, checker_policy);
152 if let Ok((saved, _)) = load(&root, &identity, true) {
153 return Ok(Self {
154 root,
155 epoch: saved.epoch,
156 identity,
157 });
158 }
159 let epoch = recover_epoch(&root)
160 .unwrap_or(0)
161 .checked_add(1)
162 .ok_or("cache epoch exhausted")?;
163 initialize(&root, epoch, &identity)?;
164 Ok(Self {
165 root,
166 epoch,
167 identity,
168 })
169 }
170 pub fn reset(&mut self, boot: &str, schema: &str, checker_policy: &str) -> Result<(), String> {
171 nonempty(&[boot, schema, checker_policy])?;
172 let epoch = self.epoch.checked_add(1).ok_or("cache epoch exhausted")?;
173 let identity = cache_identity(boot, schema, checker_policy);
174 initialize(&self.root, epoch, &identity)?;
175 (self.epoch, self.identity) = (epoch, identity);
176 Ok(())
177 }
178 pub fn root(&self) -> &Path {
179 &self.root
180 }
181 pub fn source_root(&self) -> PathBuf {
182 self.root.join("source")
183 }
184 pub fn resolution_root(&self) -> PathBuf {
185 self.root.join("resolution")
186 }
187 pub fn diagnostics_root(&self) -> PathBuf {
188 self.root.join("diagnostics")
189 }
190 pub const fn epoch(&self) -> u64 {
191 self.epoch
192 }
193 pub fn source_identity(&self) -> Result<Option<SourceIdentity>, String> {
194 Ok(source_identity(&self.state()?))
195 }
196 pub fn resolution_identity(&self) -> Result<Option<ResolutionIdentity>, String> {
197 Ok(resolution_identity(&self.state()?))
198 }
199 pub fn materialize(&self, package: &SourcePackage) -> Result<SourceIdentity, String> {
200 let (mut saved, files) = load(&self.root, &self.identity, false)?;
201 self.same_epoch(&saved)?;
202 let source = package_digest(package);
203 let tree = files_digest(package.files());
204 if let Some(current) = source_identity(&saved)
205 && current.digest == source
206 && saved.source.is_some_and(|value| value.2 == tree.0)
207 && files_digest(&files) == tree
208 {
209 return Ok(current);
210 }
211 let generation = saved
212 .source_counter
213 .checked_add(1)
214 .ok_or("source generation exhausted")?;
215 let identity = SourceIdentity {
216 generation,
217 digest: source,
218 };
219 saved.source_counter = generation;
220 saved.source = Some((generation, source.0, tree.0));
221 saved.resolution = None;
222 self.mutate(|| {
223 storage::replace_source(&self.source_root(), package.files())?;
224 storage::replace_empty(&self.resolution_root())?;
225 storage::replace_empty(&self.diagnostics_root())?;
226 storage::atomic_write(&self.root.join("control/state"), &encode_saved(&saved))?;
227 storage::remove_file(&self.root.join("control/receipt"))
228 })?;
229 Ok(identity)
230 }
231 pub fn record_resolution(
232 &self,
233 manifest: &ResolutionManifest,
234 ) -> Result<ResolutionIdentity, String> {
235 let mut saved = self.state()?;
236 if source_identity(&saved) != Some(manifest.source) {
237 return Err("resolution source identity is stale".into());
238 }
239 let bytes = encode_manifest(manifest);
240 let digest = hash(&bytes);
241 if let Some(current) = resolution_identity(&saved)
242 && current.digest == digest
243 && resolution_cursor(&saved) == manifest.projection_cursor
244 {
245 return Ok(current);
246 }
247 let generation = saved
248 .resolution_counter
249 .checked_add(1)
250 .ok_or("resolution generation exhausted")?;
251 let identity = ResolutionIdentity { generation, digest };
252 saved.resolution_counter = generation;
253 saved.resolution = Some((
254 generation,
255 digest.0,
256 manifest.projection_cursor.map(|id| *id.as_bytes()),
257 ));
258 self.mutate(|| {
259 storage::replace_manifest(&self.resolution_root(), &bytes)?;
260 storage::replace_empty(&self.diagnostics_root())?;
261 storage::atomic_write(&self.root.join("control/state"), &encode_saved(&saved))?;
262 storage::remove_file(&self.root.join("control/receipt"))
263 })?;
264 Ok(identity)
265 }
266 pub fn has_check(&self, identity: &CheckIdentity) -> Result<bool, String> {
267 let saved = self.state()?;
268 let receipt = check_receipt(identity)?;
269 if !self.current(&saved, identity) {
270 return Ok(false);
271 }
272 Ok(storage::read(&self.root.join("control/receipt"))?.as_deref() == Some(&receipt.0))
273 }
274 pub fn record_check(&self, identity: &CheckIdentity) -> Result<(), String> {
275 let saved = self.state()?;
276 let receipt = check_receipt(identity)?;
277 if !self.current(&saved, identity) {
278 return Err("check identity is stale".into());
279 }
280 self.mutate(|| storage::atomic_write(&self.root.join("control/receipt"), &receipt.0))
281 }
282 fn state(&self) -> Result<Saved, String> {
283 let (saved, _) = load(&self.root, &self.identity, true)?;
284 self.same_epoch(&saved)?;
285 Ok(saved)
286 }
287 fn same_epoch(&self, saved: &Saved) -> Result<(), String> {
288 (saved.epoch == self.epoch)
289 .then_some(())
290 .ok_or_else(|| "cache epoch changed".into())
291 }
292 fn current(&self, saved: &Saved, identity: &CheckIdentity) -> bool {
293 identity.cache_epoch == self.epoch
294 && source_identity(saved) == Some(identity.source)
295 && resolution_identity(saved) == Some(identity.resolution)
296 && resolution_cursor(saved) == identity.projection_cursor
297 }
298 fn mutate<T>(&self, operation: impl FnOnce() -> Result<T, String>) -> Result<T, String> {
299 fs::OpenOptions::new()
300 .write(true)
301 .create_new(true)
302 .open(self.root.join("control/dirty"))
303 .map_err(|e| e.to_string())?;
304 let value = operation()?;
305 fs::remove_file(self.root.join("control/dirty")).map_err(|e| e.to_string())?;
306 Ok(value)
307 }
308}
309
310fn initialize(root: &Path, epoch: u64, identity: &[u8]) -> Result<(), String> {
311 let saved = Saved {
312 epoch,
313 source_counter: 0,
314 source: None,
315 resolution_counter: 0,
316 resolution: None,
317 };
318 storage::replace_directory(root, |stage| {
319 for name in ["control", "source", "resolution", "diagnostics"] {
320 fs::create_dir(stage.join(name)).map_err(|e| e.to_string())?;
321 }
322 fs::write(stage.join("control/identity"), identity).map_err(|e| e.to_string())?;
323 fs::write(stage.join("control/state"), encode_saved(&saved)).map_err(|e| e.to_string())
324 })
325}
326fn load(root: &Path, identity: &[u8], exact: bool) -> Result<(Saved, Vec<SourceFile>), String> {
327 let files = storage::inspect(root)?;
328 let control = root.join("control");
329 if storage::read(&control.join("dirty"))?.is_some() {
330 return Err("cache mutation is incomplete".into());
331 }
332 if storage::read(&control.join("identity"))?.as_deref() != Some(identity) {
333 return Err("cache identity mismatch".into());
334 }
335 let saved =
336 decode_saved(&storage::read(&control.join("state"))?.ok_or("missing cache state")?)?;
337 if exact
338 && !saved
339 .source
340 .map_or(files.is_empty(), |value| files_digest(&files).0 == value.2)
341 {
342 return Err("materialized source mismatch".into());
343 }
344 match (
345 saved.resolution,
346 storage::read(&root.join("resolution/manifest"))?,
347 ) {
348 (Some(value), Some(bytes)) if hash(&bytes).0 == value.1 => {}
349 (None, None) => {}
350 _ => return Err("resolution manifest mismatch".into()),
351 }
352 let invalid_receipt = storage::read(&control.join("receipt"))?
353 .is_some_and(|receipt| receipt.len() != 32 || saved.resolution.is_none());
354 if invalid_receipt {
355 return Err("invalid check receipt".into());
356 }
357 Ok((saved, files))
358}
359
360fn encode_manifest(value: &ResolutionManifest) -> Vec<u8> {
361 let id = &value.candidate;
362 let cursor = value
363 .projection_cursor
364 .map_or_else(|| "-".into(), |item| item.to_string());
365 let mut text = format!(
366 "K1WEBRESOLUTION1\n{}\n{}\n{}\n{}\n{}\n{}\n",
367 id.family().authority(),
368 id.family().logical_name(),
369 id.version(),
370 value.source.generation,
371 hex(value.source.digest),
372 cursor
373 );
374 for entry in &value.entries {
375 text.push_str(&format!(
376 "{}\t{}\t{}\t{}\t{}\n",
377 entry.family.authority(),
378 entry.family.logical_name(),
379 entry.selector,
380 entry.resolved.version(),
381 entry.winning
382 ));
383 }
384 text.into_bytes()
385}
386fn encode_saved(value: &Saved) -> Vec<u8> {
387 serde_json::to_vec(value).expect("saved state serializes")
388}
389fn decode_saved(bytes: &[u8]) -> Result<Saved, String> {
390 let value: Saved = serde_json::from_slice(bytes).map_err(|_| "invalid cache state")?;
391 if encode_saved(&value) != bytes
392 || value
393 .source
394 .is_some_and(|item| item.0 == 0 || item.0 > value.source_counter)
395 || value
396 .resolution
397 .is_some_and(|item| item.0 == 0 || item.0 > value.resolution_counter)
398 {
399 return Err("noncanonical cache state".into());
400 }
401 Ok(value)
402}
403fn source_identity(value: &Saved) -> Option<SourceIdentity> {
404 value.source.map(|item| SourceIdentity {
405 generation: item.0,
406 digest: Digest(item.1),
407 })
408}
409fn resolution_identity(value: &Saved) -> Option<ResolutionIdentity> {
410 value.resolution.map(|item| ResolutionIdentity {
411 generation: item.0,
412 digest: Digest(item.1),
413 })
414}
415fn resolution_cursor(value: &Saved) -> Option<TxId> {
416 value
417 .resolution
418 .and_then(|item| item.2.map(TxId::from_bytes))
419}
420fn package_digest(package: &SourcePackage) -> Digest {
421 let mut hasher = Sha256::new();
422 append(
423 &mut hasher,
424 package
425 .id()
426 .family()
427 .authority()
428 .transaction_id()
429 .as_bytes(),
430 );
431 append(&mut hasher, package.id().family().logical_name().as_bytes());
432 append(&mut hasher, package.id().version().to_string().as_bytes());
433 append_files(&mut hasher, package.files());
434 Digest(hasher.finalize().into())
435}
436fn files_digest(files: &[SourceFile]) -> Digest {
437 let mut hasher = Sha256::new();
438 append_files(&mut hasher, files);
439 Digest(hasher.finalize().into())
440}
441fn append_files(hasher: &mut Sha256, files: &[SourceFile]) {
442 for file in files {
443 append(hasher, file.path().as_bytes());
444 append(hasher, file.bytes());
445 }
446}
447fn check_receipt(value: &CheckIdentity) -> Result<Digest, String> {
448 nonempty(&[
449 &value.chromium_version,
450 &value.route_revision,
451 &value.harness_revision,
452 &value.check_policy_revision,
453 &value.command_policy,
454 ])?;
455 let mut hasher = Sha256::new();
456 for number in [
457 value.cache_epoch,
458 value.source.generation,
459 value.resolution.generation,
460 ] {
461 append(&mut hasher, &number.to_le_bytes());
462 }
463 for digest in [
464 value.source.digest,
465 value.resolution.digest,
466 value.graph,
467 value.checker_executable,
468 value.container_image,
469 ] {
470 append(&mut hasher, &digest.0);
471 }
472 match value.projection_cursor {
473 Some(cursor) => {
474 append(&mut hasher, &[1]);
475 append(&mut hasher, cursor.as_bytes());
476 }
477 None => append(&mut hasher, &[0]),
478 }
479 for text in [
480 &value.chromium_version,
481 &value.route_revision,
482 &value.harness_revision,
483 &value.check_policy_revision,
484 &value.command_policy,
485 ] {
486 append(&mut hasher, text.as_bytes());
487 }
488 Ok(Digest(hasher.finalize().into()))
489}
490fn recover_epoch(root: &Path) -> Option<u64> {
491 serde_json::from_slice::<Saved>(&storage::read(&root.join("control/state")).ok().flatten()?)
492 .ok()
493 .map(|value| value.epoch)
494}
495fn cache_identity(boot: &str, schema: &str, policy: &str) -> Vec<u8> {
496 serde_json::to_vec(&(boot, schema, policy)).expect("cache identity serializes")
497}
498fn hash(bytes: &[u8]) -> Digest {
499 Digest(Sha256::digest(bytes).into())
500}
501fn append(hasher: &mut Sha256, bytes: &[u8]) {
502 hasher.update((bytes.len() as u64).to_le_bytes());
503 hasher.update(bytes);
504}
505fn nonempty(values: &[&str]) -> Result<(), String> {
506 (!values.iter().any(|value| value.is_empty()))
507 .then_some(())
508 .ok_or_else(|| "identity values must be nonempty".into())
509}
510fn hex(value: Digest) -> String {
511 value.0.iter().map(|byte| format!("{byte:02x}")).collect()
512}