use super::*;
use crate::bucket::{GcsBucket, PersistenceConfig, gcs::GcsClients};
use google_cloud_storage::{
appendable_object_writer::AppendableObjectWriter, model_ext::ReadRange,
};
impl RapidSnapshots {
pub(crate) async fn validate_gcs(
config: &PersistenceConfig,
clients: GcsClients,
) -> Result<()> {
let Some((buckets, archive_bucket)) = config.rapid_settings() else {
anyhow::bail!("Rapid persistence configuration required");
};
GcsBucket::with_clients(archive_bucket, clients.clone())?
.require_standard()
.await?;
for placement in buckets {
let bucket = clients
.control
.get_bucket()
.set_name(format!("projects/_/buckets/{}", placement.bucket))
.send()
.await?;
ensure!(
bucket.storage_class == "RAPID",
"Rapid bucket has the wrong storage class"
);
validate_retention(&bucket)?;
let locations = bucket
.custom_placement_config
.context("Rapid bucket zone missing")?
.data_locations;
ensure!(
locations.len() == 1 && locations[0].eq_ignore_ascii_case(&placement.zone),
"Rapid bucket does not occupy its configured zone"
);
}
Ok(())
}
pub(crate) fn gcs(
config: &PersistenceConfig,
clients: GcsClients,
stop: CancellationToken,
actor: Option<&crate::actor::ActorKey>,
) -> Result<Self> {
config.validate()?;
let PersistenceConfig::Rapid {
buckets,
archive_bucket,
archive_batch,
} = config
else {
anyhow::bail!("append-log persistence configuration required");
};
let archive = Arc::new(GcsBucket::with_clients(archive_bucket, clients.clone())?);
let snapshots = Arc::new(super::standard::GcsSnapshots::new(
archive_bucket,
clients.clone(),
)?);
let zones = buckets
.iter()
.map(|placement| {
Arc::new(GcsZone {
bucket: format!("projects/_/buckets/{}", placement.bucket),
clients: clients.clone(),
}) as Arc<dyn LogZone>
})
.collect();
let store = Self::new(
archive,
snapshots,
zones,
*archive_batch,
Arc::new(crate::litestream::compaction::RustCompactor),
stop,
)?;
if let Some(actor) = actor {
store.prepare(actor)?;
}
Ok(store)
}
}
struct GcsZone {
bucket: String,
clients: GcsClients,
}
struct GcsWriter(AppendableObjectWriter);
#[async_trait]
impl LogZone for GcsZone {
fn bucket(&self) -> &str {
&self.bucket
}
async fn open(&self, object: &str) -> Result<(Replica, Box<dyn LogWriter>)> {
let writer = self
.clients
.storage
.open_appendable_object(&self.bucket, object)
.set_if_generation_match(0)
.send()
.await?;
let replica = Replica {
bucket: self.bucket.clone(),
object: object.into(),
generation: writer.generation(),
};
Ok((replica, Box::new(GcsWriter(writer))))
}
async fn read(&self, replica: &Replica, fence: bool) -> Result<Bytes> {
let _transfer = self.clients.transfers.acquire().await?;
let mut writer = if fence {
Some(
self.clients
.storage
.reopen_appendable_object(&self.bucket, &replica.object, replica.generation)
.send()
.await?,
)
} else {
None
};
let persisted = if let Some(writer) = &mut writer {
Some(writer.flush().await?)
} else {
None
};
if persisted == Some(0) {
return Ok(Bytes::new());
}
let descriptor = self
.clients
.storage
.open_object(&self.bucket, &replica.object)
.set_generation(replica.generation)
.send()
.await?;
let size = descriptor.object().size;
ensure!(size >= 0, "invalid log segment size");
if size == 0 {
return Ok(Bytes::new());
}
let mut reader = descriptor
.read_range(ReadRange::segment(0, size as u64))
.await;
let mut download = crate::payload::Download::new();
let mut length = 0usize;
while let Some(chunk) = reader.next().await {
let chunk = chunk?;
length += chunk.len();
ensure!(length <= size as usize, "log read exceeded persisted size");
download = download.append(chunk).await?;
}
ensure!(
length == size as usize && persisted.is_none_or(|size| size as usize == length),
"incomplete fenced log read"
);
drop(writer);
download.finish().await
}
async fn read_range(&self, replica: &Replica, start: u64, length: u64) -> Result<Bytes> {
let _transfer = self.clients.transfers.acquire().await?;
let descriptor = self
.clients
.storage
.open_object(&self.bucket, &replica.object)
.set_generation(replica.generation)
.send()
.await?;
let mut reader = descriptor
.read_range(ReadRange::segment(start, length))
.await;
let mut download = crate::payload::Download::new();
let mut received = 0u64;
while let Some(chunk) = reader.next().await {
let chunk = chunk?;
received += chunk.len() as u64;
ensure!(received <= length, "Rapid range exceeded requested length");
download = download.append(chunk).await?;
}
ensure!(received == length, "incomplete Rapid range");
download.finish().await
}
async fn delete(&self, replica: &Replica) -> Result<()> {
match self
.clients
.control
.delete_object()
.set_bucket(&self.bucket)
.set_object(&replica.object)
.set_if_generation_match(replica.generation)
.send()
.await
{
Ok(_) => Ok(()),
Err(error)
if error.http_status_code() == Some(404)
|| error.status().is_some_and(|s| s.code.name() == "NOT_FOUND") =>
{
Ok(())
}
Err(error) => Err(error.into()),
}
}
}
#[async_trait]
impl LogWriter for GcsWriter {
async fn append_and_flush(&mut self, bytes: Bytes) -> Result<u64> {
use google_cloud_storage::streaming_source::StreamingSource;
let mut upload = crate::payload::Upload::new(bytes);
while let Some(chunk) = upload.next().await {
self.0.append(chunk?).await?;
}
Ok(self.0.flush().await?.try_into()?)
}
}
pub(crate) fn validate_retention(bucket: &google_cloud_storage::model::Bucket) -> Result<()> {
const PREFIX: &str = "durable-actors-v3-logs-";
if let Some(lifecycle) = &bucket.lifecycle {
for rule in &lifecycle.rule {
if rule
.action
.as_ref()
.is_none_or(|action| action.r#type != "Delete")
{
continue;
}
let prefixes = rule
.condition
.as_ref()
.map(|condition| condition.matches_prefix.as_slice())
.unwrap_or_default();
ensure!(
!prefixes.is_empty()
&& prefixes
.iter()
.all(|prefix| !PREFIX.starts_with(prefix) && !prefix.starts_with(PREFIX)),
"Rapid log objects must be excluded from automatic lifecycle deletion"
);
}
}
Ok(())
}
#[cfg(test)]
#[path = "../../../tests/unit/bucket/rapid/gcs.rs"]
mod tests;