whiskers.git / backend / worker / src / object.rs

The Durable Object: one per household, holding that household's documents in its own SQLite.

This file is the whole of the Cloudflare-and-celld-specific storage code. It reads a request, lends [documents::serve] an object's SQLite as a [DocStore], and writes back the answer.

Atomicity. serve is synchronous and runs after the request body has been read: the select, the merge and the upsert happen in one turn of the object's thread with no await between them, so no other request can be interleaved. That holds on Cloudflare (input gates) and under celld (where an awaiting event can overlap another, so a merge that awaited would not be safe). Do not add an await between the read and the write.

12use whiskers_ports::StoreError;
13use worker::{DurableObject, Env, Request, Response, Result, SqlStorage, SqlStorageValue, State, durable_object};
15use crate::documents::{DocStore, serve};
16use crate::logging;

One household's documents, in a table of name and JSON text.

19struct Sql(SqlStorage);
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}