Skip to main content

notedthat_core/
storage.rs

1//! `Storage` trait — object store abstraction.
2
3use crate::conditional::ConditionalHeaders;
4use crate::error::StorageError;
5use crate::kb::{KbManifest, ObjectMeta};
6use crate::object_path::ObjectPath;
7use crate::range::ByteRange;
8use crate::slug::KbSlug;
9use async_trait::async_trait;
10use bytes::Bytes;
11use futures::Stream;
12use std::pin::Pin;
13
14use crate::staging::StagedBody;
15
16/// Ordered chunks returned by a streaming object read.
17pub type ObjectChunkStream = Pin<Box<dyn Stream<Item = Result<Bytes, StorageError>> + Send>>;
18
19/// Metadata and ordered chunks returned without materializing the object.
20pub struct ObjectStream {
21    /// Ordered object body chunks.
22    pub chunks: ObjectChunkStream,
23    /// Associated metadata.
24    pub meta: ObjectMeta,
25    /// Backend `Content-Range` value for range reads.
26    pub content_range: Option<String>,
27}
28
29/// Preconditions and metadata for a server-side object copy.
30#[derive(Clone, Debug, Default, PartialEq, Eq)]
31pub struct CopyObjectOptions {
32    /// Required source `ETag` match when present.
33    pub source_if_match: Option<String>,
34    /// Required destination non-match condition when present.
35    pub destination_if_none_match: Option<String>,
36    /// Content type stored on the destination.
37    pub content_type: Option<String>,
38}
39
40/// The bytes and metadata returned by a GET or HEAD operation.
41pub struct ObjectRead {
42    /// The raw object bytes.
43    pub bytes: Bytes,
44    /// Associated metadata (size, content-type, last-modified).
45    pub meta: ObjectMeta,
46    /// `Content-Range: bytes start-end/total` header value from the backend
47    /// when responding to a range request, or `None` for full-body reads.
48    /// Passed through to HTTP clients on 206 responses.
49    pub content_range: Option<String>,
50}
51
52/// The result of a LIST operation.
53///
54/// # Invariant
55///
56/// `truncated == next_cursor.is_some()`. A response where `truncated=true` MUST supply a
57/// `next_cursor`; a response where `next_cursor=Some(_)` MUST also set `truncated=true`.
58#[derive(Debug)]
59pub struct ListResponse {
60    /// The matching objects, up to the requested `limit`.
61    pub objects: Vec<ObjectMeta>,
62    /// `true` if the backend indicated more objects exist beyond `limit`.
63    pub truncated: bool,
64    /// Opaque backend continuation token. Present exactly when `truncated=true`.
65    ///
66    /// Pass this value unchanged as `cursor` on the next `list_objects` call to retrieve the
67    /// next page. Clients MUST NOT parse, validate, or store this value beyond the immediate
68    /// next request. Invalid or expired tokens cause the backend to return
69    /// `StorageError::BackendUnavailable`.
70    pub next_cursor: Option<String>,
71}
72
73/// Return value from [`Storage::put_object`]. Carries the `ETag` of the stored object.
74#[derive(Debug, Clone, PartialEq, Eq)]
75pub struct PutOutcome {
76    /// `ETag` from the backend (opaque, quoted per RFC 7232 §2.3), or `None` if not returned.
77    pub etag: Option<String>,
78}
79
80/// The storage abstraction shared by all `NotedThat` components.
81///
82/// Implementations include `notedthat_storage_s3::S3Storage` (production)
83/// and `notedthat_api_http::testing::InMemoryStorage` (tests).
84///
85/// # Object safety
86///
87/// This trait is designed to be used as `Arc<dyn Storage>` from axum handlers.
88/// All methods take `&self` (not `&mut self`) to allow sharing across threads.
89#[async_trait]
90pub trait Storage: Send + Sync {
91    /// Idempotently create the S3 bucket for the given KB.
92    ///
93    /// Returns `Ok(())` if the bucket already exists (owned by this account).
94    async fn ensure_bucket(&self, kb: &KbSlug) -> Result<(), StorageError>;
95
96    /// Read the KB manifest from `.notedthat/manifest.json` in the KB's bucket.
97    ///
98    /// Returns `Err(StorageError::NotFound)` if no manifest exists yet.
99    async fn read_manifest(&self, kb: &KbSlug) -> Result<KbManifest, StorageError>;
100
101    /// Write (overwrite) the KB manifest in the KB's bucket.
102    async fn write_manifest(&self, kb: &KbSlug, manifest: &KbManifest) -> Result<(), StorageError>;
103
104    /// Return metadata for an object without fetching its body.
105    ///
106    /// `conditionals` carries raw HTTP conditional headers for the backend to evaluate.
107    /// Returns `Err(StorageError::NotFound)` if the object does not exist.
108    async fn head_object(
109        &self,
110        kb: &KbSlug,
111        path: &ObjectPath,
112        conditionals: ConditionalHeaders,
113    ) -> Result<ObjectMeta, StorageError>;
114
115    /// Fetch an object's bytes and metadata.
116    ///
117    /// `range` carries parsed byte ranges for the backend, and `conditionals`
118    /// carries raw HTTP conditional headers for the backend to evaluate.
119    /// Returns `Err(StorageError::NotFound)` if the object does not exist.
120    async fn get_object(
121        &self,
122        kb: &KbSlug,
123        path: &ObjectPath,
124        range: Option<Vec<ByteRange>>,
125        conditionals: ConditionalHeaders,
126    ) -> Result<ObjectRead, StorageError>;
127
128    /// Fetch metadata and an ordered body stream without whole-object buffering.
129    async fn get_object_stream(
130        &self,
131        kb: &KbSlug,
132        path: &ObjectPath,
133        range: Option<Vec<ByteRange>>,
134        conditionals: ConditionalHeaders,
135    ) -> Result<ObjectStream, StorageError>;
136
137    /// Store an object, overwriting any existing object at the same path.
138    ///
139    /// The `content_type` is stored with the object and echoed on GET/HEAD.
140    /// `conditionals` carries raw HTTP conditional headers for the backend to evaluate.
141    async fn put_object(
142        &self,
143        kb: &KbSlug,
144        path: &ObjectPath,
145        bytes: Bytes,
146        content_type: Option<&str>,
147        conditionals: ConditionalHeaders,
148    ) -> Result<PutOutcome, StorageError>;
149
150    /// Store an owned staged body without loading file-backed bodies into memory.
151    async fn put_staged_object(
152        &self,
153        kb: &KbSlug,
154        path: &ObjectPath,
155        body: StagedBody,
156        content_type: Option<&str>,
157        conditionals: ConditionalHeaders,
158    ) -> Result<PutOutcome, StorageError>;
159
160    /// Copy an object within one KB using backend-native conditional copy.
161    async fn copy_object(
162        &self,
163        kb: &KbSlug,
164        source: &ObjectPath,
165        destination: &ObjectPath,
166        options: CopyObjectOptions,
167    ) -> Result<PutOutcome, StorageError>;
168
169    /// Delete an object.
170    ///
171    /// `conditionals` carries raw HTTP conditional headers for the backend to evaluate.
172    /// This operation is **idempotent** — deleting a non-existent object returns
173    /// `Ok(())` (matching S3 semantics per Metis directive).
174    async fn delete_object(
175        &self,
176        kb: &KbSlug,
177        path: &ObjectPath,
178        conditionals: ConditionalHeaders,
179    ) -> Result<(), StorageError>;
180
181    /// List objects in the KB, optionally filtered by a prefix.
182    ///
183    /// Results are capped at `limit` (default 100, max 1000).
184    ///
185    /// Pass `cursor = None` on the first call. On subsequent calls, pass the opaque
186    /// `next_cursor` value from the previous [`ListResponse`] unchanged. An invalid or
187    /// expired cursor causes the backend to return `StorageError::BackendUnavailable`.
188    ///
189    /// # Invariant
190    ///
191    /// Every returned `ListResponse` satisfies `truncated == next_cursor.is_some()`.
192    async fn list_objects(
193        &self,
194        kb: &KbSlug,
195        prefix: Option<&str>,
196        limit: u32,
197        cursor: Option<&str>,
198    ) -> Result<ListResponse, StorageError>;
199}
200
201#[cfg(test)]
202mod tests {
203    fn assert_send_sync<T: Send + Sync + ?Sized>() {}
204
205    #[test]
206    fn storage_is_dyn_compatible_and_send_sync() {
207        assert_send_sync::<dyn crate::storage::Storage>();
208    }
209}