whiskers.git / backend / worker / src / object.rs
1//! The Durable Object: one per household, holding that household's documents in its own SQLite.
2//!
3//! This file is the whole of the Cloudflare-and-celld-specific storage code. It reads a request,
4//! lends [`documents::serve`] an object's SQLite as a [`DocStore`], and writes back the answer.
5//!
6//! **Atomicity.** `serve` is synchronous and runs after the request body has been read: the select,
7//! the merge and the upsert happen in one turn of the object's thread with no `await` between them,
8//! so no other request can be interleaved. That holds on Cloudflare (input gates) and under celld
9//! (where an awaiting event *can* overlap another, so a merge that awaited would not be safe). Do not
10//! add an `await` between the read and the write.
11
12use whiskers_ports::StoreError;
13use worker::{DurableObject, Env, Request, Response, Result, SqlStorage, SqlStorageValue, State, durable_object};
14
15use crate::documents::{DocStore, serve};
16use crate::logging;
17
18/// One household's documents, in a table of name and JSON text.
19struct Sql(SqlStorage);
20
21impl Sql {
22    fn open(sql: SqlStorage) -> std::result::Result<Self, StoreError> {
23        // Idempotent, so every activation of the object may run it.
24        sql.exec("CREATE TABLE IF NOT EXISTS documents (name TEXT PRIMARY KEY, body TEXT NOT NULL) WITHOUT ROWID", None)
25            .map_err(|e| {
26                log::error!("the object's table cannot be made: {e}");
27                StoreError::Unavailable
28            })?;
29        Ok(Self(sql))
30    }
31}
32
33impl DocStore for Sql {
34    fn load(&self, name: &str) -> std::result::Result<Option<String>, StoreError> {
35        #[derive(serde::Deserialize)]
36        struct Row {
37            body: String,
38        }
39        let rows: Vec<Row> = self
40            .0
41            .exec("SELECT body FROM documents WHERE name = ?", vec![SqlStorageValue::String(name.to_owned())])
42            .and_then(|c| c.to_array())
43            .map_err(|e| {
44                log::error!("{name}: cannot read the stored document: {e}");
45                StoreError::Unavailable
46            })?;
47        Ok(rows.into_iter().next().map(|r| r.body))
48    }
49
50    fn save(&self, name: &str, json: &str) -> std::result::Result<(), StoreError> {
51        self.0
52            .exec(
53                "INSERT OR REPLACE INTO documents (name, body) VALUES (?, ?)",
54                vec![SqlStorageValue::String(name.to_owned()), SqlStorageValue::String(json.to_owned())],
55            )
56            .map(|_| ())
57            .map_err(|e| {
58                log::error!("{name}: cannot write the document: {e}");
59                StoreError::Unavailable
60            })
61    }
62}
63
64#[durable_object(fetch)]
65pub struct HouseholdObject {
66    state: State,
67}
68
69impl DurableObject for HouseholdObject {
70    fn new(state: State, env: Env) -> Self {
71        logging::install(&env);
72        Self { state }
73    }
74
75    async fn fetch(&self, mut req: Request) -> Result<Response> {
76        let path = req.path();
77        // Everything that awaits comes before the synchronous section below.
78        let body = req.text().await?;
79        let answer = match Sql::open(self.state.storage().sql()) {
80            Ok(store) => serve(&store, &path, &body),
81            Err(_) => crate::documents::Answer { status: 503, body: String::new() },
82        };
83        let headers = worker::Headers::new();
84        headers.set("content-type", "application/json")?;
85        Ok(Response::from_body(worker::ResponseBody::Body(answer.body.into_bytes()))?.with_status(answer.status).with_headers(headers))
86    }
87}