Files
BlueArchiveToolkit/infrastructure/src/cas/filesystem.rs
T
nyaKazuha 786b739f99
bat-rust / Build and test Rust (push) Canceled after 0s
bat-rust / Build and test Go API (push) Canceled after 0s
fix(release): 收紧分发热路径与事务边界
2026-09-12 22:22:11 +08:00

235 lines
7.7 KiB
Rust

//! 文件系统 CAS 仓储适配层。
//!
//! 这里不实现 CAS 核心算法,只将 `bat-cas-engine` 适配到
//! `bat-core::repositories::CasRepository`。
use async_trait::async_trait;
use bat_cas_engine::hash::Hash;
use bat_cas_engine::repository as engine_repository;
use bat_core::repositories::cas_repository::{CasRepository, ObjectId};
use std::path::PathBuf;
use tokio::sync::OnceCell;
/// 文件系统 CAS Repository 适配器。
pub struct FileSystemCasRepository {
root: PathBuf,
inner: OnceCell<engine_repository::FileSystemCasRepository>,
}
impl FileSystemCasRepository {
/// 创建新的文件系统 CAS Repository 适配器。
pub fn new(root: impl Into<PathBuf>) -> Self {
Self {
root: root.into(),
inner: OnceCell::new(),
}
}
/// 初始化底层 CAS 引擎。
pub async fn init(&self) -> bat_core::Result<()> {
self.engine().await.map(|_| ())
}
/// Releases one release-owned CAS reference exactly once.
pub async fn release_reference_once(
&self,
ownership_id: &str,
ordinal: u64,
id: &ObjectId,
) -> bat_core::Result<bool> {
let hash = Self::parse_object_id(id)?;
self.engine()
.await?
.release_reference_once(ownership_id, ordinal, &hash)
.await
.map_err(Self::map_error)
}
async fn engine(&self) -> bat_core::Result<&engine_repository::FileSystemCasRepository> {
self.inner
.get_or_try_init(|| async {
engine_repository::FileSystemCasRepository::new(&self.root)
.await
.map_err(Self::map_error)
})
.await
}
fn parse_object_id(id: &ObjectId) -> bat_core::Result<Hash> {
id.parse::<Hash>()
.map_err(|error| bat_core::Error::InvalidArgument(error.to_string()))
}
fn map_error(error: bat_cas_engine::CasError) -> bat_core::Error {
match error {
bat_cas_engine::CasError::Io(error) => bat_core::Error::Io(error),
bat_cas_engine::CasError::ObjectNotFound(id) => bat_core::Error::NotFound(id),
bat_cas_engine::CasError::HashMismatch { expected, actual } => {
bat_core::Error::InvalidArgument(format!(
"Hash mismatch: expected {}, got {}",
expected, actual
))
}
bat_cas_engine::CasError::InvalidHash(hash) => {
bat_core::Error::InvalidArgument(format!("Invalid hash: {}", hash))
}
bat_cas_engine::CasError::ReferenceUnderflow(hash) => bat_core::Error::InvalidArgument(
format!("Reference count is already zero: {}", hash),
),
bat_cas_engine::CasError::Database(message) => {
bat_core::Error::Other(anyhow::anyhow!("CAS metadata error: {}", message))
}
bat_cas_engine::CasError::Other(error) => bat_core::Error::Other(error),
}
}
}
#[async_trait]
impl CasRepository for FileSystemCasRepository {
async fn store(&self, data: &[u8]) -> bat_core::Result<ObjectId> {
let hash = self
.engine()
.await?
.store(data)
.await
.map_err(Self::map_error)?;
Ok(hash.to_string())
}
async fn get(&self, id: &ObjectId) -> bat_core::Result<Vec<u8>> {
let hash = Self::parse_object_id(id)?;
let data = self
.engine()
.await?
.get(&hash)
.await
.map_err(Self::map_error)?;
Ok(data)
}
async fn exists(&self, id: &ObjectId) -> bat_core::Result<bool> {
let hash = Self::parse_object_id(id)?;
self.engine()
.await?
.exists(&hash)
.await
.map_err(Self::map_error)
}
async fn add_reference(&self, id: &ObjectId) -> bat_core::Result<u64> {
let hash = Self::parse_object_id(id)?;
self.engine()
.await?
.add_reference(&hash)
.await
.map_err(Self::map_error)
}
async fn remove_reference(&self, id: &ObjectId) -> bat_core::Result<u64> {
let hash = Self::parse_object_id(id)?;
self.engine()
.await?
.remove_reference(&hash)
.await
.map_err(Self::map_error)
}
async fn get_reference_count(&self, id: &ObjectId) -> bat_core::Result<u64> {
let hash = Self::parse_object_id(id)?;
self.engine()
.await?
.get_reference_count(&hash)
.await
.map_err(Self::map_error)
}
async fn gc(&self) -> bat_core::Result<u64> {
self.engine().await?.gc().await.map_err(Self::map_error)
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
#[tokio::test]
async fn store_get_and_reference_count() {
let temp_dir = TempDir::new().unwrap();
let repo = FileSystemCasRepository::new(temp_dir.path());
repo.init().await.unwrap();
let data = b"Hello, World!";
let id = repo.store(data).await.unwrap();
let retrieved = repo.get(&id).await.unwrap();
assert_eq!(data, retrieved.as_slice());
assert_eq!(repo.get_reference_count(&id).await.unwrap(), 1);
}
#[tokio::test]
async fn duplicate_store_increments_reference_count() {
let temp_dir = TempDir::new().unwrap();
let repo = FileSystemCasRepository::new(temp_dir.path());
repo.init().await.unwrap();
let id1 = repo.store(b"Duplicate data").await.unwrap();
let id2 = repo.store(b"Duplicate data").await.unwrap();
assert_eq!(id1, id2);
assert_eq!(repo.get_reference_count(&id1).await.unwrap(), 2);
}
#[tokio::test]
async fn gc_removes_zero_reference_objects() {
let temp_dir = TempDir::new().unwrap();
let repo = FileSystemCasRepository::new(temp_dir.path());
repo.init().await.unwrap();
let id = repo.store(b"collect me").await.unwrap();
assert_eq!(repo.remove_reference(&id).await.unwrap(), 0);
assert_eq!(repo.gc().await.unwrap(), 1);
assert!(!repo.exists(&id).await.unwrap());
let missing = repo.get(&id).await;
assert!(matches!(missing, Err(bat_core::Error::NotFound(_))));
}
#[tokio::test]
async fn invalid_object_id_is_rejected() {
let temp_dir = TempDir::new().unwrap();
let repo = FileSystemCasRepository::new(temp_dir.path());
repo.init().await.unwrap();
let error = repo.get(&"not-a-hash".to_string()).await.unwrap_err();
assert!(matches!(error, bat_core::Error::InvalidArgument(_)));
}
#[tokio::test]
async fn exists_distinguishes_missing_from_invalid_id() {
let temp_dir = TempDir::new().unwrap();
let repo = FileSystemCasRepository::new(temp_dir.path());
repo.init().await.unwrap();
// 合法但不存在的对象:Ok(false),而非静默错误。
let absent = repo.store(b"absent-seed").await.unwrap();
repo.remove_reference(&absent).await.unwrap();
repo.gc().await.unwrap();
assert!(!repo.exists(&absent).await.unwrap());
// 非法 ObjectId:返回 Err 而非 Ok(false),让调用方能区分故障。
let error = repo.exists(&"not-a-hash".to_string()).await.unwrap_err();
assert!(matches!(error, bat_core::Error::InvalidArgument(_)));
}
#[tokio::test]
async fn store_returns_blake3_hex_object_id() {
let temp_dir = TempDir::new().unwrap();
let repo = FileSystemCasRepository::new(temp_dir.path());
let hash = repo.store(b"test").await.unwrap();
assert_eq!(hash.len(), 64);
}
}