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>>,
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"),
)
})?;
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",
}
})?;
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(¤t)
.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(¤t)
.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(),
}
}