Skip to main content

kcode_k1_web_cache/
lib.rs

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}