diff --git a/crates/celld/bucket.rs b/crates/celld/bucket.rs index 8bb342ae8..9eace55a5 100644 --- a/crates/celld/bucket.rs +++ b/crates/celld/bucket.rs @@ -9,13 +9,11 @@ //! `Ok(None)` only for a clean 412/409 rejection; every other failure is //! ambiguous — the write may have committed — and surfaces as `Err`. +use crate::storage_backend::ObjectStorageConfig; use anyhow::anyhow; use anyhow::Context; use bytes::Bytes; use futures_util::StreamExt; -use object_store::aws::AmazonS3; -use object_store::aws::AmazonS3Builder; -use object_store::aws::S3ConditionalPut; use object_store::path::Path; use object_store::Attribute; use object_store::Attributes; @@ -28,31 +26,24 @@ use object_store::PutMode; use object_store::PutOptions; use object_store::PutPayload; use object_store::RetryConfig; -use object_store::UpdateVersion; use std::borrow::Cow; use std::sync::Arc; use std::time::Duration; -/// Explicit credentials for a managed installation; everything else comes -/// from the standard `AWS_*` environment. -pub struct StaticCredentials { - pub access_key_id: String, - pub secret_access_key: String, - pub session_token: Option, -} - /// One S3-compatible bucket. Cheap to clone; each `open` builds its own /// HTTP transport, so a dedicated instance also isolates its traffic. #[derive(Clone)] pub struct Bucket { - store: Arc, + store: Arc, /// Conditional writes only, built with retries OFF: a retried CAS put - /// can land on the first attempt's own etag change and report a clean - /// 412 — converting "may have committed" into a false rejection. The - /// ambiguity must surface as `Err` so the caller reconciles. - cas_store: Arc, + /// can observe the first attempt's object-version change and report a + /// definite precondition rejection — converting "may have committed" into + /// a false rejection. The ambiguity must surface as `Err` so the caller + /// reconciles. + cas_store: Arc, /// Bucket name, for messages — the store is already bound to it. pub name: String, + storage_config: ObjectStorageConfig, } impl Bucket { @@ -60,10 +51,7 @@ impl Bucket { /// AppName format, `app/`), keeping e.g. the lease safety lane /// observable in black-box storage traces. pub fn open( - bucket: &str, - endpoint: Option<&str>, - region: &str, - credentials: Option, + storage_config: ObjectStorageConfig, app: Option<&str>, ) -> anyhow::Result { // These bounds mirror the aws-sdk TimeoutConfig they replace @@ -80,65 +68,46 @@ impl Bucket { .context("app user agent")?, ); } - let mut builder = AmazonS3Builder::from_env() - .with_bucket_name(bucket) - .with_region(region) - .with_conditional_put(S3ConditionalPut::ETagMatch) - .with_retry(RetryConfig { - max_retries: 2, - retry_timeout: Duration::from_secs(30), - ..RetryConfig::default() - }) - .with_client_options(options); - if let Some(endpoint) = endpoint { - // Path-style against explicit S3-compatible endpoints, exactly - // as the aws client's force_path_style(endpoint.is_some()). - builder = builder - .with_endpoint(endpoint) - .with_virtual_hosted_style_request(false); - } else { - builder = builder.with_virtual_hosted_style_request(true); - } - if let Some(credentials) = credentials { - builder = builder - .with_access_key_id(credentials.access_key_id) - .with_secret_access_key(credentials.secret_access_key); - if let Some(token) = credentials.session_token { - builder = builder.with_token(token); - } - } - let cas_builder = builder.clone().with_retry(RetryConfig { - max_retries: 0, + let retry = |max_retries| RetryConfig { + max_retries, retry_timeout: Duration::from_secs(30), ..RetryConfig::default() - }); + }; + let ordinary_retry = retry(2); + let cas_retry = retry(0); + let (store, cas_store) = + storage_config.build_bucket_stores(options, ordinary_retry, cas_retry)?; Ok(Bucket { - store: Arc::new(builder.build().context("build s3 client")?), - cas_store: Arc::new(cas_builder.build().context("build s3 cas client")?), - name: bucket.to_string(), + store, + cas_store, + name: storage_config.bucket().to_string(), + storage_config, }) } - /// Body and etag, or `None` when the key does not exist. + /// Body and object version, or `None` when the key does not exist. pub async fn get(&self, key: &str) -> anyhow::Result> { match self.store.get(&Path::from(key)).await { Ok(result) => { - let etag = result.meta.e_tag.clone().unwrap_or_default(); + let version = self.storage_config.object_version(&result.meta); let bytes = result .bytes() .await .with_context(|| format!("read body s3://{}/{key}", self.name))?; - Ok(Some((bytes, etag))) + Ok(Some((bytes, version))) } Err(Error::NotFound { .. }) => Ok(None), Err(error) => Err(anyhow!(error).context(format!("read s3://{}/{key}", self.name))), } } - /// Size and etag, or `None` when the key does not exist. + /// Size and object version, or `None` when the key does not exist. pub async fn head(&self, key: &str) -> anyhow::Result> { match self.store.head(&Path::from(key)).await { - Ok(meta) => Ok(Some((meta.size as u64, meta.e_tag.unwrap_or_default()))), + Ok(meta) => { + let version = self.storage_config.object_version(&meta); + Ok(Some((meta.size as u64, version))) + } Err(Error::NotFound { .. }) => Ok(None), Err(error) => Err(anyhow!(error).context(format!("head s3://{}/{key}", self.name))), } @@ -201,29 +170,29 @@ impl Bucket { Ok(()) } - /// Conditional write. `etag: None` requires the key to be absent - /// (If-None-Match: *); `Some` requires the current etag (If-Match). - /// `Ok(Some(new_etag))` applied, `Ok(None)` cleanly rejected; any other + /// Conditional write. `version: None` requires the key to be absent; + /// `Some` requires the current provider object version. + /// `Ok(Some(new_version))` applied, `Ok(None)` cleanly rejected; any other /// failure is ambiguous and stays an error. pub async fn put_cas( &self, key: &str, body: impl Into, - etag: Option<&str>, + version: Option<&str>, ) -> anyhow::Result> { - let mode = match etag { + let mode = match version { None => PutMode::Create, - Some(etag) => PutMode::Update(UpdateVersion { - e_tag: Some(etag.to_string()), - version: None, - }), + Some(version) => PutMode::Update(self.storage_config.update_version(version)), }; match self .cas_store .put_opts(&Path::from(key), body.into(), PutOptions::from(mode)) .await { - Ok(result) => Ok(Some(result.e_tag.unwrap_or_default())), + Ok(result) => { + let new_version = self.storage_config.put_result_version(result); + Ok(Some(new_version)) + } Err(Error::Precondition { .. } | Error::AlreadyExists { .. }) => Ok(None), Err(error) => Err(anyhow!(error).context(format!( "conditional write s3://{}/{key} may have committed", @@ -308,7 +277,8 @@ pub fn is_unauthorized(error: &anyhow::Error) -> bool { #[cfg(test)] mod live_cas { - use super::{Bucket, StaticCredentials}; + use super::Bucket; + use crate::storage_backend::{ObjectStorageConfig, StaticCredentials}; // Live CAS contract against a real S3-compatible bucket (R2). Gated on // CELLD_CAS_LIVE=1 so it never runs in CI; a mock cannot answer whether @@ -324,19 +294,15 @@ mod live_cas { let name = std::env::var("CELLD_CAS_BUCKET").expect("CELLD_CAS_BUCKET"); let endpoint = std::env::var("CELLD_CAS_ENDPOINT").ok(); let region = std::env::var("AWS_REGION").unwrap_or_else(|_| "auto".into()); - let creds = StaticCredentials { + let credentials = StaticCredentials { access_key_id: std::env::var("AWS_ACCESS_KEY_ID").unwrap(), secret_access_key: std::env::var("AWS_SECRET_ACCESS_KEY").unwrap(), session_token: std::env::var("AWS_SESSION_TOKEN").ok(), }; - let bucket = Bucket::open( - &name, - endpoint.as_deref(), - ®ion, - Some(creds), - Some("cas-test"), - ) - .expect("open bucket"); + let storage = ObjectStorageConfig::from_bucket_uri(&name, endpoint.as_deref(), ®ion) + .expect("normalize storage") + .with_credentials(credentials); + let bucket = Bucket::open(storage, Some("cas-test")).expect("open bucket"); let nanos = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .unwrap() diff --git a/crates/celld/dead_node_gc.rs b/crates/celld/dead_node_gc.rs index 1564a4e2d..3134833ef 100644 --- a/crates/celld/dead_node_gc.rs +++ b/crates/celld/dead_node_gc.rs @@ -293,7 +293,7 @@ async fn gc_markers( async fn retire_dead_node(bucket: &Bucket, node: &str, now_ms: u64) -> anyhow::Result { let key = format!("nodes/{node}.json"); - let Some((record, etag)) = read_node(bucket, &key).await? else { + let Some((record, version)) = read_node(bucket, &key).await? else { return Ok(true); }; if !celld_logic::dead_node_reconciliation::node_record_is_dead( @@ -312,7 +312,7 @@ async fn retire_dead_node(bucket: &Bucket, node: &str, now_ms: u64) -> anyhow::R expires_ms: 0, ..record })?; - match bucket.put_cas(&key, tombstone, Some(&etag)).await? { + match bucket.put_cas(&key, tombstone, Some(&version)).await? { Some(_) => { bucket.delete(&key).await?; Ok(true) @@ -322,10 +322,10 @@ async fn retire_dead_node(bucket: &Bucket, node: &str, now_ms: u64) -> anyhow::R } async fn read_node(bucket: &Bucket, key: &str) -> anyhow::Result> { - let Some((bytes, etag)) = bucket.get(key).await? else { + let Some((bytes, version)) = bucket.get(key).await? else { return Ok(None); }; let record = serde_json::from_slice(&bytes) .with_context(|| format!("decode s3://{}/{key}", bucket.name))?; - Ok(Some((record, etag))) + Ok(Some((record, version))) } diff --git a/crates/celld/fleet.rs b/crates/celld/fleet.rs index 53d387903..dda587a48 100644 --- a/crates/celld/fleet.rs +++ b/crates/celld/fleet.rs @@ -5,6 +5,7 @@ use crate::bucket::Bucket; use crate::deploy; use crate::js::WorkerConfigOptions; +use crate::storage_backend::{ObjectStorageConfig, StaticCredentials}; use crate::protocol::{DeployPointer, Manifest}; use anyhow::{bail, Context}; use serde::Deserialize; @@ -12,8 +13,59 @@ use std::collections::BTreeMap; use std::time::Duration; use tracing::info; -pub fn s3_client(bucket: &str, endpoint: Option<&str>, region: &str) -> anyhow::Result { - s3_client_with_credentials(bucket, endpoint, region, None) +pub fn normalize_storage( + bucket: &str, + endpoint: Option<&str>, + region: &str, + managed: Option<&crate::control_plane::ManagedStorageConfig>, +) -> anyhow::Result { + if let Some(managed) = managed { + let name = bucket.trim_start_matches("s3://"); + return ObjectStorageConfig::managed( + name, + managed.region.clone(), + managed.endpoint.clone(), + StaticCredentials { + access_key_id: managed.access_key_id.clone(), + secret_access_key: managed.secret_access_key.clone(), + session_token: managed.session_token.clone(), + }, + ); + } + + ObjectStorageConfig::from_bucket_uri(bucket, endpoint, region) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn managed_storage_reaches_ltx_config() { + let managed = crate::control_plane::ManagedStorageConfig { + bucket: "managed-bucket".into(), + endpoint: "https://managed.example".into(), + region: "managed-region".into(), + access_key_id: "managed-access-key".into(), + secret_access_key: "managed-secret-key".into(), + session_token: Some("managed-session-token".into()), + }; + let storage = normalize_storage(&managed.bucket, None, "ignored", Some(&managed)).unwrap(); + let replica = storage.replica_config("replicas/epoch".into()); + + assert_eq!(replica.bucket, "managed-bucket"); + assert_eq!(replica.path, "replicas/epoch"); + assert_eq!(replica.endpoint, "https://managed.example"); + assert_eq!(replica.region, "managed-region"); + assert_eq!(replica.access_key_id, "managed-access-key"); + assert_eq!(replica.secret_access_key, "managed-secret-key"); + assert_eq!(replica.session_token, "managed-session-token"); + assert!(replica.force_path_style); + } +} + +pub fn s3_client(backend: &ObjectStorageConfig) -> anyhow::Result { + Bucket::open(backend.clone(), None) } /// Build the authority-heartbeat client on its own HTTP connection pool. @@ -23,36 +75,9 @@ pub fn s3_client(bucket: &str, endpoint: Option<&str>, region: &str) -> anyhow:: /// dedicated instance keeps the safety lane isolated, and the `celld-lease` /// app tag labels it in black-box storage traces. pub fn s3_lease_client_with_credentials( - bucket: &str, - endpoint: Option<&str>, - region: &str, - managed: Option<&crate::control_plane::ManagedStorageConfig>, -) -> anyhow::Result { - open(bucket, endpoint, region, managed, Some("celld-lease")) -} - -pub fn s3_client_with_credentials( - bucket: &str, - endpoint: Option<&str>, - region: &str, - managed: Option<&crate::control_plane::ManagedStorageConfig>, -) -> anyhow::Result { - open(bucket, endpoint, region, managed, None) -} - -fn open( - bucket: &str, - endpoint: Option<&str>, - region: &str, - managed: Option<&crate::control_plane::ManagedStorageConfig>, - app: Option<&str>, + backend: &ObjectStorageConfig, ) -> anyhow::Result { - let credentials = managed.map(|managed| crate::bucket::StaticCredentials { - access_key_id: managed.access_key_id.clone(), - secret_access_key: managed.secret_access_key.clone(), - session_token: managed.session_token.clone(), - }); - Bucket::open(bucket, endpoint, region, credentials, app) + Bucket::open(backend.clone(), Some("celld-lease")) } pub async fn validate_bucket(bucket: &Bucket) -> anyhow::Result<()> { @@ -323,7 +348,8 @@ pub async fn run_deploy(arguments: Vec) -> anyhow::Result<()> { .or_else(|| env("AWS_REGION")) .or_else(|| env("AWS_DEFAULT_REGION")) .unwrap_or_else(|| "us-east-1".to_string()); - let store = s3_client(&bucket, options.endpoint.as_deref(), ®ion)?; + let storage = normalize_storage(&bucket, options.endpoint.as_deref(), ®ion, None)?; + let store = s3_client(&storage)?; validate_bucket(&store).await?; let started = std::time::Instant::now(); deploy::write(&store, &built).await?; diff --git a/crates/celld/lib.rs b/crates/celld/lib.rs index 2ee46f72d..3b4661c80 100644 --- a/crates/celld/lib.rs +++ b/crates/celld/lib.rs @@ -20,6 +20,8 @@ pub mod deploy; /// for it. #[cfg(all(test, celld_internal_tests))] mod fault; +mod storage_backend; + pub mod fleet; pub mod js; pub mod ltx_repl; diff --git a/crates/celld/ltx_repl.rs b/crates/celld/ltx_repl.rs index 02dcb8922..cb5aa265d 100644 --- a/crates/celld/ltx_repl.rs +++ b/crates/celld/ltx_repl.rs @@ -21,12 +21,12 @@ use std::sync::Arc; use std::sync::Mutex; use std::time::Duration; +use crate::storage_backend::ObjectStorageConfig; use anyhow::anyhow; use celld_ltx::object_store::ObjectStore; use celld_ltx::replica; use celld_ltx::Db; use celld_ltx::ObjectStoreClient; -use celld_ltx::ObjectStoreConfig; use celld_ltx::Replica; use celld_ltx::TXID; use tokio::sync::Notify; @@ -39,7 +39,6 @@ use crate::replication::sqlite_snapshot; use crate::replication::ActivationOptions; use crate::replication::ActivationResult; use crate::replication::RestoredSnapshot; -use crate::replication::StorageCredentials; use crate::replication::SyncWait; /// Max cells uploading concurrently across the node. Caps blocking-pool threads @@ -70,10 +69,7 @@ type CellHandle = Arc; pub struct LtxRepl { /// Local root: cell dbs live at `watch//ltx/e/db.sqlite`. watch: PathBuf, - bucket: String, - endpoint: Option, - region: String, - credentials: Option, + storage: ObjectStorageConfig, /// One connection pool for the whole node, shared by every cell client. store: Arc, cells: Arc>>, @@ -83,15 +79,9 @@ pub struct LtxRepl { } impl LtxRepl { - pub fn start( - watch: &Path, - bucket: String, - endpoint: Option, - region: String, - credentials: Option, - ) -> anyhow::Result { - let store = node_config(&bucket, endpoint.as_deref(), ®ion, credentials.as_ref()) - .build_store() + pub fn start(watch: &Path, storage: ObjectStorageConfig) -> anyhow::Result { + let store = storage + .build_ltx_store() .map_err(|error| anyhow!("build shared object store: {error}"))?; let cells: Arc>> = Arc::default(); let dirty = Arc::new(Notify::new()); @@ -101,10 +91,7 @@ impl LtxRepl { tokio::spawn(sync_loop(cells.clone(), dirty.clone(), slots)); Ok(Self { watch: watch.to_path_buf(), - bucket, - endpoint, - region, - credentials, + storage, store, cells, dirty, @@ -123,13 +110,9 @@ impl LtxRepl { /// prefix. `cells//ltx/e` matches [`Self::db_path`]'s remote /// twin so the same coordinates address local and replica state. fn client_for(&self, cell: &str, epoch: u64) -> ObjectStoreClient { - let mut config = node_config( - &self.bucket, - self.endpoint.as_deref(), - &self.region, - self.credentials.as_ref(), - ); - config.path = format!("cells/{cell}/ltx/e{epoch}"); + let config = self + .storage + .replica_config(format!("cells/{cell}/ltx/e{epoch}")); ObjectStoreClient::with_store(config, self.store.clone()) } @@ -484,51 +467,3 @@ async fn sync_loop( } } } - -/// Node-level object-store config (no per-cell prefix). `build_store` on this -/// yields the one shared client; per-cell clients set only `path`. -fn node_config( - bucket: &str, - endpoint: Option<&str>, - region: &str, - credentials: Option<&StorageCredentials>, -) -> ObjectStoreConfig { - let endpoint = endpoint.unwrap_or_default().to_string(); - // Static credentials come from the managed control plane when present, - // else the `AWS_*` env the node already carries. Without this, - // `build_store` sees empty keys and object_store falls back to the - // instance credential provider, which off-EC2 sends unsigned requests (R2 - // answers "404 page not found"). - let env = |key: &str| std::env::var(key).ok().filter(|value| !value.is_empty()); - let access_key_id = credentials - .map(|c| c.access_key_id.clone()) - .filter(|value| !value.is_empty()) - .or_else(|| env("AWS_ACCESS_KEY_ID")) - .unwrap_or_default(); - let secret_access_key = credentials - .map(|c| c.secret_access_key.clone()) - .filter(|value| !value.is_empty()) - .or_else(|| env("AWS_SECRET_ACCESS_KEY")) - .unwrap_or_default(); - // Temporary R2/STS credentials require the session token, or signing fails. - let session_token = credentials - .and_then(|c| c.session_token.clone()) - .filter(|value| !value.is_empty()) - .or_else(|| env("AWS_SESSION_TOKEN")) - .unwrap_or_default(); - ObjectStoreConfig { - bucket: bucket.to_string(), - path: String::new(), - region: region.to_string(), - // A custom endpoint (R2/MinIO) uses path-style addressing, matching - // `ObjectStoreConfig::from_url`'s default for non-AWS hosts. - force_path_style: !endpoint.is_empty(), - endpoint, - access_key_id, - secret_access_key, - session_token, - skip_verify: false, - part_size: 0, - concurrency: 0, - } -} diff --git a/crates/celld/main.rs b/crates/celld/main.rs index 4f3b6670a..87fe4912d 100644 --- a/crates/celld/main.rs +++ b/crates/celld/main.rs @@ -70,7 +70,7 @@ struct MemoryOwnership { node: String, owners: BTreeMap, leases: BTreeMap, - next_etag: u64, + next_version: u64, } #[derive(Clone)] @@ -136,18 +136,18 @@ impl Ownership { CasGuard::Match(expected) => memory .owners .get(cell) - .is_some_and(|owner| owner.etag == expected), + .is_some_and(|owner| owner.version == expected), }; if allowed { - let etag = format!("e{}", memory.next_etag); - memory.next_etag += 1; + let version = format!("e{}", memory.next_version); + memory.next_version += 1; let node = memory.node.clone(); memory.owners.insert( cell.into(), OwnerRecord { node: Some(node), epoch, - etag, + version, }, ); } @@ -157,12 +157,14 @@ impl Ownership { CasOutcome::Rejected }) } - Self::S3(s3) => s3.cas_owner(cell, guard, epoch).await.map_err(|error| { - // Any transport or 5xx failure may have happened after S3 - // committed. The core reconciles by reading the owner again. - eprintln!("celld ownership CAS ambiguous: {error:#}"); - Failure::Ambiguous - }), + Self::S3(s3) => { + s3.cas_owner(cell, guard, epoch).await.map_err(|error| { + // Any transport or 5xx failure may have happened after S3 + // committed. The core reconciles by reading the owner again. + eprintln!("celld ownership CAS ambiguous: {error:#}"); + Failure::Ambiguous + }) + } } } @@ -175,14 +177,14 @@ impl Ownership { owner.node.as_deref() == Some(node.as_str()) && owner.epoch == epoch }); if releasable { - let etag = format!("e{}", memory.next_etag); - memory.next_etag += 1; + let version = format!("e{}", memory.next_version); + memory.next_version += 1; memory.owners.insert( cell.into(), OwnerRecord { node: None, epoch, - etag, + version, }, ); } @@ -215,22 +217,24 @@ impl Ownership { let allowed = match guard { CasGuard::Absent => current.is_none(), CasGuard::Match(expected) => { - current.is_some_and(|lease| lease.etag == expected) + current.is_some_and(|lease| lease.version == expected) } }; if !allowed { return Ok(LeaseCasOutcome::Rejected); } - let etag = format!("e{}", memory.next_etag); - memory.next_etag += 1; - record.etag = etag.clone(); + let version = format!("e{}", memory.next_version); + memory.next_version += 1; + record.version = version.clone(); memory.leases.insert(record.node.clone(), record); - Ok(LeaseCasOutcome::Applied { etag }) + Ok(LeaseCasOutcome::Applied { version }) + } + Self::S3(s3) => { + s3.cas_node_lease(guard, &record).await.map_err(|error| { + eprintln!("celld node-lease CAS ambiguous: {error:#}"); + Failure::Ambiguous + }) } - Self::S3(s3) => s3.cas_node_lease(guard, &record).await.map_err(|error| { - eprintln!("celld node-lease CAS ambiguous: {error:#}"); - Failure::Ambiguous - }), } } @@ -830,7 +834,7 @@ impl Actor { node: node.clone(), owners: BTreeMap::new(), leases: BTreeMap::new(), - next_etag: 1, + next_version: 1, }))) }; let live_load = match &ownership { @@ -4455,12 +4459,13 @@ async fn async_main() -> anyhow::Result<()> { let bucket = settings .bucket .ok_or_else(|| anyhow::anyhow!("celld diagnose requires --bucket"))?; - let client = fleet::s3_client_with_credentials( + let backend = fleet::normalize_storage( &bucket, settings.endpoint.as_deref(), &settings.region, managed_storage.as_ref(), )?; + let client = fleet::s3_client(&backend)?; return fleet::diagnose(&client, peers, settings.unsafe_public_advertise).await; } Action::Run(settings) => settings, @@ -4518,14 +4523,18 @@ async fn async_main() -> anyhow::Result<()> { } else { None }; - let storage_credentials = - managed_storage - .as_ref() - .map(|storage| celld::replication::StorageCredentials { - access_key_id: storage.access_key_id.clone(), - secret_access_key: storage.secret_access_key.clone(), - session_token: storage.session_token.clone(), - }); + let storage_backend = settings + .bucket + .as_deref() + .map(|bucket| { + fleet::normalize_storage( + bucket, + settings.endpoint.as_deref(), + &settings.region, + managed_storage.as_ref(), + ) + }) + .transpose()?; let (tx, rx) = mpsc::unbounded_channel(); let sample_tx = tx.clone(); let alarm_tx = tx.clone(); @@ -4558,25 +4567,18 @@ async fn async_main() -> anyhow::Result<()> { .map(std::path::PathBuf::from) .unwrap_or_else(|| std::env::temp_dir().join(format!("celld-{}", std::process::id()))); let mut deploy_agent = None; + let load_backend = storage_backend + .as_ref() + .filter(|_| settings.load_deployment); let (runtime, ownership, peer_key, wake_scan, assets, asset_script) = - if let Some(bucket) = settings.bucket.clone().filter(|_| settings.load_deployment) { - let client = fleet::s3_client_with_credentials( - &bucket, - settings.endpoint.as_deref(), - &settings.region, - managed_storage.as_ref(), - )?; + if let Some(backend) = load_backend { + let client = fleet::s3_client(backend)?; if settings.control_plane { fleet::validate_managed_bucket(&client).await?; } else { fleet::validate_bucket(&client).await?; } - let lease_client = fleet::s3_lease_client_with_credentials( - &bucket, - settings.endpoint.as_deref(), - &settings.region, - managed_storage.as_ref(), - )?; + let lease_client = fleet::s3_lease_client_with_credentials(backend)?; if settings.control_plane { celld::control_plane::wait_for_initial_deployment(&client).await?; deploy_agent = Some(client.clone()); @@ -4637,13 +4639,7 @@ async fn async_main() -> anyhow::Result<()> { // of today's manifest. Start it even for a stateless deployment // so a later deployment can introduce cells without changing the // durability contract underneath the node. - let replication = Some(Replication::start( - client.clone(), - &data_dir, - settings.endpoint.clone(), - settings.region.clone(), - storage_credentials.clone(), - )?); + let replication = Some(Replication::start(backend.clone(), &data_dir)?); let asset_script = Some(Arc::::from(primary_script)); let assets = Arc::new(asset_resolvers); let runtime = RuntimeManager::start(RuntimeOptions { @@ -4706,20 +4702,10 @@ async fn async_main() -> anyhow::Result<()> { text: Vec::new(), compat: Compat::default(), }; - let (ownership, peer_key, wake, wake_scan) = match settings.bucket.clone() { - Some(bucket) => { - let client = fleet::s3_client_with_credentials( - &bucket, - settings.endpoint.as_deref(), - &settings.region, - managed_storage.as_ref(), - )?; - let lease_client = fleet::s3_lease_client_with_credentials( - &bucket, - settings.endpoint.as_deref(), - &settings.region, - managed_storage.as_ref(), - )?; + let (ownership, peer_key, wake, wake_scan) = match storage_backend.as_ref() { + Some(backend) => { + let client = fleet::s3_client(backend)?; + let lease_client = fleet::s3_lease_client_with_credentials(backend)?; let peer_key = peer_auth::load_or_create(&client).await?; let wake = Arc::new(celld::wake::WakeFlusher::new()); celld::js::set_arm_gate(ArmGate { @@ -4766,20 +4752,10 @@ async fn async_main() -> anyhow::Result<()> { None, ) } else { - let (ownership, peer_key) = match settings.bucket.clone() { - Some(bucket) => { - let client = fleet::s3_client_with_credentials( - &bucket, - settings.endpoint.as_deref(), - &settings.region, - managed_storage.as_ref(), - )?; - let lease_client = fleet::s3_lease_client_with_credentials( - &bucket, - settings.endpoint.as_deref(), - &settings.region, - managed_storage.as_ref(), - )?; + let (ownership, peer_key) = match storage_backend.as_ref() { + Some(backend) => { + let client = fleet::s3_client(backend)?; + let lease_client = fleet::s3_lease_client_with_credentials(backend)?; let peer_key = peer_auth::load_or_create(&client).await?; ( Some(Ownership::S3(Arc::new( diff --git a/crates/celld/ownership_store.rs b/crates/celld/ownership_store.rs index 63f1fa367..cafeffbea 100644 --- a/crates/celld/ownership_store.rs +++ b/crates/celld/ownership_store.rs @@ -1,6 +1,6 @@ // Copyright 2026 Deno Land Inc. Apache-2.0 license. -//! S3-compatible ownership effect adapter. +//! Object-store ownership effect adapter. //! //! This module deliberately contains serialization, wall-clock sampling, SDK //! configuration and error classification only. Ownership decisions remain in @@ -140,10 +140,12 @@ impl S3Ownership { .ok() .or_else(|| std::env::var("AWS_DEFAULT_REGION").ok()) .unwrap_or_else(|| "us-east-1".into()); - Ok(Self::new( - Bucket::open(&bucket, endpoint.as_deref(), ®ion, None, None)?, - node, - )) + let storage = crate::storage_backend::ObjectStorageConfig::from_bucket_uri( + &bucket, + endpoint.as_deref(), + ®ion, + )?; + Ok(Self::new(Bucket::open(storage, None)?, node)) } /// The lease lifetime this fleet renews on, used to decide which node @@ -165,13 +167,13 @@ impl S3Ownership { pub async fn read_owner(&self, cell: &str) -> anyhow::Result> { let key = format!("cells/{cell}/own.json"); - let Some((owner, etag)) = self.read_json::(&key).await? else { + let Some((owner, version)) = self.read_json::(&key).await? else { return Ok(None); }; Ok(Some(OwnerRecord { node: (!owner.node.is_empty()).then_some(owner.node), epoch: owner.epoch, - etag, + version, })) } @@ -196,13 +198,13 @@ impl S3Ownership { Ok(self .read_json_with::(bucket, &key) .await? - .map(|(lease, etag)| NodeLeaseRecord { + .map(|(lease, version)| NodeLeaseRecord { node: lease.node, addr: lease.addr, expires_ms: lease.expires_ms, peer_protocol: lease.peer_protocol, generation: lease.generation, - etag, + version, })) } @@ -284,7 +286,11 @@ impl S3Ownership { } let key = format!("cells/{cell}/own.json"); let body = serde_json::to_vec(&OwnerWire { node: "", epoch })?; - match self.bucket.put_cas(&key, body, Some(¤t.etag)).await? { + match self + .bucket + .put_cas(&key, body, Some(¤t.version)) + .await? + { Some(_) => Ok(CasOutcome::Applied), None => Ok(CasOutcome::Rejected), } @@ -301,11 +307,11 @@ impl S3Ownership { node: &self.node, epoch, })?; - let etag = match &guard { + let version = match &guard { CasGuard::Absent => None, - CasGuard::Match(etag) => Some(etag.as_str()), + CasGuard::Match(version) => Some(version.as_str()), }; - match self.bucket.put_cas(&key, body, etag).await? { + match self.bucket.put_cas(&key, body, version).await? { Some(_) => Ok(CasOutcome::Applied), None => Ok(CasOutcome::Rejected), } @@ -331,12 +337,12 @@ impl S3Ownership { generation: record.generation.clone(), load: process_load(&self.live), })?; - let etag = match &guard { + let version = match &guard { CasGuard::Absent => None, - CasGuard::Match(etag) => Some(etag.as_str()), + CasGuard::Match(version) => Some(version.as_str()), }; - match self.lease_bucket.put_cas(&key, body, etag).await? { - Some(etag) => Ok(LeaseCasOutcome::Applied { etag }), + match self.lease_bucket.put_cas(&key, body, version).await? { + Some(version) => Ok(LeaseCasOutcome::Applied { version }), None => Ok(LeaseCasOutcome::Rejected), } } @@ -353,12 +359,12 @@ impl S3Ownership { bucket: &Bucket, key: &str, ) -> anyhow::Result> { - let Some((bytes, etag)) = bucket.get(key).await? else { + let Some((bytes, version)) = bucket.get(key).await? else { return Ok(None); }; let value = serde_json::from_slice(&bytes) .with_context(|| format!("decode s3://{}/{key}", bucket.name))?; - Ok(Some((value, etag))) + Ok(Some((value, version))) } } diff --git a/crates/celld/runtime.rs b/crates/celld/runtime.rs index 7f135eaae..8786b191a 100644 --- a/crates/celld/runtime.rs +++ b/crates/celld/runtime.rs @@ -11,7 +11,7 @@ use crate::js::{ self, CellJob, FetchRequest, HttpResponse, Worker, WorkerConfig, WorkerConfigOptions, }; use crate::ltx_repl::LtxRepl; -use crate::replication::{ActivationOptions, StorageCredentials, SyncWait}; +use crate::replication::{ActivationOptions, SyncWait}; use crate::storage; use crate::wake::WakeFlusher; use anyhow::{anyhow, Context}; @@ -106,20 +106,11 @@ pub struct Replication { impl Replication { pub fn start( - bucket: crate::bucket::Bucket, + storage: crate::storage_backend::ObjectStorageConfig, watch: &Path, - endpoint: Option, - region: String, - credentials: Option, ) -> anyhow::Result { Ok(Self { - ltx: Arc::new(LtxRepl::start( - watch, - bucket.name, - endpoint, - region, - credentials, - )?), + ltx: Arc::new(LtxRepl::start(watch, storage)?), }) } diff --git a/crates/celld/storage_backend.rs b/crates/celld/storage_backend.rs new file mode 100644 index 000000000..5de15fd42 --- /dev/null +++ b/crates/celld/storage_backend.rs @@ -0,0 +1,197 @@ +//! Celld's daemon-wide object-storage policy. + +use anyhow::Context; +use object_store::aws::{AmazonS3Builder, S3ConditionalPut}; +use object_store::{ClientOptions, ObjectMeta, ObjectStore, PutResult, RetryConfig}; +use std::sync::Arc; + +#[derive(Clone, PartialEq, Eq)] +pub(crate) struct StaticCredentials { + pub(crate) access_key_id: String, + pub(crate) secret_access_key: String, + pub(crate) session_token: Option, +} + +#[derive(Clone, PartialEq, Eq)] +pub struct ObjectStorageConfig { + bucket: String, + region: String, + endpoint: Option, + credentials: Option, + force_path_style: bool, +} + +impl ObjectStorageConfig { + pub(crate) fn from_bucket_uri( + bucket: &str, + endpoint: Option<&str>, + region: &str, + ) -> anyhow::Result { + Self::s3(bucket, endpoint, region, None) + } + + pub(crate) fn s3( + bucket: &str, + endpoint: Option<&str>, + region: &str, + credentials: Option, + ) -> anyhow::Result { + let bucket = bucket.trim_start_matches("s3://"); + Ok(Self { + bucket: bucket.into(), + region: region.into(), + endpoint: endpoint.map(Into::into), + credentials, + force_path_style: endpoint.is_some(), + }) + } + + pub(crate) fn managed( + bucket: &str, + region: String, + endpoint: String, + credentials: StaticCredentials, + ) -> anyhow::Result { + Self::s3(bucket, Some(&endpoint), ®ion, Some(credentials)) + } + + #[cfg(test)] + pub(crate) fn with_credentials(mut self, credentials: StaticCredentials) -> Self { + self.credentials = Some(credentials); + self + } + + pub(crate) fn bucket(&self) -> &str { + &self.bucket + } + + fn runtime_builder(&self) -> AmazonS3Builder { + let mut builder = AmazonS3Builder::from_env() + .with_bucket_name(&self.bucket) + .with_region(&self.region) + .with_virtual_hosted_style_request(!self.force_path_style); + if let Some(endpoint) = self.endpoint.as_deref() { + builder = builder.with_endpoint(endpoint); + } + if let Some(StaticCredentials { + access_key_id, + secret_access_key, + session_token, + }) = &self.credentials + { + builder = builder + .with_access_key_id(access_key_id) + .with_secret_access_key(secret_access_key); + if let Some(token) = session_token.as_deref() { + builder = builder.with_token(token); + } + } + builder + } + + pub(crate) fn build_ltx_store(&self) -> anyhow::Result> { + self.replica_config(String::new()) + .build_store() + .map_err(anyhow::Error::from) + } + + pub(crate) fn build_bucket_stores( + &self, + options: ClientOptions, + ordinary: RetryConfig, + cas: RetryConfig, + ) -> anyhow::Result<(Arc, Arc)> { + let builder = self + .runtime_builder() + .with_client_options(options) + .with_conditional_put(S3ConditionalPut::ETagMatch); + let store = builder + .clone() + .with_retry(ordinary) + .build() + .context("build s3 client")?; + let cas_store = builder + .with_retry(cas) + .build() + .context("build s3 cas client")?; + Ok((Arc::new(store), Arc::new(cas_store))) + } + + pub(crate) fn object_version(&self, meta: &ObjectMeta) -> String { + meta.e_tag.clone().unwrap_or_default() + } + pub(crate) fn put_result_version(&self, result: PutResult) -> String { + result.e_tag.unwrap_or_default() + } + pub(crate) fn update_version(&self, version: &str) -> object_store::UpdateVersion { + object_store::UpdateVersion { + e_tag: Some(version.into()), + version: None, + } + } + + pub(crate) fn replica_config(&self, path: String) -> celld_ltx::ObjectStoreConfig { + let env = |name| std::env::var(name).ok().filter(|value| !value.is_empty()); + let (access_key_id, secret_access_key, session_token) = match &self.credentials { + None => ( + env("AWS_ACCESS_KEY_ID").unwrap_or_default(), + env("AWS_SECRET_ACCESS_KEY").unwrap_or_default(), + env("AWS_SESSION_TOKEN").unwrap_or_default(), + ), + Some(StaticCredentials { + access_key_id, + secret_access_key, + session_token, + }) => ( + (!access_key_id.is_empty()) + .then(|| access_key_id.clone()) + .or_else(|| env("AWS_ACCESS_KEY_ID")) + .unwrap_or_default(), + (!secret_access_key.is_empty()) + .then(|| secret_access_key.clone()) + .or_else(|| env("AWS_SECRET_ACCESS_KEY")) + .unwrap_or_default(), + session_token + .clone() + .filter(|value| !value.is_empty()) + .or_else(|| env("AWS_SESSION_TOKEN")) + .unwrap_or_default(), + ), + }; + celld_ltx::ObjectStoreConfig { + bucket: self.bucket.clone(), + path, + region: self.region.clone(), + endpoint: self.endpoint.clone().unwrap_or_default(), + access_key_id, + secret_access_key, + session_token, + force_path_style: self + .endpoint + .as_deref() + .is_some_and(|value| !value.is_empty()), + skip_verify: false, + part_size: 0, + concurrency: 0, + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + #[test] + fn parses_s3_and_bare_identically() { + assert!( + ObjectStorageConfig::from_bucket_uri("bucket", None, "r").unwrap() + == ObjectStorageConfig::from_bucket_uri("s3://bucket", None, "r").unwrap() + ); + } + #[test] + fn maps_etag_version() { + let storage = ObjectStorageConfig::from_bucket_uri("b", None, "r").unwrap(); + let update = storage.update_version("tag"); + assert_eq!(update.e_tag.as_deref(), Some("tag")); + assert!(update.version.is_none()); + } +} diff --git a/crates/celld/wake.rs b/crates/celld/wake.rs index f38b07ddd..789b27d35 100644 --- a/crates/celld/wake.rs +++ b/crates/celld/wake.rs @@ -76,7 +76,7 @@ pub async fn try_hold_waker(bucket: &Bucket, node: &str, now_ms: i64, ttl_ms: i6 bucket.put_cas(KEY, body(now_ms + ttl_ms), None).await, Ok(Some(_)) ), - Ok(Some((bytes, etag))) => { + Ok(Some((bytes, version))) => { let text = String::from_utf8_lossy(&bytes); let held_by_us = text.contains(&format!("\"node\":{node:?}")); let expires = text @@ -87,7 +87,7 @@ pub async fn try_hold_waker(bucket: &Bucket, node: &str, now_ms: i64, ttl_ms: i6 if celld_logic::wake::waker_may_claim(held_by_us, expires, now_ms) { matches!( bucket - .put_cas(KEY, body(now_ms + ttl_ms), Some(&etag)) + .put_cas(KEY, body(now_ms + ttl_ms), Some(&version)) .await, Ok(Some(_)) ) diff --git a/crates/logic/lib.rs b/crates/logic/lib.rs index 3b7a46963..27bafbfa0 100644 --- a/crates/logic/lib.rs +++ b/crates/logic/lib.rs @@ -140,7 +140,7 @@ pub struct OwnerRecord { /// `None` is a deliberately released, fenced record. Epochs never reset. pub node: Option, pub epoch: Epoch, - pub etag: String, + pub version: String, } /// The routing and authority fields read from `nodes/.json`. @@ -157,7 +157,7 @@ pub struct NodeLeaseRecord { /// `ownership_index_generation`. pub generation: String, /// Object version observed by the read. Empty only in synthetic events. - pub etag: String, + pub version: String, } /// One advisory fleet-capacity observation returned by the storage shell. @@ -298,7 +298,7 @@ pub struct RestoredAlarm { #[derive(Clone, Debug, PartialEq, Eq)] pub enum LeaseCasOutcome { - Applied { etag: String }, + Applied { version: String }, Rejected, } @@ -1794,7 +1794,7 @@ impl State { expires_ms: now_ms.saturating_add(spec.ttl_ms), peer_protocol: spec.peer_protocol, generation: spec.generation.clone(), - etag: String::new(), + version: String::new(), }; let pending = PendingNodeLease { spec, @@ -1942,7 +1942,7 @@ impl State { .is_some_and(|prior| same_node_lease(&record, &prior.record)) => { let mut prior = pending.prior.expect("checked above"); - prior.record.etag = record.etag; + prior.record.version = record.version; self.resume_node_lease_after_failure(prior, now_mono_ms, effects); } Ok(record) if pending.prior.is_some() => { @@ -1962,7 +1962,7 @@ impl State { // node replaces its prior process generation immediately. The // ETag still serializes competing replacements, and a process // which loses that CAS never becomes authoritative. - self.begin_node_lease_write(pending, CasGuard::Match(record.etag), effects); + self.begin_node_lease_write(pending, CasGuard::Match(record.version), effects); } Ok(None) => self.begin_node_lease_write(pending, CasGuard::Absent, effects), } @@ -1992,9 +1992,9 @@ impl State { return; } match result { - Ok(LeaseCasOutcome::Applied { etag }) => { + Ok(LeaseCasOutcome::Applied { version }) => { let mut record = pending.desired; - record.etag = etag; + record.version = version; self.hold_node_lease(pending.spec, record, pending.prior, now_mono_ms, effects); } Ok(LeaseCasOutcome::Rejected) if pending.prior.is_some() => { @@ -2223,9 +2223,9 @@ impl State { expires_ms: now_ms.saturating_add(spec.ttl_ms), peer_protocol: spec.peer_protocol, generation: spec.generation.clone(), - etag: String::new(), + version: String::new(), }; - let guard = CasGuard::Match(prior.record.etag.clone()); + let guard = CasGuard::Match(prior.record.version.clone()); self.begin_node_lease_write( PendingNodeLease { spec, @@ -2814,7 +2814,7 @@ impl State { &id, &mut cell, Activation::Claim(Claim { - guard: CasGuard::Match(record.etag), + guard: CasGuard::Match(record.version), epoch, takeover: false, reconciles, @@ -2828,7 +2828,7 @@ impl State { &id, &mut cell, Claim { - guard: CasGuard::Match(record.etag), + guard: CasGuard::Match(record.version), epoch, takeover: true, reconciles, @@ -3073,7 +3073,7 @@ impl State { id, cell, Activation::Claim(Claim { - guard: CasGuard::Match(record.etag), + guard: CasGuard::Match(record.version), epoch, takeover: true, reconciles: 0,