haematite 0.7.0

Content-addressed, branchable, actor-native storage engine
Documentation
use wasm_bindgen::JsCast;
use wasm_bindgen_futures::JsFuture;
use web_sys::{DedicatedWorkerGlobalScope, FileSystemDirectoryHandle, WorkerGlobalScope};

use super::codec::{
    decode_store_format, default_browser_create_policy, encode_store_format, generation_id,
    record_file_name,
};
use super::js::{haem_acquire_owner, haem_move_file};
use super::opfs_io::{
    get_directory, get_file, js_source, list_files, read_access, read_optional, read_required,
    sync_handle, write_access, write_flushed,
};
use super::owner::OwnerLock;
use super::source::{node_publication_failure, record_entry_failure};
use super::{BrowserLocalError, BrowserSourceError};
use crate::branch::{preserve_node_publication_fence, preserve_record_entry_fence};
use crate::tree::TreePolicy;

const ROOT_NAME: &str = "haematite-browser-local-v1";
const FORMAT_FILE: &str = "format";
const GENERATION_FILE: &str = "generation";
const NODES_FILE: &str = "nodes";
const RECORDS_DIRECTORY: &str = "records";
const EVIDENCE: &str = "docs/design/evidence/HAEMATITE-WASM-MINIMAL-RUNG-PROBE-RELOAD-1.json — sha256:36036e22a9230548ed3198a78666024b9af4a3fbaa58bdb31a7f8db984524140";

pub(super) struct OpenedBackend {
    pub(super) backend: BrowserBackend,
    pub(super) generation: String,
    pub(super) records: Vec<Vec<u8>>,
    pub(super) nodes: Option<Vec<u8>>,
    /// The chunking [`TreePolicy`] resolved from this store's persisted `format`
    /// marker — v2 from birth for a freshly created store, the stamped v2 policy
    /// for an existing v2 store, or [`TreePolicy::V1_DEFAULT`] for a legacy v1
    /// store. The SOLE policy source for the browser mutation path (§4.1).
    pub(super) policy: TreePolicy,
}

pub(super) struct BrowserBackend {
    directory: FileSystemDirectoryHandle,
    records: FileSystemDirectoryHandle,
    owner: OwnerLock,
}

async fn initialize_fresh(
    directory: FileSystemDirectoryHandle,
    owner: OwnerLock,
    store_name: &str,
) -> Result<OpenedBackend, BrowserLocalError> {
    let records = get_directory(&directory, RECORDS_DIRECTORY, true)
        .await
        .map_err(|error| record_error("create records directory", error))?
        .ok_or_else(|| {
            record_error(
                "create records directory",
                BrowserSourceError::backend("create records directory", "directory unavailable"),
            )
        })?;
    // A freshly created browser store is chunking-v2 from birth (§5), stamped
    // with the ratified policy in its `format` marker — the browser analog of
    // native `Database::create`. The stamp is the SOLE policy source at open.
    let policy = default_browser_create_policy();
    write_flushed(&directory, FORMAT_FILE, &encode_store_format(policy))
        .await
        .map_err(|error| record_error("write format", error))?;
    let generation = generation_id(store_name);
    write_flushed(&directory, GENERATION_FILE, generation.as_bytes())
        .await
        .map_err(|error| record_error("write generation", error))?;
    Ok(OpenedBackend {
        backend: BrowserBackend {
            directory,
            records,
            owner,
        },
        generation,
        records: Vec::new(),
        nodes: None,
        policy,
    })
}

async fn open_existing(
    directory: FileSystemDirectoryHandle,
    owner: OwnerLock,
) -> Result<OpenedBackend, BrowserLocalError> {
    let records = get_directory(&directory, RECORDS_DIRECTORY, false)
        .await
        .map_err(|error| record_error("open records directory", error))?
        .ok_or(BrowserLocalError::StoreGenerationLost {
            component: "records directory",
        })?;
    let format = read_required(&directory, FORMAT_FILE).await.map_err(|_| {
        BrowserLocalError::StoreGenerationLost {
            component: "format",
        }
    })?;
    // Resolve the persisted format marker to this store's chunking policy: a
    // legacy v1 store opens under V1_DEFAULT forever, a v2 store under its
    // stamped targets, and a zero/wrong-length v2 stamp is a typed refusal
    // (§4.1 dispatch, browser seam; the engine never guesses a target).
    let policy = decode_store_format(&format)?;
    let generation = String::from_utf8(read_required(&directory, GENERATION_FILE).await.map_err(
        |_| BrowserLocalError::StoreGenerationLost {
            component: "generation",
        },
    )?)
    .map_err(|_| BrowserLocalError::StoreGenerationLost {
        component: "generation encoding",
    })?;
    let names = list_files(&records)
        .await
        .map_err(|error| record_error("list/read at open", error))?;
    let mut loaded = Vec::with_capacity(names.len());
    for name in names {
        let path = std::path::Path::new(&name);
        if name.starts_with(".branch-")
            && path
                .extension()
                .is_some_and(|extension| extension.eq_ignore_ascii_case("tmp"))
        {
            JsFuture::from(records.remove_entry(&name))
                .await
                .map_err(|error| {
                    record_error("sweep record temp", js_source("remove temp entry", &error))
                })?;
            let temp_absent = preserve_record_entry_fence(
                read_optional(&records, &name)
                    .await
                    .map(|entry| entry.is_none()),
            )
            .map_err(|error| record_entry_failure("temp unlink entry fence", error))?;
            if !temp_absent {
                let source = BrowserSourceError::backend(
                    "temp unlink entry fence",
                    "unlinked temp entry remained visible",
                );
                return Err(record_entry_failure(
                    "temp unlink entry fence",
                    crate::branch::FenceOperationError::RecordEntry(source),
                ));
            }
            continue;
        }
        if path
            .extension()
            .is_some_and(|extension| extension.eq_ignore_ascii_case("hbr"))
        {
            loaded.push(
                read_required(&records, &name)
                    .await
                    .map_err(|error| record_error("read record at open", error))?,
            );
        }
    }
    let nodes = read_optional(&directory, NODES_FILE)
        .await
        .map_err(|source| BrowserLocalError::NodeFailure {
            operation: "read nodes",
            source,
        })?;
    Ok(OpenedBackend {
        backend: BrowserBackend {
            directory,
            records,
            owner,
        },
        generation,
        records: loaded,
        nodes,
        policy,
    })
}

impl BrowserBackend {
    pub(super) async fn open(store_name: &str) -> Result<OpenedBackend, BrowserLocalError> {
        let global = js_sys::global();
        if !global.is_instance_of::<DedicatedWorkerGlobalScope>() {
            return Err(unsupported("not a dedicated Worker"));
        }
        let scope: WorkerGlobalScope = global.unchecked_into();
        let browser = scope
            .navigator()
            .user_agent()
            .unwrap_or_else(|_| "unknown".to_owned());
        if !(browser.contains("Chrome/150") && browser.contains("Macintosh")) {
            return Err(unsupported(&browser));
        }
        let root: FileSystemDirectoryHandle =
            JsFuture::from(scope.navigator().storage().get_directory())
                .await
                .map_err(|error| unsupported(&js_source("open OPFS root", &error).to_string()))?
                .dyn_into()
                .map_err(|_| unsupported("OPFS root cast failed"))?;
        let base = get_directory(&root, ROOT_NAME, true)
            .await
            .map_err(|error| record_error("open browser-local root", error))?
            .ok_or_else(|| unsupported("OPFS root unavailable"))?;
        let existing = get_directory(&base, store_name, false)
            .await
            .map_err(|error| record_error("open browser-local store", error))?;
        let fresh = existing.is_none();
        let directory = match existing {
            Some(directory) => directory,
            None => get_directory(&base, store_name, true)
                .await
                .map_err(|error| record_error("create browser-local store", error))?
                .ok_or_else(|| unsupported("OPFS store unavailable"))?,
        };
        let owner = OwnerLock::new(
            JsFuture::from(haem_acquire_owner(&format!(
                "haematite-browser-local:{store_name}"
            )))
            .await
            .map_err(|error| unsupported(&js_source("acquire owner lock", &error).to_string()))?,
        );
        if fresh {
            return initialize_fresh(directory, owner, store_name).await;
        }
        open_existing(directory, owner).await
    }

    pub(super) async fn write_nodes(&self, bytes: &[u8]) -> Result<(), BrowserLocalError> {
        preserve_node_publication_fence(write_flushed(&self.directory, NODES_FILE, bytes).await)
            .map_err(|error| node_publication_failure("node publication fence", error))
    }

    pub(super) async fn create_exclusive(
        &self,
        name: &str,
        bytes: &[u8],
    ) -> Result<Option<Vec<u8>>, BrowserLocalError> {
        let file_name = record_file_name(name);
        let file = get_file(&self.records, &file_name, true)
            .await
            .map_err(|error| record_error("create-exclusive", error))?;
        let access = sync_handle(&file)
            .await
            .map_err(|error| record_error("create-exclusive", error))?;
        let size = access.get_size().map_err(|error| {
            record_error("create-exclusive", js_source("get file size", &error))
        })?;
        if size != 0.0 {
            let existing =
                read_access(&access).map_err(|error| record_error("create-exclusive", error))?;
            access.close();
            return Ok(Some(existing));
        }
        write_access(&access, bytes).map_err(|error| record_error("create-exclusive", error))?;
        preserve_record_entry_fence(
            access
                .flush()
                .map_err(|error| js_source("flush created record", &error)),
        )
        .map_err(|error| record_entry_failure("create entry fence", error))?;
        access.close();
        self.entry_fence_present(&file_name, bytes, "create entry fence")
            .await?;
        Ok(None)
    }

    pub(super) async fn cas_replace_install(
        &self,
        name: &str,
        expected_created: u64,
        expected_seq: u64,
        bytes: &[u8],
    ) -> Result<(), BrowserLocalError> {
        let target = record_file_name(name);
        let current = read_optional(&self.records, &target)
            .await
            .map_err(|error| record_error("CAS read", error))?
            .ok_or_else(|| record_error("CAS replace-install", "branch removed"))?;
        let (_, record) = super::codec::decode_record(&current)
            .map_err(|error| record_error("CAS decode", error))?;
        if record.created != expected_created {
            return Err(record_error(
                "CAS replace-install",
                format!(
                    "branch recreated: expected {expected_created}, found {}",
                    record.created
                ),
            ));
        }
        if record.seq != expected_seq {
            return Err(record_error(
                "CAS replace-install",
                format!(
                    "stale sequence: expected {expected_seq}, found {}",
                    record.seq
                ),
            ));
        }
        let temporary = format!(".branch-{expected_created}-{expected_seq}.tmp");
        let temp = get_file(&self.records, &temporary, true)
            .await
            .map_err(|error| record_error("CAS temp", error))?;
        let access = sync_handle(&temp)
            .await
            .map_err(|error| record_error("CAS temp", error))?;
        write_access(&access, bytes).map_err(|error| record_error("CAS temp", error))?;
        access
            .flush()
            .map_err(|error| BrowserLocalError::FenceFailure {
                operation: "record temp flush",
                source: js_source("flush record temp", &error),
            })?;
        access.close();
        JsFuture::from(haem_move_file(&temp, &target))
            .await
            .map_err(|error| {
                record_error("CAS replace-install", js_source("move record file", &error))
            })?;
        self.entry_fence_present(&target, bytes, "replace entry fence")
            .await
    }

    pub(super) async fn cas_check(
        &self,
        name: &str,
        expected_created: u64,
        expected_seq: u64,
    ) -> Result<(), BrowserLocalError> {
        let current = read_optional(&self.records, &record_file_name(name))
            .await
            .map_err(|error| record_error("CAS read", error))?
            .ok_or_else(|| record_error("CAS check", "branch removed"))?;
        let (_, record) = super::codec::decode_record(&current)
            .map_err(|error| record_error("CAS decode", error))?;
        if record.created != expected_created || record.seq != expected_seq {
            return Err(record_error(
                "CAS check",
                "branch generation or sequence changed",
            ));
        }
        Ok(())
    }

    async fn entry_fence_present(
        &self,
        name: &str,
        expected: &[u8],
        operation: &'static str,
    ) -> Result<(), BrowserLocalError> {
        let observed = preserve_record_entry_fence(read_required(&self.records, name).await)
            .map_err(|error| record_entry_failure(operation, error))?;
        if observed != expected {
            let source = BrowserSourceError::backend(operation, "entry bytes mismatched");
            return Err(record_entry_failure(
                operation,
                crate::branch::FenceOperationError::RecordEntry(source),
            ));
        }
        Ok(())
    }

    pub(super) async fn close(&self) -> Result<(), BrowserLocalError> {
        self.owner.release().await
    }
}

fn unsupported(browser: &str) -> BrowserLocalError {
    BrowserLocalError::UnsupportedSubstrate {
        probe: "PROBE-RELOAD-1/committed-state",
        browser: browser.to_owned(),
        evidence: EVIDENCE,
    }
}
fn record_error(
    operation: &'static str,
    source: impl Into<BrowserSourceError>,
) -> BrowserLocalError {
    BrowserLocalError::RecordFailure {
        operation,
        source: source.into(),
    }
}