use std::io::{Cursor, Read, Write};
use std::path::{Path, PathBuf};
use base64::{engine::general_purpose::URL_SAFE_NO_PAD, Engine as _};
use ring::signature::{UnparsedPublicKey, ED25519};
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use sha2::{Digest, Sha256};
use crate::fsx::write_atomic;
use crate::store::Store;
pub const HUB_URL_ENV: &str = "DBMD_HUB_URL";
pub const HUB_KEY_ENV: &str = "DBMD_HUB_KEY";
pub const CONFIG_REL_PATH: &str = ".dbmd/config";
const MAX_RESPONSE_BYTES: u64 = 256 * 1024 * 1024;
const MAX_PUSH_BYTES: usize = 4 * 1024 * 1024;
const MAX_PUSH_FILES: usize = 100_000;
const MAX_STORE_BYTES: u64 = 512 * 1024 * 1024;
const MAX_PACK_BYTES: u64 = 256 * 1024 * 1024;
pub const MAX_PROPOSE_BYTES: u64 = 16 * 1024;
const CONNECT_TIMEOUT_SECS: u64 = 10;
const READ_TIMEOUT_SECS: u64 = 120;
#[derive(Debug, thiserror::Error)]
pub enum LinkError {
#[error(
"no hub configured — pass --hub <URL>, set {HUB_URL_ENV}, or add `hub = <URL>` to {CONFIG_REL_PATH}"
)]
NoHub,
#[error("no hub credential — set {HUB_KEY_ENV} (credentials never live in {CONFIG_REL_PATH})")]
NoCredential,
#[error(
"the hub credential in {HUB_KEY_ENV} contains whitespace or non-ASCII characters — re-copy it (the key is not shown here on purpose)"
)]
BadKey,
#[error("refusing non-HTTPS hub {hub} — the credential would travel in cleartext (localhost is exempt)")]
UnsafeHub {
hub: String,
},
#[error("hub unreachable at {hub}: {message}")]
Transport {
hub: String,
message: String,
},
#[error("{what} failed (HTTP {status}): {message}")]
Http {
what: &'static str,
status: u16,
message: String,
code: Option<String>,
},
#[error("{what}: the hub answered HTTP {status} with a non-JSON body — check the hub URL")]
NotJson {
what: &'static str,
status: u16,
},
#[error("hub response exceeded {} MB — refusing to buffer it", MAX_RESPONSE_BYTES / (1024 * 1024))]
ResponseTooLarge,
#[error("invalid address `{given}`: {reason}")]
BadAddress {
given: String,
reason: String,
},
#[error(
"invalid grant id `{given}` — grant ids come from `grant list` (lowercase letters, digits, hyphens)"
)]
BadGrantId {
given: String,
},
#[error("refusing unsafe path from the hub: `{path}`")]
UnsafePath {
path: String,
},
#[error(
"push too large ({detail}) — one snapshot caps at {} MB uncompressed, {} MB compressed, and {MAX_PUSH_FILES} files",
MAX_STORE_BYTES / (1024 * 1024),
MAX_PACK_BYTES / (1024 * 1024)
)]
PushTooLarge {
detail: String,
},
#[error(
"propose body too large ({bytes} bytes) — the hub's inbox caps one submission at {} KB",
MAX_PROPOSE_BYTES / 1024
)]
ProposeTooLarge {
bytes: u64,
},
#[error("store file `{path}` is not valid UTF-8 — the JSON push path carries text only")]
NotUtf8 {
path: String,
},
#[error("invalid store pack: {message}")]
InvalidPack {
message: String,
},
#[error("invalid signed feed: {message}")]
InvalidFeed {
message: String,
},
#[error(transparent)]
Io(#[from] std::io::Error),
#[error(transparent)]
Store(#[from] crate::StoreError),
}
pub type LinkResult<T> = std::result::Result<T, LinkError>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AddressTarget {
Id(String),
Path(String),
}
const BAD_BRAIN_REASON: &str =
"the brain reference must be a brain id (lowercase ULID) or a slug (lowercase letters, digits, hyphens)";
const BAD_TARGET_REASON: &str =
"the part after `/` must be a record id (lowercase ULID) or a store-relative `.md` path";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Address {
pub brain: String,
pub target: Option<AddressTarget>,
}
impl Address {
pub fn parse(raw: &str) -> LinkResult<Address> {
let bad = |reason: &str| LinkError::BadAddress {
given: raw.to_string(),
reason: reason.to_string(),
};
let trimmed = raw.trim();
let body = trimmed.strip_prefix('@').unwrap_or(trimmed);
if body.is_empty() {
return Err(bad("empty address"));
}
let (brain, rest) = match body.split_once('/') {
Some((b, r)) => (b, Some(r)),
None => (body, None),
};
if brain.is_empty() {
return Err(bad("missing brain reference before `/`"));
}
if !is_safe_ref(brain) {
return Err(bad(BAD_BRAIN_REASON));
}
let target = match rest {
None => None,
Some("") => return Err(bad("trailing `/` with no record id or path")),
Some(r) if crate::ulid::is_ulid(r) => Some(AddressTarget::Id(r.to_string())),
Some(r) => {
if !safe_store_rel_path(r) || !r.ends_with(".md") {
return Err(bad(BAD_TARGET_REASON));
}
Some(AddressTarget::Path(r.to_string()))
}
};
Ok(Address {
brain: brain.to_string(),
target,
})
}
}
fn is_safe_ref(s: &str) -> bool {
!s.is_empty()
&& s.len() <= 64
&& s.bytes()
.all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-')
}
pub fn is_valid_handle(s: &str) -> bool {
is_safe_ref(s)
}
pub fn safe_store_rel_path(p: &str) -> bool {
if p.is_empty() || p.len() > 512 || p.starts_with('/') {
return false;
}
if !p
.bytes()
.all(|b| b.is_ascii_alphanumeric() || matches!(b, b'.' | b'_' | b'-' | b'/'))
{
return false;
}
p.split('/')
.all(|seg| !seg.is_empty() && seg != "." && seg != ".." && !seg.starts_with('.'))
}
fn require_safe_ref(brain: &str) -> LinkResult<()> {
if is_safe_ref(brain) {
Ok(())
} else {
Err(LinkError::BadAddress {
given: brain.to_string(),
reason: BAD_BRAIN_REASON.to_string(),
})
}
}
fn require_valid_handle(handle: &str) -> LinkResult<()> {
if is_valid_handle(handle) {
Ok(())
} else {
Err(LinkError::BadAddress {
given: handle.to_string(),
reason: "the site handle must be lowercase letters, digits, hyphens".to_string(),
})
}
}
fn require_safe_grant_id(id: &str) -> LinkResult<()> {
if is_safe_ref(id) {
Ok(())
} else {
Err(LinkError::BadGrantId {
given: id.to_string(),
})
}
}
#[derive(Debug, Clone)]
pub struct HubConfig {
pub hub: String,
pub key: Option<String>,
}
impl HubConfig {
pub fn require_key(&self) -> LinkResult<&str> {
self.key.as_deref().ok_or(LinkError::NoCredential)
}
}
pub fn hub_config(flag_hub: Option<&str>, dir: &Path) -> LinkResult<HubConfig> {
let hub = flag_hub
.map(str::to_string)
.or_else(|| env_nonempty(HUB_URL_ENV))
.or_else(|| config_file_hub(&dir.join(CONFIG_REL_PATH)))
.ok_or(LinkError::NoHub)?;
let hub = hub.trim().trim_end_matches('/').to_string();
assert_safe_hub(&hub)?;
let key = match env_nonempty(HUB_KEY_ENV) {
Some(raw) => Some(clean_key(&raw)?),
None => None,
};
Ok(HubConfig { hub, key })
}
fn env_nonempty(name: &str) -> Option<String> {
std::env::var(name).ok().filter(|v| !v.trim().is_empty())
}
fn config_file_hub(path: &Path) -> Option<String> {
let text = std::fs::read_to_string(path).ok()?;
for line in text.lines() {
let line = line.trim();
if line.is_empty() || line.starts_with('#') {
continue;
}
if let Some((k, v)) = line.split_once('=') {
if k.trim() == "hub" {
let v = v.trim();
if !v.is_empty() {
return Some(v.to_string());
}
}
}
}
None
}
fn assert_safe_hub(hub: &str) -> LinkResult<()> {
let parsed = url::Url::parse(hub).map_err(|_| LinkError::UnsafeHub {
hub: hub.to_string(),
})?;
if !parsed.username().is_empty()
|| parsed.password().is_some()
|| parsed.query().is_some()
|| parsed.fragment().is_some()
{
return Err(LinkError::UnsafeHub {
hub: hub.to_string(),
});
}
let loopback = match parsed.host() {
Some(url::Host::Domain(host)) => host.eq_ignore_ascii_case("localhost"),
Some(url::Host::Ipv4(ip)) => ip.is_loopback(),
Some(url::Host::Ipv6(ip)) => ip.is_loopback(),
None => false,
};
if parsed.scheme().eq_ignore_ascii_case("https") || loopback {
Ok(())
} else {
Err(LinkError::UnsafeHub {
hub: hub.to_string(),
})
}
}
fn clean_key(raw: &str) -> LinkResult<String> {
let k = raw.trim();
if k.is_empty() || k.bytes().any(|b| !(0x21..=0x7e).contains(&b)) {
return Err(LinkError::BadKey);
}
Ok(k.to_string())
}
#[derive(Debug)]
pub struct HubResponse {
pub status: u16,
pub body: Option<Value>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Auth {
Required,
None,
}
fn agent() -> ureq::Agent {
ureq::AgentBuilder::new()
.user_agent(concat!("dbmd/", env!("CARGO_PKG_VERSION")))
.redirects(0)
.timeout_connect(std::time::Duration::from_secs(CONNECT_TIMEOUT_SECS))
.timeout_read(std::time::Duration::from_secs(READ_TIMEOUT_SECS))
.build()
}
fn request(
cfg: &HubConfig,
method: &str,
path: &str,
body: Option<&Value>,
auth: Auth,
) -> LinkResult<HubResponse> {
let url = format!("{}{}", cfg.hub, path);
let mut req = agent().request(method, &url);
if auth == Auth::Required {
req = req.set("authorization", &format!("Bearer {}", cfg.require_key()?));
}
let result = match body {
Some(v) => req
.set("content-type", "application/json")
.send_string(&v.to_string()),
None => req.call(),
};
let resp = match result {
Ok(resp) => resp,
Err(ureq::Error::Status(_, resp)) => resp,
Err(ureq::Error::Transport(t)) => {
return Err(LinkError::Transport {
hub: cfg.hub.clone(),
message: t.to_string(),
})
}
};
let status = resp.status();
let mut buf = Vec::new();
resp.into_reader()
.take(MAX_RESPONSE_BYTES + 1)
.read_to_end(&mut buf)?;
if buf.len() as u64 > MAX_RESPONSE_BYTES {
return Err(LinkError::ResponseTooLarge);
}
let parsed: Option<Value> = serde_json::from_slice(&buf).ok();
Ok(HubResponse {
status,
body: parsed,
})
}
fn assert_safe_presigned_url(raw: &str) -> LinkResult<()> {
let parsed = url::Url::parse(raw).map_err(|_| LinkError::InvalidPack {
message: "the hub returned an invalid object-store URL".to_string(),
})?;
if !parsed.scheme().eq_ignore_ascii_case("https")
|| !parsed.username().is_empty()
|| parsed.password().is_some()
|| parsed.fragment().is_some()
{
return Err(LinkError::InvalidPack {
message: "the hub returned an unsafe object-store URL".to_string(),
});
}
Ok(())
}
fn put_presigned(raw: &str, headers: &Value, bytes: &[u8]) -> LinkResult<()> {
assert_safe_presigned_url(raw)?;
let mut req = agent().put(raw);
if let Some(map) = headers.as_object() {
for (name, value) in map {
if let Some(value) = value.as_str() {
req = req.set(name, value);
}
}
}
match req.send_bytes(bytes) {
Ok(resp) if resp.status() < 300 => Ok(()),
Ok(resp) | Err(ureq::Error::Status(_, resp)) => Err(LinkError::Http {
what: "pack upload",
status: resp.status(),
message: "object store rejected the upload".to_string(),
code: None,
}),
Err(ureq::Error::Transport(err)) => Err(LinkError::Transport {
hub: "the object store".to_string(),
message: err.to_string(),
}),
}
}
fn get_presigned(raw: &str) -> LinkResult<Vec<u8>> {
assert_safe_presigned_url(raw)?;
let resp = match agent().get(raw).call() {
Ok(resp) => resp,
Err(ureq::Error::Status(_, resp)) => {
return Err(LinkError::Http {
what: "pack download",
status: resp.status(),
message: "object store rejected the download".to_string(),
code: None,
})
}
Err(ureq::Error::Transport(err)) => {
return Err(LinkError::Transport {
hub: "the object store".to_string(),
message: err.to_string(),
})
}
};
let mut bytes = Vec::new();
resp.into_reader()
.take(MAX_PACK_BYTES + 1)
.read_to_end(&mut bytes)?;
if bytes.len() as u64 > MAX_PACK_BYTES {
return Err(LinkError::InvalidPack {
message: "download exceeds the compressed-size limit".to_string(),
});
}
Ok(bytes)
}
fn ensure_ok(r: HubResponse, what: &'static str) -> LinkResult<Value> {
if r.status >= 400 {
let message = r
.body
.as_ref()
.and_then(|b| b.get("error"))
.and_then(Value::as_str)
.unwrap_or("unknown error")
.to_string();
let code = r
.body
.as_ref()
.and_then(|b| b.get("code"))
.and_then(Value::as_str)
.map(str::to_string);
return Err(LinkError::Http {
what,
status: r.status,
message,
code,
});
}
r.body.ok_or(LinkError::NotJson {
what,
status: r.status,
})
}
pub fn resolve(cfg: &HubConfig, addr: &Address) -> LinkResult<Value> {
require_safe_ref(&addr.brain)?;
if let Some(target) = &addr.target {
let (given, ok) = match target {
AddressTarget::Id(id) => (id, crate::ulid::is_ulid(id)),
AddressTarget::Path(p) => (p, safe_store_rel_path(p) && p.ends_with(".md")),
};
if !ok {
return Err(LinkError::BadAddress {
given: given.clone(),
reason: BAD_TARGET_REASON.to_string(),
});
}
}
let path = match &addr.target {
None => format!("/api/hub/brains/{}", addr.brain),
Some(AddressTarget::Id(id)) => {
format!("/api/hub/brains/{}/resolve?id={id}", addr.brain)
}
Some(AddressTarget::Path(p)) => {
format!("/api/hub/brains/{}/resolve?path={p}", addr.brain)
}
};
ensure_ok(request(cfg, "GET", &path, None, Auth::Required)?, "resolve")
}
#[derive(Debug, serde::Serialize)]
pub struct PullReport {
pub brain: String,
pub slug: String,
#[serde(rename = "headSeq")]
pub head_seq: u64,
pub files: usize,
pub dest: String,
#[serde(rename = "extraLocal")]
pub extra_local: Vec<String>,
}
pub fn sync_pull(cfg: &HubConfig, brain: &str, out: Option<&Path>) -> LinkResult<PullReport> {
require_safe_ref(brain)?;
let path = format!("/api/hub/brains/{brain}/export?format=pack");
let body = ensure_ok(
request(cfg, "GET", &path, None, Auth::Required)?,
"sync pull",
)?;
let remote_slug = body
.get("slug")
.and_then(Value::as_str)
.filter(|slug| is_safe_slug(slug));
let slug = remote_slug
.or_else(|| is_safe_slug(brain).then_some(brain))
.unwrap_or("brain")
.to_string();
let brain_id = body
.get("brain")
.and_then(Value::as_str)
.unwrap_or(brain)
.to_string();
let head_seq = body.get("headSeq").and_then(Value::as_u64).unwrap_or(0);
let dest: PathBuf = match out {
Some(p) => p.to_path_buf(),
None => PathBuf::from(&slug),
};
let entries =
if let Some(url) = body.get("url").and_then(Value::as_str) {
let expected = body
.get("sha256")
.and_then(Value::as_str)
.filter(|hash| {
hash.len() == 64
&& hash
.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
})
.ok_or_else(|| LinkError::InvalidPack {
message: "the hub returned an invalid SHA-256".to_string(),
})?;
let bytes = get_presigned(url)?;
let actual = format!("{:x}", Sha256::digest(&bytes));
if actual != expected {
return Err(LinkError::InvalidPack {
message: "SHA-256 verification failed".to_string(),
});
}
parse_store_pack(bytes)?
} else {
let files = body.get("files").and_then(Value::as_array).ok_or_else(|| {
LinkError::InvalidPack {
message: "the hub returned neither a pack nor a file manifest".to_string(),
}
})?;
let mut entries = Vec::with_capacity(files.len());
for file in files {
let path = file.get("path").and_then(Value::as_str).ok_or_else(|| {
LinkError::InvalidPack {
message: "a file entry has no string path".to_string(),
}
})?;
let content = file.get("content").and_then(Value::as_str).ok_or_else(|| {
LinkError::InvalidPack {
message: format!("file `{path}` has no string content"),
}
})?;
entries.push((path.to_string(), content.as_bytes().to_vec()));
}
entries
};
let mut seen = std::collections::HashSet::new();
for (path, _) in &entries {
if !safe_store_rel_path(path) {
return Err(LinkError::UnsafePath { path: path.clone() });
}
if !seen.insert(path) {
return Err(LinkError::InvalidPack {
message: format!("duplicate path `{path}`"),
});
}
}
std::fs::create_dir_all(&dest)?;
let real_dest = std::fs::canonicalize(&dest)?;
for (p, content) in &entries {
let abs = dest.join(p);
if let Some(parent) = abs.parent() {
std::fs::create_dir_all(parent)?;
let real_parent = std::fs::canonicalize(parent)?;
if !real_parent.starts_with(&real_dest) {
return Err(LinkError::UnsafePath { path: p.clone() });
}
}
if std::fs::symlink_metadata(&abs).is_ok_and(|meta| meta.file_type().is_symlink()) {
return Err(LinkError::UnsafePath { path: p.clone() });
}
write_atomic(&abs, content)?;
}
let pulled: std::collections::BTreeSet<&str> =
entries.iter().map(|(p, _)| p.as_str()).collect();
let mut extra_local = Vec::new();
if let Ok(store) = Store::open(&dest) {
if let Ok(walked) = store.walk() {
for rel in walked {
let rel_str = rel.to_string_lossy().replace('\\', "/");
if !pulled.contains(rel_str.as_str()) {
extra_local.push(rel_str);
}
}
}
}
Ok(PullReport {
brain: brain_id,
slug,
head_seq,
files: entries.len(),
dest: dest.to_string_lossy().into_owned(),
extra_local,
})
}
fn is_safe_slug(slug: &str) -> bool {
!slug.is_empty()
&& slug.len() <= 63
&& !slug.starts_with('-')
&& !slug.ends_with('-')
&& slug
.bytes()
.all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-')
}
fn parse_store_pack(bytes: Vec<u8>) -> LinkResult<Vec<(String, Vec<u8>)>> {
let mut archive =
zip::ZipArchive::new(Cursor::new(bytes)).map_err(|err| LinkError::InvalidPack {
message: format!("ZIP parse failed: {err}"),
})?;
if archive.is_empty() || archive.len() > MAX_PUSH_FILES {
return Err(LinkError::InvalidPack {
message: format!("invalid file count {}", archive.len()),
});
}
let mut total = 0u64;
let mut entries = Vec::with_capacity(archive.len());
for index in 0..archive.len() {
let mut file = archive
.by_index(index)
.map_err(|err| LinkError::InvalidPack {
message: format!("ZIP entry failed: {err}"),
})?;
if file.is_dir() {
continue;
}
let path = file.name().to_string();
if file.enclosed_name().is_none() || !safe_store_rel_path(&path) {
return Err(LinkError::UnsafePath { path });
}
if file
.unix_mode()
.is_some_and(|mode| !matches!(mode & 0o170000, 0 | 0o100000))
{
return Err(LinkError::InvalidPack {
message: format!("non-file entry `{path}`"),
});
}
total = total.saturating_add(file.size());
if total > MAX_STORE_BYTES {
return Err(LinkError::InvalidPack {
message: "expanded content exceeds the 512 MB limit".to_string(),
});
}
let mut content = Vec::new();
file.read_to_end(&mut content)
.map_err(|err| LinkError::InvalidPack {
message: format!("could not decompress `{path}`: {err}"),
})?;
if content.len() as u64 != file.size() {
return Err(LinkError::InvalidPack {
message: format!("length mismatch for `{path}`"),
});
}
entries.push((path, content));
}
if entries.is_empty() {
return Err(LinkError::InvalidPack {
message: "pack contains no files".to_string(),
});
}
Ok(entries)
}
pub fn collect_push_files(store: &Store) -> LinkResult<Vec<(String, String)>> {
let mut out: Vec<(String, String)> = Vec::new();
let read_text = |rel: &str| -> LinkResult<String> {
let abs = store.root.join(rel);
std::fs::read(&abs)
.map_err(LinkError::from)
.and_then(|bytes| {
String::from_utf8(bytes).map_err(|_| LinkError::NotUtf8 {
path: rel.to_string(),
})
})
};
out.push(("DB.md".to_string(), read_text("DB.md")?));
if store.root.join("assets.jsonl").is_file() {
out.push(("assets.jsonl".to_string(), read_text("assets.jsonl")?));
}
for rel in store.walk()? {
let rel_str = rel.to_string_lossy().replace('\\', "/");
if !safe_store_rel_path(&rel_str) {
return Err(LinkError::UnsafePath { path: rel_str });
}
let content = read_text(&rel_str)?;
out.push((rel_str, content));
}
out.sort_by(|a, b| a.0.cmp(&b.0));
Ok(out)
}
pub fn sync_push(cfg: &HubConfig, brain: &str, files: &[(String, String)]) -> LinkResult<Value> {
require_safe_ref(brain)?;
if files.len() > MAX_PUSH_FILES {
return Err(LinkError::PushTooLarge {
detail: format!("{} files", files.len()),
});
}
let raw_total: u64 = files.iter().map(|(_, content)| content.len() as u64).sum();
if raw_total > MAX_STORE_BYTES {
return Err(LinkError::PushTooLarge {
detail: format!("{raw_total} uncompressed bytes"),
});
}
let body = json!({
"files": files
.iter()
.map(|(p, c)| json!({ "path": p, "content": c }))
.collect::<Vec<_>>(),
});
if body.to_string().len() <= MAX_PUSH_BYTES {
let path = format!("/api/hub/brains/{brain}/push");
return ensure_ok(
request(cfg, "POST", &path, Some(&body), Auth::Required)?,
"sync push",
);
}
let pack = build_store_pack(files)?;
if pack.len() as u64 > MAX_PACK_BYTES {
return Err(LinkError::PushTooLarge {
detail: format!("{} compressed bytes", pack.len()),
});
}
let sha256 = format!("{:x}", Sha256::digest(&pack));
let meta = json!({ "sha256": sha256, "bytes": pack.len() });
let presigned = ensure_ok(
request(
cfg,
"POST",
&format!("/api/hub/brains/{brain}/packs/presign"),
Some(&meta),
Auth::Required,
)?,
"prepare pack upload",
)?;
let url = presigned
.get("url")
.and_then(Value::as_str)
.ok_or_else(|| LinkError::InvalidPack {
message: "the hub returned no upload URL".to_string(),
})?;
put_presigned(url, presigned.get("headers").unwrap_or(&Value::Null), &pack)?;
ensure_ok(
request(
cfg,
"POST",
&format!("/api/hub/brains/{brain}/packs/commit"),
Some(&meta),
Auth::Required,
)?,
"commit pack",
)
}
fn build_store_pack(files: &[(String, String)]) -> LinkResult<Vec<u8>> {
let mut sorted: Vec<_> = files.iter().collect();
sorted.sort_by(|a, b| a.0.cmp(&b.0));
let mut writer = zip::ZipWriter::new(Cursor::new(Vec::new()));
let options = zip::write::SimpleFileOptions::default()
.compression_method(zip::CompressionMethod::Deflated)
.last_modified_time(zip::DateTime::default())
.unix_permissions(0o600);
for (path, content) in sorted {
writer
.start_file(path, options)
.map_err(|err| LinkError::InvalidPack {
message: format!("could not create ZIP entry `{path}`: {err}"),
})?;
writer.write_all(content.as_bytes())?;
}
writer
.finish()
.map(Cursor::into_inner)
.map_err(|err| LinkError::InvalidPack {
message: format!("could not finish ZIP: {err}"),
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Capability {
Read,
Write,
}
impl Capability {
pub fn as_str(self) -> &'static str {
match self {
Capability::Read => "read",
Capability::Write => "write",
}
}
}
pub fn grant_issue(
cfg: &HubConfig,
brain: &str,
grantee: &str,
can: Capability,
scope: Option<&str>,
until: Option<&str>,
) -> LinkResult<Value> {
require_safe_ref(brain)?;
let mut body = json!({
"email": grantee,
"capability": can.as_str(),
});
if let Some(s) = scope {
body["scopePrefix"] = json!(s);
}
if let Some(u) = until {
body["expiresAt"] = json!(u);
}
let path = format!("/api/hub/brains/{brain}/grants");
ensure_ok(
request(cfg, "POST", &path, Some(&body), Auth::Required)?,
"grant issue",
)
}
pub fn grant_list(cfg: &HubConfig, brain: &str) -> LinkResult<Value> {
require_safe_ref(brain)?;
let path = format!("/api/hub/brains/{brain}/grants");
ensure_ok(
request(cfg, "GET", &path, None, Auth::Required)?,
"grant list",
)
}
pub fn grant_revoke(cfg: &HubConfig, brain: &str, grant_id: &str) -> LinkResult<Value> {
require_safe_ref(brain)?;
require_safe_grant_id(grant_id)?;
let path = format!("/api/hub/brains/{brain}/grants/{grant_id}");
ensure_ok(
request(cfg, "DELETE", &path, None, Auth::Required)?,
"grant revoke",
)
}
pub fn propose(cfg: &HubConfig, handle: &str, app: &str, body: &str) -> LinkResult<Value> {
require_valid_handle(handle)?;
if body.len() as u64 > MAX_PROPOSE_BYTES {
return Err(LinkError::ProposeTooLarge {
bytes: body.len() as u64,
});
}
let payload = json!({ "app": app, "body": body });
let path = format!("/api/hub/sites/{handle}/inbox");
ensure_ok(
request(cfg, "POST", &path, Some(&payload), Auth::None)?,
"propose",
)
}
#[derive(Debug, serde::Serialize)]
pub struct Head {
pub brain: String,
pub seq: u64,
#[serde(rename = "updatedAt", skip_serializing_if = "Option::is_none")]
pub updated_at: Option<String>,
#[serde(rename = "feedHash", skip_serializing_if = "Option::is_none")]
pub feed_hash: Option<String>,
pub verified: bool,
}
#[derive(Debug, Deserialize, Serialize)]
struct FeedFile {
path: String,
sha256: String,
bytes: u64,
}
#[derive(Debug, Deserialize, Serialize)]
struct FeedEntry {
v: u8,
seq: u64,
ts: String,
brain: String,
public_key: String,
kind: String,
op: String,
pack_sha256: String,
files: Vec<FeedFile>,
removed: Vec<String>,
prev_entry_hash: Option<String>,
sig: String,
}
#[derive(Serialize)]
struct UnsignedFeedEntry<'a> {
v: u8,
seq: u64,
ts: &'a str,
brain: &'a str,
public_key: &'a str,
kind: &'a str,
op: &'a str,
pack_sha256: &'a str,
files: &'a [FeedFile],
removed: &'a [String],
prev_entry_hash: &'a Option<String>,
}
#[derive(Debug, Deserialize)]
struct FeedItem {
hash: String,
entry: FeedEntry,
}
#[derive(Debug, Deserialize)]
struct FeedIdentity {
fingerprint: String,
#[serde(rename = "publicKeySpki")]
public_key_spki: String,
}
#[derive(Debug, Deserialize)]
struct FeedResponse {
#[serde(rename = "headSeq")]
head_seq: u64,
#[serde(rename = "feedHash")]
feed_hash: Option<String>,
identity: Option<FeedIdentity>,
entries: Vec<FeedItem>,
#[serde(rename = "scopeLimited")]
scope_limited: bool,
}
fn invalid_feed(message: impl Into<String>) -> LinkError {
LinkError::InvalidFeed {
message: message.into(),
}
}
fn verify_feed_item(item: &FeedItem, identity: &FeedIdentity) -> LinkResult<()> {
const ED25519_SPKI_PREFIX: &[u8] = &[
0x30, 0x2a, 0x30, 0x05, 0x06, 0x03, 0x2b, 0x65, 0x70, 0x03, 0x21, 0x00,
];
let entry = &item.entry;
let public_der = URL_SAFE_NO_PAD
.decode(&entry.public_key)
.map_err(|_| invalid_feed("public key is not base64url"))?;
if public_der.len() != ED25519_SPKI_PREFIX.len() + 32
|| !public_der.starts_with(ED25519_SPKI_PREFIX)
|| identity.public_key_spki != entry.public_key
{
return Err(invalid_feed("public key does not match the brain card"));
}
let fingerprint = URL_SAFE_NO_PAD.encode(Sha256::digest(&public_der));
if fingerprint != identity.fingerprint || entry.brain != format!("ed25519:{fingerprint}") {
return Err(invalid_feed(
"brain fingerprint does not match its public key",
));
}
let unsigned = UnsignedFeedEntry {
v: entry.v,
seq: entry.seq,
ts: &entry.ts,
brain: &entry.brain,
public_key: &entry.public_key,
kind: &entry.kind,
op: &entry.op,
pack_sha256: &entry.pack_sha256,
files: &entry.files,
removed: &entry.removed,
prev_entry_hash: &entry.prev_entry_hash,
};
let message =
serde_json::to_vec(&unsigned).map_err(|_| invalid_feed("could not canonicalize entry"))?;
let signature = URL_SAFE_NO_PAD
.decode(&entry.sig)
.map_err(|_| invalid_feed("signature is not base64url"))?;
UnparsedPublicKey::new(&ED25519, &public_der[ED25519_SPKI_PREFIX.len()..])
.verify(&message, &signature)
.map_err(|_| invalid_feed("Ed25519 signature verification failed"))?;
let mut exact = serde_json::to_vec(entry).map_err(|_| invalid_feed("could not hash entry"))?;
exact.push(b'\n');
let actual_hash = format!("{:x}", Sha256::digest(&exact));
if actual_hash != item.hash {
return Err(invalid_feed("entry SHA-256 does not match"));
}
Ok(())
}
pub fn head(cfg: &HubConfig, brain: &str) -> LinkResult<Head> {
require_safe_ref(brain)?;
let path = format!("/api/hub/brains/{brain}");
let body = ensure_ok(
request(cfg, "GET", &path, None, Auth::Required)?,
"subscribe",
)?;
let resolved_brain = body
.get("id")
.and_then(Value::as_str)
.unwrap_or(brain)
.to_string();
let seq = body.get("headSeq").and_then(Value::as_u64).unwrap_or(0);
let advertised_hash = body
.get("feedHash")
.and_then(Value::as_str)
.map(str::to_string);
let updated_at = body
.get("updatedAt")
.and_then(Value::as_str)
.map(str::to_string);
if seq == 0 {
return Ok(Head {
brain: resolved_brain,
seq,
updated_at,
feed_hash: None,
verified: true,
});
}
let feed_value = ensure_ok(
request(
cfg,
"GET",
&format!("/api/hub/brains/{brain}/feed?after={}&limit=1", seq - 1),
None,
Auth::Required,
)?,
"subscribe feed",
)?;
let feed: FeedResponse = serde_json::from_value(feed_value)
.map_err(|_| invalid_feed("hub returned an invalid feed shape"))?;
if feed.head_seq != seq || feed.feed_hash != advertised_hash {
return Err(invalid_feed("brain card and feed head disagree"));
}
if feed.scope_limited {
return Ok(Head {
brain: resolved_brain,
seq,
updated_at,
feed_hash: advertised_hash,
verified: false,
});
}
let identity = feed
.identity
.as_ref()
.ok_or_else(|| invalid_feed("feed has no brain identity"))?;
let item = feed
.entries
.first()
.ok_or_else(|| invalid_feed("feed head entry is missing"))?;
if item.entry.seq != seq || Some(&item.hash) != advertised_hash.as_ref() {
return Err(invalid_feed(
"advertised feed hash does not address the head entry",
));
}
verify_feed_item(item, identity)?;
Ok(Head {
brain: resolved_brain,
seq,
updated_at,
feed_hash: advertised_hash,
verified: true,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn signed_feed_item_verifies_identity_hash_and_signature() {
use ring::rand::SystemRandom;
use ring::signature::{Ed25519KeyPair, KeyPair};
const PREFIX: &[u8] = &[
0x30, 0x2a, 0x30, 0x05, 0x06, 0x03, 0x2b, 0x65, 0x70, 0x03, 0x21, 0x00,
];
let pkcs8 = Ed25519KeyPair::generate_pkcs8(&SystemRandom::new()).unwrap();
let pair = Ed25519KeyPair::from_pkcs8(pkcs8.as_ref()).unwrap();
let mut spki = PREFIX.to_vec();
spki.extend_from_slice(pair.public_key().as_ref());
let public_key = URL_SAFE_NO_PAD.encode(&spki);
let fingerprint = URL_SAFE_NO_PAD.encode(Sha256::digest(&spki));
let mut entry = FeedEntry {
v: 1,
seq: 1,
ts: "2026-07-14T00:00:00.000Z".to_string(),
brain: format!("ed25519:{fingerprint}"),
public_key: public_key.clone(),
kind: "push".to_string(),
op: "snapshot".to_string(),
pack_sha256: "a".repeat(64),
files: vec![FeedFile {
path: "DB.md".to_string(),
sha256: "b".repeat(64),
bytes: 3,
}],
removed: vec![],
prev_entry_hash: None,
sig: String::new(),
};
let unsigned = UnsignedFeedEntry {
v: entry.v,
seq: entry.seq,
ts: &entry.ts,
brain: &entry.brain,
public_key: &entry.public_key,
kind: &entry.kind,
op: &entry.op,
pack_sha256: &entry.pack_sha256,
files: &entry.files,
removed: &entry.removed,
prev_entry_hash: &entry.prev_entry_hash,
};
entry.sig =
URL_SAFE_NO_PAD.encode(pair.sign(&serde_json::to_vec(&unsigned).unwrap()).as_ref());
let mut exact = serde_json::to_vec(&entry).unwrap();
exact.push(b'\n');
let item = FeedItem {
hash: format!("{:x}", Sha256::digest(&exact)),
entry,
};
let identity = FeedIdentity {
fingerprint,
public_key_spki: public_key,
};
assert!(verify_feed_item(&item, &identity).is_ok());
let mut tampered = item;
tampered.entry.pack_sha256 = "c".repeat(64);
assert!(verify_feed_item(&tampered, &identity).is_err());
}
#[test]
fn address_bare_brain_with_and_without_sigil() {
for raw in ["@acme-ops", "acme-ops"] {
let a = Address::parse(raw).expect(raw);
assert_eq!(a.brain, "acme-ops");
assert_eq!(a.target, None);
}
}
#[test]
fn address_ulid_target_parses_as_id() {
let a = Address::parse("@acme/01j5qc3v9k4ym8rwbn2tqe6f7d").unwrap();
assert_eq!(a.brain, "acme");
assert_eq!(
a.target,
Some(AddressTarget::Id("01j5qc3v9k4ym8rwbn2tqe6f7d".to_string()))
);
}
#[test]
fn address_md_path_target_parses_as_path() {
let a = Address::parse("@acme/records/clients/lumio.md").unwrap();
assert_eq!(
a.target,
Some(AddressTarget::Path("records/clients/lumio.md".to_string()))
);
}
#[test]
fn address_rejects_malformed_forms() {
for raw in [
"",
"@",
"@/x",
"@acme/",
"@acme/../etc/passwd",
"@acme/records/.hidden.md",
"@ACME", "@acme/notes/x.txt", "@a b", ] {
assert!(Address::parse(raw).is_err(), "should reject {raw:?}");
}
}
#[test]
fn safe_paths_accept_store_shapes_and_reject_escapes() {
for ok in [
"DB.md",
"assets.jsonl",
"records/clients/lumio.md",
"sources/emails/2026/07/x.md",
] {
assert!(safe_store_rel_path(ok), "should accept {ok:?}");
}
for bad in [
"",
"/etc/passwd",
"../up.md",
"records/../../up.md",
"records//x.md",
".dbmd/config",
"records/.hidden/x.md",
"records/a b.md",
"records\\win.md",
] {
assert!(!safe_store_rel_path(bad), "should reject {bad:?}");
}
}
#[test]
fn hub_config_flag_beats_file_and_requires_some_source() {
let dir = tempfile::tempdir().unwrap();
std::fs::create_dir_all(dir.path().join(".dbmd")).unwrap();
std::fs::write(
dir.path().join(CONFIG_REL_PATH),
"# toolkit state\nhub = https://file.example.com\nunknown = ignored\n",
)
.unwrap();
let from_flag = hub_config(Some("https://flag.example.com/"), dir.path()).unwrap();
assert_eq!(from_flag.hub, "https://flag.example.com");
let from_file = hub_config(None, dir.path()).unwrap();
assert_eq!(from_file.hub, "https://file.example.com");
let none = hub_config(None, tempfile::tempdir().unwrap().path());
assert!(matches!(none, Err(LinkError::NoHub)));
}
#[test]
fn https_guard_allows_loopback_only_for_plain_http() {
assert!(assert_safe_hub("https://hub.example.com").is_ok());
assert!(assert_safe_hub("http://localhost:3000").is_ok());
assert!(assert_safe_hub("http://127.0.0.1:3000").is_ok());
assert!(assert_safe_hub("http://[::1]:3000").is_ok());
assert!(matches!(
assert_safe_hub("http://hub.example.com"),
Err(LinkError::UnsafeHub { .. })
));
assert!(matches!(
assert_safe_hub("hub.example.com"),
Err(LinkError::UnsafeHub { .. })
));
assert!(matches!(
assert_safe_hub("http://localhost:80@127.0.0.1:1"),
Err(LinkError::UnsafeHub { .. })
));
assert!(matches!(
assert_safe_hub("https://hub.example.com@attacker.example"),
Err(LinkError::UnsafeHub { .. })
));
}
#[test]
fn https_guard_matches_the_scheme_case_insensitively() {
assert!(assert_safe_hub("HTTPS://hub.example.com").is_ok());
assert!(assert_safe_hub("Https://hub.example.com").is_ok());
assert!(matches!(
assert_safe_hub("HTTP://hub.example.com"),
Err(LinkError::UnsafeHub { .. })
));
}
#[test]
fn clean_key_refuses_paste_artifacts_without_echoing() {
assert_eq!(clean_key(" vc_account_abc ").unwrap(), "vc_account_abc");
for bad in ["vc account", "vc\naccount", "ключ", ""] {
let err = clean_key(bad).unwrap_err();
assert!(matches!(err, LinkError::BadKey));
assert!(
!err.to_string().contains(bad.trim()) || bad.trim().is_empty(),
"error must not echo the key"
);
}
}
fn dead_hub() -> HubConfig {
HubConfig {
hub: "http://127.0.0.1:9".to_string(),
key: Some("k".to_string()),
}
}
#[test]
fn verb_entry_gates_accept_the_hub_ref_shapes() {
for ok in ["acme-ops", "a", "01j5qc3v9k4ym8rwbn2tqe6f7d"] {
assert!(require_safe_ref(ok).is_ok(), "brain ref {ok:?}");
assert!(require_valid_handle(ok).is_ok(), "handle {ok:?}");
assert!(require_safe_grant_id(ok).is_ok(), "grant id {ok:?}");
}
}
#[test]
fn raw_ref_verbs_refuse_url_reshaping_brain_refs_before_any_request() {
let cfg = dead_hub();
for bad in ["../up", "a/b", "a?x=1", "a#frag", "a%2e%2e", "A", "a b", ""] {
assert!(
matches!(
sync_pull(&cfg, bad, None),
Err(LinkError::BadAddress { .. })
),
"sync_pull must refuse {bad:?}"
);
assert!(
matches!(sync_push(&cfg, bad, &[]), Err(LinkError::BadAddress { .. })),
"sync_push must refuse {bad:?}"
);
assert!(
matches!(
grant_issue(&cfg, bad, "maya@example.com", Capability::Read, None, None),
Err(LinkError::BadAddress { .. })
),
"grant_issue must refuse {bad:?}"
);
assert!(
matches!(grant_list(&cfg, bad), Err(LinkError::BadAddress { .. })),
"grant_list must refuse {bad:?}"
);
assert!(
matches!(
grant_revoke(&cfg, bad, "01j5qc3v9k4ym8rwbn2tqe6f7f"),
Err(LinkError::BadAddress { .. })
),
"grant_revoke must refuse brain {bad:?}"
);
assert!(
matches!(head(&cfg, bad), Err(LinkError::BadAddress { .. })),
"head must refuse {bad:?}"
);
}
}
#[test]
fn grant_revoke_refuses_url_reshaping_grant_ids() {
let cfg = dead_hub();
for bad in ["../01j", "a/b", "id?x=1", "id#frag", "ID", ""] {
assert!(
matches!(
grant_revoke(&cfg, "acme", bad),
Err(LinkError::BadGrantId { .. })
),
"grant_revoke must refuse grant id {bad:?}"
);
}
}
#[test]
fn propose_refuses_url_reshaping_handles_and_oversize_bodies_before_upload() {
let cfg = dead_hub();
for bad in ["../up", "a/b", "a?x=1", "a#frag", "A", ""] {
assert!(
matches!(
propose(&cfg, bad, "intake", "hi"),
Err(LinkError::BadAddress { .. })
),
"propose must refuse handle {bad:?}"
);
}
let oversize = "a".repeat(MAX_PROPOSE_BYTES as usize + 1);
assert!(matches!(
propose(&cfg, "acme-site", "intake", &oversize),
Err(LinkError::ProposeTooLarge { .. })
));
assert!(matches!(
propose(&cfg, "acme-site", "intake", "hi"),
Err(LinkError::Transport { .. })
));
}
#[test]
fn resolve_refuses_a_hand_built_unsafe_address() {
let cfg = dead_hub();
for brain in ["../up", "a/b", "a?x", "a#f"] {
let addr = Address {
brain: brain.to_string(),
target: None,
};
assert!(
matches!(resolve(&cfg, &addr), Err(LinkError::BadAddress { .. })),
"resolve must refuse brain {brain:?}"
);
}
for target in [
AddressTarget::Id("01j5qc3v9k4ym8rwbn2tqe6f7d?id=other".to_string()),
AddressTarget::Id("01J5QC3V9K4YM8RWBN2TQE6F7D".to_string()), AddressTarget::Path("../up.md".to_string()),
AddressTarget::Path("records/x.md#frag".to_string()),
] {
let addr = Address {
brain: "acme".to_string(),
target: Some(target.clone()),
};
assert!(
matches!(resolve(&cfg, &addr), Err(LinkError::BadAddress { .. })),
"resolve must refuse target {target:?}"
);
}
}
}