diff --git a/crates/credentials-module/src/bin/cli_support/opencode_files.rs b/crates/credentials-module/src/bin/cli_support/opencode_files.rs index 09c6cc1..be5727d 100644 --- a/crates/credentials-module/src/bin/cli_support/opencode_files.rs +++ b/crates/credentials-module/src/bin/cli_support/opencode_files.rs @@ -4,15 +4,115 @@ use std::{ fs::{self, File, OpenOptions}, io::Write, path::{Path, PathBuf}, - sync::atomic::{AtomicU64, Ordering}, + sync::{ + atomic::{AtomicBool, AtomicU64, Ordering}, + mpsc, Arc, + }, + thread, + time::{Duration, Instant, SystemTime, UNIX_EPOCH}, }; +use ring::rand::{SecureRandom, SystemRandom}; use serde::{Deserialize, Serialize}; use serde_json::Value; static TEMP_SEQ: AtomicU64 = AtomicU64::new(0); const AUTH_FILE_MAX_BYTES: u64 = 1024 * 1024; const HANDLE_FILE_MAX_BYTES: u64 = 256 * 1024; +const MANIFEST_LOCK_TTL_MS: u64 = 30_000; +const MANIFEST_LOCK_RENEW_EVERY_MS: u64 = 10_000; +const MANIFEST_LOCK_OWNER_KEYS: [&str; 4] = ["tenant", "pid", "claimed_at_ms", "nonce"]; +const MANIFEST_LOCK_STALE_TARGET_PATTERN: &str = r"^\.lock\.stale-\d+-[A-Za-z0-9_-]+$"; +const OPENCODE_CLAUSTRUM_TENANT: &str = "opencode-claustrum"; +type BeforeManifestRename = Arc; + +#[cfg(test)] +static LEASE_LOST_WARNINGS: AtomicU64 = AtomicU64::new(0); + +#[derive(Clone)] +struct ManifestLockOptions { + ttl: Duration, + renew_every: Duration, + retry_min: Duration, + retry_max: Duration, + after_claim: Option>, + before_evict: Option>, + #[cfg(test)] + after_evict_rename_attempt: Option>, + after_evict: Option>, + before_manifest_rename: Option, + // Manifest lock staleness is judged against the contender's clock at each observation (not at claim start); claim deadline expiry remains monotonic so the production bound is still exercised. + now_override_ms: Option, + #[cfg(test)] + now_sequence_ms: Option>, +} + +impl Default for ManifestLockOptions { + fn default() -> Self { + Self { + ttl: Duration::from_millis(MANIFEST_LOCK_TTL_MS), + renew_every: Duration::from_millis(MANIFEST_LOCK_RENEW_EVERY_MS), + retry_min: Duration::from_millis(25), + retry_max: Duration::from_millis(75), + after_claim: None, + before_evict: None, + #[cfg(test)] + after_evict_rename_attempt: None, + after_evict: None, + before_manifest_rename: None, + now_override_ms: None, + #[cfg(test)] + now_sequence_ms: None, + } + } +} + +struct ManifestLease { + lock: PathBuf, + nonce: String, + ttl: Duration, + renewal_failed: Arc, + stop_tx: Option>, + renewal: Option>, +} + +impl ManifestLease { + fn stop_renewal(&mut self) { + if let Some(stop_tx) = self.stop_tx.take() { + let _ = stop_tx.send(()); + } + if let Some(renewal) = self.renewal.take() { + let _ = renewal.join(); + } + } + + fn commit(&mut self) -> Result<(), OpenCodeFilesError> { + self.stop_renewal(); + let owner = read_lock_owner(&self.lock.join("owner")).ok(); + let ours_and_fresh = owner.is_some_and(|owner| { + owner.nonce == self.nonce + && current_time_ms().is_ok_and(|now| { + now.saturating_sub(owner.claimed_at_ms) < self.ttl.as_millis() as u64 + }) + }); + if self.renewal_failed.load(Ordering::SeqCst) || !ours_and_fresh { + return Err(OpenCodeFilesError::Invalid( + "manifest lock renewal failed; write aborted".into(), + )); + } + Ok(()) + } +} + +#[derive(Clone, Serialize, Deserialize)] +struct ManifestLockOwner { + #[serde(default)] + tenant: Value, + #[serde(default)] + pid: Value, + claimed_at_ms: u64, + nonce: String, +} #[derive(Debug)] pub enum OpenCodeFilesError { @@ -226,20 +326,430 @@ pub fn read_handle_file(path: &Path) -> Result { } pub fn write_handle_file(path: &Path, file: &HandleFile) -> Result<(), OpenCodeFilesError> { - validate_handle_file(file)?; - let bytes = serde_json::to_vec(file).map_err(OpenCodeFilesError::Json)?; - write_atomic(path, &bytes, true) + write_handle_file_for_tenant( + path, + OPENCODE_CLAUSTRUM_TENANT, + file, + ManifestLockOptions::default(), + ) } pub fn verify_handle_written(path: &Path, expected: &HandleFile) -> Result<(), OpenCodeFilesError> { - if &read_handle_file(path)? != expected { + validate_handle_file(expected)?; + let written = read_handle_file(path)?; + let expected_owned: Vec<_> = expected + .providers + .iter() + .filter(|provider| provider.serve == OPENCODE_CLAUSTRUM_TENANT) + .cloned() + .collect(); + let written_owned: Vec<_> = written + .providers + .iter() + .filter(|provider| provider.serve == OPENCODE_CLAUSTRUM_TENANT) + .cloned() + .collect(); + if written_owned != expected_owned { return Err(OpenCodeFilesError::Invalid( - "handle file did not persist exactly".into(), + "handle file tenant block did not persist exactly".into(), )); } Ok(()) } +fn write_handle_file_for_tenant( + path: &Path, + tenant: &str, + desired: &HandleFile, + options: ManifestLockOptions, +) -> Result<(), OpenCodeFilesError> { + validate_handle_file(desired)?; + let parent = path + .parent() + .filter(|parent| !parent.as_os_str().is_empty()) + .ok_or_else(|| OpenCodeFilesError::Invalid("file path has no parent".into()))?; + fs::create_dir_all(parent).map_err(|source| io_error("create parent directory", source))?; + validate_secure_parent(parent)?; + set_mode(parent, 0o700)?; + let before_manifest_rename = options.before_manifest_rename.clone(); + with_manifest_lock_with_options(path, tenant, options, |lease| { + let current = match fs::symlink_metadata(path) { + Ok(_) => read_handle_file(path)?, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => HandleFile { + version: 1, + providers: Vec::new(), + }, + Err(error) => return Err(io_error("stat handle file", error)), + }; + let before_foreign: Vec> = current + .providers + .iter() + .filter(|provider| provider.serve != tenant) + .map(serde_json::to_vec) + .collect::>() + .map_err(OpenCodeFilesError::Json)?; + let mut providers: Vec<_> = current + .providers + .into_iter() + .filter(|provider| provider.serve != tenant) + .collect(); + providers.extend( + desired + .providers + .iter() + .filter(|provider| provider.serve == tenant) + .cloned(), + ); + let next = HandleFile { + version: 1, + providers, + }; + validate_handle_file(&next)?; + let bytes = serde_json::to_vec(&next).map_err(OpenCodeFilesError::Json)?; + write_atomic_guarded(path, &bytes, true, || { + if let Some(before_manifest_rename) = &before_manifest_rename { + before_manifest_rename(&lock_path(path)); + } + lease.commit() + })?; + let readback = read_handle_file(path)?; + if readback != next { + return Err(OpenCodeFilesError::Invalid( + "handle file readback did not persist exactly".into(), + )); + } + let after_foreign: Vec> = readback + .providers + .iter() + .filter(|provider| provider.serve != tenant) + .map(serde_json::to_vec) + .collect::>() + .map_err(OpenCodeFilesError::Json)?; + if after_foreign != before_foreign { + return Err(OpenCodeFilesError::Invalid( + "handle file readback changed another tenant block".into(), + )); + } + Ok(()) + }) +} + +fn lock_path(path: &Path) -> PathBuf { + let mut lock = path.as_os_str().to_os_string(); + lock.push(".lock"); + PathBuf::from(lock) +} + +fn current_time_ms() -> Result { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|duration| duration.as_millis() as u64) + .map_err(|_| OpenCodeFilesError::Invalid("system clock is before UNIX epoch".into())) +} + +fn resolve_now_ms(options: &ManifestLockOptions) -> Result { + #[cfg(test)] + if let Some(clock) = &options.now_sequence_ms { + return Ok(clock.load(Ordering::SeqCst)); + } + match options.now_override_ms { + Some(fixed) => Ok(fixed), + None => current_time_ms(), + } +} + +fn random_nonce() -> Result { + let mut bytes = [0_u8; 16]; + SystemRandom::new() + .fill(&mut bytes) + .map_err(|_| OpenCodeFilesError::Invalid("generate manifest lock nonce failed".into()))?; + Ok(bytes.iter().map(|byte| format!("{byte:02x}")).collect()) +} + +fn io_error(action: &'static str, source: std::io::Error) -> OpenCodeFilesError { + OpenCodeFilesError::Io { action, source } +} + +fn read_lock_owner(path: &Path) -> Result { + let source = + fs::read_to_string(path).map_err(|source| io_error("read manifest lock owner", source))?; + let owner: ManifestLockOwner = serde_json::from_str(&source) + .map_err(|_| OpenCodeFilesError::Invalid("manifest lock owner invalid".into()))?; + let stale_target = format!(".lock.stale-{}-{}", owner.claimed_at_ms, owner.nonce); + if !stale_target_matches(&stale_target) { + return Err(OpenCodeFilesError::Invalid( + "manifest lock owner invalid".into(), + )); + } + Ok(owner) +} + +fn write_lock_owner(lock: &Path, owner: &ManifestLockOwner) -> Result<(), OpenCodeFilesError> { + let owner_path = lock.join("owner"); + let temporary = lock.join(format!( + "owner.{}.{}.tmp", + std::process::id(), + random_nonce()? + )); + let result = (|| -> Result<(), OpenCodeFilesError> { + #[cfg(unix)] + let mut file = { + use std::os::unix::fs::OpenOptionsExt; + OpenOptions::new() + .write(true) + .create_new(true) + .mode(0o600) + .open(&temporary) + .map_err(|source| io_error("create manifest lock owner", source))? + }; + #[cfg(not(unix))] + let mut file = OpenOptions::new() + .write(true) + .create_new(true) + .open(&temporary) + .map_err(|source| io_error("create manifest lock owner", source))?; + set_mode(&temporary, 0o600)?; + serde_json::to_writer(&mut file, owner).map_err(OpenCodeFilesError::Json)?; + file.write_all(b"\n") + .map_err(|source| io_error("write manifest lock owner", source))?; + file.sync_all() + .map_err(|source| io_error("sync manifest lock owner", source))?; + drop(file); + fs::rename(&temporary, &owner_path) + .map_err(|source| io_error("rename manifest lock owner", source)) + })(); + if result.is_err() { + let _ = fs::remove_file(&temporary); + } + result +} + +fn stale_target_matches(value: &str) -> bool { + let Some(rest) = value.strip_prefix(".lock.stale-") else { + return false; + }; + let Some((claimed, random)) = rest.split_once('-') else { + return false; + }; + !claimed.is_empty() + && claimed.bytes().all(|byte| byte.is_ascii_digit()) + && !random.is_empty() + && random + .bytes() + .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-')) +} + +fn warn_lease_lost(path: &Path) { + #[cfg(test)] + LEASE_LOST_WARNINGS.fetch_add(1, Ordering::SeqCst); + eprintln!( + "manifest lock lease lost, not releasing: {}", + path.display() + ); +} + +fn jitter(options: &ManifestLockOptions) -> Duration { + let min = options.retry_min.as_millis() as u64; + let max = options.retry_max.as_millis() as u64; + if max <= min { + return Duration::from_millis(min); + } + let mut bytes = [0_u8; 8]; + if SystemRandom::new().fill(&mut bytes).is_err() { + return Duration::from_millis(min); + } + Duration::from_millis(min + u64::from_le_bytes(bytes) % (max - min + 1)) +} + +fn release_manifest_lock( + path: &Path, + lock: &Path, + nonce: &str, + ttl: Duration, +) -> Result<(), OpenCodeFilesError> { + let owner = match read_lock_owner(&lock.join("owner")) { + Ok(owner) => owner, + Err(_) => { + warn_lease_lost(path); + return Ok(()); + } + }; + let now = current_time_ms()?; + if owner.nonce != nonce || now.saturating_sub(owner.claimed_at_ms) >= ttl.as_millis() as u64 { + warn_lease_lost(path); + return Ok(()); + } + let release = PathBuf::from(format!("{}.release-{nonce}", lock.display())); + if fs::rename(lock, &release).is_err() { + warn_lease_lost(path); + return Ok(()); + } + let moved = read_lock_owner(&release.join("owner")).ok(); + let moved_is_ours = moved.is_some_and(|owner| { + owner.nonce == nonce + && current_time_ms() + .is_ok_and(|now| now.saturating_sub(owner.claimed_at_ms) < ttl.as_millis() as u64) + }); + if !moved_is_ours { + let _ = fs::rename(&release, lock); + warn_lease_lost(path); + return Ok(()); + } + fs::remove_dir_all(&release).map_err(|source| io_error("remove manifest lock", source)) +} + +fn with_manifest_lock_with_options( + path: &Path, + tenant: &str, + options: ManifestLockOptions, + operation: F, +) -> Result +where + F: FnOnce(&mut ManifestLease) -> Result, +{ + let lock = lock_path(path); + let owner_path = lock.join("owner"); + let nonce = random_nonce()?; + let deadline = Instant::now() + options.ttl; + loop { + match fs::create_dir(&lock) { + Ok(()) => { + set_mode(&lock, 0o700)?; + let owner = ManifestLockOwner { + tenant: Value::String(tenant.into()), + pid: Value::from(std::process::id()), + claimed_at_ms: resolve_now_ms(&options)?, + nonce: nonce.clone(), + }; + if let Err(error) = write_lock_owner(&lock, &owner) { + let _ = fs::remove_dir_all(&lock); + return Err(error); + } + if let Some(after_claim) = &options.after_claim { + after_claim(); + } + break; + } + Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {} + Err(error) => return Err(io_error("create manifest lock", error)), + } + + let owner_read_error = match read_lock_owner(&owner_path) { + Ok(observed) => { + if resolve_now_ms(&options)?.saturating_sub(observed.claimed_at_ms) + >= options.ttl.as_millis() as u64 + { + if let Some(before_evict) = &options.before_evict { + before_evict(); + } + let stale = PathBuf::from(format!( + "{}.stale-{}-{}", + lock.display(), + observed.claimed_at_ms, + observed.nonce + )); + let rename_result = fs::rename(&lock, &stale); + #[cfg(test)] + if let Some(after_evict_rename_attempt) = &options.after_evict_rename_attempt { + after_evict_rename_attempt(); + } + match rename_result { + Ok(()) => { + let moved = read_lock_owner(&stale.join("owner")).ok(); + if moved.is_some_and(|owner| { + owner.nonce == observed.nonce + && owner.claimed_at_ms == observed.claimed_at_ms + }) { + if let Some(after_evict) = &options.after_evict { + after_evict(); + } + continue; + } + let _ = fs::rename(&stale, &lock); + } + Err(error) + if matches!( + error.kind(), + std::io::ErrorKind::NotFound + | std::io::ErrorKind::AlreadyExists + | std::io::ErrorKind::DirectoryNotEmpty + ) => {} + Err(error) => return Err(io_error("rename stale manifest lock", error)), + } + } + None + } + Err(error) => Some(error), + }; + if Instant::now() >= deadline { + if matches!(owner_read_error, Some(OpenCodeFilesError::Invalid(_))) { + return Err(OpenCodeFilesError::Invalid( + "manifest lock owner invalid".into(), + )); + } + return Err(OpenCodeFilesError::Invalid("manifest lock busy".into())); + } + thread::sleep(jitter(&options).min(deadline.saturating_duration_since(Instant::now()))); + } + + let (stop_tx, stop_rx) = mpsc::channel::<()>(); + let renewal_lock = lock.clone(); + let renewal_nonce = nonce.clone(); + let renewal_ttl = options.ttl; + let renewal_every = options.renew_every; + let renewal_failed = Arc::new(AtomicBool::new(false)); + let renewal_failed_thread = Arc::clone(&renewal_failed); + let renewal = thread::spawn(move || loop { + match stop_rx.recv_timeout(renewal_every) { + Ok(()) | Err(mpsc::RecvTimeoutError::Disconnected) => break, + Err(mpsc::RecvTimeoutError::Timeout) => { + let owner_path = renewal_lock.join("owner"); + let Ok(mut owner) = read_lock_owner(&owner_path) else { + renewal_failed_thread.store(true, Ordering::SeqCst); + break; + }; + let Ok(now) = current_time_ms() else { + renewal_failed_thread.store(true, Ordering::SeqCst); + break; + }; + if owner.nonce != renewal_nonce + || now.saturating_sub(owner.claimed_at_ms) >= renewal_ttl.as_millis() as u64 + { + renewal_failed_thread.store(true, Ordering::SeqCst); + break; + } + owner.claimed_at_ms = now; + if write_lock_owner(&renewal_lock, &owner).is_err() { + renewal_failed_thread.store(true, Ordering::SeqCst); + break; + } + } + } + }); + let mut lease = ManifestLease { + lock: lock.clone(), + nonce: nonce.clone(), + ttl: options.ttl, + renewal_failed, + stop_tx: Some(stop_tx), + renewal: Some(renewal), + }; + let result = operation(&mut lease); + lease.stop_renewal(); + let result = match result { + Ok(_) if lease.renewal_failed.load(Ordering::SeqCst) => Err(OpenCodeFilesError::Invalid( + "manifest lock renewal failed; write aborted".into(), + )), + other => other, + }; + let release = release_manifest_lock(path, &lock, &nonce, options.ttl); + match (result, release) { + (Err(error), _) => Err(error), + (Ok(_), Err(error)) => Err(error), + (Ok(value), Ok(())) => Ok(value), + } +} + fn validate_auth_entry(entry: &Value) -> Result<(), OpenCodeFilesError> { let object = entry .as_object() @@ -343,6 +853,18 @@ fn validate_identifier(value: &str, kind: &str) -> Result<(), OpenCodeFilesError } fn write_atomic(path: &Path, bytes: &[u8], secure_parent: bool) -> Result<(), OpenCodeFilesError> { + write_atomic_guarded(path, bytes, secure_parent, || Ok(())) +} + +fn write_atomic_guarded( + path: &Path, + bytes: &[u8], + secure_parent: bool, + before_rename: F, +) -> Result<(), OpenCodeFilesError> +where + F: FnOnce() -> Result<(), OpenCodeFilesError>, +{ let parent = path .parent() .filter(|parent| !parent.as_os_str().is_empty()) @@ -397,6 +919,7 @@ fn write_atomic(path: &Path, bytes: &[u8], secure_parent: bool) -> Result<(), Op action: "sync temporary file", source, })?; + before_rename()?; fs::rename(&temp, path).map_err(|source| OpenCodeFilesError::Io { action: "rename temporary file", source, @@ -533,3 +1056,333 @@ fn current_uid() -> Result { .map_err(|_| OpenCodeFilesError::Invalid("current uid was invalid".into())) }) } + +#[cfg(test)] +mod manifest_lock_aba_regression { + use super::*; + use std::{ + os::unix::fs::PermissionsExt, + sync::{Arc, Barrier}, + thread, + time::{Duration, SystemTime, UNIX_EPOCH}, + }; + + fn now_ms() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_millis() as u64 + } + + #[test] + fn aba_observation_cannot_rename_a_replacement_lock() { + let root = std::env::temp_dir().join(format!( + "claustrum-manifest-lock-aba-{}-{}", + std::process::id(), + TEMP_SEQ.fetch_add(1, Ordering::Relaxed) + )); + fs::create_dir_all(&root).unwrap(); + let path = root.join("opencode-handles.json"); + let lock = lock_path(&path); + let now = now_ms(); + fs::create_dir(&lock).unwrap(); + fs::set_permissions(&lock, fs::Permissions::from_mode(0o700)).unwrap(); + fs::write( + lock.join("owner"), + format!( + "{{\"tenant\":\"other-tenant\",\"pid\":41,\"claimed_at_ms\":{},\"nonce\":\"0123456789abcdef0123456789abcdef\"}}\n", + now - 501 + ), + ) + .unwrap(); + fs::set_permissions(lock.join("owner"), fs::Permissions::from_mode(0o600)).unwrap(); + + let loser_observed = Arc::new(Barrier::new(2)); + let allow_loser_rename = Arc::new(Barrier::new(2)); + let replacement_claimed = Arc::new(Barrier::new(2)); + let allow_replacement_release = Arc::new(Barrier::new(2)); + let rename_attempted = Arc::new(Barrier::new(2)); + let allow_attempt_completion = Arc::new(Barrier::new(2)); + + let loser_path = path.clone(); + let loser = thread::spawn({ + let loser_observed = Arc::clone(&loser_observed); + let allow_loser_rename = Arc::clone(&allow_loser_rename); + let rename_attempted = Arc::clone(&rename_attempted); + let allow_attempt_completion = Arc::clone(&allow_attempt_completion); + move || { + with_manifest_lock_with_options( + &loser_path, + "loser", + ManifestLockOptions { + ttl: Duration::from_millis(500), + renew_every: Duration::from_secs(1), + retry_min: Duration::from_millis(2), + retry_max: Duration::from_millis(3), + before_evict: Some(Arc::new(move || { + loser_observed.wait(); + allow_loser_rename.wait(); + })), + after_evict_rename_attempt: Some(Arc::new(move || { + rename_attempted.wait(); + allow_attempt_completion.wait(); + })), + now_override_ms: Some(now), + ..ManifestLockOptions::default() + }, + |_| Ok(()), + ) + } + }); + + loser_observed.wait(); + let replacement_path = path.clone(); + let replacement = thread::spawn({ + let replacement_claimed = Arc::clone(&replacement_claimed); + let allow_replacement_release = Arc::clone(&allow_replacement_release); + move || { + with_manifest_lock_with_options( + &replacement_path, + "replacement", + ManifestLockOptions { + ttl: Duration::from_millis(500), + renew_every: Duration::from_secs(1), + retry_min: Duration::from_millis(2), + retry_max: Duration::from_millis(3), + after_claim: Some(Arc::new(move || { + replacement_claimed.wait(); + allow_replacement_release.wait(); + })), + now_override_ms: Some(now), + ..ManifestLockOptions::default() + }, + |_| Ok(()), + ) + } + }); + + replacement_claimed.wait(); + allow_loser_rename.wait(); + rename_attempted.wait(); + allow_replacement_release.wait(); + replacement.join().unwrap().unwrap(); + allow_attempt_completion.wait(); + loser.join().unwrap().unwrap(); + assert!(!lock.exists()); + let _ = fs::remove_dir_all(root); + } + + #[test] + fn two_stale_evictors_create_exactly_one_quarantine_directory() { + let root = std::env::temp_dir().join(format!( + "claustrum-manifest-lock-quarantine-{}-{}", + std::process::id(), + TEMP_SEQ.fetch_add(1, Ordering::Relaxed) + )); + fs::create_dir_all(&root).unwrap(); + let path = root.join("opencode-handles.json"); + let lock = lock_path(&path); + let now = now_ms(); + fs::create_dir(&lock).unwrap(); + fs::set_permissions(&lock, fs::Permissions::from_mode(0o700)).unwrap(); + fs::write( + lock.join("owner"), + format!( + "{{\"tenant\":\"other-tenant\",\"pid\":41,\"claimed_at_ms\":{},\"nonce\":\"0123456789abcdef0123456789abcdef\"}}\n", + now - 501 + ), + ) + .unwrap(); + fs::set_permissions(lock.join("owner"), fs::Permissions::from_mode(0o600)).unwrap(); + let ready = Arc::new(Barrier::new(2)); + let evictions = Arc::new(AtomicU64::new(0)); + let mut joins = Vec::new(); + for tenant in ["anthropic-auth", "openai-auth"] { + let path = path.clone(); + let ready = Arc::clone(&ready); + let evictions = Arc::clone(&evictions); + joins.push(thread::spawn(move || { + with_manifest_lock_with_options( + &path, + tenant, + ManifestLockOptions { + ttl: Duration::from_millis(500), + renew_every: Duration::from_secs(1), + retry_min: Duration::from_millis(2), + retry_max: Duration::from_millis(3), + before_evict: Some(Arc::new(move || { + ready.wait(); + })), + after_evict: Some(Arc::new(move || { + evictions.fetch_add(1, Ordering::SeqCst); + })), + now_override_ms: Some(now), + ..ManifestLockOptions::default() + }, + |_| Ok(()), + ) + })); + } + for join in joins { + join.join().unwrap().unwrap(); + } + let stale = fs::read_dir(&root) + .unwrap() + .filter_map(Result::ok) + .filter(|entry| entry.file_name().to_string_lossy().contains(".lock.stale-")) + .count(); + assert_eq!(evictions.load(Ordering::SeqCst), 1); + assert_eq!(stale, 1); + let _ = fs::remove_dir_all(root); + } + + fn seed_owner(path: &Path, owner: &str) -> PathBuf { + let lock = lock_path(path); + fs::create_dir(&lock).unwrap(); + fs::set_permissions(&lock, fs::Permissions::from_mode(0o700)).unwrap(); + fs::write(lock.join("owner"), owner).unwrap(); + fs::set_permissions(lock.join("owner"), fs::Permissions::from_mode(0o600)).unwrap(); + lock + } + + #[test] + fn unknown_owner_keys_are_tolerated_and_evictable_once_stale() { + let root = std::env::temp_dir().join(format!( + "claustrum-manifest-lock-unknown-key-{}-{}", + std::process::id(), + TEMP_SEQ.fetch_add(1, Ordering::Relaxed) + )); + fs::create_dir_all(&root).unwrap(); + let path = root.join("opencode-handles.json"); + let now = now_ms(); + seed_owner( + &path, + &format!( + "{{\"tenant\":\"other\",\"pid\":41,\"claimed_at_ms\":{},\"nonce\":\"0123456789abcdef0123456789abcdef\",\"host\":\"x\"}}\n", + now - 501 + ), + ); + let result = with_manifest_lock_with_options( + &path, + "claimant", + ManifestLockOptions { + ttl: Duration::from_millis(500), + now_override_ms: Some(now), + ..ManifestLockOptions::default() + }, + |_| Ok(()), + ); + assert!(result.is_ok()); + assert!(!lock_path(&path).exists()); + let _ = fs::remove_dir_all(root); + } + + #[test] + fn malformed_diagnostic_owner_fields_are_tolerated_and_evictable_once_stale() { + let root = std::env::temp_dir().join(format!( + "claustrum-manifest-lock-malformed-diagnostic-{}-{}", + std::process::id(), + TEMP_SEQ.fetch_add(1, Ordering::Relaxed) + )); + fs::create_dir_all(&root).unwrap(); + let path = root.join("opencode-handles.json"); + let now = now_ms(); + seed_owner( + &path, + &format!( + "{{\"pid\":\"not-a-number\",\"claimed_at_ms\":{},\"nonce\":\"0123456789abcdef0123456789abcdef\"}}\n", + now - 501 + ), + ); + let result = with_manifest_lock_with_options( + &path, + "claimant", + ManifestLockOptions { + ttl: Duration::from_millis(500), + now_override_ms: Some(now), + ..ManifestLockOptions::default() + }, + |_| Ok(()), + ); + assert!(result.is_ok()); + assert!(!lock_path(&path).exists()); + let _ = fs::remove_dir_all(root); + } + + #[test] + fn missing_owner_nonce_fails_with_owner_invalid_at_deadline() { + let root = std::env::temp_dir().join(format!( + "claustrum-manifest-lock-owner-invalid-{}-{}", + std::process::id(), + TEMP_SEQ.fetch_add(1, Ordering::Relaxed) + )); + fs::create_dir_all(&root).unwrap(); + let path = root.join("opencode-handles.json"); + let now = now_ms(); + seed_owner( + &path, + &format!( + "{{\"tenant\":\"other\",\"pid\":41,\"claimed_at_ms\":{}}}\n", + now - 501 + ), + ); + let result = with_manifest_lock_with_options( + &path, + "claimant", + ManifestLockOptions { + ttl: Duration::from_millis(20), + retry_min: Duration::from_millis(2), + retry_max: Duration::from_millis(3), + now_override_ms: Some(now), + ..ManifestLockOptions::default() + }, + |_| Ok(()), + ); + assert_eq!( + result.unwrap_err().to_string(), + "manifest lock owner invalid" + ); + let _ = fs::remove_dir_all(root); + } + + #[test] + fn owner_that_becomes_stale_during_retry_window_is_evicted() { + let root = std::env::temp_dir().join(format!( + "claustrum-manifest-lock-observation-clock-{}-{}", + std::process::id(), + TEMP_SEQ.fetch_add(1, Ordering::Relaxed) + )); + fs::create_dir_all(&root).unwrap(); + let path = root.join("opencode-handles.json"); + let now = now_ms(); + seed_owner( + &path, + &format!( + "{{\"tenant\":\"other\",\"pid\":41,\"claimed_at_ms\":{},\"nonce\":\"0123456789abcdef0123456789abcdef\"}}\n", + now - 80 + ), + ); + let clock = Arc::new(AtomicU64::new(now)); + let advancing_clock = Arc::clone(&clock); + let advance = thread::spawn(move || { + thread::sleep(Duration::from_millis(20)); + advancing_clock.store(now + 100, Ordering::SeqCst); + }); + let result = with_manifest_lock_with_options( + &path, + "claimant", + ManifestLockOptions { + ttl: Duration::from_millis(100), + retry_min: Duration::from_millis(50), + retry_max: Duration::from_millis(50), + now_sequence_ms: Some(clock), + ..ManifestLockOptions::default() + }, + |_| Ok(()), + ); + advance.join().unwrap(); + assert!(result.is_ok()); + assert!(!lock_path(&path).exists()); + let _ = fs::remove_dir_all(root); + } +} diff --git a/packages/client/README.md b/packages/client/README.md index 0ecdcfb..3ad6b97 100644 --- a/packages/client/README.md +++ b/packages/client/README.md @@ -21,3 +21,10 @@ credential caching, refresh, account scheduling, and retry policy. The client sends `consumerIdentity: null` for every managed request so inherited `SUBC_MODULE_ID` and `SUBC_LAUNCH_NONCE` cannot impersonate a supervising host. + +## OpenCode handle-file exports + +`defaultHandleFilePath`, `readHandleFile`, `handleFileRevision`, and `parseHandleFile` +are exported for tenant plugins consuming the OpenCode handle manifest. The shared +`HANDLE_FILE_CONTRACT` is pinned to `maxBytes: 262144`, `mode: 0o600`, +`labelRe: /^[a-z0-9][a-z0-9._-]{0,63}$/`, and `handleRe: /^ckh_[A-Za-z0-9_-]{43}$/`. diff --git a/packages/client/src/handles.ts b/packages/client/src/handles.ts new file mode 100644 index 0000000..836e93b --- /dev/null +++ b/packages/client/src/handles.ts @@ -0,0 +1,223 @@ +import { createHash } from 'node:crypto' +import { constants } from 'node:fs' +import { lstat as nodeLstat, open as nodeOpen, readFile, stat as nodeStat } from 'node:fs/promises' +import { userInfo } from 'node:os' +import { dirname, join } from 'node:path' + +export const HANDLE_FILE_CONTRACT = { + maxBytes: 256 * 1024, + mode: 0o600, + labelRe: /^[a-z0-9][a-z0-9._-]{0,63}$/, + handleRe: /^ckh_[A-Za-z0-9_-]{43}$/, +} as const + +const FORBIDDEN_IDENTIFIERS = new Set(['__proto__', 'constructor', 'prototype']) + +export class HandleFileValidationError extends Error { + override name = 'HandleFileValidationError' +} + +export type HandleAccount = { + label: string + handle: string + credential_id: string + superseded?: string[] +} +export type HandleProvider = { + provider: string + shape: 'api' | 'oauth' + serve: string + accounts: HandleAccount[] +} +export type OpenCodeHandleFileV1 = { version: 1; providers: HandleProvider[] } + +function isAccount(value: unknown): value is HandleAccount { + if (!value || typeof value !== 'object') return false + const account = value as Record + return typeof account.label === 'string' && typeof account.handle === 'string' && + typeof account.credential_id === 'string' && + (account.superseded === undefined || + (Array.isArray(account.superseded) && account.superseded.every((handle) => typeof handle === 'string'))) +} + +function handleIsValid(handle: unknown): handle is string { + return typeof handle === 'string' && HANDLE_FILE_CONTRACT.handleRe.test(handle) +} + +function identifierIsValid(value: unknown): value is string { + return typeof value === 'string' && HANDLE_FILE_CONTRACT.labelRe.test(value) && !FORBIDDEN_IDENTIFIERS.has(value) +} + +function invalid(message: string): never { + throw new HandleFileValidationError(message) +} + +export function parseHandleFile(value: unknown): OpenCodeHandleFileV1 { + if (!value || typeof value !== 'object') invalid('handle file must be an object') + const file = value as Record + if (file.version !== 1 || !Array.isArray(file.providers)) { + invalid('handle file must have version 1 and providers') + } + const providerIds = new Set() + const providers = file.providers.map((provider, index): HandleProvider => { + if (!provider || typeof provider !== 'object') invalid(`provider ${index} must be an object`) + const item = provider as Record + if (!identifierIsValid(item.provider)) invalid(`provider ${index} has invalid provider`) + if (providerIds.has(item.provider)) invalid(`provider ${index} duplicates provider ${item.provider}`) + providerIds.add(item.provider) + if (item.shape !== 'api' && item.shape !== 'oauth') invalid(`provider ${index} has invalid shape`) + if (typeof item.serve !== 'string' || !item.serve) invalid(`provider ${index} requires serve`) + if (!Array.isArray(item.accounts) || item.accounts.length === 0 || !item.accounts.every(isAccount)) { + invalid(`provider ${index} has invalid accounts`) + } + const labels = new Set() + for (const account of item.accounts) { + if (!identifierIsValid(account.label)) invalid(`provider ${index} has an invalid account label`) + if (labels.has(account.label)) invalid(`provider ${index} duplicates account label ${account.label}`) + labels.add(account.label) + if (!handleIsValid(account.handle)) invalid(`provider ${index} account ${account.label} has invalid handle`) + if (!account.credential_id) invalid(`provider ${index} account ${account.label} has invalid credential id`) + if (account.superseded?.some((handle) => !handleIsValid(handle))) { + invalid(`provider ${index} account ${account.label} has invalid superseded handle`) + } + } + return { + provider: item.provider, + shape: item.shape, + serve: item.serve, + accounts: item.accounts.map((account) => ({ + ...account, + ...(account.superseded === undefined ? {} : { superseded: account.superseded }), + })), + } + }) + return { version: 1, providers } +} + +type HandleFileStat = { + isFile(): boolean + isDirectory?(): boolean + isSymbolicLink?(): boolean + mode: number + size?: number + uid?: number + mtimeMs?: number +} +type HandleFileDescriptor = { + read?(buffer: Uint8Array, offset: number, length: number, position: number): Promise<{ bytesRead: number }> | { bytesRead: number } + stat(): Promise + readFile(options: { encoding: 'utf8' }): Promise + close(): Promise +} +export type HandleFileIo = { + stat?: (path: string) => Promise + lstat?: (path: string) => Promise + readFile?: (path: string, encoding: 'utf8') => Promise + open?: (path: string) => Promise + currentUid?: () => number | undefined +} + +export function defaultHandleFilePath(env: NodeJS.ProcessEnv = process.env): string { + if (env.CLAUSTRUM_OPENCODE_HANDLES) return env.CLAUSTRUM_OPENCODE_HANDLES + const configHome = env.XDG_CONFIG_HOME || (env.HOME ? join(env.HOME, '.config') : '.config') + return join(configHome, 'cortexkit', 'opencode-handles.json') +} + +function currentUid(): number | undefined { + return process.getuid?.() ?? userInfo().uid +} + +type HandleFileSnapshot = { + file: OpenCodeHandleFileV1 + source?: string + mtimeMs?: number +} + +async function readBounded(descriptor: HandleFileDescriptor, cap: number): Promise<{ buffer: Buffer; bytes: number }> { + if (!descriptor.read) throw new Error('readBounded requires a descriptor exposing read()') + const buffer = Buffer.alloc(cap + 1) + let total = 0 + while (total < cap + 1) { + const chunk = await descriptor.read(buffer, total, buffer.length - total, total) + if (chunk.bytesRead === 0) break + total += chunk.bytesRead + } + return { buffer, bytes: total > cap ? -1 : total } +} + +async function readHandleSnapshot(path = defaultHandleFilePath(), io: HandleFileIo = {}): Promise { + const stat = io.stat ?? nodeStat + const lstat = io.lstat ?? nodeLstat + const read = io.readFile ?? readFile + const openFd = io.open ?? ((candidate: string) => nodeOpen(candidate, constants.O_RDONLY | constants.O_NOFOLLOW)) + let descriptor: HandleFileDescriptor | undefined + try { + let metadata: HandleFileStat + try { + if (io.lstat || io.readFile) { + metadata = await lstat(path) + } else { + descriptor = await openFd(path) + metadata = await descriptor.stat() + } + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return { file: { version: 1, providers: [] } } + if ((error as NodeJS.ErrnoException).code === 'ELOOP') invalid('handle file must not be a symlink') + invalid(`cannot stat handle file: ${error instanceof Error ? error.message : String(error)}`) + } + if (metadata.isSymbolicLink?.()) invalid('handle file must not be a symlink') + if (!metadata.isFile()) invalid('handle file must be a regular file') + if ((metadata.size ?? 0) > HANDLE_FILE_CONTRACT.maxBytes) invalid('handle file exceeds 256 KiB') + if ((metadata.mode & 0o777) !== HANDLE_FILE_CONTRACT.mode) invalid('handle file mode must be exactly 0600') + const uid = io.currentUid ?? currentUid + const expectedUid = uid() + if (expectedUid !== undefined && metadata.uid !== undefined && metadata.uid !== expectedUid) { + invalid('handle file is not owned by the current uid') + } + let parent: HandleFileStat + try { + parent = await stat(dirname(path)) + } catch (error) { + invalid(`cannot stat handle file parent: ${error instanceof Error ? error.message : String(error)}`) + } + if (!parent.isDirectory?.()) invalid('handle file parent must be a directory') + if (expectedUid !== undefined && parent.uid !== undefined && parent.uid !== expectedUid) { + invalid('handle file parent is not owned by the current uid') + } + if ((parent.mode & 0o002) !== 0 && (parent.mode & 0o1000) === 0) { + invalid('handle file parent is world-writable without sticky bit') + } + let source: string + try { + if (descriptor) { + const { buffer, bytes } = await readBounded(descriptor, HANDLE_FILE_CONTRACT.maxBytes) + if (bytes === -1) invalid('handle file exceeds 256 KiB') + source = buffer.subarray(0, bytes).toString('utf8') + } else { + source = await read(path, 'utf8') + } + } catch (error) { + if (error instanceof HandleFileValidationError) throw error + invalid(`cannot read handle file: ${error instanceof Error ? error.message : String(error)}`) + } + let value: unknown + try { + value = JSON.parse(source) + } catch { + invalid('handle file contains invalid JSON') + } + return { file: parseHandleFile(value), source, mtimeMs: metadata.mtimeMs } + } finally { + await descriptor?.close() + } +} + +export async function readHandleFile(path = defaultHandleFilePath(), io: HandleFileIo = {}): Promise { + return (await readHandleSnapshot(path, io)).file +} + +export async function handleFileRevision(path = defaultHandleFilePath(), io: HandleFileIo = {}): Promise { + const snapshot = await readHandleSnapshot(path, io) + if (snapshot.source === undefined) invalid('cannot revise absent handle file') + return `${snapshot.mtimeMs ?? 0}:${createHash('sha256').update(snapshot.source).digest('hex')}` +} diff --git a/packages/client/src/index.ts b/packages/client/src/index.ts index c4f0789..85268d8 100644 --- a/packages/client/src/index.ts +++ b/packages/client/src/index.ts @@ -6,6 +6,16 @@ export { type ClaustrumEndpoint, } from './detect.js' export { storeIdentity, storageFingerprint } from './identity.js' +export { + MANIFEST_LOCK, + withManifestLock, + writeHandleFileLocked, + type ManifestHandleAccount, + type ManifestHandleFile, + type ManifestHandleProvider, + type ManifestLockError, + type ManifestLockErrorCode, +} from './manifest-lock.js' export { ClaustrumCredentialError, credentialErrorAction, @@ -21,3 +31,15 @@ export { type CredentialStatus, type ServedCredential, } from './wire.js' +export { + HANDLE_FILE_CONTRACT, + HandleFileValidationError, + defaultHandleFilePath, + handleFileRevision, + parseHandleFile, + readHandleFile, + type HandleAccount, + type HandleFileIo, + type HandleProvider, + type OpenCodeHandleFileV1, +} from './handles.js' diff --git a/packages/client/src/manifest-lock.ts b/packages/client/src/manifest-lock.ts new file mode 100644 index 0000000..445df0e --- /dev/null +++ b/packages/client/src/manifest-lock.ts @@ -0,0 +1,117 @@ +import { constants as fsConstants } from 'node:fs' +import { chmod, lstat, mkdir, open, readFile, rename, rm, stat, unlink } from 'node:fs/promises' +import { randomBytes, randomInt } from 'node:crypto' +import { dirname, join } from 'node:path' +import { HANDLE_FILE_CONTRACT, parseHandleFile, type OpenCodeHandleFileV1 } from './handles.js' + +export const MANIFEST_LOCK = { ttlMs: 30_000, renewEveryMs: 10_000, ownerKeys: ['tenant', 'pid', 'claimed_at_ms', 'nonce'] as const, staleTargetRe: /^\.lock\.stale-\d+-(?!\.{1,2}$)(?!.*[. ]$)[^/\\\x00-\x1f:*?"<>|]{1,128}$/, errorCodes: ['lock_busy', 'owner_invalid', 'renewal_failed'] as const } + +/** Thrown by the lock. Branch on `code`; the message is diagnostic and may be reworded. */ +export type ManifestLockErrorCode = (typeof MANIFEST_LOCK.errorCodes)[number] +export type ManifestLockError = Error & { code: ManifestLockErrorCode } +export type ManifestHandleAccount = OpenCodeHandleFileV1['providers'][number]['accounts'][number] +export type ManifestHandleProvider = OpenCodeHandleFileV1['providers'][number] +export type ManifestHandleFile = OpenCodeHandleFileV1 +type Owner = { tenant: string; pid: number; claimed_at_ms: number; nonce: string } +type TestOptions = { ttlMs?: number; renewEveryMs?: number; retryMinMs?: number; retryMaxMs?: number; afterClaim?: () => Promise | void; beforeEvict?: () => Promise; afterEvictRenameAttempt?: () => Promise; afterEvict?: () => void; beforeManifestRename?: (path: string) => Promise } +let testOptions: TestOptions | undefined +export function __setManifestLockTestOptions(options?: TestOptions): void { testOptions = options } +const token = () => randomBytes(16).toString('base64url') +const code = (error: unknown) => (error as NodeJS.ErrnoException | undefined)?.code +const sleep = async (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)) +// A tenant classifying a failure must not string-match our prose: a copy-edit would +// silently reclassify a retryable busy-lock as an unknown error, with nothing failing +// loudly. The code is the contract; the message is free to change. +const lockError = (code: ManifestLockErrorCode, message: string): ManifestLockError => Object.assign(new Error(message), { code }) + +function parseOwner(source: string): Owner { + let value: unknown + try { value = JSON.parse(source) as unknown } catch { throw lockError('owner_invalid', 'manifest lock owner invalid') } + if (!value || typeof value !== 'object') throw lockError('owner_invalid', 'manifest lock owner invalid') + const owner = value as Record + // Widen the nonce alphabet only after every tenant has this path-safe reader; an older + // allowlist reader can otherwise wedge forever on the first owner using the new alphabet. + const staleTarget = `.lock.stale-${owner.claimed_at_ms}-${owner.nonce}` + if (MANIFEST_LOCK.ownerKeys.some((key) => !Object.hasOwn(owner, key)) || typeof owner.claimed_at_ms !== 'number' || !Number.isInteger(owner.claimed_at_ms) || owner.claimed_at_ms < 0 || typeof owner.nonce !== 'string' || !MANIFEST_LOCK.staleTargetRe.test(staleTarget)) throw lockError('owner_invalid', 'manifest lock owner invalid') + return owner as Owner +} +const readOwner = async (path: string) => parseOwner(await readFile(path, 'utf8')) +async function writeOwner(lock: string, owner: Owner): Promise { + const target = join(lock, 'owner'); const temporary = join(lock, `owner.${process.pid}.${token()}.tmp`) + let file: Awaited> | undefined + try { + file = await open(temporary, fsConstants.O_CREAT | fsConstants.O_EXCL | fsConstants.O_WRONLY, 0o600) + await file.chmod(0o600); await file.writeFile(`${JSON.stringify(owner)}\n`); await file.sync(); await file.close(); file = undefined + await rename(temporary, target) + } finally { await file?.close().catch(() => {}); await unlink(temporary).catch(() => {}) } +} + +export async function withManifestLock(path: string, tenant: string, fn: () => Promise | T): Promise { + return withLockCommit(path, tenant, async () => fn()) +} +async function withLockCommit(path: string, tenant: string, fn: (commit: () => Promise) => Promise | T): Promise { + const lock = `${path}.lock`, ownerPath = join(lock, 'owner'), ttl = testOptions?.ttlMs ?? MANIFEST_LOCK.ttlMs, renewEvery = testOptions?.renewEveryMs ?? MANIFEST_LOCK.renewEveryMs, retryMin = testOptions?.retryMinMs ?? 25, retryMax = testOptions?.retryMaxMs ?? 75 + const nonce = token(), started = Date.now(), deadline = started + ttl + while (true) { + try { await mkdir(lock, { mode: 0o700 }); await writeOwner(lock, { tenant, pid: process.pid, claimed_at_ms: Date.now(), nonce }); await testOptions?.afterClaim?.(); break } catch (error) { + if (code(error) !== 'EEXIST') { if (code(error) !== 'ENOENT') await rm(lock, { recursive: true, force: true }).catch(() => {}); throw error } + } + let observed: Owner | undefined, ownerReadError: unknown + try { observed = await readOwner(ownerPath) } catch (error) { ownerReadError = error; if (code(error) !== 'ENOENT' && Date.now() >= deadline) throw code(error) === 'owner_invalid' ? error : lockError('lock_busy', 'manifest lock busy') } + if (observed && Date.now() - observed.claimed_at_ms >= ttl) { + await testOptions?.beforeEvict?.() + const stale = `${lock}.stale-${observed.claimed_at_ms}-${observed.nonce}` + let renameError: unknown + try { await rename(lock, stale) } catch (error) { renameError = error } + await testOptions?.afterEvictRenameAttempt?.() + if (renameError === undefined) { + const moved = await readOwner(join(stale, 'owner')).catch(() => undefined) + if (moved?.nonce === observed.nonce && moved.claimed_at_ms === observed.claimed_at_ms) { testOptions?.afterEvict?.(); continue } + await rename(stale, lock).catch(() => {}) + } else if (!['ENOENT', 'EEXIST', 'ENOTEMPTY'].includes(code(renameError) ?? '')) throw renameError + } + if (Date.now() >= deadline) throw code(ownerReadError) === 'owner_invalid' ? ownerReadError : lockError('lock_busy', 'manifest lock busy') + await sleep(Math.min(randomInt(retryMin, retryMax + 1), Math.max(1, deadline - Date.now()))) + } + let renewal = Promise.resolve(), failed = false, stopped = false + const timer = setInterval(() => { renewal = renewal.then(async () => { if (failed) return; try { const current = await readOwner(ownerPath); if (current.nonce !== nonce || Date.now() - current.claimed_at_ms >= ttl) throw new Error('lease lost'); await writeOwner(lock, { ...current, claimed_at_ms: Date.now() }) } catch { failed = true } }) }, renewEvery) + timer.unref?.() + const commit = async () => { if (!stopped) { stopped = true; clearInterval(timer); await renewal }; const current = await readOwner(ownerPath).catch(() => undefined); if (failed || !current || current.nonce !== nonce || Date.now() - current.claimed_at_ms >= ttl) throw lockError('renewal_failed', 'manifest lock renewal failed; write aborted') } + try { const result = await fn(commit); if (failed) throw lockError('renewal_failed', 'manifest lock renewal failed; write aborted'); return result } finally { + if (!stopped) clearInterval(timer); await renewal + const current = await readOwner(ownerPath).catch(() => undefined) + if (!current || current.nonce !== nonce || Date.now() - current.claimed_at_ms >= ttl) console.warn('manifest lock lease lost, not releasing', { path, tenant }) + else { + const release = `${lock}.release-${nonce}` + try { await rename(lock, release); const moved = await readOwner(join(release, 'owner')).catch(() => undefined); if (!moved || moved.nonce !== nonce || Date.now() - moved.claimed_at_ms >= ttl) { await rename(release, lock).catch(() => {}); console.warn('manifest lock lease lost, not releasing', { path, tenant }) } else await rm(release, { recursive: true, force: true }) } catch { console.warn('manifest lock lease lost, not releasing', { path, tenant }) } + } + } +} + +async function readManifest(path: string): Promise { + let metadata: Awaited> + try { metadata = await lstat(path) } catch (error) { if (code(error) === 'ENOENT') return { version: 1, providers: [] }; throw error } + if (metadata.isSymbolicLink() || !metadata.isFile()) throw new Error('handle file must be a regular file') + if ((metadata.mode & 0o777) !== HANDLE_FILE_CONTRACT.mode) throw new Error('handle file mode must be exactly 0600') + const source = await readFile(path); if (source.byteLength > HANDLE_FILE_CONTRACT.maxBytes) throw new Error('handle file exceeds 256 KiB') + return parseHandleFile(JSON.parse(source.toString('utf8'))) +} +const foreign = (file: ManifestHandleFile, tenant: string) => file.providers.filter((provider) => provider.serve !== tenant).map((provider) => JSON.stringify(provider)) +async function prepareParent(path: string): Promise { const parent = dirname(path); await mkdir(parent, { recursive: true, mode: 0o700 }); const metadata = await stat(parent); if (!metadata.isDirectory()) throw new Error('handle file parent must be a directory'); if ((metadata.mode & 0o002) !== 0 && (metadata.mode & 0o1000) === 0) throw new Error('handle file parent is world-writable without sticky bit'); if ((metadata.mode & 0o022) !== 0) throw new Error('handle file parent must not be group- or other-writable') } +async function writeAtomic(path: string, file: ManifestHandleFile, commit: () => Promise): Promise { + const bytes = Buffer.from(JSON.stringify(file)); if (bytes.byteLength > HANDLE_FILE_CONTRACT.maxBytes) throw new Error('handle file exceeds 256 KiB') + const temporary = join(dirname(path), `.${path.split('/').pop()}.${process.pid}.${token()}.tmp`); let handle: Awaited> | undefined + try { handle = await open(temporary, fsConstants.O_CREAT | fsConstants.O_EXCL | fsConstants.O_WRONLY, 0o600); await handle.chmod(0o600); await handle.writeFile(bytes); await handle.sync(); await handle.close(); handle = undefined; await chmod(temporary, 0o600); await testOptions?.beforeManifestRename?.(`${path}.lock`); await commit(); await rename(temporary, path) } finally { await handle?.close().catch(() => {}); await unlink(temporary).catch(() => {}) } +} +export async function writeHandleFileLocked(path: string, tenant: string, mutate: (file: ManifestHandleFile) => void | ManifestHandleFile | Promise): Promise { + await prepareParent(path) + await withLockCommit(path, tenant, async (commit) => { + const before = await readManifest(path), working = structuredClone(before), result = await mutate(working), next = parseHandleFile(result ?? working), beforeForeign = foreign(before, tenant) + if (JSON.stringify(foreign(next, tenant)) !== JSON.stringify(beforeForeign)) throw new Error('manifest mutation changed another tenant block') + await writeAtomic(path, next, commit) + if (((await lstat(path)).mode & 0o777) !== HANDLE_FILE_CONTRACT.mode) throw new Error('manifest readback mode is not 0600') + const readback = await readManifest(path) + if (JSON.stringify(readback) !== JSON.stringify(next)) throw new Error('manifest readback differs') + if (JSON.stringify(foreign(readback, tenant)) !== JSON.stringify(beforeForeign)) throw new Error('manifest readback changed another tenant block') + }) +} diff --git a/packages/client/src/tests/handles.test.ts b/packages/client/src/tests/handles.test.ts new file mode 100644 index 0000000..f1d90d5 --- /dev/null +++ b/packages/client/src/tests/handles.test.ts @@ -0,0 +1,131 @@ +import { afterEach, describe, expect, test } from 'bun:test' +import { chmod, mkdir, rm, symlink, writeFile } from 'node:fs/promises' +import { join } from 'node:path' +import { + HANDLE_FILE_CONTRACT, + defaultHandleFilePath, + handleFileRevision, + parseHandleFile, + readHandleFile, + type OpenCodeHandleFileV1, +} from '../handles.js' + +const root = '/tmp/claustrum-client-handles-tests' +const handle = `ckh_${'a'.repeat(43)}` + +afterEach(() => rm(root, { recursive: true, force: true })) + +function validFile(): OpenCodeHandleFileV1 { + return { version: 1, providers: [{ provider: 'deepseek', shape: 'api', serve: 'opencode-claustrum', accounts: [{ label: 'main', handle, credential_id: 'apikey:deepseek:main' }] }] } +} + +describe('client handle-file contract', () => { + test('exports the pinned resolver and contract constants', () => { + expect(defaultHandleFilePath({ CLAUSTRUM_OPENCODE_HANDLES: '/tmp/custom.json' })).toBe('/tmp/custom.json') + expect(HANDLE_FILE_CONTRACT.maxBytes).toBe(262144) + expect(HANDLE_FILE_CONTRACT.mode).toBe(0o600) + expect(HANDLE_FILE_CONTRACT.labelRe.test('main.account-1')).toBe(true) + expect(HANDLE_FILE_CONTRACT.handleRe.test(handle)).toBe(true) + }) + + test('reads a valid owned 0600 manifest and computes its revision', async () => { + await mkdir(root, { recursive: true, mode: 0o700 }) + const path = join(root, 'handles.json') + await writeFile(path, `${JSON.stringify(validFile())}\n`, { mode: 0o600 }) + expect(await readHandleFile(path)).toEqual(validFile()) + expect(await handleFileRevision(path)).toMatch(/^\d+(\.\d+)?:\w{64}$/) + }) + + test('rejects insecure mode and world-writable parents', async () => { + await mkdir(root, { recursive: true, mode: 0o777 }) + const path = join(root, 'handles.json') + await writeFile(path, JSON.stringify(validFile()), { mode: 0o600 }) + await chmod(root, 0o777) + await expect(readHandleFile(path)).rejects.toThrow('world-writable without sticky bit') + await chmod(root, 0o700) + await chmod(path, 0o640) + await expect(readHandleFile(path)).rejects.toThrow('exactly 0600') + }) + + test('rejects invalid labels, handles, prototype keys, and oversized input', async () => { + expect(() => parseHandleFile({ version: 1, providers: [{ ...validFile().providers[0], accounts: [{ ...validFile().providers[0].accounts[0], label: '__proto__' }] }] })).toThrow('invalid account label') + expect(() => parseHandleFile({ version: 1, providers: [{ ...validFile().providers[0], accounts: [{ ...validFile().providers[0].accounts[0], handle: 'ckh_short' }] }] })).toThrow('invalid handle') + await mkdir(root, { recursive: true, mode: 0o700 }) + const path = join(root, 'handles.json') + await writeFile(path, 'x'.repeat(262145), { mode: 0o600 }) + await expect(readHandleFile(path)).rejects.toThrow('exceeds 256 KiB') + }) + + test('preserves the historical parser fixture outcomes and exact messages', () => { + const provider = validFile().providers[0] + const account = provider.accounts[0] + const fixtures: Array<[unknown, string]> = [ + [null, 'handle file must be an object'], + [{ version: 2, providers: [] }, 'handle file must have version 1 and providers'], + [{ version: 1, providers: [null] }, 'provider 0 must be an object'], + [{ version: 1, providers: [{ ...provider, provider: '__proto__' }] }, 'provider 0 has invalid provider'], + [{ version: 1, providers: [provider, provider] }, 'provider 1 duplicates provider deepseek'], + [{ version: 1, providers: [{ ...provider, shape: 'other' }] }, 'provider 0 has invalid shape'], + [{ version: 1, providers: [{ ...provider, serve: '' }] }, 'provider 0 requires serve'], + [{ version: 1, providers: [{ ...provider, accounts: [{ ...account, credential_id: 3 }] }] }, 'provider 0 has invalid accounts'], + [{ version: 1, providers: [{ ...provider, accounts: [{ ...account, label: '__proto__' }] }] }, 'provider 0 has an invalid account label'], + [{ version: 1, providers: [{ ...provider, accounts: [account, account] }] }, 'provider 0 duplicates account label main'], + [{ version: 1, providers: [{ ...provider, accounts: [{ ...account, handle: 'ckh_short' }] }] }, 'provider 0 account main has invalid handle'], + [{ version: 1, providers: [{ ...provider, accounts: [{ ...account, credential_id: '' }] }] }, 'provider 0 account main has invalid credential id'], + [{ version: 1, providers: [{ ...provider, accounts: [{ ...account, superseded: ['ckh_short'] }] }] }, 'provider 0 account main has invalid superseded handle'], + ] + for (const [fixture, message] of fixtures) expect(() => parseHandleFile(fixture)).toThrow(message) + }) + + test('preserves historical reader validation order and normalized failures', async () => { + const source = JSON.stringify(validFile()) + const regular = { isFile: () => true, mode: 0o100600, uid: 1000, size: source.length } + const parent = { isFile: () => false, isDirectory: () => true, mode: 0o040755, uid: 1000 } + await expect(readHandleFile('/tmp/handles.json', { + currentUid: () => 1000, + lstat: async () => regular, + stat: async () => { throw new Error('parent denied') }, + readFile: async () => source, + })).rejects.toThrow('cannot stat handle file parent: parent denied') + await expect(readHandleFile('/tmp/handles.json', { + currentUid: () => 1000, + lstat: async () => regular, + stat: async () => parent, + readFile: async () => { throw new Error('read denied') }, + })).rejects.toThrow('cannot read handle file: read denied') + await expect(readHandleFile('/tmp/handles.json', { + currentUid: () => 1000, + lstat: async () => ({ ...regular, mode: 0o100640, uid: 1001 }), + stat: async () => { throw new Error('must not reach parent') }, + readFile: async () => { throw new Error('must not read') }, + })).rejects.toThrow('handle file mode must be exactly 0600') + }) + + test('preserves the historical symlink and grow-after-fstat fixture outcomes', async () => { + await mkdir(root, { recursive: true, mode: 0o700 }) + const target = join(root, 'target.json') + const link = join(root, 'handles.json') + await writeFile(target, JSON.stringify(validFile()), { mode: 0o600 }) + await symlink(target, link) + await expect(readHandleFile(link)).rejects.toThrow('handle file must not be a symlink') + + const chunk = Buffer.alloc(HANDLE_FILE_CONTRACT.maxBytes + 2, 0x78) + await expect(readHandleFile('/tmp/handles.json', { + currentUid: () => 1000, + stat: async () => ({ isFile: () => false, isDirectory: () => true, mode: 0o040755, uid: 1000 }), + open: async () => ({ + stat: async () => ({ isFile: () => true, mode: 0o100600, uid: 1000, size: 256 }), + readFile: async () => chunk.toString('utf8'), + read: (buffer, offset, length, position) => { + const remaining = chunk.length - position + if (remaining <= 0) return { bytesRead: 0 } + const slice = chunk.subarray(position, position + Math.min(length, remaining)) + buffer.set(slice, offset) + return { bytesRead: slice.length } + }, + close: async () => {}, + }), + })).rejects.toThrow('handle file exceeds 256 KiB') + }) + +}) diff --git a/packages/client/src/tests/manifest-lock.test.ts b/packages/client/src/tests/manifest-lock.test.ts new file mode 100644 index 0000000..a87bf17 --- /dev/null +++ b/packages/client/src/tests/manifest-lock.test.ts @@ -0,0 +1,505 @@ +import { afterEach, describe, expect, test } from 'bun:test' +import { chmod, lstat, mkdir, mkdtemp, readFile, readdir, rename, rm, stat, symlink, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { basename, join } from 'node:path' + +import { + __setManifestLockTestOptions, + MANIFEST_LOCK, + withManifestLock, + writeHandleFileLocked, +} from '../manifest-lock' +import type { ManifestLockError, ManifestLockErrorCode } from '../index' + +const roots: string[] = [] +const handle = (letter: string) => `ckh_${letter.repeat(43)}` +const sleep = async (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)) + +async function manifestPath(): Promise { + const root = await mkdtemp(join(tmpdir(), 'claustrum-manifest-lock-')) + roots.push(root) + return join(root, 'opencode-handles.json') +} + +function provider(provider: string, tenant: string) { + return { + provider, + shape: 'api' as const, + serve: tenant, + accounts: [{ + label: 'main', + handle: handle(provider[0] ?? 'A'), + credential_id: `apikey:${provider}:main`, + }], + } +} + +async function owner(path: string, claimedAtMs: number, tenant = 'other-tenant'): Promise { + const lockPath = `${path}.lock` + await mkdir(lockPath, { mode: 0o700 }) + await writeFile(join(lockPath, 'owner'), `${JSON.stringify({ + tenant, + pid: 41, + claimed_at_ms: claimedAtMs, + nonce: '0123456789abcdef0123456789abcdef', + })}\n`, { mode: 0o600 }) + await chmod(join(lockPath, 'owner'), 0o600) +} + +afterEach(async () => { + __setManifestLockTestOptions() + await Promise.all(roots.splice(0).map((root) => rm(root, { recursive: true, force: true }))) +}) + +describe('manifest writer lock', () => { + test('two concurrent tenant writers preserve both provider blocks', async () => { + const path = await manifestPath() + const firstEntered = Promise.withResolvers() + const releaseFirst = Promise.withResolvers() + + const first = writeHandleFileLocked(path, 'anthropic-auth', async (file) => { + firstEntered.resolve() + await releaseFirst.promise + file.providers.push(provider('anthropic', 'anthropic-auth')) + }) + await firstEntered.promise + const second = writeHandleFileLocked(path, 'openai-auth', (file) => { + file.providers.push(provider('openai', 'openai-auth')) + }) + releaseFirst.resolve() + await Promise.all([first, second]) + + const written = JSON.parse(await readFile(path, 'utf8')) as { providers: Array<{ provider: string }> } + expect(written.providers.map((entry) => entry.provider).sort()).toEqual(['anthropic', 'openai']) + }) + + test('stale owner is evicted by rename and retained as a quarantine directory', async () => { + const path = await manifestPath() + await owner(path, Date.now() - MANIFEST_LOCK.ttlMs - 1) + + await withManifestLock(path, 'anthropic-auth', async () => {}) + + const suffixes = (await readdir(join(path, '..'))) + .filter((name) => name.startsWith(`${basename(path)}.lock.stale-`)) + .map((name) => name.slice(basename(path).length)) + expect(suffixes).toHaveLength(1) + expect(MANIFEST_LOCK.staleTargetRe.test(suffixes[0]!)).toBe(true) + }) + + test('owner that becomes stale during the retry window is evicted', async () => { + const path = await manifestPath() + __setManifestLockTestOptions({ ttlMs: 80, renewEveryMs: 1_000, retryMinMs: 2, retryMaxMs: 3 }) + await owner(path, Date.now() - 30) + + await withManifestLock(path, 'anthropic-auth', async () => {}) + + const suffixes = (await readdir(join(path, '..'))).filter((name) => name.startsWith(`${basename(path)}.lock.stale-`)) + expect(suffixes).toHaveLength(1) + }) + + test('owner nonce containing a path traversal is invalid and never renamed', async () => { + for (const nonce of ['../escape', 'a/b', 'a:b', 'a*b', 'a?b', 'a|b']) { + const path = await manifestPath() + const lockPath = `${path}.lock` + __setManifestLockTestOptions({ ttlMs: 30, renewEveryMs: 10, retryMinMs: 2, retryMaxMs: 3 }) + await mkdir(lockPath, { mode: 0o700 }) + await writeFile(join(lockPath, 'owner'), `${JSON.stringify({ tenant: 'squatter', pid: 41, claimed_at_ms: Date.now() - 31, nonce })}\n`, { mode: 0o600 }) + + const error = (await withManifestLock(path, 'anthropic-auth', async () => {}).catch((caught) => caught)) as Error & { code?: string } + + expect(error.code).toBe('owner_invalid') + expect((await stat(lockPath)).isDirectory()).toBe(true) + expect((await readdir(join(path, '..'))).some((name) => name.includes('.lock.stale-'))).toBe(false) + await expect(stat(join(path, '..', 'escape'))).rejects.toMatchObject({ code: 'ENOENT' }) + } + }) + + test('path-safe unfamiliar nonce alphabets remain evictable', async () => { + for (const nonce of ['abc.def', 'AAAA====']) { + const path = await manifestPath() + const lockPath = `${path}.lock` + __setManifestLockTestOptions({ ttlMs: 30, renewEveryMs: 10, retryMinMs: 2, retryMaxMs: 3 }) + await mkdir(lockPath, { mode: 0o700 }) + await writeFile(join(lockPath, 'owner'), `${JSON.stringify({ tenant: 'newer-writer', pid: 41, claimed_at_ms: Date.now() - 31, nonce })}\n`, { mode: 0o600 }) + + await withManifestLock(path, 'anthropic-auth', async () => {}) + + expect((await readdir(join(path, '..'))).some((name) => name.startsWith(`${basename(path)}.lock.stale-`))).toBe(true) + } + }) + + test('unknown owner keys are busy while fresh and evictable once stale', async () => { + const path = await manifestPath() + const lockPath = `${path}.lock` + __setManifestLockTestOptions({ ttlMs: 30, renewEveryMs: 10, retryMinMs: 2, retryMaxMs: 3 }) + await mkdir(lockPath, { mode: 0o700 }) + const record = { tenant: 'newer-writer', pid: 41, claimed_at_ms: Date.now() + 1_000, nonce: 'newer_writer_nonce', generation: 2 } + await writeFile(join(lockPath, 'owner'), `${JSON.stringify(record)}\n`, { mode: 0o600 }) + + const freshError = (await withManifestLock(path, 'anthropic-auth', async () => {}).catch((caught) => caught)) as Error & { code?: string } + expect(freshError.code).toBe('lock_busy') + + record.claimed_at_ms = Date.now() - 31 + await writeFile(join(lockPath, 'owner'), `${JSON.stringify(record)}\n`, { mode: 0o600 }) + await withManifestLock(path, 'anthropic-auth', async () => {}) + expect((await readdir(join(path, '..'))).some((name) => name.startsWith(`${basename(path)}.lock.stale-`))).toBe(true) + }) + + test('malformed diagnostic owner fields do not prevent stale eviction', async () => { + const path = await manifestPath() + const lockPath = `${path}.lock` + __setManifestLockTestOptions({ ttlMs: 30, renewEveryMs: 10, retryMinMs: 2, retryMaxMs: 3 }) + await mkdir(lockPath, { mode: 0o700 }) + await writeFile(join(lockPath, 'owner'), `${JSON.stringify({ tenant: 41, pid: 'unknown', claimed_at_ms: Date.now() - 31, nonce: 'valid_nonce' })}\n`, { mode: 0o600 }) + + await withManifestLock(path, 'anthropic-auth', async () => {}) + + expect((await readdir(join(path, '..'))).some((name) => name.startsWith(`${basename(path)}.lock.stale-`))).toBe(true) + }) + + test('renewing owner fails loudly after the bounded retry window', async () => { + const path = await manifestPath() + __setManifestLockTestOptions({ ttlMs: 40, renewEveryMs: 10, retryMinMs: 2, retryMaxMs: 3 }) + await owner(path, Date.now()) + + const started = Date.now() + const renewal = setInterval(async () => { + const ownerPath = join(`${path}.lock`, 'owner') + const current = JSON.parse(await readFile(ownerPath, 'utf8')) as Record + current.claimed_at_ms = Date.now() + await writeFile(ownerPath, `${JSON.stringify(current)}\n`, { mode: 0o600 }) + }, 5) + try { + await expect(withManifestLock(path, 'anthropic-auth', async () => {})).rejects.toThrow('manifest lock busy') + expect(Date.now() - started).toBeGreaterThanOrEqual(35) + } finally { clearInterval(renewal) } + }) + + test('owner file exists while held and disappears with the lock after release', async () => { + const path = await manifestPath() + const lockPath = `${path}.lock` + + await withManifestLock(path, 'anthropic-auth', async () => { + const parsed = JSON.parse(await readFile(join(lockPath, 'owner'), 'utf8')) as Record + expect(Object.keys(parsed).sort()).toEqual([...MANIFEST_LOCK.ownerKeys].sort()) + expect(parsed.tenant).toBe('anthropic-auth') + }) + + await expect(stat(lockPath)).rejects.toMatchObject({ code: 'ENOENT' }) + }) + + test('two stale evictors produce one eviction winner and never overlap holders', async () => { + const path = await manifestPath() + await owner(path, Date.now() - MANIFEST_LOCK.ttlMs - 1) + let waiting = 0 + const bothReady = Promise.withResolvers() + let evictionWins = 0 + __setManifestLockTestOptions({ + beforeEvict: async () => { + waiting += 1 + if (waiting === 2) bothReady.resolve() + await bothReady.promise + }, + afterEvict: () => { evictionWins += 1 }, + retryMinMs: 2, + retryMaxMs: 3, + }) + let active = 0 + let maxActive = 0 + const hold = async () => { + active += 1 + maxActive = Math.max(maxActive, active) + await Bun.sleep(15) + active -= 1 + } + + await Promise.all([ + withManifestLock(path, 'anthropic-auth', hold), + withManifestLock(path, 'openai-auth', hold), + ]) + + expect(evictionWins).toBe(1) + expect(maxActive).toBe(1) + }) + + test('an evictor that observed stale owner cannot rename a replacement lock', async () => { + const path = await manifestPath() + const now = Date.now() + await owner(path, now - 501) + const loserObserved = Promise.withResolvers() + const allowLoserRename = Promise.withResolvers() + const replacementClaimed = Promise.withResolvers() + const allowReplacementRelease = Promise.withResolvers() + const renameAttempted = Promise.withResolvers() + const allowAttemptCompletion = Promise.withResolvers() + let beforeEvictCalls = 0 + let evictRenameAttempts = 0 + + __setManifestLockTestOptions({ + ttlMs: 500, + renewEveryMs: 1_000, + retryMinMs: 2, + retryMaxMs: 3, + beforeEvict: async () => { + beforeEvictCalls += 1 + if (beforeEvictCalls === 1) { + loserObserved.resolve() + await allowLoserRename.promise + } + }, + afterClaim: async () => { + replacementClaimed.resolve() + await allowReplacementRelease.promise + }, + afterEvictRenameAttempt: async () => { + evictRenameAttempts += 1 + if (evictRenameAttempts !== 2) return + renameAttempted.resolve() + await allowAttemptCompletion.promise + }, + }) + + const loser = withManifestLock(path, 'loser', async () => {}) + await loserObserved.promise + const replacement = withManifestLock(path, 'replacement', async () => {}) + await replacementClaimed.promise + allowLoserRename.resolve() + await renameAttempted.promise + allowReplacementRelease.resolve() + await replacement + allowAttemptCompletion.resolve() + await loser + + await expect(stat(`${path}.lock`)).rejects.toMatchObject({ code: 'ENOENT' }) + }) + + test('an expired holder does not release its directory and logs the lost lease', async () => { + const path = await manifestPath() + const lockPath = `${path}.lock` + __setManifestLockTestOptions({ ttlMs: 40, renewEveryMs: 1_000, retryMinMs: 2, retryMaxMs: 3 }) + const warnings: unknown[][] = [] + const originalWarn = console.warn + console.warn = (...args: unknown[]) => { warnings.push(args) } + try { + await withManifestLock(path, 'anthropic-auth', async () => { + const ownerPath = join(lockPath, 'owner') + const parsed = JSON.parse(await readFile(ownerPath, 'utf8')) as Record + parsed.claimed_at_ms = Date.now() - 41 + await writeFile(ownerPath, `${JSON.stringify(parsed)}\n`, { mode: 0o600 }) + }) + } finally { + console.warn = originalWarn + } + + await expect(stat(lockPath)).resolves.toBeDefined() + expect(warnings.some((args) => args.includes('manifest lock lease lost, not releasing'))).toBe(true) + }) + + test('atomic publication remains 0600 under umask 022', async () => { + const path = await manifestPath() + const previous = process.umask(0o022) + try { + await writeHandleFileLocked(path, 'anthropic-auth', (file) => { + file.providers.push(provider('anthropic', 'anthropic-auth')) + }) + } finally { + process.umask(previous) + } + expect((await stat(path)).mode & 0o777).toBe(0o600) + }) + + test('creates a missing manifest parent before claiming its colocated lock', async () => { + const root = await mkdtemp(join(tmpdir(), 'claustrum-manifest-parent-')) + roots.push(root) + const path = join(root, 'nested', 'opencode-handles.json') + + await writeHandleFileLocked(path, 'anthropic-auth', (file) => { + file.providers.push(provider('anthropic', 'anthropic-auth')) + }) + + expect((await stat(path)).mode & 0o777).toBe(0o600) + }) + + test('leaves the mode of a pre-existing benign parent unchanged', async () => { + const path = await manifestPath() + const parent = join(path, '..') + await chmod(parent, 0o755) + const before = (await stat(parent)).mode & 0o777 + + await writeHandleFileLocked(path, 'anthropic-auth', (file) => { + file.providers.push(provider('anthropic', 'anthropic-auth')) + }) + + expect((await stat(parent)).mode & 0o777).toBe(before) + }) + + test('refuses a group-writable manifest parent without changing its mode', async () => { + const path = await manifestPath() + const parent = join(path, '..') + await chmod(parent, 0o770) + + await expect(writeHandleFileLocked(path, 'anthropic-auth', () => {})).rejects.toThrow('handle file parent must not be group- or other-writable') + expect((await stat(parent)).mode & 0o777).toBe(0o770) + }) + + test('pins the shared lock constants and renewal bound', () => { + expect(MANIFEST_LOCK.ttlMs).toBe(30_000) + expect(MANIFEST_LOCK.renewEveryMs).toBe(10_000) + expect(MANIFEST_LOCK.ownerKeys).toEqual(['tenant', 'pid', 'claimed_at_ms', 'nonce']) + for (const [target, accepted] of [ + ['.lock.stale-1-nonce_2', true], + ['.lock.stale-1-abc.def', true], + ['.lock.stale-1-AAAA====', true], + ['.lock.stale-1.bad', false], + ['.lock.stale-1-a/b', false], + ['.lock.stale-1-..', false], + ['.lock.stale-1-a:b', false], + ['.lock.stale-1-a*b', false], + ['.lock.stale-1-a?b', false], + ['.lock.stale-1-a|b', false], + // Windows aliases trailing dots and spaces, collapsing distinct nonces onto one ABA target. + ['.lock.stale-1-abc.', false], + ['.lock.stale-1-abc ', false], + ] as const) expect(MANIFEST_LOCK.staleTargetRe.test(target)).toBe(accepted) + expect(MANIFEST_LOCK.renewEveryMs * 3).toBeLessThanOrEqual(MANIFEST_LOCK.ttlMs) + }) + + test('reads a Rust-shaped owner fixture using the shared field contract', async () => { + const path = await manifestPath() + __setManifestLockTestOptions({ ttlMs: 30, renewEveryMs: 10, retryMinMs: 2, retryMaxMs: 3 }) + await owner(path, Date.now() + 1_000, 'opencode-claustrum') + await expect(withManifestLock(path, 'anthropic-auth', async () => {})).rejects.toThrow('manifest lock busy') + }) + + test('refuses a dangling manifest symlink without replacing it', async () => { + const path = await manifestPath() + const target = join(path, '..', 'missing-target.json') + await symlink(target, path) + + await expect(writeHandleFileLocked(path, 'anthropic-auth', (file) => { + file.providers.push(provider('anthropic', 'anthropic-auth')) + })).rejects.toThrow('handle file must be a regular file') + + expect((await lstat(path)).isSymbolicLink()).toBe(true) + await expect(stat(target)).rejects.toMatchObject({ code: 'ENOENT' }) + }) + + test('aborts before manifest rename when renewal loses the original lock path', async () => { + const path = await manifestPath() + await writeHandleFileLocked(path, 'anthropic-auth', (file) => { + file.providers.push(provider('anthropic', 'anthropic-auth')) + }) + const before = await readFile(path, 'utf8') + __setManifestLockTestOptions({ + ttlMs: 100, + renewEveryMs: 2, + retryMinMs: 2, + retryMaxMs: 3, + beforeManifestRename: async (lockPath: string) => { + await rename(lockPath, `${lockPath}.vanished`) + await Bun.sleep(10) + }, + } as never) + + await expect(writeHandleFileLocked(path, 'anthropic-auth', (file) => { + file.providers[0]!.accounts.push({ + label: 'backup', + handle: handle('Z'), + credential_id: 'apikey:anthropic:backup', + }) + })).rejects.toThrow('manifest lock renewal failed; write aborted') + + expect(await readFile(path, 'utf8')).toBe(before) + }) + + test('pins missing and unparseable owner records without eviction', async () => { + const path = await manifestPath() + const lockPath = `${path}.lock` + __setManifestLockTestOptions({ ttlMs: 25, renewEveryMs: 8, retryMinMs: 2, retryMaxMs: 3 }) + for (const ownerSource of [undefined, '{']) { + await mkdir(lockPath, { mode: 0o700 }) + if (ownerSource !== undefined) { + await writeFile(join(lockPath, 'owner'), ownerSource, { mode: 0o600 }) + } + + const error = (await withManifestLock(path, 'anthropic-auth', async () => {}).catch((caught) => caught)) as Error & { code?: string } + expect(error.code).toBe(ownerSource === undefined ? 'lock_busy' : 'owner_invalid') + expect((await lstat(lockPath)).isDirectory()).toBe(true) + expect((await readdir(join(path, '..'))).some((name) => name.includes('.lock.stale-'))).toBe(false) + await rm(lockPath, { recursive: true }) + } + }) +}) + +describe('thrown errors carry a stable code', () => { + // A tenant classifying "retry later" vs "the artefact is wrong" vs "the write was + // abandoned" had only the message text to branch on, so any copy-edit here silently + // reclassified a busy lock as an unknown error. openai-auth asked for this before + // writing its conformance suite, which is the cheap moment to add it. + test('a busy lock throws code lock_busy', async () => { + const path = await manifestPath() + __setManifestLockTestOptions({ ttlMs: 400, retryMinMs: 5, retryMaxMs: 10 }) + await mkdir(`${path}.lock`, { mode: 0o700 }) + await writeFile(join(`${path}.lock`, 'owner'), `${JSON.stringify({ tenant: 'squatter', pid: 1, claimed_at_ms: Date.now() + 1_000, nonce: 'n'.repeat(22) })}\n`, { mode: 0o600 }) + const error = (await withManifestLock(path, 'probe', async () => {}).catch((e) => e)) as Error & { code?: string } + expect(error.code).toBe('lock_busy') + expect(MANIFEST_LOCK.errorCodes).toContain('lock_busy') + }) + + test('the package entrypoint exports the lock error types', () => { + const classify = (error: ManifestLockError): ManifestLockErrorCode => error.code + expect(classify(Object.assign(new Error('busy'), { code: 'lock_busy' as const }))).toBe('lock_busy') + }) + + test('an unparseable owner is invalid, never evicted, and keeps that code', async () => { + const path = await manifestPath() + __setManifestLockTestOptions({ ttlMs: 400, retryMinMs: 5, retryMaxMs: 10 }) + await mkdir(`${path}.lock`, { mode: 0o700 }) + await writeFile(join(`${path}.lock`, 'owner'), 'not json at all\n', { mode: 0o600 }) + const error = (await withManifestLock(path, 'probe', async () => {}).catch((e) => e)) as Error & { code?: string } + expect(error.code).toBe('owner_invalid') + // the squatter's lock must still be standing: unreadable owner is never evicted + expect((await stat(`${path}.lock`)).isDirectory()).toBe(true) + expect((await readFile(join(`${path}.lock`, 'owner'), 'utf8')).trim()).toBe('not json at all') + }) + + test('a throwing callback releases the lock and re-raises the original error unwrapped', async () => { + const path = await manifestPath() + class EnrollRefusal extends Error { + constructor() { + super('identity mismatch') + this.name = 'EnrollRefusal' + } + } + const error = await withManifestLock(path, 'probe', async () => { + throw new EnrollRefusal() + }).catch((e) => e) + expect(error).toBeInstanceOf(EnrollRefusal) + expect((error as Error).message).toBe('identity mismatch') + await expect(stat(`${path}.lock`)).rejects.toThrow() + // The consumer-visible consequence: the next claimant is not stalled for a TTL. + // Bounded by a race rather than by awaiting and measuring afterwards -- if release + // regressed, the await itself would block for the full 30s TTL and the suite would + // report a timeout with no attribution, which is indistinguishable from a slow box + // or a hang anywhere else. Failing fast with the property named beats hanging. + const started = Date.now() + const outcome = await Promise.race([ + withManifestLock(path, 'probe', async () => 'reacquired' as const), + sleep(1_000).then(() => 'lock not released on throw: re-acquire exceeded 1000ms' as const), + ]) + expect(outcome).toBe('reacquired') + expect(Date.now() - started).toBeLessThan(1_000) + }) + + test('distinct manifest paths do not contend', async () => { + const first = await manifestPath() + const second = await manifestPath() + let bothInside = false + await withManifestLock(first, 'tenant-a', async () => { + await withManifestLock(second, 'tenant-b', async () => { + bothInside = true + }) + }) + expect(bothInside).toBe(true) + }) +}) diff --git a/packages/opencode/src/handles.ts b/packages/opencode/src/handles.ts index f2b91ef..5e3e3a4 100644 --- a/packages/opencode/src/handles.ts +++ b/packages/opencode/src/handles.ts @@ -1,221 +1,34 @@ -import { constants } from "node:fs"; -import { lstat as nodeLstat, open as nodeOpen, readFile, stat as nodeStat } from "node:fs/promises"; -import { userInfo } from "node:os"; -import { dirname, join } from "node:path"; -import { createHash } from "node:crypto"; - -import { boundedBytesText, readBounded, type BoundedReadDescriptor } from "./bounded-read"; -import { HandleFileValidationError } from "./errors"; -import { parseSecretJson, SecretJsonParseError } from "./secret-json"; - -export const OUR_PLUGIN_ID = "opencode-claustrum"; -const HANDLE_FILE_MAX_BYTES = 256 * 1024; -const PROVIDER_ID = /^[a-z0-9][a-z0-9._-]{0,63}$/; -const FORBIDDEN_IDENTIFIERS = new Set(["__proto__", "constructor", "prototype"]); - -export type HandleAccount = { - label: string; - handle: string; - credential_id: string; - superseded?: string[]; -}; -export type HandleProvider = { - provider: string; - shape: "api" | "oauth"; - serve: string; - accounts: HandleAccount[]; -}; -export type OpenCodeHandleFileV1 = { version: 1; providers: HandleProvider[] }; - -function isAccount(value: unknown): value is HandleAccount { - if (!value || typeof value !== "object") return false; - const account = value as Record; - return typeof account.label === "string" && typeof account.handle === "string" && - typeof account.credential_id === "string" && - (account.superseded === undefined || - (Array.isArray(account.superseded) && account.superseded.every((handle) => typeof handle === "string"))); -} - -function handleIsValid(handle: unknown): handle is string { - return typeof handle === "string" && /^ckh_[A-Za-z0-9_-]{43}$/.test(handle); -} - -function identifierIsValid(value: unknown): value is string { - return typeof value === "string" && PROVIDER_ID.test(value) && !FORBIDDEN_IDENTIFIERS.has(value); -} - -function invalid(message: string): never { - throw new HandleFileValidationError(message); -} - -export function parseHandleFile(value: unknown): OpenCodeHandleFileV1 { - if (!value || typeof value !== "object") invalid("handle file must be an object"); - const file = value as Record; - if (file.version !== 1 || !Array.isArray(file.providers)) { - invalid("handle file must have version 1 and providers"); +import { + defaultHandleFilePath, + handleFileRevision as clientHandleFileRevision, + parseHandleFile as clientParseHandleFile, + readHandleFile as clientReadHandleFile, + type HandleFileIo, + type OpenCodeHandleFileV1, +} from '@cortexkit/claustrum-client' +import { HandleFileValidationError } from './errors' + +export { defaultHandleFilePath } +export type { HandleFileIo, OpenCodeHandleFileV1 } +export const OUR_PLUGIN_ID = 'opencode-claustrum' + +function preserveError(operation: () => T): T { + try { return operation() } catch (error) { + if (error instanceof Error && error.name === 'HandleFileValidationError') throw new HandleFileValidationError(error.message) + throw error } - const providerIds = new Set(); - const providers = file.providers.map((provider, index): HandleProvider => { - if (!provider || typeof provider !== "object") invalid(`provider ${index} must be an object`); - const item = provider as Record; - if (!identifierIsValid(item.provider)) invalid(`provider ${index} has invalid provider`); - if (providerIds.has(item.provider)) invalid(`provider ${index} duplicates provider ${item.provider}`); - providerIds.add(item.provider); - if (item.shape !== "api" && item.shape !== "oauth") invalid(`provider ${index} has invalid shape`); - if (typeof item.serve !== "string" || !item.serve) invalid(`provider ${index} requires serve`); - if (!Array.isArray(item.accounts) || item.accounts.length === 0 || !item.accounts.every(isAccount)) { - invalid(`provider ${index} has invalid accounts`); - } - const labels = new Set(); - for (const account of item.accounts) { - if (!identifierIsValid(account.label)) invalid(`provider ${index} has an invalid account label`); - if (labels.has(account.label)) invalid(`provider ${index} duplicates account label ${account.label}`); - labels.add(account.label); - if (!handleIsValid(account.handle)) invalid(`provider ${index} account ${account.label} has invalid handle`); - if (!account.credential_id) invalid(`provider ${index} account ${account.label} has invalid credential id`); - if (account.superseded?.some((handle) => !handleIsValid(handle))) { - invalid(`provider ${index} account ${account.label} has invalid superseded handle`); - } - } - return { - provider: item.provider, - shape: item.shape, - serve: item.serve, - accounts: item.accounts.map((account) => ({ - ...account, - ...(account.superseded === undefined ? {} : { superseded: account.superseded }), - })), - }; - }); - return { version: 1, providers }; -} - -type HandleFileStat = { - isFile(): boolean; - isDirectory?(): boolean; - isSymbolicLink?(): boolean; - mode: number; - size?: number; - uid?: number; - mtimeMs?: number; -}; -type HandleFileDescriptor = BoundedReadDescriptor & { - stat(): Promise; - readFile(options: { encoding: "utf8" }): Promise; - close(): Promise; -}; -export type HandleFileIo = { - stat?: (path: string) => Promise; - lstat?: (path: string) => Promise; - readFile?: (path: string, encoding: "utf8") => Promise; - // Injectable descriptor: when supplied, the handle reader uses a bounded read into a - // cap+1 buffer instead of `readFile()`. The cap check then catches a TOCTOU write that - // grows the file between fstat and read; the unbounded path is preserved for callers - // that pre-trust the source. Mirrors `ConfigHookDependencies.authReader`'s `openFile` - // so a grow-after-fstat handle test can exercise the same bounded read path. - open?: (path: string) => Promise; - currentUid?: () => number | undefined; -}; - -export function defaultHandleFilePath(env: NodeJS.ProcessEnv = process.env): string { - if (env.CLAUSTRUM_OPENCODE_HANDLES) return env.CLAUSTRUM_OPENCODE_HANDLES; - const configHome = env.XDG_CONFIG_HOME || (env.HOME ? join(env.HOME, ".config") : ".config"); - return join(configHome, "cortexkit", "opencode-handles.json"); -} - -function currentUid(): number | undefined { - return process.getuid?.() ?? userInfo().uid; } -type HandleFileSnapshot = { - file: OpenCodeHandleFileV1; - source?: string; - mtimeMs?: number; -}; - -async function readHandleSnapshot(path = defaultHandleFilePath(), io: HandleFileIo = {}): Promise { - const stat = io.stat ?? nodeStat; - const lstat = io.lstat ?? nodeLstat; - const read = io.readFile ?? readFile; - const openFd = io.open ?? ((candidate: string) => nodeOpen(candidate, constants.O_RDONLY | constants.O_NOFOLLOW)); - let descriptor: HandleFileDescriptor | undefined; - try { - let metadata: HandleFileStat; - try { - if (io.lstat || io.readFile) { - metadata = await lstat(path); - } else { - descriptor = await openFd(path); - metadata = await descriptor.stat(); - } - } catch (error) { - if ((error as NodeJS.ErrnoException).code === "ENOENT") return { file: { version: 1, providers: [] } }; - if ((error as NodeJS.ErrnoException).code === "ELOOP") invalid("handle file must not be a symlink"); - invalid(`cannot stat handle file: ${error instanceof Error ? error.message : String(error)}`); - } - if (metadata.isSymbolicLink?.()) invalid("handle file must not be a symlink"); - if (!metadata.isFile()) invalid("handle file must be a regular file"); - if ((metadata.size ?? 0) > HANDLE_FILE_MAX_BYTES) invalid("handle file exceeds 256 KiB"); - if ((metadata.mode & 0o777) !== 0o600) invalid("handle file mode must be exactly 0600"); - const uid = io.currentUid ?? currentUid; - const expectedUid = uid(); - if (expectedUid !== undefined && metadata.uid !== undefined && metadata.uid !== expectedUid) { - invalid("handle file is not owned by the current uid"); - } - let parent: HandleFileStat; - try { - parent = await stat(dirname(path)); - } catch (error) { - invalid(`cannot stat handle file parent: ${error instanceof Error ? error.message : String(error)}`); - } - if (!parent.isDirectory?.()) invalid("handle file parent must be a directory"); - if (expectedUid !== undefined && parent.uid !== undefined && parent.uid !== expectedUid) { - invalid("handle file parent is not owned by the current uid"); - } - if ((parent.mode & 0o002) !== 0 && (parent.mode & 0o1000) === 0) { - invalid("handle file parent is world-writable without sticky bit"); - } - let source: string; - try { - if (descriptor) { - // Bounded read on the already-fstat'd descriptor closes the TOCTOU window a - // size-only check leaves open: a writer that grows the file between fstat and - // readFile would otherwise drive the read past the cap. The shared helper - // allocates the cap+1 buffer and reports bytes > cap so the message is uniform - // across auth and handle paths. - const { buffer, bytes } = await readBounded(descriptor, HANDLE_FILE_MAX_BYTES); - if (bytes === -1) invalid("handle file exceeds 256 KiB"); - source = boundedBytesText(bytes, buffer); - } else { - source = await read(path, "utf8"); - } - } catch (error) { - if (error instanceof HandleFileValidationError) throw error; - invalid(`cannot read handle file: ${error instanceof Error ? error.message : String(error)}`); - } - let value: unknown; - try { - value = parseSecretJson(source, "handle file"); - } catch (error) { - if (error instanceof SecretJsonParseError) invalid("handle file contains invalid JSON"); - throw error; - } - return { - file: parseHandleFile(value), - source, - mtimeMs: metadata.mtimeMs, - }; - } finally { - await descriptor?.close(); +export function parseHandleFile(value: unknown): OpenCodeHandleFileV1 { return preserveError(() => clientParseHandleFile(value)) } +export async function readHandleFile(path?: string, io?: HandleFileIo): Promise { + try { return await clientReadHandleFile(path, io) } catch (error) { + if (error instanceof Error && error.name === 'HandleFileValidationError') throw new HandleFileValidationError(error.message) + throw error } } - -export async function readHandleFile(path = defaultHandleFilePath(), io: HandleFileIo = {}): Promise { - return (await readHandleSnapshot(path, io)).file; -} - -export async function handleFileRevision(path = defaultHandleFilePath(), io: HandleFileIo = {}): Promise { - const snapshot = await readHandleSnapshot(path, io); - if (snapshot.source === undefined) invalid("cannot revise absent handle file"); - return `${snapshot.mtimeMs ?? 0}:${createHash("sha256").update(snapshot.source).digest("hex")}`; +export async function handleFileRevision(path?: string, io?: HandleFileIo): Promise { + try { return await clientHandleFileRevision(path, io) } catch (error) { + if (error instanceof Error && error.name === 'HandleFileValidationError') throw new HandleFileValidationError(error.message) + throw error + } } diff --git a/scripts/gate.sh b/scripts/gate.sh index c5220b5..bce225b 100755 --- a/scripts/gate.sh +++ b/scripts/gate.sh @@ -235,16 +235,23 @@ stream and pass the arm without ever seeing it skip." # follows it), and any gap between the floor and the real count is how many can go # before anyone is told. Measured 402 across the workspace's suites at the time of # writing; an earlier floor of 200 left a third of them free to disappear. -# The current measured total is 578 (debug profile, the same -# `cargo test --locked --workspace` this arm runs). The latest five tests pin the resolved -# credential id on get/status, its omission on unresolved shapes, and the get response key -# set in both directions. +# The current measured total is 590 (debug profile, the same +# `cargo test --locked --workspace` this arm runs). The latest twelve tests pin the Rust +# half of the manifest writer lock: the ABA observation that cannot rename a replacement, +# one quarantine directory per stale owner, unknown and malformed owner fields tolerated +# but still evictable, a missing nonce surfacing as owner_invalid at the deadline, and an +# owner that goes stale inside the retry window. +# +# MEASURED, NOT ARITHMETIC: this branch bumped 564 to 568 when it was written, master +# reached 578 independently, and the sum of those edits (582) is WRONG -- later commits on +# this branch added Rust tests without touching the floor. Re-run the arm and read the +# number rather than adding the two deltas. # Do not measure it with `--release`: one login test deliberately fails there and the # pipeline still prints a number, 32 short of the truth. # # Raise this when tests are added. A failure here is normally that, not a defect -- # but it should be a deliberate edit rather than a number nobody revisits. -run_expect 578 "workspace unit + integration" \ +run_expect 590 "workspace unit + integration" \ cargo test --locked --workspace # Two independent defences, because each catches what the other misses: