ocre/runtime/storage.rs
1use axum::{body::Body, http::HeaderMap, response::Response};
2use futures_util::StreamExt as _;
3use worker::{
4 Bucket, Conditional, Env, FixedLengthStream, HttpMetadata, Include, Object, Range, ResponseBody,
5 send::{SendFuture, SendWrapper},
6 web_sys,
7};
8
9use super::Ctx;
10use crate::{
11 Error, IntoParam as _, Result,
12 storage::{
13 Attachment, ByteRange, DirectUpload, DirectUploadRequest, Disposition, Fetch, KEYS_PER_QUERY, Listing,
14 MAX_LIST, Purged, Rules, S3Endpoint, STORAGE_BINDING, STORAGE_PUBLIC_URL, StoredObject, Upload,
15 attachment_from_head, essence, file_response, join_public_url, missing_binding, not_modified, presign_get_url,
16 r2_endpoint, redirect_response, referenced_keys_sql, stale_keys, unsatisfiable, upload_secret, verify_key,
17 },
18};
19
20pub(super) fn bucket(env: &Env) -> Result<Bucket> {
21 env.bucket(STORAGE_BINDING).map_err(|err| missing_binding(&err.to_string()))
22}
23
24async fn put(env: &Env, attachment: &Attachment, data: worker::Data) -> Result<()> {
25 let metadata = HttpMetadata { content_type: Some(attachment.content_type.clone()), ..Default::default() };
26 bucket(env)?
27 .put(&attachment.key, data)
28 .http_metadata(metadata)
29 .custom_metadata([("filename".to_owned(), attachment.filename.clone())])
30 .execute()
31 .await
32 .map_err(|err| Error::internal(format!("R2 put of `{}` failed: {err}", attachment.key)))?;
33 Ok(())
34}
35
36/// Stores an [`Upload`] in R2 under a new random key starting with `prefix`, and returns the [`Attachment`] to save.
37///
38/// The key is `<prefix>/<22 random characters>` (`posts/avatar/...`),
39/// never derived from the file name; the file name and content type are
40/// cleaned up like [`store_bytes`](crate::storage::store_bytes)'s. Check
41/// the upload with [`Validator::file`](crate::Validator::file) first, and
42/// save the attachment with [`columns`](crate::storage::columns); if saving
43/// the row fails, [`delete`](crate::storage::delete) the new key.
44///
45/// Free plan: one R2 class A operation (1M free per month); copying 10 MB
46/// to R2 takes about 0.15 ms of CPU. Stored files count towards the 10
47/// GB-month of free storage.
48///
49/// # Errors
50///
51/// [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
52/// binding is missing (the message shows the `STORAGE` entry to add
53/// to cloudflare.config.ts) or the R2 write fails.
54///
55/// # Examples
56///
57/// ```no_run
58/// use axum::extract::State;
59/// use ocre::storage::{self, Multipart, Rules};
60/// use ocre::{Ctx, Error, Result, Validator};
61///
62/// const DOCUMENT: Rules = Rules { max_bytes: 10 * 1024 * 1024, content_types: &["application/pdf"] };
63///
64/// // PUT /documents with `curl -X PUT -F file=@spec.pdf`.
65/// async fn upload(State(ctx): State<Ctx>, Multipart(mut form): Multipart<{ 11 * 1024 * 1024 }>) -> Result<String> {
66/// let file = form.file("file").ok_or_else(|| Error::bad_request("Send the file as `file`"))?;
67/// Validator::new().file("file", &file, &DOCUMENT).finish()?;
68/// let document = storage::store(&ctx, "documents/file", file).await?;
69/// Ok(document.key)
70/// }
71/// ```
72pub fn store(ctx: &Ctx, prefix: &str, upload: Upload) -> impl Future<Output = Result<Attachment>> + Send + use<> {
73 store_bytes(ctx, prefix, &upload.filename, &upload.content_type, Vec::from(upload.bytes))
74}
75
76/// Stores bytes built by the app (a generated report, an export) like [`store`](crate::storage::store).
77///
78/// The object gets a new random key under `prefix`; `filename` is cleaned
79/// up (no directories or control characters, at most 200 characters) and
80/// `content_type` normalized (lowercase, no parameters; empty becomes
81/// `application/octet-stream`). The returned [`Attachment`] records all
82/// four, ready to save with [`columns`](crate::storage::columns).
83///
84/// Free plan: one R2 class A operation (1M free per month); copying 10 MB
85/// to R2 takes about 0.15 ms of CPU.
86///
87/// # Errors
88///
89/// [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
90/// binding is missing (the message shows the `STORAGE` entry to add
91/// to cloudflare.config.ts) or the R2 write fails.
92///
93/// # Examples
94///
95/// ```no_run
96/// use axum::extract::State;
97/// use ocre::{Ctx, Result, storage};
98/// use serde::Deserialize;
99///
100/// #[derive(Deserialize)]
101/// struct Post {
102/// id: i64,
103/// title: String,
104/// }
105///
106/// // POST /exports: writes every post to a CSV kept in R2.
107/// async fn export(State(ctx): State<Ctx>) -> Result<String> {
108/// let posts: Vec<Post> = ctx.db()?.all("SELECT id, title FROM posts ORDER BY id", vec![]).await?;
109/// let mut csv = String::from("id,title\n");
110/// for post in posts {
111/// csv.push_str(&format!("{},\"{}\"\n", post.id, post.title.replace('"', "\"\"")));
112/// }
113/// let file = storage::store_bytes(&ctx, "exports", "posts.csv", "text/csv", csv.into_bytes()).await?;
114/// Ok(file.key)
115/// }
116/// ```
117pub fn store_bytes(
118 ctx: &Ctx,
119 prefix: &str,
120 filename: &str,
121 content_type: &str,
122 bytes: Vec<u8>,
123) -> impl Future<Output = Result<Attachment>> + Send + use<> {
124 let env = ctx.env().clone();
125 let attachment = Attachment::prepare(prefix, filename, content_type, bytes.len() as u64);
126 SendFuture::new(async move {
127 put(&env, &attachment, bytes.into()).await?;
128 Ok(attachment)
129 })
130}
131
132/// Streams `body` (exactly `size` bytes) into R2 without holding it in memory, and returns its [`Attachment`].
133///
134/// For a raw request body with a `Content-Length` (`curl -T big.zip`),
135/// where [`Multipart`](crate::storage::Multipart) would hold the whole body
136/// in memory. R2 needs the length up front: a body of another length makes
137/// the write fail and nothing is stored. The file name and content type are
138/// cleaned up like [`store_bytes`]'s; no [`Rules`](crate::storage::Rules)
139/// are checked, so check `size` and `content_type` first.
140///
141/// Free plan: one R2 class A operation (1M free per month). Chunks still
142/// pass through WebAssembly, but memory stays flat; Cloudflare refuses
143/// request bodies over 100 MB on the Free plan.
144///
145/// # Errors
146///
147/// [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
148/// binding is missing (the message shows the `STORAGE` entry to add
149/// to cloudflare.config.ts), the body is not `size` bytes long, or the R2 write fails.
150///
151/// # Examples
152///
153/// ```no_run
154/// use axum::{body::Body, extract::State, http::{HeaderMap, header}};
155/// use ocre::{Ctx, Error, Result, storage};
156///
157/// const MAX: u64 = 50 * 1024 * 1024;
158///
159/// // PUT /backups with `curl -T backup.zip`.
160/// async fn upload(State(ctx): State<Ctx>, headers: HeaderMap, body: Body) -> Result<String> {
161/// let size: u64 = headers
162/// .get(header::CONTENT_LENGTH)
163/// .and_then(|value| value.to_str().ok()?.parse().ok())
164/// .ok_or_else(|| Error::bad_request("Content-Length required"))?;
165/// if size > MAX {
166/// return Err(Error::PayloadTooLarge("The backup is too large (maximum is 50 MB)".into()));
167/// }
168/// let backup = storage::store_body(&ctx, "backups", "backup.zip", "application/zip", size, body).await?;
169/// Ok(backup.key)
170/// }
171/// ```
172pub fn store_body(
173 ctx: &Ctx,
174 prefix: &str,
175 filename: &str,
176 content_type: &str,
177 size: u64,
178 body: Body,
179) -> impl Future<Output = Result<Attachment>> + Send + use<> {
180 let env = ctx.env().clone();
181 let attachment = Attachment::prepare(prefix, filename, content_type, size);
182 SendFuture::new(async move {
183 let chunks = body
184 .into_data_stream()
185 .map(|chunk| chunk.map(|bytes| bytes.to_vec()).map_err(|err| worker::Error::RustError(err.to_string())));
186 put(&env, &attachment, FixedLengthStream::wrap(chunks, size).into()).await?;
187 Ok(attachment)
188 })
189}
190
191/// Reads a whole object into memory, or returns `None` when `key` does not exist.
192///
193/// Meant for files the Worker itself processes (parsing an uploaded CSV,
194/// attaching a file to an email). To send a file to a browser use
195/// [`serve`](crate::storage::serve), which streams it without copying it
196/// into WebAssembly and handles 304 and `Range`.
197///
198/// Free plan: one R2 class B operation (10M free per month); the bytes are
199/// copied into WebAssembly memory (a Worker has 128 MB).
200///
201/// # Errors
202///
203/// [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
204/// binding is missing (the message shows the `STORAGE` entry to add
205/// to cloudflare.config.ts) or R2 fails.
206///
207/// # Examples
208///
209/// ```no_run
210/// use axum::extract::{Path, State};
211/// use ocre::{Ctx, OptionExt, Result, storage};
212///
213/// // Counts the lines of an uploaded CSV.
214/// async fn line_count(State(ctx): State<Ctx>, Path(key): Path<String>) -> Result<String> {
215/// let bytes = storage::read(&ctx, &key).await?.or_404()?;
216/// Ok(bytes.split(|&b| b == b'\n').filter(|line| !line.is_empty()).count().to_string())
217/// }
218/// ```
219pub fn read(ctx: &Ctx, key: &str) -> impl Future<Output = Result<Option<Vec<u8>>>> + Send + use<> {
220 let env = ctx.env().clone();
221 let key = key.to_owned();
222 SendFuture::new(async move {
223 let Some(object) = bucket(&env)?.get(&key).execute().await? else { return Ok(None) };
224 match object.body() {
225 Some(body) => Ok(Some(body.bytes().await?)),
226 None => Ok(Some(Vec::new())),
227 }
228 })
229}
230
231/// Deletes the object at `key`.
232///
233/// Deleting a missing key is not an error, so a retried cleanup is safe.
234/// Delete the object only once no row points to it any more (after the
235/// database write succeeded). R2 deletes are free.
236///
237/// # Errors
238///
239/// [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
240/// binding is missing (the message shows the `STORAGE` entry to add
241/// to cloudflare.config.ts) or R2 fails.
242///
243/// # Examples
244///
245/// ```no_run
246/// use axum::{extract::{Path, State}, http::StatusCode};
247/// use ocre::{Ctx, Result, params, storage};
248///
249/// // DELETE /exports/{key}: forget a generated export.
250/// async fn destroy(State(ctx): State<Ctx>, Path(key): Path<String>) -> Result<StatusCode> {
251/// ctx.db()?.execute("DELETE FROM exports WHERE file_key = ?1", params![key.as_str()]).await?;
252/// storage::delete(&ctx, &key).await?;
253/// Ok(StatusCode::NO_CONTENT)
254/// }
255/// ```
256pub fn delete(ctx: &Ctx, key: &str) -> impl Future<Output = Result<()>> + Send + use<> {
257 let env = ctx.env().clone();
258 let key = key.to_owned();
259 SendFuture::new(async move { Ok(bucket(&env)?.delete(key).await?) })
260}
261
262/// Deletes the objects of every `Some` attachment in one R2 call.
263///
264/// `None` entries (optional attachments left empty) are skipped, and
265/// nothing is sent when all are `None`. Generated models call it after a
266/// record is deleted, or after an update replaced or removed its files, so
267/// a failed database write never loses a file that a row still points to.
268/// Missing keys are not an error. R2 deletes are free.
269///
270/// # Errors
271///
272/// [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
273/// binding is missing (the message shows the `STORAGE` entry to add
274/// to cloudflare.config.ts) or R2 fails.
275///
276/// # Examples
277///
278/// ```no_run
279/// use axum::{extract::{Path, State}, http::StatusCode};
280/// use ocre::storage::{self, Attachment};
281/// use ocre::{Ctx, Result, params};
282///
283/// async fn destroy(State(ctx): State<Ctx>, Path(id): Path<i64>) -> Result<StatusCode> {
284/// let db = ctx.db()?;
285/// let image: Option<Attachment> = db
286/// .first(
287/// "SELECT image_key AS key, image_filename AS filename, image_content_type AS content_type, \
288/// image_size AS size FROM photos WHERE id = ?1 AND image_key IS NOT NULL",
289/// params![id],
290/// )
291/// .await?;
292/// db.execute("DELETE FROM photos WHERE id = ?1", params![id]).await?;
293/// storage::delete_attachments(&ctx, &[image]).await?;
294/// Ok(StatusCode::NO_CONTENT)
295/// }
296/// ```
297pub fn delete_attachments(
298 ctx: &Ctx,
299 attachments: &[Option<Attachment>],
300) -> impl Future<Output = Result<()>> + Send + use<> {
301 let env = ctx.env().clone();
302 let keys: Vec<String> = attachments.iter().flatten().map(|attachment| attachment.key.clone()).collect();
303 SendFuture::new(async move {
304 if !keys.is_empty() {
305 bucket(&env)?.delete_multiple(keys).await?;
306 }
307 Ok(())
308 })
309}
310
311/// The [`StoredObject`] of an R2 object (from `head`, `get` or `list` with metadata).
312fn stored(object: &Object) -> StoredObject {
313 StoredObject {
314 key: object.key(),
315 size: object.size(),
316 content_type: essence(object.http_metadata().content_type.as_deref().unwrap_or_default()),
317 etag: object.etag(),
318 uploaded_at: i64::try_from(object.uploaded().as_millis() / 1000).unwrap_or(i64::MAX),
319 filename: object.custom_metadata().ok().and_then(|mut metadata| metadata.remove("filename")),
320 }
321}
322
323/// The S3 endpoint of the `R2_*` variables and secrets.
324fn endpoint(env: &Env) -> Result<S3Endpoint> {
325 r2_endpoint(&|name| env.var(name).ok().map(|value| value.to_string()))
326}
327
328/// How long the URL of a [`direct_upload`] can be used to start the `PUT`: 10 minutes.
329const DIRECT_UPLOAD_EXPIRES_IN: u64 = 600;
330
331/// Describes the object at `key` without reading it: size, content type, ETag, upload time; `None` when it does not exist.
332///
333/// Active Storage's `blob.byte_size`/`service.exist?` in one call. Direct
334/// uploads are checked with it ([`attach_direct_upload`]), and purges read
335/// the upload time.
336///
337/// Free plan: one R2 class B operation (10M free per month).
338///
339/// # Errors
340///
341/// [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
342/// binding is missing (the message shows the `STORAGE` entry to add
343/// to cloudflare.config.ts) or R2 fails.
344///
345/// # Examples
346///
347/// ```no_run
348/// use axum::extract::{Path, State};
349/// use ocre::{Ctx, OptionExt, Result, storage};
350///
351/// // GET /files/{*key}/size
352/// async fn size(State(ctx): State<Ctx>, Path(key): Path<String>) -> Result<String> {
353/// let object = storage::head(&ctx, &key).await?.or_404()?;
354/// Ok(format!("{} ({}), uploaded at {}", storage::human_size(object.size), object.content_type, object.uploaded_at))
355/// }
356/// ```
357pub fn head(ctx: &Ctx, key: &str) -> impl Future<Output = Result<Option<StoredObject>>> + Send + use<> {
358 let env = ctx.env().clone();
359 let key = key.to_owned();
360 SendFuture::new(async move { Ok(bucket(&env)?.head(key).await?.as_ref().map(stored)) })
361}
362
363/// Returns whether an object exists at `key` (Active Storage's `exist?`); a [`head`] without the details.
364///
365/// Free plan: one R2 class B operation (10M free per month).
366///
367/// # Errors
368///
369/// [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
370/// binding is missing or R2 fails.
371///
372/// # Examples
373///
374/// ```no_run
375/// use axum::extract::{Path, State};
376/// use ocre::{Ctx, Result, storage};
377///
378/// async fn check(State(ctx): State<Ctx>, Path(key): Path<String>) -> Result<String> {
379/// Ok(storage::exists(&ctx, &key).await?.to_string())
380/// }
381/// ```
382pub fn exists(ctx: &Ctx, key: &str) -> impl Future<Output = Result<bool>> + Send + use<> {
383 let object = head(ctx, key);
384 async move { Ok(object.await?.is_some()) }
385}
386
387/// Lists one page of the objects whose key starts with `prefix`, in key order, with their size, type and upload time.
388///
389/// Pass `None` as `cursor` for the first page, then the returned
390/// [`Listing::cursor`] until it is `None`. `limit` is clamped to 1..=1000;
391/// a page may hold fewer objects even when more follow (R2 caps the
392/// metadata it returns), so rely on the cursor, not the count.
393///
394/// Free plan: one R2 class A operation per call, like an upload (1M free
395/// per month), and decoding 1,000 entries takes a few ms of the 10 ms of
396/// CPU: list in scheduled tasks or admin pages, never on every request.
397///
398/// # Errors
399///
400/// [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
401/// binding is missing or R2 fails.
402///
403/// # Examples
404///
405/// ```no_run
406/// use ocre::{Ctx, Result, storage};
407///
408/// // Total size of the exports, 1,000 keys (one class A operation) at a time.
409/// async fn exports_size(ctx: &Ctx) -> Result<u64> {
410/// let (mut total, mut cursor) = (0, None);
411/// loop {
412/// let page = storage::list(ctx, "exports/", cursor.as_deref(), 1000).await?;
413/// total += page.objects.iter().map(|object| object.size).sum::<u64>();
414/// cursor = page.cursor;
415/// if cursor.is_none() {
416/// return Ok(total);
417/// }
418/// }
419/// }
420/// ```
421pub fn list(
422 ctx: &Ctx,
423 prefix: &str,
424 cursor: Option<&str>,
425 limit: u32,
426) -> impl Future<Output = Result<Listing>> + Send + use<> {
427 let env = ctx.env().clone();
428 let prefix = prefix.to_owned();
429 let cursor = cursor.map(str::to_owned);
430 let limit = limit.clamp(1, MAX_LIST);
431 SendFuture::new(async move {
432 let bucket = bucket(&env)?;
433 let mut request =
434 bucket.list().prefix(prefix).limit(limit).include(vec![Include::HttpMetadata, Include::CustomMetadata]);
435 if let Some(cursor) = cursor {
436 request = request.cursor(cursor);
437 }
438 let page = request.execute().await?;
439 let objects = page.objects().iter().map(stored).collect();
440 Ok(Listing { objects, cursor: if page.truncated() { page.cursor() } else { None } })
441 })
442}
443
444/// Reads the first `length` bytes of an object (all of it when shorter), or `None` when `key` does not exist.
445///
446/// For [`analyze`](crate::storage::analyze) on a file uploaded directly:
447/// the signature and image size are in the first few KB (use 256 KB for
448/// JPEGs with large EXIF blocks), so the Worker never holds the whole file.
449///
450/// Free plan: one R2 class B operation (10M free per month).
451///
452/// # Errors
453///
454/// [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
455/// binding is missing or R2 fails.
456///
457/// # Examples
458///
459/// ```no_run
460/// use ocre::{Ctx, Result, storage};
461///
462/// async fn dimensions(ctx: &Ctx, key: &str) -> Result<Option<(u32, u32)>> {
463/// let Some(start) = storage::read_first(ctx, key, 64 * 1024).await? else { return Ok(None) };
464/// let analysis = storage::analyze(&start);
465/// Ok(analysis.width.zip(analysis.height))
466/// }
467/// ```
468pub fn read_first(ctx: &Ctx, key: &str, length: u64) -> impl Future<Output = Result<Option<Vec<u8>>>> + Send + use<> {
469 let env = ctx.env().clone();
470 let key = key.to_owned();
471 let length = length.max(1);
472 SendFuture::new(async move {
473 let bucket = bucket(&env)?;
474 let Some(object) = bucket.get(&key).range(Range::Prefix { length }).execute().await? else { return Ok(None) };
475 match object.body() {
476 Some(body) => Ok(Some(body.bytes().await?)),
477 None => Ok(Some(Vec::new())),
478 }
479 })
480}
481
482/// Presigns a `GET` of an attachment on R2's S3 API, valid `expires_in` seconds (at most 7 days).
483///
484/// The browser downloads straight from R2 (`https://<account>.r2.cloudflarestorage.com/...`),
485/// not through the Worker. R2 answers with the same safe `Content-Type`
486/// and `Content-Disposition` as [`serve`] (the URL carries
487/// `response-content-type` and `response-content-disposition`). Anyone with
488/// the URL can download the file until it expires: authorize before
489/// presigning and keep lifetimes short. [`serve_redirect`] answers a
490/// redirect to it.
491///
492/// Reads the [`R2_ACCOUNT_ID`](crate::storage::R2_ACCOUNT_ID) and
493/// [`R2_BUCKET`](crate::storage::R2_BUCKET) variables and the
494/// [`R2_ACCESS_KEY_ID`](crate::storage::R2_ACCESS_KEY_ID) and
495/// [`R2_SECRET_ACCESS_KEY`](crate::storage::R2_SECRET_ACCESS_KEY) secrets.
496/// Signing is local: no R2 operation, microseconds of CPU; the download
497/// is one class B operation. In `ocre dev` the URL points to the real
498/// bucket, not the local simulation.
499///
500/// # Errors
501///
502/// [`Error::Internal`](crate::Error::Internal) (500) when one of the four
503/// settings is missing (the message names them and how to set them), or
504/// `expires_in` is 0 or over 7 days.
505///
506/// # Examples
507///
508/// ```no_run
509/// use axum::{Json, extract::{Path, State}};
510/// use ocre::storage::{self, Attachment, Disposition};
511/// use ocre::{Ctx, OptionExt, Result, params};
512///
513/// // GET /api/documents/{id}/download_url: a link valid 5 minutes.
514/// async fn download_url(State(ctx): State<Ctx>, Path(id): Path<i64>) -> Result<Json<String>> {
515/// let sql = "SELECT file_key AS key, file_filename AS filename, file_content_type AS content_type, \
516/// file_size AS size FROM documents WHERE id = ?1 AND file_key IS NOT NULL";
517/// let file: Attachment = ctx.db()?.first(sql, params![id]).await?.or_404()?;
518/// Ok(Json(storage::presign_get(&ctx, &file, Disposition::Download, 300)?))
519/// }
520/// ```
521pub fn presign_get(ctx: &Ctx, attachment: &Attachment, disposition: Disposition, expires_in: u64) -> Result<String> {
522 presign_get_url(&endpoint(ctx.env())?, attachment, disposition, crate::now(), expires_in)
523}
524
525/// Presigns a `PUT` of exactly `size` bytes of type `content_type` at `key`, valid `expires_in` seconds.
526///
527/// The URL signs `Content-Type` (normalized) and `Content-Length`: R2
528/// refuses a `PUT` with another type or size. For browser uploads prefer
529/// [`direct_upload`], which also picks a new key, checks [`Rules`] and
530/// signs the key for [`attach_direct_upload`]. Settings, costs and errors
531/// as [`presign_get`]; the upload is one class A operation.
532///
533/// # Examples
534///
535/// ```no_run
536/// use ocre::{Ctx, Result, storage};
537///
538/// // A URL a backup script can `curl -T` to, valid one hour.
539/// fn backup_url(ctx: &Ctx, size: u64) -> Result<String> {
540/// storage::presign_put(ctx, "backups/latest.tar", "application/x-tar", size, 3600)
541/// }
542/// ```
543pub fn presign_put(ctx: &Ctx, key: &str, content_type: &str, size: u64, expires_in: u64) -> Result<String> {
544 let content_type = essence(content_type);
545 let size = size.to_string();
546 let headers = [("content-type", content_type.as_str()), ("content-length", size.as_str())];
547 endpoint(ctx.env())?.presign("PUT", key, &headers, &[], crate::now(), expires_in)
548}
549
550/// Starts a direct upload: checks what the browser declares against `rules` and presigns a `PUT` of a new key under `prefix`.
551///
552/// Rails' `DirectUploadsController#create`. The browser `PUT`s the file to
553/// [`DirectUpload::url`] with [`DirectUpload::headers`] (within 10
554/// minutes), then submits [`DirectUpload::signed_key`] with its form;
555/// [`attach_direct_upload`] checks the stored object and returns its
556/// [`Attachment`]. The file never passes through the Worker: no 100 MB
557/// request limit, no Worker memory or CPU. The key is
558/// `<prefix>/<22 random characters>`; use a prefix of its own (`uploads/photos`)
559/// so [`purge_unattached`] can find abandoned uploads.
560///
561/// Settings as [`presign_get`]; the bucket also needs a CORS rule allowing
562/// `PUT` from the app's origin (see the file storage guide). Signing costs
563/// no R2 operation; the upload is one class A operation.
564///
565/// # Errors
566///
567/// - [`Error::Invalid`](crate::Error::Invalid) (422) on `field` when the
568/// declared size or type breaks `rules`.
569/// - [`Error::Internal`](crate::Error::Internal) (500) when an `R2_*`
570/// setting is missing.
571///
572/// # Examples
573///
574/// ```no_run
575/// use axum::{Json, extract::State};
576/// use ocre::storage::{self, DirectUpload, DirectUploadRequest, Rules};
577/// use ocre::{ApiResult, Ctx};
578///
579/// const VIDEO: Rules = Rules { max_bytes: 500 * 1024 * 1024, content_types: &["video/mp4"] };
580///
581/// // POST /videos/uploads with {"filename": "a.mp4", "content_type": "video/mp4", "size": 1234}
582/// async fn start(State(ctx): State<Ctx>, Json(request): Json<DirectUploadRequest>) -> ApiResult<Json<DirectUpload>> {
583/// Ok(Json(storage::direct_upload(&ctx, "uploads/videos", "video", &request, &VIDEO)?))
584/// }
585/// ```
586pub fn direct_upload(
587 ctx: &Ctx,
588 prefix: &str,
589 field: &str,
590 request: &DirectUploadRequest,
591 rules: &Rules,
592) -> Result<DirectUpload> {
593 DirectUpload::sign(&endpoint(ctx.env())?, prefix, field, request, rules, crate::now(), DIRECT_UPLOAD_EXPIRES_IN)
594}
595
596/// Finishes a direct upload (or a [`multipart_uploads`](crate::storage::multipart_uploads) one): checks the object behind `signed_key` against `rules` and returns its [`Attachment`].
597///
598/// Rails' `attach(signed_blob_id)`. `signed_key` must come from
599/// [`direct_upload`] (a key this app signed, so a client cannot claim
600/// another record's file); the object must exist, and its size and
601/// content type (as R2 recorded them) must pass `rules`. A refused object
602/// is deleted. `filename` is the name the form sends (cleaned up). Save
603/// the attachment with [`columns`](crate::storage::columns); if that write
604/// fails, [`delete`] the key.
605///
606/// Free plan: one R2 class B operation (`head`), plus a free delete when
607/// the object is refused.
608///
609/// # Errors
610///
611/// - [`Error::Invalid`](crate::Error::Invalid) (422) on `field`: "is not a
612/// valid upload" (bad signature), "was not uploaded" (no object), or the
613/// messages of [`Validator::file`](crate::Validator::file).
614/// - [`Error::Internal`](crate::Error::Internal) (500) when neither
615/// `R2_SECRET_ACCESS_KEY` nor `SECRET_KEY_BASE` is set (the secret that
616/// signs upload keys), the `STORAGE` binding is missing, or R2 fails.
617///
618/// # Examples
619///
620/// ```no_run
621/// use axum::{Form, extract::{Path, State}};
622/// use ocre::storage::{self, Rules};
623/// use ocre::{Ctx, IntoParam, Result};
624/// use serde::Deserialize;
625///
626/// const VIDEO: Rules = Rules { max_bytes: 500 * 1024 * 1024, content_types: &["video/mp4"] };
627///
628/// #[derive(Deserialize)]
629/// struct VideoForm {
630/// video_key: String,
631/// video_filename: String,
632/// }
633///
634/// async fn update(State(ctx): State<Ctx>, Path(id): Path<i64>, Form(form): Form<VideoForm>) -> Result<String> {
635/// let video = storage::attach_direct_upload(&ctx, "video", &form.video_key, &form.video_filename, &VIDEO).await?;
636/// let mut values = Vec::from(storage::columns(Some(&video)));
637/// values.push(id.into_param());
638/// let sql = "UPDATE lessons SET video_key = ?1, video_filename = ?2, video_content_type = ?3, video_size = ?4 \
639/// WHERE id = ?5";
640/// if let Err(err) = ctx.db()?.execute(sql, values).await {
641/// storage::delete(&ctx, &video.key).await?;
642/// return Err(err);
643/// }
644/// Ok(video.key)
645/// }
646/// ```
647pub fn attach_direct_upload(
648 ctx: &Ctx,
649 field: &str,
650 signed_key: &str,
651 filename: &str,
652 rules: &Rules,
653) -> impl Future<Output = Result<Attachment>> + Send + use<> {
654 let env = ctx.env().clone();
655 let (field, signed_key, filename, rules) = (field.to_owned(), signed_key.to_owned(), filename.to_owned(), *rules);
656 SendFuture::new(async move {
657 let key = verify_key(&field, &upload_secret(&|name| env.var(name).ok().map(|v| v.to_string()))?, &signed_key)?;
658 let bucket = bucket(&env)?;
659 let object = bucket.head(key).await?.as_ref().map(stored);
660 let attached = attachment_from_head(&field, object.as_ref(), &filename, &rules);
661 if let (Err(_), Some(refused)) = (&attached, object) {
662 bucket.delete(refused.key).await?;
663 }
664 attached
665 })
666}
667
668/// Answers `302 Found` to a presigned `GET` of the attachment (Active Storage's redirect mode).
669///
670/// For large or popular files: the browser downloads from R2 directly,
671/// so the Worker only signs a URL. The redirect may be cached privately
672/// for half of `expires_in`. Authorize before calling it, as with
673/// [`serve`]: whoever gets the URL can use it until it expires. Settings,
674/// costs and errors as [`presign_get`].
675///
676/// # Examples
677///
678/// ```no_run
679/// use axum::{extract::{Path, State}, response::Response};
680/// use ocre::storage::{self, Attachment, Disposition};
681/// use ocre::{Ctx, OptionExt, Result, params};
682///
683/// // GET /videos/{id}/file: a redirect valid one hour.
684/// async fn file(State(ctx): State<Ctx>, Path(id): Path<i64>) -> Result<Response> {
685/// let sql = "SELECT file_key AS key, file_filename AS filename, file_content_type AS content_type, \
686/// file_size AS size FROM videos WHERE id = ?1";
687/// let video: Attachment = ctx.db()?.first(sql, params![id]).await?.or_404()?;
688/// storage::serve_redirect(&ctx, &video, Disposition::Inline, 3600)
689/// }
690/// ```
691pub fn serve_redirect(
692 ctx: &Ctx,
693 attachment: &Attachment,
694 disposition: Disposition,
695 expires_in: u64,
696) -> Result<Response> {
697 Ok(redirect_response(&presign_get(ctx, attachment, disposition, expires_in)?, expires_in))
698}
699
700/// The permanent public URL of `key` in a public bucket: `<STORAGE_PUBLIC_URL>/<key>` (Active Storage's `public: true`).
701///
702/// Requires public access on the bucket (an `r2.dev` URL, rate-limited and
703/// meant for development, or a custom domain) and the
704/// [`STORAGE_PUBLIC_URL`](crate::storage::STORAGE_PUBLIC_URL) variable.
705/// Every object of the bucket is then public forever to whoever has its
706/// key, with no authorization and no expiry: use it for avatars and
707/// product images, never for private documents. Pure string building: no
708/// R2 operation; downloads are class B operations on R2's side.
709///
710/// # Errors
711///
712/// [`Error::Internal`](crate::Error::Internal) (500) when
713/// `STORAGE_PUBLIC_URL` is not set (the message says how).
714///
715/// # Examples
716///
717/// ```no_run
718/// use ocre::{Ctx, Result, storage::{self, Attachment}};
719///
720/// fn avatar_src(ctx: &Ctx, avatar: &Attachment) -> Result<String> {
721/// storage::public_url(ctx, &avatar.key) // https://files.example.com/users/avatar/2u1Vd0zJ8sQqS6rJq0rVmA
722/// }
723/// ```
724pub fn public_url(ctx: &Ctx, key: &str) -> Result<String> {
725 join_public_url(ctx.env().var(STORAGE_PUBLIC_URL).ok().map(|value| value.to_string()), key)
726}
727
728/// Deletes the objects under `prefix` older than `max_age` seconds that no `table.column` row references, one listing page per call.
729///
730/// Direct uploads that were never attached (a closed tab, a failed form)
731/// would stay in R2 forever: Rails' `ActiveStorage::Blob.unattached`
732/// purge. Each call lists one page of up to 1,000 keys from `cursor`
733/// (`None` to start), keeps those uploaded more than `max_age` seconds
734/// ago, looks them up with `SELECT <column> FROM <table> WHERE <column> IN
735/// (...)` (100 keys per query) and deletes the others in one call. Pass
736/// the returned [`Purged::cursor`] to the next call until it is `None`.
737/// `table` and `column` are identifiers from your code (letters, digits,
738/// `_`); give the column an index (`CREATE UNIQUE INDEX ... ON
739/// lessons(video_key)`) so each lookup reads only matching rows. Keep
740/// `max_age` well above 10 minutes (the lifetime of an upload URL) plus
741/// the time a form stays open: a day is safe.
742///
743/// Free plan: one R2 class A operation per call (the listing), one D1
744/// query per 100 old keys, deletes free; a few ms of CPU per 1,000 keys.
745/// Run it from a scheduled task, a few pages per run.
746///
747/// # Errors
748///
749/// [`Error::Internal`](crate::Error::Internal) (500) when `table` or
750/// `column` is not a plain identifier, the `STORAGE` or `DB` binding is
751/// missing, or R2 or D1 fails (nothing is deleted for the page then).
752///
753/// # Examples
754///
755/// ```no_run
756/// use ocre::{Ctx, Result, storage};
757///
758/// // src/schedules/purge_uploads.rs, run every day: at most 5 pages (5 class A operations).
759/// pub async fn run(ctx: &Ctx) -> Result<()> {
760/// let mut cursor = None;
761/// for _ in 0..5 {
762/// let purged = storage::purge_unattached(ctx, "uploads/videos", "lessons", "video_key", 86_400, cursor.as_deref()).await?;
763/// println!("purge_uploads: {} unattached uploads deleted", purged.deleted.len());
764/// cursor = purged.cursor;
765/// if cursor.is_none() {
766/// break;
767/// }
768/// }
769/// Ok(())
770/// }
771/// ```
772pub fn purge_unattached(
773 ctx: &Ctx,
774 prefix: &str,
775 table: &str,
776 column: &str,
777 max_age: i64,
778 cursor: Option<&str>,
779) -> impl Future<Output = Result<Purged>> + Send + use<> {
780 #[derive(serde::Deserialize)]
781 struct Row {
782 key: String,
783 }
784 let ctx = ctx.clone();
785 let (table, column) = (table.to_owned(), column.to_owned());
786 let page = list(&ctx, prefix, cursor, MAX_LIST);
787 SendFuture::new(async move {
788 referenced_keys_sql(&table, &column, 0)?;
789 let page = page.await?;
790 let mut stale = stale_keys(&page.objects, crate::now() - max_age);
791 if !stale.is_empty() {
792 let db = ctx.db()?;
793 let mut referenced = std::collections::HashSet::new();
794 for keys in stale.chunks(KEYS_PER_QUERY) {
795 let sql = referenced_keys_sql(&table, &column, keys.len())?;
796 let params = keys.iter().map(|key| key.as_str().into_param()).collect();
797 let rows: Vec<Row> = db.all(&sql, params).await?;
798 referenced.extend(rows.into_iter().map(|row| row.key));
799 }
800 stale.retain(|key| !referenced.contains(key));
801 if !stale.is_empty() {
802 bucket(ctx.env())?.delete_multiple(stale.iter().map(String::as_str).collect()).await?;
803 }
804 }
805 Ok(Purged { deleted: stale, cursor: page.cursor })
806 })
807}
808
809/// Streams a stored file to the client, with the headers a browser needs for caching, seeking and saving it.
810///
811/// The response carries `Content-Type`, `Content-Length`,
812/// `Content-Disposition` (the original file name; `inline` only for types
813/// safe to display, see [`Disposition`](crate::storage::Disposition)),
814/// `ETag` and `Cache-Control` ([`CACHE_CONTROL`](crate::storage::CACHE_CONTROL)).
815/// A matching `If-None-Match` answers `304 Not Modified` without a body; a
816/// single `Range` (video seeking, resumed downloads) answers 206 with
817/// `Content-Range`, or 416 when it lies outside the file (without calling
818/// R2). Several ranges, other units and malformed values send the whole
819/// file, as RFC 9110 allows. The headers come from `attachment`, so it must
820/// be the row saved for that key.
821///
822/// Check that the user may see the record before calling it: the route is
823/// the only protection.
824///
825/// Free plan: one R2 class B operation per call, 304s included (10M free
826/// per month). The bytes never pass through WebAssembly: the R2 stream is
827/// attached to the response and [`crate::serve`] answers with it directly,
828/// so a download costs almost no CPU whatever its size.
829///
830/// # Errors
831///
832/// - [`Error::NotFound`](crate::Error::NotFound) (404) when no object has
833/// `attachment.key`.
834/// - [`Error::Internal`](crate::Error::Internal) (500) when the `STORAGE`
835/// binding is missing (the message shows the `STORAGE` entry to add
836/// to cloudflare.config.ts) or R2 fails.
837///
838/// # Examples
839///
840/// ```no_run
841/// use axum::{extract::{Path, Query, State}, http::HeaderMap, response::Response};
842/// use ocre::storage::{self, Attachment, Disposition};
843/// use ocre::{Ctx, OptionExt, Result, params};
844/// use serde::Deserialize;
845///
846/// #[derive(Deserialize)]
847/// struct Download {
848/// #[serde(default)]
849/// download: bool,
850/// }
851///
852/// // GET /photos/{id}/image, or /photos/{id}/image?download=true for "Save as".
853/// async fn image(
854/// State(ctx): State<Ctx>,
855/// Path(id): Path<i64>,
856/// Query(query): Query<Download>,
857/// headers: HeaderMap,
858/// ) -> Result<Response> {
859/// let sql = "SELECT image_key AS key, image_filename AS filename, image_content_type AS content_type, \
860/// image_size AS size FROM photos WHERE id = ?1";
861/// let image: Attachment = ctx.db()?.first(sql, params![id]).await?.or_404()?;
862/// let disposition = if query.download { Disposition::Download } else { Disposition::Inline };
863/// storage::serve(&ctx, &image, &headers, disposition).await
864/// }
865/// ```
866pub fn serve(
867 ctx: &Ctx,
868 attachment: &Attachment,
869 headers: &HeaderMap,
870 disposition: Disposition,
871) -> impl Future<Output = Result<Response>> + Send + use<> {
872 let env = ctx.env().clone();
873 let attachment = attachment.clone();
874 let fetch = Fetch::from_headers(headers, attachment.size);
875 SendFuture::new(async move {
876 if fetch.range == ByteRange::Unsatisfiable {
877 return Ok(unsatisfiable(attachment.size));
878 }
879 let bucket = bucket(&env)?;
880 let mut get = bucket.get(&attachment.key);
881 if let Some(etag) = fetch.if_none_match {
882 get = get.only_if(Conditional { etag_does_not_match: Some(etag), ..Default::default() });
883 }
884 if let ByteRange::Partial { offset, length } = fetch.range {
885 get = get.range(Range::OffsetWithLength { offset, length });
886 }
887 let object = get.execute().await?.ok_or(Error::NotFound)?;
888 let etag = object.http_etag();
889 let Some(body) = object.body() else { return Ok(not_modified(&etag)) };
890 let mut response = file_response(&attachment, disposition, &etag, fetch.range, Body::empty());
891 match body.response_body()? {
892 // Handed to `ocre::serve`, which answers with this stream itself.
893 ResponseBody::Stream(stream) => {
894 response.extensions_mut().insert(R2Stream(SendWrapper::new(stream)));
895 }
896 ResponseBody::Body(bytes) => *response.body_mut() = Body::from(bytes),
897 ResponseBody::Empty => {}
898 }
899 Ok(response)
900 })
901}
902
903/// An R2 object body attached to a response by [`serve`]: `ocre::serve`
904/// sends it as the JavaScript response body, so the bytes never pass
905/// through WebAssembly and workerd keeps `Content-Length`.
906#[derive(Clone)]
907pub(crate) struct R2Stream(SendWrapper<web_sys::ReadableStream>);
908
909/// The JavaScript response for `response`'s status and headers, with `stream` as its body.
910pub(crate) fn into_js_response(response: Response, stream: R2Stream) -> worker::Result<web_sys::Response> {
911 let headers = web_sys::Headers::new()?;
912 for (name, value) in response.headers() {
913 headers.append(name.as_str(), value.to_str().unwrap_or_default())?;
914 }
915 let init = web_sys::ResponseInit::new();
916 init.set_status(response.status().as_u16());
917 init.set_headers(&headers);
918 Ok(web_sys::Response::new_with_opt_readable_stream_and_init(Some(&stream.0), &init)?)
919}