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}