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.
69 ///
70 /// What an *invalid* token does is backend-defined and must not be relied on. A backend
71 /// that detects one returns `StorageError::BackendUnavailable`; some S3 implementations
72 /// accept an unrecognised token and answer with a page instead. Since the token is
73 /// opaque and backend-issued, a client has no legitimate reason to send one that did
74 /// not come from this API.
75 pub next_cursor: Option<String>,
76}
77
78/// Return value from [`Storage::put_object`]. Carries the `ETag` of the stored object.
79#[derive(Debug, Clone, PartialEq, Eq)]
80pub struct PutOutcome {
81 /// `ETag` from the backend (opaque, quoted per RFC 7232 ยง2.3), or `None` if not returned.
82 pub etag: Option<String>,
83}
84
85/// The storage abstraction shared by all `NotedThat` components.
86///
87/// Implementations include `notedthat_storage_s3::S3Storage` (production)
88/// and `notedthat_api_http::testing::InMemoryStorage` (tests).
89///
90/// # Object safety
91///
92/// This trait is designed to be used as `Arc<dyn Storage>` from axum handlers.
93/// All methods take `&self` (not `&mut self`) to allow sharing across threads.
94#[async_trait]
95pub trait Storage: Send + Sync {
96 /// Idempotently create the S3 bucket for the given KB.
97 ///
98 /// Returns `Ok(())` if the bucket already exists (owned by this account).
99 async fn ensure_bucket(&self, kb: &KbSlug) -> Result<(), StorageError>;
100
101 /// Confirm the backend is reachable and `kb`'s bucket still exists.
102 ///
103 /// A read-only check for `/readyz`: it never creates or lists anything.
104 /// Returns `Err(StorageError::BucketNotFound)` when the bucket is gone and
105 /// `Err(StorageError::BackendUnavailable)` when the backend cannot answer.
106 /// Callers bound it with a timeout; the probe itself carries none.
107 async fn probe(&self, kb: &KbSlug) -> Result<(), StorageError>;
108
109 /// Read the KB manifest from `.notedthat/manifest.json` in the KB's bucket.
110 ///
111 /// Returns `Err(StorageError::NotFound)` if no manifest exists yet.
112 async fn read_manifest(&self, kb: &KbSlug) -> Result<KbManifest, StorageError>;
113
114 /// Write (overwrite) the KB manifest in the KB's bucket.
115 async fn write_manifest(&self, kb: &KbSlug, manifest: &KbManifest) -> Result<(), StorageError>;
116
117 /// Return metadata for an object without fetching its body.
118 ///
119 /// `conditionals` carries raw HTTP conditional headers for the backend to evaluate.
120 /// Returns `Err(StorageError::NotFound)` if the object does not exist.
121 async fn head_object(
122 &self,
123 kb: &KbSlug,
124 path: &ObjectPath,
125 conditionals: ConditionalHeaders,
126 ) -> Result<ObjectMeta, StorageError>;
127
128 /// Fetch an object's bytes and metadata.
129 ///
130 /// `range` carries the single requested byte range for the backend, and
131 /// `conditionals` carries raw HTTP conditional headers for the backend to
132 /// evaluate.
133 /// Returns `Err(StorageError::NotFound)` if the object does not exist.
134 async fn get_object(
135 &self,
136 kb: &KbSlug,
137 path: &ObjectPath,
138 range: Option<ByteRange>,
139 conditionals: ConditionalHeaders,
140 ) -> Result<ObjectRead, StorageError>;
141
142 /// Fetch metadata and an ordered body stream without whole-object buffering.
143 async fn get_object_stream(
144 &self,
145 kb: &KbSlug,
146 path: &ObjectPath,
147 range: Option<ByteRange>,
148 conditionals: ConditionalHeaders,
149 ) -> Result<ObjectStream, StorageError>;
150
151 /// Store an object, overwriting any existing object at the same path.
152 ///
153 /// The `content_type` is stored with the object and echoed on GET/HEAD.
154 /// `conditionals` carries raw HTTP conditional headers for the backend to evaluate.
155 async fn put_object(
156 &self,
157 kb: &KbSlug,
158 path: &ObjectPath,
159 bytes: Bytes,
160 content_type: Option<&str>,
161 conditionals: ConditionalHeaders,
162 ) -> Result<PutOutcome, StorageError>;
163
164 /// Store an owned staged body without loading file-backed bodies into memory.
165 async fn put_staged_object(
166 &self,
167 kb: &KbSlug,
168 path: &ObjectPath,
169 body: StagedBody,
170 content_type: Option<&str>,
171 conditionals: ConditionalHeaders,
172 ) -> Result<PutOutcome, StorageError>;
173
174 /// Copy an object within one KB using backend-native conditional copy.
175 async fn copy_object(
176 &self,
177 kb: &KbSlug,
178 source: &ObjectPath,
179 destination: &ObjectPath,
180 options: CopyObjectOptions,
181 ) -> Result<PutOutcome, StorageError>;
182
183 /// Delete an object.
184 ///
185 /// `conditionals` carries raw HTTP conditional headers for the backend to evaluate.
186 /// This operation is **idempotent** โ deleting a non-existent object returns
187 /// `Ok(())` (matching S3 semantics per Metis directive).
188 async fn delete_object(
189 &self,
190 kb: &KbSlug,
191 path: &ObjectPath,
192 conditionals: ConditionalHeaders,
193 ) -> Result<(), StorageError>;
194
195 /// List objects in the KB, optionally filtered by a prefix.
196 ///
197 /// Results are capped at `limit` (default 100, max 1000).
198 ///
199 /// Pass `cursor = None` on the first call. On subsequent calls, pass the opaque
200 /// `next_cursor` value from the previous [`ListResponse`] unchanged. What an invalid
201 /// cursor does is backend-defined โ see [`ListResponse::next_cursor`].
202 ///
203 /// # Invariant
204 ///
205 /// Every returned `ListResponse` satisfies `truncated == next_cursor.is_some()`.
206 async fn list_objects(
207 &self,
208 kb: &KbSlug,
209 prefix: Option<&str>,
210 limit: u32,
211 cursor: Option<&str>,
212 ) -> Result<ListResponse, StorageError>;
213}
214
215#[cfg(test)]
216mod tests {
217 fn assert_send_sync<T: Send + Sync + ?Sized>() {}
218
219 #[test]
220 fn storage_is_dyn_compatible_and_send_sync() {
221 assert_send_sync::<dyn crate::storage::Storage>();
222 }
223}