From 499ea212e470286a835b008c120008efddfee66b Mon Sep 17 00:00:00 2001 From: Reinaldy Rafli Date: Sat, 1 Aug 2026 20:06:35 +0700 Subject: [PATCH 1/3] feat(keeper): integrate filesystem backend with keeper trait --- objectstore-server/src/config.rs | 5 +++ objectstore-service/src/backend/local_fs.rs | 50 +++++++++++++++++++-- objectstore-service/src/backend/mod.rs | 6 ++- objectstore-service/src/keeper/mod.rs | 34 +++++++++++++- 4 files changed, 87 insertions(+), 8 deletions(-) diff --git a/objectstore-server/src/config.rs b/objectstore-server/src/config.rs index 18b7b62e..0d787aea 100644 --- a/objectstore-server/src/config.rs +++ b/objectstore-server/src/config.rs @@ -41,6 +41,7 @@ use std::time::Duration; use anyhow::Result; use figment::providers::{Env, Format, Serialized, Yaml}; use objectstore_service::backend::local_fs::FileSystemConfig; +use objectstore_service::keeper::{KeeperBackend, KeeperConfig}; use objectstore_types::auth::Permission; use secrecy::{CloneableSecret, SecretBox, SerializableSecret, zeroize::Zeroize}; use serde::{Deserialize, Serialize}; @@ -586,6 +587,10 @@ impl Default for Config { storage: StorageConfig::FileSystem(FileSystemConfig { path: PathBuf::from("data"), + keeper: KeeperConfig { + backend: KeeperBackend::Sqlite, + connection_url: "sqlite:///opt/objectstore/keeper.db".into(), + }, }), runtime: Runtime::default(), diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index af0b233d..ef368176 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -17,6 +17,8 @@ use crate::backend::common::{ }; use crate::error::{Error, Result}; use crate::id::ObjectId; +use crate::keeper::sqlite_backed::SqliteBackedKeeper; +use crate::keeper::{Keeper, KeeperBackend, KeeperConfig}; use crate::multipart::{ AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse, ListPartsResponse, Part, PartNumber, UploadId, UploadPartResponse, @@ -34,6 +36,9 @@ use crate::stream::{self, ClientStream}; /// storage: /// type: filesystem /// path: /data +/// keeper: +/// backend: sqlite +/// connection_url: sqlite://data/keeper.db /// ``` #[derive(Debug, Clone, serde::Deserialize, serde::Serialize)] pub struct FileSystemConfig { @@ -51,18 +56,31 @@ pub struct FileSystemConfig { /// - `OS__STORAGE__TYPE=filesystem` /// - `OS__STORAGE__PATH=/path/to/storage` pub path: PathBuf, + + /// Configuration for the keeper backend. + /// See [`KeeperConfig`] for details. + pub keeper: KeeperConfig, } /// Local filesystem backend for development and testing. #[derive(Debug)] pub struct LocalFsBackend { path: PathBuf, + keeper: Box, } impl LocalFsBackend { /// Creates a new [`LocalFsBackend`] rooted at the directory in `config`. - pub fn new(config: FileSystemConfig) -> Self { - Self { path: config.path } + pub async fn new(config: FileSystemConfig) -> Result { + let keeper = match config.keeper.backend { + KeeperBackend::Sqlite => { + Box::new(SqliteBackedKeeper::new(&config.keeper.connection_url).await?) + } + }; + Ok(Self { + path: config.path, + keeper, + }) } } @@ -115,6 +133,8 @@ impl Backend for LocalFsBackend { file.sync_data().await?; drop(file); + self.keeper.keep(id, metadata.expiration_policy).await?; + Ok(()) } @@ -131,6 +151,7 @@ impl Backend for LocalFsBackend { } err => err?, }; + self.keeper.mark_accessed(id).await?; let mut reader = BufReader::new(file); let mut metadata_line = String::new(); @@ -174,6 +195,7 @@ impl Backend for LocalFsBackend { { objectstore_log::debug!("Object not found"); } + self.keeper.remove(id).await?; Ok(result?) } } @@ -432,6 +454,8 @@ impl MultipartUploadBackend for LocalFsBackend { file.sync_data().await?; drop(file); + self.keeper.keep(id, metadata.expiration_policy).await?; + // Clean up multipart state tokio::fs::remove_dir_all(dir).await?; @@ -451,6 +475,7 @@ mod tests { use super::*; use crate::id::ObjectContext; + use crate::keeper::sqlite_backed::SqliteBackedKeeper; use crate::stream; #[tokio::test] @@ -458,6 +483,10 @@ mod tests { let tempdir = tempfile::tempdir().unwrap(); let backend = LocalFsBackend::new(FileSystemConfig { path: tempdir.path().to_path_buf(), + keeper: KeeperConfig { + backend: KeeperBackend::Sqlite, + connection_url: "sqlite::memory:".into(), + }, }); let id = ObjectId::random(ObjectContext { @@ -499,6 +528,10 @@ mod tests { let tempdir = tempfile::tempdir().unwrap(); let backend = LocalFsBackend::new(FileSystemConfig { path: tempdir.path().to_path_buf(), + keeper: KeeperConfig { + backend: KeeperBackend::Sqlite, + connection_url: "sqlite::memory:".into(), + }, }); let id = ObjectId::random(ObjectContext { @@ -533,6 +566,10 @@ mod tests { let tempdir = tempfile::tempdir().unwrap(); let backend = LocalFsBackend::new(FileSystemConfig { path: tempdir.path().to_path_buf(), + keeper: KeeperConfig { + backend: KeeperBackend::Sqlite, + connection_url: "sqlite::memory:".into(), + }, }); let id = ObjectId::random(ObjectContext { @@ -551,11 +588,16 @@ mod tests { }) } - fn make_backend() -> (tempfile::TempDir, LocalFsBackend) { + async fn make_backend() -> (tempfile::TempDir, LocalFsBackend) { let tempdir = tempfile::tempdir().unwrap(); let backend = LocalFsBackend::new(FileSystemConfig { path: tempdir.path().to_path_buf(), - }); + keeper: KeeperConfig { + backend: KeeperBackend::Sqlite, + connection_url: "sqlite::memory:".into(), + }, + }) + .await?; (tempdir, backend) } diff --git a/objectstore-service/src/backend/mod.rs b/objectstore-service/src/backend/mod.rs index 4d265139..52d8d710 100644 --- a/objectstore-service/src/backend/mod.rs +++ b/objectstore-service/src/backend/mod.rs @@ -75,7 +75,7 @@ pub async fn from_config(config: StorageConfig) -> Result Result> { Ok(match config { - StorageConfig::FileSystem(c) => Box::new(local_fs::LocalFsBackend::new(c)), + StorageConfig::FileSystem(c) => Box::new(local_fs::LocalFsBackend::new(c).await?), StorageConfig::S3Compatible(c) => { Box::new(s3_compatible::S3CompatibleBackend::without_token(c)) } @@ -127,7 +127,9 @@ async fn lt_from_config( config: MultipartUploadStorageConfig, ) -> Result> { Ok(match config { - MultipartUploadStorageConfig::FileSystem(c) => Box::new(local_fs::LocalFsBackend::new(c)), + MultipartUploadStorageConfig::FileSystem(c) => { + Box::new(local_fs::LocalFsBackend::new(c).await?) + } MultipartUploadStorageConfig::Gcs(c) => Box::new(gcs::GcsBackend::new(c).await?), }) } diff --git a/objectstore-service/src/keeper/mod.rs b/objectstore-service/src/keeper/mod.rs index 07338f67..982d5623 100644 --- a/objectstore-service/src/keeper/mod.rs +++ b/objectstore-service/src/keeper/mod.rs @@ -1,9 +1,12 @@ //! Keeper defines a trait to support extending an object retention for backends //! that does not support them natively (e.g., S3-compatible backend and filesystem). +use std::str::FromStr; + use objectstore_types::metadata::ExpirationPolicy; +use serde::{Deserialize, Serialize}; -use crate::error::Result; +use crate::error::{Error, Result}; use crate::id::ObjectId; /// SQLite-backed implementation of the [`Keeper`] trait. @@ -18,7 +21,7 @@ pub struct ObjectExpiry<'a> { /// Object retention keeper trait. #[async_trait::async_trait] -pub trait Keeper: Send + Sync { +pub trait Keeper: Send + Sync + std::fmt::Debug { /// Keep is the first step in the object retention lifecycle. /// Practically speaking, it would not be kept if the `expiration_policy` is `Manual`. /// The `expiration_policy` is set by the client at upload time via the @@ -33,3 +36,30 @@ pub trait Keeper: Send + Sync { /// extend the object retention. async fn mark_accessed(&self, id: &ObjectId) -> Result<()>; } + +/// Keeper backend of choice. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub enum KeeperBackend { + /// SQLite-backed keeper. + Sqlite, +} + +impl FromStr for KeeperBackend { + type Err = Error; + fn from_str(s: &str) -> Result { + match s { + "sqlite" => Ok(KeeperBackend::Sqlite), + _ => Err(Error::generic(format!("unknown backend {}", s))), + } + } +} + +/// Configuration for the keeper backend. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct KeeperConfig { + /// Specifies the backend to use. + pub backend: KeeperBackend, + /// Connection URL for the backend. + /// Refer to each backend's documentation for details. + pub connection_url: String, +} From 9db675f09fd352ab73af2fe77a2e74041237add8 Mon Sep 17 00:00:00 2001 From: Reinaldy Rafli Date: Sat, 1 Aug 2026 20:31:27 +0700 Subject: [PATCH 2/3] fix(keeper): tests lints --- objectstore-service/src/backend/local_fs.rs | 34 +++++++++++---------- 1 file changed, 18 insertions(+), 16 deletions(-) diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index ef368176..681c8d3b 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -475,7 +475,6 @@ mod tests { use super::*; use crate::id::ObjectContext; - use crate::keeper::sqlite_backed::SqliteBackedKeeper; use crate::stream; #[tokio::test] @@ -487,7 +486,8 @@ mod tests { backend: KeeperBackend::Sqlite, connection_url: "sqlite::memory:".into(), }, - }); + }) + .await?; let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -532,7 +532,8 @@ mod tests { backend: KeeperBackend::Sqlite, connection_url: "sqlite::memory:".into(), }, - }); + }) + .await?; let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -570,7 +571,8 @@ mod tests { backend: KeeperBackend::Sqlite, connection_url: "sqlite::memory:".into(), }, - }); + }) + .await?; let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -588,7 +590,7 @@ mod tests { }) } - async fn make_backend() -> (tempfile::TempDir, LocalFsBackend) { + async fn make_backend() -> Result<(tempfile::TempDir, LocalFsBackend)> { let tempdir = tempfile::tempdir().unwrap(); let backend = LocalFsBackend::new(FileSystemConfig { path: tempdir.path().to_path_buf(), @@ -598,12 +600,12 @@ mod tests { }, }) .await?; - (tempdir, backend) + Ok((tempdir, backend)) } #[tokio::test] async fn multipart_single_part() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await?; let id = make_id(); let metadata = Metadata { content_type: "text/plain".into(), @@ -655,7 +657,7 @@ mod tests { #[tokio::test] async fn multipart_multiple_parts() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await?; let id = make_id(); let metadata = Metadata::default(); @@ -729,7 +731,7 @@ mod tests { #[tokio::test] async fn multipart_list_parts() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await?; let id = make_id(); let metadata = Metadata::default(); @@ -792,7 +794,7 @@ mod tests { #[tokio::test] async fn get_object_range_bounded() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await?; let id = make_id(); let metadata = Metadata::default(); @@ -819,7 +821,7 @@ mod tests { #[tokio::test] async fn get_object_range_from() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await?; let id = make_id(); let metadata = Metadata::default(); @@ -846,7 +848,7 @@ mod tests { #[tokio::test] async fn get_object_range_last() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await?; let id = make_id(); let metadata = Metadata::default(); @@ -873,7 +875,7 @@ mod tests { #[tokio::test] async fn get_object_range_unsatisfiable() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await?; let id = make_id(); let metadata = Metadata::default(); @@ -891,7 +893,7 @@ mod tests { #[tokio::test] async fn multipart_abort() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await?; let id = make_id(); let metadata = Metadata::default(); @@ -917,7 +919,7 @@ mod tests { #[tokio::test] async fn multipart_invalid_etag() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await?; let id = make_id(); let metadata = Metadata::default(); @@ -966,7 +968,7 @@ mod tests { #[tokio::test] async fn multipart_missing_part() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await?; let id = make_id(); let metadata = Metadata::default(); From bad756589293399530a58c6fcde58cd45f49a291 Mon Sep 17 00:00:00 2001 From: Reinaldy Rafli Date: Sun, 2 Aug 2026 05:57:26 +0700 Subject: [PATCH 3/3] test(localfs): use unwrap --- objectstore-service/src/backend/local_fs.rs | 29 ++++++++++++--------- 1 file changed, 16 insertions(+), 13 deletions(-) diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index 681c8d3b..b281a3fe 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -487,7 +487,8 @@ mod tests { connection_url: "sqlite::memory:".into(), }, }) - .await?; + .await + .unwrap(); let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -533,7 +534,8 @@ mod tests { connection_url: "sqlite::memory:".into(), }, }) - .await?; + .await + .unwrap(); let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -572,7 +574,8 @@ mod tests { connection_url: "sqlite::memory:".into(), }, }) - .await?; + .await + .unwrap(); let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -605,7 +608,7 @@ mod tests { #[tokio::test] async fn multipart_single_part() { - let (_tempdir, backend) = make_backend().await?; + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata { content_type: "text/plain".into(), @@ -657,7 +660,7 @@ mod tests { #[tokio::test] async fn multipart_multiple_parts() { - let (_tempdir, backend) = make_backend().await?; + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -731,7 +734,7 @@ mod tests { #[tokio::test] async fn multipart_list_parts() { - let (_tempdir, backend) = make_backend().await?; + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -794,7 +797,7 @@ mod tests { #[tokio::test] async fn get_object_range_bounded() { - let (_tempdir, backend) = make_backend().await?; + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -821,7 +824,7 @@ mod tests { #[tokio::test] async fn get_object_range_from() { - let (_tempdir, backend) = make_backend().await?; + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -848,7 +851,7 @@ mod tests { #[tokio::test] async fn get_object_range_last() { - let (_tempdir, backend) = make_backend().await?; + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -875,7 +878,7 @@ mod tests { #[tokio::test] async fn get_object_range_unsatisfiable() { - let (_tempdir, backend) = make_backend().await?; + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -893,7 +896,7 @@ mod tests { #[tokio::test] async fn multipart_abort() { - let (_tempdir, backend) = make_backend().await?; + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -919,7 +922,7 @@ mod tests { #[tokio::test] async fn multipart_invalid_etag() { - let (_tempdir, backend) = make_backend().await?; + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -968,7 +971,7 @@ mod tests { #[tokio::test] async fn multipart_missing_part() { - let (_tempdir, backend) = make_backend().await?; + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default();