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..b281a3fe 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?; @@ -458,7 +482,13 @@ 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(), + }, + }) + .await + .unwrap(); let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -499,7 +529,13 @@ 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(), + }, + }) + .await + .unwrap(); let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -533,7 +569,13 @@ 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(), + }, + }) + .await + .unwrap(); let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -551,17 +593,22 @@ mod tests { }) } - 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(), - }); - (tempdir, backend) + keeper: KeeperConfig { + backend: KeeperBackend::Sqlite, + connection_url: "sqlite::memory:".into(), + }, + }) + .await?; + Ok((tempdir, backend)) } #[tokio::test] async fn multipart_single_part() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata { content_type: "text/plain".into(), @@ -613,7 +660,7 @@ mod tests { #[tokio::test] async fn multipart_multiple_parts() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -687,7 +734,7 @@ mod tests { #[tokio::test] async fn multipart_list_parts() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -750,7 +797,7 @@ mod tests { #[tokio::test] async fn get_object_range_bounded() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -777,7 +824,7 @@ mod tests { #[tokio::test] async fn get_object_range_from() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -804,7 +851,7 @@ mod tests { #[tokio::test] async fn get_object_range_last() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -831,7 +878,7 @@ mod tests { #[tokio::test] async fn get_object_range_unsatisfiable() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -849,7 +896,7 @@ mod tests { #[tokio::test] async fn multipart_abort() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -875,7 +922,7 @@ mod tests { #[tokio::test] async fn multipart_invalid_etag() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); @@ -924,7 +971,7 @@ mod tests { #[tokio::test] async fn multipart_missing_part() { - let (_tempdir, backend) = make_backend(); + let (_tempdir, backend) = make_backend().await.unwrap(); let id = make_id(); let metadata = Metadata::default(); 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 85661629..c946dc3d 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 supplied Metadata. @@ -32,3 +35,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, +}