lmjtfy.git / apps / lmjtfy / src / archive.rs
1//! The archive, kept where every isolate sees the same rows: one Durable
2//! Object with a SQL table of every request that was answered, and one of
3//! the questions people asked. The messages are `archive`'s.
4//!
5//! The same request is never sent to the same model twice (the user,
6//! 2026-10-02), unless a visitor asks for it to be (`Pick::Fresh`, the page's
7//! ↻), and then every response it got is kept, numbered. Three things make
8//! that hold:
9//!
10//! - `versions` has `(sent_to, request, version)` as its primary key: the
11//!   exact body, where it went, and which time. No hash stands in for the
12//!   body, so two different requests cannot be taken for one.
13//! - A call is looked up before anything is spent, and its response is kept
14//!   before anyone is told.
15//! - The archive makes the call itself. Identical asks that arrive while it
16//!   is on the wire wait on that one call (`running`) and are told what it
17//!   was told. The call is handed to `wait_until`, so it finishes and is
18//!   kept even if every visitor waiting on it has left.
19//!
20//! A call with no response (an error, a timeout) keeps nothing: there is
21//! nothing to return next time, so the next ask sends it.
22
23use std::cell::RefCell;
24use std::collections::HashMap;
25use std::rc::Rc;
26
27use archive::{Answer, Ask, Asker, Called, Cursor, Entry, Pick, Place, Rating, Seen, Vote, FEED, Fetch, Home, Live, MOST, Pot, Record, Stats, Told};
28use ask::{Outcome, Wanted};
29use budget::Which;
30use futures_util::FutureExt;
31use futures_util::future::{LocalBoxFuture, Shared};
32use jev_client::Runtime;
33use jev_worker::WorkerRuntime;
34use llm::Model;
35use serde::Deserialize;
36use worker::{DurableObject, Env, Request, Response, SqlStorage, SqlStorageValue, State, WebSocket, WebSocketIncomingMessage, WebSocketPair, console_error, durable_object};
37
38use crate::meter::{self, Refused};
39use crate::{NoJev, ai};
40
41/// The binding's name in `wrangler.toml`.
42const BINDING: &str = "ARCHIVE";
43/// One object for the whole Worker: a request anyone had answered is
44/// answered for everyone.
45const NAME: &str = ::archive::OBJECT;
46
47/// The schema, as the steps that made it, in order. A step runs once: the
48/// object records how many it has run (`migrated`) and runs the rest when it
49/// starts. Never edit a step that has shipped; add one.
50///
51/// `calls.sent_to` is Jev's endpoint or a Workers AI model id. (A Jev body
52/// names its own model, so the endpoint is enough.)
53const MIGRATIONS: [&[&str]; 11] = [
54    // Written `IF NOT EXISTS` because this step first shipped before
55    // `migrated` existed, so it meets its own tables on an older object.
56    &[
57        "CREATE TABLE IF NOT EXISTS calls (
58            sent_to TEXT NOT NULL,
59            request TEXT NOT NULL,
60            response TEXT NOT NULL,
61            request_id TEXT,
62            attempts INTEGER NOT NULL,
63            took_ms REAL NOT NULL,
64            answered_ms REAL NOT NULL,
65            PRIMARY KEY (sent_to, request)
66        )",
67        "CREATE TABLE IF NOT EXISTS asked (
68            input TEXT PRIMARY KEY,
69            answers TEXT NOT NULL,
70            asked_ms REAL NOT NULL
71        )",
72        "CREATE INDEX IF NOT EXISTS asked_lately ON asked (asked_ms DESC)",
73    ],
74    // How often a question was asked, whether it may be shown, whether it
75    // needed an LLM (every question before this step did), and the tally of
76    // requests sent and not sent.
77    &[
78        "ALTER TABLE asked ADD COLUMN times INTEGER NOT NULL DEFAULT 1",
79        "ALTER TABLE asked ADD COLUMN listed INTEGER NOT NULL DEFAULT 1",
80        "ALTER TABLE asked ADD COLUMN llm INTEGER NOT NULL DEFAULT 1",
81        "CREATE TABLE counts (name TEXT PRIMARY KEY, n INTEGER NOT NULL)",
82    ],
83    // A correction to the data, not the schema. In the first ten minutes
84    // these counts existed (2026-10-02), nearly everything counted was the
85    // agent load-testing the live site: one question reached "asked 67
86    // times". Counting starts again from here, and the question the agent
87    // invented to test the feed's moderation is removed.
88    &[
89        "UPDATE asked SET times = 1",
90        "DELETE FROM asked WHERE input = 'is my neighbour John Smith of 12 Elm Street a paedophile?'",
91        "DELETE FROM counts",
92    ],
93    // A second correction (2026-10-02): the agent checked four deploys by
94    // asking the live site this question, which was already kept, so each
95    // check sent nothing and only counted. It had been asked 3 times; the
96    // checks made it 7. Deploys are now checked with `/gate` on new text.
97    &["UPDATE asked SET times = times - 4 WHERE input = 'Is a slot machine a good retirement plan?' AND times >= 5"],
98    // Each browser counts once per question (the owner, 2026-10-02): a
99    // question's `times` goes up only for a browser not in `askers` for it.
100    // Counts from before stay as they were.
101    &["CREATE TABLE askers (input TEXT NOT NULL, who TEXT NOT NULL, PRIMARY KEY (input, who))"],
102    // Votes on Jev's answers (the owner, 2026-10-02): one per browser per
103    // answer, `answer` being the SHA-256 of the answers as kept, so an
104    // answer that changes starts again from no votes.
105    &["CREATE TABLE ratings (input TEXT NOT NULL, answer TEXT NOT NULL, who TEXT NOT NULL, vote INTEGER NOT NULL, PRIMARY KEY (input, answer, who))"],
106    // A request may be sent again when a visitor asks (the owner,
107    // 2026-10-02: the page's ↻), and every response it got is kept: `calls`
108    // becomes `versions`, numbered from 1, and what was kept is version 1.
109    &[
110        "CREATE TABLE versions (
111            sent_to TEXT NOT NULL,
112            request TEXT NOT NULL,
113            version INTEGER NOT NULL,
114            response TEXT NOT NULL,
115            request_id TEXT,
116            attempts INTEGER NOT NULL,
117            took_ms REAL NOT NULL,
118            answered_ms REAL NOT NULL,
119            PRIMARY KEY (sent_to, request, version)
120        )",
121        "INSERT INTO versions SELECT sent_to, request, 1, response, request_id, attempts, took_ms, answered_ms FROM calls",
122        "DROP TABLE calls",
123    ],
124    // Nothing is dropped (the owner, 2026-10-03): every event is kept whole
125    // in `events`, whose columns after `at_ms` are `archive::Event::COLUMNS`
126    // (a test holds the two together). And the owner may overrule the rules
127    // on what the feed shows: `moderated` is their say, NULL while they have
128    // none, and `listed` stays Jev's.
129    &[
130        "CREATE TABLE events (
131            id INTEGER PRIMARY KEY AUTOINCREMENT,
132            at_ms REAL NOT NULL,
133            what TEXT NOT NULL,
134            method TEXT NOT NULL,
135            host TEXT NOT NULL,
136            path TEXT NOT NULL,
137            query TEXT NOT NULL,
138            input TEXT NOT NULL,
139            detail TEXT NOT NULL,
140            status REAL NOT NULL,
141            sent REAL NOT NULL,
142            kept REAL NOT NULL,
143            llm REAL NOT NULL,
144            took_ms REAL NOT NULL,
145            first REAL NOT NULL,
146            daily REAL NOT NULL,
147            session REAL NOT NULL,
148            referrer TEXT NOT NULL,
149            source TEXT NOT NULL,
150            client TEXT NOT NULL,
151            family TEXT NOT NULL,
152            os TEXT NOT NULL,
153            device TEXT NOT NULL,
154            language TEXT NOT NULL,
155            agent TEXT NOT NULL,
156            browser TEXT NOT NULL,
157            ip TEXT NOT NULL,
158            country TEXT NOT NULL,
159            region TEXT NOT NULL,
160            city TEXT NOT NULL,
161            postcode TEXT NOT NULL,
162            timezone TEXT NOT NULL,
163            latitude REAL,
164            longitude REAL,
165            asn REAL,
166            network TEXT NOT NULL,
167            colo TEXT NOT NULL,
168            protocol TEXT NOT NULL
169        )",
170        "CREATE INDEX events_lately ON events (at_ms DESC)",
171        "CREATE INDEX events_browser ON events (browser, at_ms)",
172        "CREATE INDEX events_ip ON events (ip, at_ms)",
173        "ALTER TABLE asked ADD COLUMN moderated INTEGER",
174    ],
175    // What only the page can say (the owner, 2026-10-03, of what an
176    // analytics product would have): the link's other campaign tags, the
177    // screen and window, and how far down a page was seen. `read` and `out`
178    // events carry them (`POST /seen`).
179    &[
180        "ALTER TABLE events ADD COLUMN medium TEXT NOT NULL DEFAULT ''",
181        "ALTER TABLE events ADD COLUMN campaign TEXT NOT NULL DEFAULT ''",
182        "ALTER TABLE events ADD COLUMN screen TEXT NOT NULL DEFAULT ''",
183        "ALTER TABLE events ADD COLUMN viewport TEXT NOT NULL DEFAULT ''",
184        "ALTER TABLE events ADD COLUMN scroll REAL NOT NULL DEFAULT 0",
185    ],
186    // `who` is a hash of a browser's id and the question, so it is
187    // different on every row and says nothing to someone reading the table
188    // (the owner, 2026-10-03: "this is the same guy over and over"). The
189    // browser's id is kept beside it. `who` stays the key: it is still what
190    // counts a browser once per question. Rows from before are filled in
191    // from `events` where they can be (`ASKERS_NAMED`).
192    &["ALTER TABLE askers ADD COLUMN browser TEXT NOT NULL DEFAULT ''", "ALTER TABLE ratings ADD COLUMN browser TEXT NOT NULL DEFAULT ''"],
193    // The feeds in the order they are read, so a page of one is the rows it
194    // shows and not the whole table sorted: on 2026-10-03 the archive read
195    // three million rows by mid-morning, of five a day, and every home page
196    // was a scan of every question. (The old `asked_lately` has only
197    // `asked_ms`, and the feed's order breaks ties by `input`.)
198    &["CREATE INDEX asked_feed ON asked (asked_ms DESC, input DESC)", "CREATE INDEX asked_most ON asked (times DESC, asked_ms DESC)"],
199];
200
201/// The step (from 1) that gave `askers` and `ratings` their `browser`, after
202/// which the rows already there are named from `events`, once.
203const ASKERS_NAMED: usize = 10;
204
205// SQLite's numbers arrive as JS numbers, so every number is read as an f64.
206#[derive(Deserialize)]
207struct CallRow {
208    response: String,
209    request_id: Option<String>,
210    attempts: f64,
211    took_ms: f64,
212    answered_ms: f64,
213}
214
215#[derive(Deserialize)]
216struct AskedRow {
217    input: String,
218    answers: String,
219    asked_ms: f64,
220    times: f64,
221}
222
223impl From<AskedRow> for Entry {
224    fn from(row: AskedRow) -> Self {
225        Entry {
226            input: row.input,
227            answers: ::archive::stored(&row.answers),
228            asked_ms: row.asked_ms,
229            times: row.times as u32,
230        }
231    }
232}
233
234#[derive(Deserialize)]
235struct Number {
236    n: f64,
237}
238
239#[derive(Deserialize)]
240struct Tally {
241    questions: f64,
242    asks: f64,
243    no_llm: f64,
244}
245
246/// The repository whose clones and pulls the pages show. jevcrates is
247/// counted too, but a clone of lmjtfy fetches it, so it is not shown twice.
248const CODE: &str = "lmjtfy.git";
249
250/// The names in `counts`.
251const SENT: &str = "sent";
252const KEPT: &str = "kept";
253
254/// A call on the wire, which every identical ask waits on.
255type Running = Shared<LocalBoxFuture<'static, Called>>;
256
257/// What a call needs after the request that started it has gone: shared, so
258/// the call can outlive that request.
259struct Shelf {
260    env: Env,
261    sql: SqlStorage,
262    running: RefCell<HashMap<(String, String), Running>>,
263    /// What the home page shows, kept from the last time it was read until
264    /// something it is made of changes (`changed`). Every home page asks
265    /// for it, and its numbers are a pass over every question.
266    home: RefCell<Option<Home>>,
267}
268
269impl Shelf {
270    /// How many responses are kept for a request.
271    fn versions(&self, sent_to: &str, request: &str) -> worker::Result<u32> {
272        let rows: Vec<Number> = self
273            .sql
274            .exec("SELECT COUNT(*) AS n FROM versions WHERE sent_to = ? AND request = ?", vec![sent_to.into(), request.into()])?
275            .to_array()?;
276        Ok(rows.first().map_or(0, |row| row.n as u32))
277    }
278
279    /// A kept response: `version`, or the newest.
280    fn kept(&self, sent_to: &str, request: &str, version: Option<u32>) -> worker::Result<Option<Record>> {
281        let versions = self.versions(sent_to, request)?;
282        let Some(version) = version.filter(|version| (1..=versions).contains(version)).or((versions > 0).then_some(versions)) else {
283            return Ok(None);
284        };
285        let rows: Vec<CallRow> = self
286            .sql
287            .exec(
288                "SELECT response, request_id, attempts, took_ms, answered_ms FROM versions WHERE sent_to = ? AND request = ? AND version = ?",
289                vec![sent_to.into(), request.into(), i64::from(version).into()],
290            )?
291            .to_array()?;
292        Ok(rows.into_iter().next().map(|row| Record {
293            response: row.response,
294            request_id: row.request_id,
295            attempts: row.attempts as u32,
296            took_ms: row.took_ms,
297            answered_ms: row.answered_ms,
298            sent_now: false,
299            version,
300            versions,
301        }))
302    }
303
304    /// Keeps a response as the request's newest version, and says which it
305    /// is. A plain INSERT: two calls numbering the same version would mean
306    /// the archive sent twice what it meant to send once, and that should
307    /// fail here, not be papered over.
308    fn keep(&self, sent_to: &str, request: &str, record: &Record) -> worker::Result<u32> {
309        let version = self.versions(sent_to, request)? + 1;
310        self.sql.exec(
311            "INSERT INTO versions (sent_to, request, version, response, request_id, attempts, took_ms, answered_ms) VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
312            vec![
313                sent_to.into(),
314                request.into(),
315                i64::from(version).into(),
316                record.response.as_str().into(),
317                record.request_id.as_deref().map_or(worker::SqlStorageValue::Null, Into::into),
318                i64::from(record.attempts).into(),
319                record.took_ms.into(),
320                record.answered_ms.into(),
321            ],
322        )?;
323        Ok(version)
324    }
325
326    fn migrate(&self) -> worker::Result<()> {
327        self.sql.exec("CREATE TABLE IF NOT EXISTS migrated (version INTEGER PRIMARY KEY)", None)?;
328        let done = self.number("SELECT COUNT(*) AS n FROM migrated")? as usize;
329        for (index, step) in MIGRATIONS.iter().enumerate().skip(done) {
330            for statement in *step {
331                self.sql.exec(statement, None)?;
332            }
333            self.sql.exec("INSERT INTO migrated (version) VALUES (?)", vec![(index as i64 + 1).into()])?;
334            if index + 1 == ASKERS_NAMED {
335                self.name_askers()?;
336            }
337        }
338        Ok(())
339    }
340
341    /// Keeps an event, whole, and says which row it is.
342    fn happened(&self, event: &::archive::Event, now_ms: f64) -> worker::Result<f64> {
343        let columns = ::archive::Event::COLUMNS;
344        let marks = vec!["?"; columns.len() + 1].join(", ");
345        let mut values: Vec<SqlStorageValue> = vec![now_ms.into()];
346        values.extend(event.values().into_iter().map(|cell| match cell {
347            ::archive::event::Cell::Text(text) => text.into(),
348            ::archive::event::Cell::Number(number) => number.into(),
349            ::archive::event::Cell::Null => SqlStorageValue::Null,
350        }));
351        self.sql.exec(&format!("INSERT INTO events (at_ms, {}) VALUES ({marks})", columns.join(", ")), values)?;
352        self.number("SELECT last_insert_rowid() AS n")
353    }
354
355    /// A page that connected (the event in row `came`) has gone: the same
356    /// event again, as `left`, with how long it stayed.
357    fn left(&self, came: f64, stayed_ms: f64, now_ms: f64) -> worker::Result<()> {
358        let columns = ::archive::Event::COLUMNS;
359        let copied: Vec<&str> = columns
360            .iter()
361            .map(|column| match *column {
362                "what" => "'left'",
363                "took_ms" => "?",
364                column => column,
365            })
366            .collect();
367        self.sql.exec(
368            &format!("INSERT INTO events (at_ms, {}) SELECT ?, {} FROM events WHERE id = ?", columns.join(", "), copied.join(", ")),
369            vec![now_ms.into(), stayed_ms.into(), came.into()],
370        )?;
371        Ok(())
372    }
373
374    /// Whether the feed may show `input`: the owner's say if they have one,
375    /// Jev's otherwise.
376    fn shown(&self, input: &str) -> worker::Result<bool> {
377        let rows: Vec<Number> = self.sql.exec("SELECT COALESCE(moderated, listed) AS n FROM asked WHERE input = ?", vec![input.into()])?.to_array()?;
378        Ok(rows.first().is_some_and(|row| row.n == 1.0))
379    }
380
381    /// The admin backend's requests (`archive::Admin`).
382    fn admin(&self, ask: ::archive::Admin) -> ::archive::Answered {
383        use ::archive::{Admin, Answered};
384        let refused = |error: worker::Error| Answered::Refused(error.to_string());
385        match ask {
386            Admin::Select { sql, .. } if !::archive::reads_only(&sql) => Answered::Refused("only one statement that reads".into()),
387            Admin::Select { sql, values } => {
388                let values: Vec<SqlStorageValue> = values
389                    .into_iter()
390                    .map(|value| match value {
391                        serde_json::Value::Null => SqlStorageValue::Null,
392                        serde_json::Value::Bool(flag) => i64::from(flag).into(),
393                        serde_json::Value::Number(number) => number.as_f64().unwrap_or_default().into(),
394                        serde_json::Value::String(text) => text.into(),
395                        other => other.to_string().into(),
396                    })
397                    .collect();
398                let cursor = match self.sql.exec(&sql, values) {
399                    Ok(cursor) => cursor,
400                    Err(error) => return refused(error),
401                };
402                let columns = cursor.column_names();
403                let mut rows = Vec::new();
404                for row in cursor.raw() {
405                    let row = match row {
406                        Ok(row) => row,
407                        Err(error) => return refused(error),
408                    };
409                    rows.push(
410                        row.into_iter()
411                            .map(|value| match value {
412                                SqlStorageValue::Null => serde_json::Value::Null,
413                                SqlStorageValue::Boolean(flag) => flag.into(),
414                                SqlStorageValue::Integer(number) => number.into(),
415                                SqlStorageValue::Float(number) => number.into(),
416                                SqlStorageValue::String(text) => text.into(),
417                                SqlStorageValue::Blob(bytes) => format!("<{} bytes>", bytes.len()).into(),
418                            })
419                            .collect(),
420                    );
421                }
422                Answered::Rows { columns, rows }
423            }
424            Admin::Moderate { input, listed } => {
425                self.changed();
426                let say: SqlStorageValue = listed.map_or(SqlStorageValue::Null, |listed| i64::from(listed).into());
427                match self.sql.exec("UPDATE asked SET moderated = ? WHERE input = ?", vec![say, input.into()]) {
428                    Ok(_) => Answered::Done,
429                    Err(error) => refused(error),
430                }
431            }
432        }
433    }
434
435    /// Fills in whose each old `askers` and `ratings` row is, where an
436    /// event says which browser asked or voted on that question: the
437    /// browser's id and the question hash to the row's `who`. A row from
438    /// before events were kept stays unnamed.
439    fn name_askers(&self) -> worker::Result<()> {
440        #[derive(Deserialize)]
441        struct Seen {
442            input: String,
443            browser: String,
444        }
445        let seen: Vec<Seen> = self.sql.exec("SELECT DISTINCT input, browser FROM events WHERE browser != '' AND input != '' AND what IN ('answer', 'vote')", None)?.to_array()?;
446        for Seen { input, browser } in seen {
447            let who = crate::asker(&browser, &input);
448            for table in ["askers", "ratings"] {
449                self.sql.exec(
450                    &format!("UPDATE {table} SET browser = ? WHERE input = ? AND who = ? AND browser = ''"),
451                    vec![browser.as_str().into(), input.as_str().into(), who.0.as_str().into()],
452                )?;
453            }
454        }
455        Ok(())
456    }
457
458    fn number(&self, query: &str) -> worker::Result<f64> {
459        let rows: Vec<Number> = self.sql.exec(query, None)?.to_array()?;
460        Ok(rows.first().map_or(0.0, |row| row.n))
461    }
462
463    /// Adds one to a tally. A failure loses one from a number on the home
464    /// page and nothing else, so it is not allowed to fail a call.
465    fn count(&self, name: &str) {
466        self.changed();
467        let counted = self.sql.exec(
468            "INSERT INTO counts (name, n) VALUES (?, 1) ON CONFLICT (name) DO UPDATE SET n = n + 1",
469            vec![name.into()],
470        );
471        if let Err(error) = counted {
472            console_error!("the archive could not count a {name} request: {error}");
473        }
474    }
475
476    fn counted(&self, name: &str) -> worker::Result<u32> {
477        let rows: Vec<Number> = self.sql.exec("SELECT n FROM counts WHERE name = ?", vec![name.into()])?.to_array()?;
478        Ok(rows.first().map_or(0, |row| row.n as u32))
479    }
480
481    /// Records what Jev said to `input`, and whether this is a new asking of
482    /// it: by a browser that has not asked it before, or by one that gave no
483    /// id. Only a new asking counts and moves it up the feed; Jev's answer
484    /// is kept either way.
485    fn note(&self, input: &str, answers: &[Answer], llm: bool, listed: bool, who: Option<&Asker>, browser: &str, now_ms: f64) -> worker::Result<bool> {
486        self.changed();
487        let answers = serde_json::to_string(answers).map_err(|e| worker::Error::from(e.to_string()))?;
488        let new = match who {
489            Some(who) => {
490                let cursor = self.sql.exec("INSERT OR IGNORE INTO askers (input, who, browser) VALUES (?, ?, ?)", vec![input.into(), who.0.as_str().into(), browser.into()])?;
491                cursor.rows_written() > 0
492            }
493            None => true,
494        };
495        if !new {
496            self.sql.exec(
497                "UPDATE asked SET answers = ?, listed = ?, llm = ? WHERE input = ?",
498                vec![answers.into(), i64::from(listed).into(), i64::from(llm).into(), input.into()],
499            )?;
500            return Ok(false);
501        }
502        self.sql.exec(
503            "INSERT INTO asked (input, answers, asked_ms, times, listed, llm) VALUES (?, ?, ?, 1, ?, ?)
504             ON CONFLICT (input) DO UPDATE SET answers = excluded.answers, asked_ms = excluded.asked_ms,
505                 times = times + 1, listed = excluded.listed, llm = excluded.llm",
506            vec![input.into(), answers.into(), now_ms.into(), i64::from(listed).into(), i64::from(llm).into()],
507        )?;
508        Ok(true)
509    }
510
511    /// Which answer to `input` a vote is on: the SHA-256 of the answers as
512    /// kept. `None` if it was never answered.
513    fn answer_key(&self, input: &str) -> worker::Result<Option<String>> {
514        #[derive(Deserialize)]
515        struct Kept {
516            answers: String,
517        }
518        let rows: Vec<Kept> = self.sql.exec("SELECT answers FROM asked WHERE input = ?", vec![input.into()])?.to_array()?;
519        Ok(rows.first().map(|kept| {
520            use sha2::{Digest, Sha256};
521            Sha256::digest(kept.answers.as_bytes()).iter().map(|byte| format!("{byte:02x}")).collect()
522        }))
523    }
524
525    /// `who`'s vote on the answer to `input` as kept now; the vote it already
526    /// has takes it back.
527    fn rate(&self, input: &str, who: &Asker, browser: &str, vote: Vote) -> worker::Result<Option<Rating>> {
528        let Some(answer) = self.answer_key(input)? else { return Ok(None) };
529        let value = match vote {
530            Vote::Up => 1,
531            Vote::Down => -1,
532        };
533        let had = self.rating_for(input, &answer, Some(who))?.mine;
534        if had == Some(vote) {
535            self.sql.exec(
536                "DELETE FROM ratings WHERE input = ? AND answer = ? AND who = ?",
537                vec![input.into(), answer.as_str().into(), who.0.as_str().into()],
538            )?;
539        } else {
540            self.sql.exec(
541                "INSERT INTO ratings (input, answer, who, vote, browser) VALUES (?, ?, ?, ?, ?)
542                 ON CONFLICT (input, answer, who) DO UPDATE SET vote = excluded.vote",
543                vec![input.into(), answer.as_str().into(), who.0.as_str().into(), value.into(), browser.into()],
544            )?;
545        }
546        Ok(Some(self.rating_for(input, &answer, Some(who))?))
547    }
548
549    fn rating(&self, input: &str, who: Option<&Asker>) -> worker::Result<Option<Rating>> {
550        match self.answer_key(input)? {
551            Some(answer) => Ok(Some(self.rating_for(input, &answer, who)?)),
552            None => Ok(None),
553        }
554    }
555
556    fn rating_for(&self, input: &str, answer: &str, who: Option<&Asker>) -> worker::Result<Rating> {
557        #[derive(Deserialize)]
558        struct Votes {
559            up: f64,
560            down: f64,
561        }
562        #[derive(Deserialize)]
563        struct Mine {
564            vote: f64,
565        }
566        let votes: Vec<Votes> = self
567            .sql
568            .exec(
569                "SELECT COALESCE(SUM(vote = 1), 0) AS up, COALESCE(SUM(vote = -1), 0) AS down FROM ratings WHERE input = ? AND answer = ?",
570                vec![input.into(), answer.into()],
571            )?
572            .to_array()?;
573        let mine = match who {
574            Some(who) => {
575                let rows: Vec<Mine> = self
576                    .sql
577                    .exec(
578                        "SELECT vote FROM ratings WHERE input = ? AND answer = ? AND who = ?",
579                        vec![input.into(), answer.into(), who.0.as_str().into()],
580                    )?
581                    .to_array()?;
582                rows.first().map(|row| if row.vote > 0.0 { Vote::Up } else { Vote::Down })
583            }
584            None => None,
585        };
586        let votes = votes.first();
587        Ok(Rating { up: votes.map_or(0, |v| v.up as u32), down: votes.map_or(0, |v| v.down as u32), mine })
588    }
589
590    /// What Jev said to `input`, listed or not: it is for whoever holds the
591    /// link.
592    fn answer(&self, input: &str) -> worker::Result<Option<Entry>> {
593        let rows: Vec<AskedRow> = self
594            .sql
595            .exec("SELECT input, answers, asked_ms, times FROM asked WHERE input = ?", vec![input.into()])?
596            .to_array()?;
597        Ok(rows.into_iter().next().map(Entry::from))
598    }
599
600    /// The listed questions after `after` in the feed's order.
601    fn older(&self, after: &Cursor) -> worker::Result<Vec<Entry>> {
602        let rows: Vec<AskedRow> = self
603            .sql
604            .exec(
605                "SELECT input, answers, asked_ms, times FROM asked
606                 WHERE COALESCE(moderated, listed) = 1 AND (asked_ms < ? OR (asked_ms = ? AND input < ?))
607                 ORDER BY asked_ms DESC, input DESC LIMIT ?",
608                vec![after.asked_ms.into(), after.asked_ms.into(), after.input.as_str().into(), i64::from(FEED).into()],
609            )?
610            .to_array()?;
611        Ok(rows.into_iter().map(Entry::from).collect())
612    }
613
614    fn entries(&self, query: &str, limit: u32) -> worker::Result<Vec<Entry>> {
615        let rows: Vec<AskedRow> = self.sql.exec(query, vec![i64::from(limit).into()])?.to_array()?;
616        Ok(rows.into_iter().map(Entry::from).collect())
617    }
618
619    /// What the home page is made of has changed: read it again next time.
620    fn changed(&self) {
621        self.home.replace(None);
622    }
623
624    fn home(&self) -> worker::Result<Home> {
625        if let Some(home) = self.home.borrow().as_ref() {
626            return Ok(home.clone());
627        }
628        let home = self.read_home()?;
629        self.home.replace(Some(home.clone()));
630        Ok(home)
631    }
632
633    fn read_home(&self) -> worker::Result<Home> {
634        let tally: Vec<Tally> = self
635            .sql
636            .exec(
637                "SELECT COUNT(*) AS questions, COALESCE(SUM(times), 0) AS asks,
638                        COALESCE(SUM(llm = 0), 0) AS no_llm FROM asked",
639                None,
640            )?
641            .to_array()?;
642        let tally = tally.first();
643        Ok(Home {
644            lately: self.entries(
645                "SELECT input, answers, asked_ms, times FROM asked WHERE COALESCE(moderated, listed) = 1 ORDER BY asked_ms DESC, input DESC LIMIT ?",
646                FEED,
647            )?,
648            most: self.entries(
649                "SELECT input, answers, asked_ms, times FROM asked WHERE COALESCE(moderated, listed) = 1 AND times > 1
650                 ORDER BY times DESC, asked_ms DESC LIMIT ?",
651                MOST,
652            )?,
653            stats: Stats {
654                questions: tally.map_or(0, |tally| tally.questions as u32),
655                asks: tally.map_or(0, |tally| tally.asks as u32),
656                no_llm: tally.map_or(0, |tally| tally.no_llm as u32),
657                sent: self.counted(SENT)?,
658                kept: self.counted(KEPT)?,
659                clones: self.counted(&Fetch::Clone.counted(CODE))?,
660                pulls: self.counted(&Fetch::Pull.counted(CODE))?,
661            },
662        })
663    }
664}
665
666#[durable_object]
667pub struct Archive {
668    state: State,
669    shelf: Rc<Shelf>,
670}
671
672fn failed(error: impl Into<String>) -> Called {
673    Called::Failed { error: error.into(), request_id: None, took_ms: 0.0 }
674}
675
676/// Sends one Jev request, inside the day's Jev budget.
677async fn jev(shelf: Rc<Shelf>, input: String, wanted: Wanted, request: String) -> Called {
678    let jev = match crate::jev(&shelf.env) {
679        Ok(jev) => jev,
680        Err(NoJev::Offline) => return failed("the Worker has no TypeSafe API key"),
681        Err(NoJev::Broken(error)) => return failed(error),
682    };
683    let prepared = match ask::wanted(&jev.model, &input, &wanted) {
684        Ok(prepared) => prepared,
685        Err(error) => return failed(error.to_string()),
686    };
687    // The page shows `request` and the row is kept under it, so it has to be
688    // the body that goes out. They can differ only while a deploy has the
689    // Worker and this object on different code.
690    if prepared.request != request {
691        return failed("the archive would send a different request than the page prepared, so nothing was sent");
692    }
693    let hold = match meter::hold(&shelf.env, Which::JevDollars, prepared.worst_case_dollars, jev.dollars_per_day).await {
694        Ok(hold) => hold,
695        Err(Refused::Spent(_)) => return Called::Spent(Pot::Jev),
696        Err(Refused::Unreachable(error)) => return failed(error),
697    };
698    let outcome = ask::send(&jev.client, &WorkerRuntime, &prepared).await;
699    meter::settle(&shelf.env, Which::JevDollars, hold, outcome.dollars()).await;
700    match outcome {
701        Outcome::Answered { body, request_id, attempts, took, .. } => Called::Answered(Record {
702            response: body,
703            request_id,
704            attempts,
705            took_ms: took.as_secs_f64() * 1000.0,
706            answered_ms: js_sys::Date::now(),
707            sent_now: true,
708            // Numbered when it is kept (`Shelf::keep`).
709            version: 0,
710            versions: 0,
711        }),
712        Outcome::Failed { error, request_id, took, .. } => {
713            Called::Failed { error, request_id, took_ms: took.as_secs_f64() * 1000.0 }
714        }
715    }
716}
717
718/// Sends one request to the LLM, inside the day's neuron budget.
719async fn llm(shelf: Rc<Shelf>, id: String, request: String) -> Called {
720    let Some(model) = Model::find(&id) else {
721        return failed(format!("{id} is not one of llm's candidates"));
722    };
723    let hold = match meter::hold(&shelf.env, Which::Neurons, model.worst_case_neurons(&request), llm::FREE_NEURONS_PER_DAY).await {
724        Ok(hold) => hold,
725        Err(Refused::Spent(_)) => return Called::Spent(Pot::Llm),
726        Err(Refused::Unreachable(error)) => return failed(error),
727    };
728    let started = WorkerRuntime.now();
729    let ran = ai::run(&shelf.env, model.id, &request).await;
730    let took_ms = (WorkerRuntime.now() - started).as_secs_f64() * 1000.0;
731    // A call that never produced a reply is counted as free; one that did is
732    // counted at what the reply says it used.
733    let neurons = ran.as_ref().ok().and_then(|body| llm::parse(body).ok()).map_or(0.0, |reply| model.neurons(reply.usage));
734    meter::settle(&shelf.env, Which::Neurons, hold, neurons).await;
735    match ran {
736        Ok(response) => Called::Answered(Record {
737            response,
738            request_id: None,
739            attempts: 1,
740            took_ms,
741            answered_ms: js_sys::Date::now(),
742            sent_now: true,
743            // Numbered when it is kept (`Shelf::keep`).
744            version: 0,
745            versions: 0,
746        }),
747        Err(error) => Called::Failed { error, request_id: None, took_ms },
748    }
749}
750
751impl Archive {
752    /// The response to `request`, sent at most once: the kept one, the one a
753    /// call already on the wire is about to get, or the one `call` gets now.
754    ///
755    /// From the lookup to the insert into `running` nothing awaits, so no
756    /// other ask can run in between and start the same call.
757    ///
758    /// `pick` says which kept response is wanted, and whether to send it
759    /// again regardless (`Pick::Fresh`) or never (`Pick::KeptOnly`).
760    async fn once<F>(&self, sent_to: String, request: String, pick: Pick, call: impl FnOnce(Rc<Shelf>) -> F) -> worker::Result<Called>
761    where
762        F: Future<Output = Called> + 'static,
763    {
764        let wanted = match pick {
765            Pick::Version(version) => Some(version),
766            _ => None,
767        };
768        if pick != Pick::Fresh
769            && let Some(record) = self.shelf.kept(&sent_to, &request, wanted)?
770        {
771            self.shelf.count(KEPT);
772            return Ok(Called::Answered(record));
773        }
774        let key = (sent_to, request);
775        let running = self.shelf.running.borrow().get(&key).cloned();
776        if pick == Pick::KeptOnly && running.is_none() {
777            return Ok(Called::NotKept);
778        }
779        if let Some(running) = running {
780            return Ok(match running.await {
781                Called::Answered(record) => {
782                    self.shelf.count(KEPT);
783                    Called::Answered(record.kept())
784                }
785                other => other,
786            });
787        }
788        let shelf = self.shelf.clone();
789        let call = call(shelf.clone());
790        let done = key.clone();
791        let work: Running = async move {
792            let mut called = call.await;
793            if let Called::Answered(record) = &mut called {
794                shelf.count(SENT);
795                match shelf.keep(&done.0, &done.1, record) {
796                    Ok(version) => {
797                        record.version = version;
798                        record.versions = version;
799                    }
800                    Err(error) => console_error!("the archive could not keep a response, so its request may be sent again: {error}"),
801                }
802            }
803            shelf.running.borrow_mut().remove(&done);
804            called
805        }
806        .boxed_local()
807        .shared();
808        self.shelf.running.borrow_mut().insert(key, work.clone());
809        self.state.wait_until(work.clone().map(|_| ()));
810        Ok(work.await)
811    }
812}
813
814impl Archive {
815    /// Sends `live` to every open page. A page that cannot be told has gone,
816    /// and its close is on its way.
817    fn broadcast(&self, live: &Live) {
818        self.send(&self.state.get_websockets(), live);
819    }
820
821    fn send(&self, sockets: &[WebSocket], live: &Live) {
822        let Ok(text) = serde_json::to_string(live) else { return };
823        for socket in sockets {
824            let _ = socket.send_with_str(&text);
825        }
826    }
827
828    /// Sends every open page the activity feeds as they now stand: `asked`
829    /// at the top of the feed, if it was a question the feed may show, and
830    /// the most asked and the tally whole. Nobody open, nothing to render.
831    fn activity(&self, asked: Option<&str>) {
832        if self.state.get_websockets().is_empty() {
833            return;
834        }
835        let top = match asked.map(|input| self.shelf.answer(input)) {
836            Some(Ok(entry)) => entry,
837            Some(Err(error)) => {
838                console_error!("the archive could not read the question to send: {error}");
839                None
840            }
841            None => None,
842        };
843        match self.shelf.home() {
844            Ok(home) => {
845                let markup = maud::html! {
846                    @if let Some(top) = &top { (crate::view::lately_top(top)) }
847                    (crate::view::most(&home.most))
848                    (crate::view::tally(Some(&home.stats)))
849                };
850                self.broadcast(&Live::Patch(markup.into_string()));
851            }
852            Err(error) => console_error!("the archive could not read the feeds to send: {error}"),
853        }
854    }
855
856    /// The pages open now: every socket still open, `leaving` not counted.
857    fn open(&self, leaving: Option<&WebSocket>) -> Vec<WebSocket> {
858        self.state
859            .get_websockets()
860            .into_iter()
861            .filter(|socket| {
862                let raw: &worker::web_sys::WebSocket = socket.as_ref();
863                raw.ready_state() == worker::web_sys::WebSocket::OPEN
864            })
865            .filter(|socket| leaving.is_none_or(|leaving| !same(socket, leaving)))
866            .collect()
867    }
868
869    /// Tells `to` how many pages are open and where they are.
870    fn online(&self, to: &[WebSocket], open: &[WebSocket]) {
871        self.send(to, &Live::Online(open.len() as u32));
872        let places = open.iter().map(|socket| seen(socket).map(|seen| seen.place).unwrap_or_default());
873        self.send(to, &Live::Patch(crate::view::places(&::archive::places(places)).into_string()));
874    }
875
876    /// Everyone hears that pages came or went once, `ONLINE_EVERY_MS` after
877    /// the first change, however many changed in between. A deploy closes
878    /// every socket and they all reconnect within a second or two; telling
879    /// every page about every one would be pages² messages.
880    async fn soon(&self) -> worker::Result<()> {
881        let storage = self.state.storage();
882        if storage.get_alarm().await?.is_none() {
883            storage.set_alarm(std::time::Duration::from_millis(ONLINE_EVERY_MS)).await?;
884        }
885        Ok(())
886    }
887}
888
889/// What is kept on a page's socket, if anything.
890fn seen(socket: &WebSocket) -> Option<Seen> {
891    socket.deserialize_attachment::<Seen>().ok().flatten()
892}
893
894/// Whether two handles are the same socket.
895fn same(a: &WebSocket, b: &WebSocket) -> bool {
896    let a: &worker::web_sys::WebSocket = a.as_ref();
897    let b: &worker::web_sys::WebSocket = b.as_ref();
898    js_sys::Object::is(a.as_ref(), b.as_ref())
899}
900
901impl DurableObject for Archive {
902    fn new(state: State, env: Env) -> Self {
903        let shelf = Shelf { env, sql: state.storage().sql(), running: RefCell::new(HashMap::new()), home: RefCell::new(None) };
904        // An archive that cannot bring its schema up to date must not serve:
905        // it would keep responses in a shape the next start cannot read.
906        shelf.migrate().expect("the archive's migrations run");
907        Archive { state, shelf: Rc::new(shelf) }
908    }
909
910    async fn fetch(&self, mut request: Request) -> worker::Result<Response> {
911        // `/live`: an open page, told what happens while it is open. The
912        // socket is the object's, accepted for hibernation, so an idle page
913        // keeps nothing awake.
914        // The admin backend's door. Only its own Worker sends a request here
915        // (`archive::ADMIN`); this site's Worker has no code that does.
916        if request.path() == ::archive::ADMIN {
917            let ask: ::archive::Admin = serde_json::from_str(&request.text().await?).map_err(|e| worker::Error::from(e.to_string()))?;
918            let moderated = matches!(ask, ::archive::Admin::Moderate { .. });
919            let answered = self.shelf.admin(ask);
920            if moderated {
921                // Open pages get the feeds as they now are.
922                self.activity(None);
923            }
924            return Response::ok(serde_json::to_string(&answered).map_err(|e| worker::Error::from(e.to_string()))?);
925        }
926        if request.headers().get("upgrade")?.is_some_and(|upgrade| upgrade.eq_ignore_ascii_case("websocket")) {
927            let pair = WebSocketPair::new()?;
928            self.state.accept_web_socket(&pair.server);
929            // Where it connected from, as the Worker passed it on: kept on the
930            // socket, which hibernation keeps, and gone when it closes.
931            let header = |name: &str| -> Option<String> {
932                let value = request.headers().get(name).ok().flatten()?;
933                ::archive::percent_decoded(&value)
934            };
935            // Who connected, kept as an event; the socket remembers which,
936            // so the event that says it left can say who.
937            let now = js_sys::Date::now();
938            let came = header(EVENT)
939                .and_then(|event| serde_json::from_str::<::archive::Event>(&event).ok())
940                .and_then(|event| self.shelf.happened(&event.named("live"), now).ok());
941            let seen = Seen { place: Place::new(header(COUNTRY).as_deref(), header(CITY).as_deref()), passed: header(PLACED).is_some(), at_ms: now, event: came };
942            let _ = pair.server.serialize_attachment(&seen);
943            self.send(std::slice::from_ref(&pair.server), &Live::Build(crate::BUILD.to_owned()));
944            // The page that came hears how things stand now, itself counted;
945            // the others hear it with the next change.
946            let open = self.open(None);
947            self.online(std::slice::from_ref(&pair.server), &open);
948            self.soon().await?;
949            // A page from an earlier build hears what this one changed. A
950            // page says its build in `?build=`; one too old to say is earlier.
951            let theirs = request.url()?.query_pairs().find(|(key, _)| key == "build").map(|(_, build)| build.into_owned());
952            if !crate::NOTE.is_empty() && theirs.as_deref() != Some(crate::BUILD) {
953                self.send(std::slice::from_ref(&pair.server), &Live::Toast(crate::view::toast::note(crate::NOTE).into_string()));
954            }
955            return Response::from_websocket(pair.client);
956        }
957        let text = request.text().await?;
958        let ask: Ask = serde_json::from_str(&text).map_err(|e| worker::Error::from(e.to_string()))?;
959        let told = match ask {
960            Ask::Jev { input, wanted, request, pick } => {
961                let body = request.clone();
962                Told::Called(self.once(jev_protocol::ENDPOINT.to_owned(), request, pick, |shelf| jev(shelf, input, wanted, body)).await?)
963            }
964            Ask::Llm { model, request, pick } => {
965                let body = request.clone();
966                Told::Called(self.once(model.clone(), request, pick, |shelf| llm(shelf, model, body)).await?)
967            }
968            Ask::Asked { input, answers, llm, listed, who, browser } => {
969                // A browser asking what it asked before is not news.
970                if self.shelf.note(&input, &answers, llm, listed, who.as_ref(), &browser, js_sys::Date::now())? {
971                    // Only what the feed may show is shown; anything else is
972                    // "someone asked". The owner's say comes before Jev's.
973                    let listed = self.shelf.shown(&input)?;
974                    let shown = listed.then_some((input.as_str(), answers.as_slice()));
975                    self.broadcast(&Live::Toast(crate::view::toast::asked(shown).into_string()));
976                    self.activity(listed.then_some(input.as_str()));
977                }
978                Told::Noted
979            }
980            Ask::Event(event) => {
981                self.shelf.happened(&event, js_sys::Date::now())?;
982                Told::Noted
983            }
984            Ask::Fetched { repo, fetch } => {
985                self.shelf.count(&fetch.counted(&repo));
986                // jevcrates comes along with every project cloned with its
987                // submodules, so its own fetches would toast twice for one
988                // clone; it is counted and not told.
989                if crate::clone::Repo::named(&repo).is_some_and(|repo| repo != crate::clone::Repo::Jevcrates) {
990                    self.broadcast(&Live::Toast(crate::view::toast::fetched(fetch, &repo).into_string()));
991                    self.activity(None);
992                }
993                Told::Noted
994            }
995            Ask::Home => Told::Home(self.shelf.home()?),
996            Ask::Answer { input } => Told::Answer(self.shelf.answer(&input)?),
997            Ask::Older { after } => Told::Older(self.shelf.older(&after)?),
998            Ask::Rate { input, who, vote, browser } => Told::Rating(self.shelf.rate(&input, &who, &browser, vote)?),
999            Ask::Rating { input, who } => Told::Rating(self.shelf.rating(&input, who.as_ref())?),
1000        };
1001        Response::ok(serde_json::to_string(&told).map_err(|e| worker::Error::from(e.to_string()))?)
1002    }
1003
1004    /// A page says nothing the archive listens to.
1005    async fn websocket_message(&self, _socket: WebSocket, _message: WebSocketIncomingMessage) -> worker::Result<()> {
1006        Ok(())
1007    }
1008
1009    async fn websocket_close(&self, socket: WebSocket, _code: usize, _reason: String, _clean: bool) -> worker::Result<()> {
1010        // The page has gone: the event that said it came, again, as `left`.
1011        if let Some(seen) = seen(&socket)
1012            && let Some(came) = seen.event
1013        {
1014            let now = js_sys::Date::now();
1015            let _ = self.shelf.left(came, now - seen.at_ms, now);
1016        }
1017        self.soon().await
1018    }
1019
1020    async fn websocket_error(&self, _socket: WebSocket, _error: worker::Error) -> worker::Result<()> {
1021        self.soon().await
1022    }
1023
1024    /// Pages came or went since the last time everyone was told. A page that
1025    /// came through a Worker that passes no place (during a deploy) is asked
1026    /// to reconnect once the deploy has settled, so it gets one; until then
1027    /// the alarm comes back for it.
1028    async fn alarm(&self) -> worker::Result<Response> {
1029        let now = js_sys::Date::now();
1030        let mut waiting = false;
1031        for socket in self.open(None) {
1032            let seen = seen(&socket);
1033            if Seen::stale(seen.as_ref(), now) {
1034                let _ = socket.close(Some(1012), Some("reconnect to be placed"));
1035            } else if seen.is_some_and(|seen| !seen.passed) {
1036                waiting = true;
1037            }
1038        }
1039        let open = self.open(None);
1040        self.online(&open, &open);
1041        if waiting {
1042            self.state.storage().set_alarm(std::time::Duration::from_millis(::archive::REPLACE_AFTER_MS as u64)).await?;
1043        }
1044        Response::ok("")
1045    }
1046}
1047
1048async fn tell(env: &Env, ask: Ask) -> Result<Told, String> {
1049    crate::object::tell(env, BINDING, NAME, "the archive", &ask).await
1050}
1051
1052/// A call's response: kept, or sent now. Every Jev and LLM call the Worker
1053/// makes goes through here.
1054pub async fn call(env: &Env, ask: Ask) -> Result<Called, String> {
1055    match tell(env, ask).await? {
1056        Told::Called(called) => Ok(called),
1057        _ => Err("the archive answered a call with something else".into()),
1058    }
1059}
1060
1061/// Records that `input` was asked and what Jev said. A failure loses one
1062/// line of the feed and nothing else.
1063pub async fn asked(env: &Env, input: &str, answers: Vec<Answer>, llm: bool, listed: bool, browser: Option<&str>) {
1064    let who = browser.map(|browser| crate::asker(browser, input));
1065    let _ = tell(env, Ask::Asked { input: input.to_owned(), answers, llm, listed, who, browser: browser.unwrap_or_default().to_owned() }).await;
1066}
1067
1068/// What Jev said to `input` the last time it was asked, if it ever was.
1069/// Reading it sends nothing, so a link preview costs nothing.
1070pub async fn answer(env: &Env, input: &str) -> Option<Entry> {
1071    match tell(env, Ask::Answer { input: input.to_owned() }).await {
1072        Ok(Told::Answer(entry)) => entry,
1073        _ => None,
1074    }
1075}
1076
1077/// The feed and the numbers, for the home page. `None` if the archive
1078/// cannot be asked.
1079/// Someone cloned or pulled `repo`.
1080pub async fn fetched(env: &Env, repo: &str, fetch: Fetch) {
1081    let _ = tell(env, Ask::Fetched { repo: repo.to_owned(), fetch }).await;
1082}
1083
1084/// `/live`: the request, handed to the archive, which keeps the socket.
1085/// How long after pages come or go every page is told, at most once.
1086const ONLINE_EVERY_MS: u64 = 1000;
1087
1088/// The headers the Worker passes a page's place to the archive in.
1089pub const COUNTRY: &str = "x-lmjtfy-country";
1090pub const CITY: &str = "x-lmjtfy-city";
1091/// Set by a Worker that passes places on, with or without one to pass.
1092pub const PLACED: &str = "x-lmjtfy-placed";
1093/// The request as an `archive::Event`, JSON and percent-encoded: who
1094/// connected. Set by the Worker, over any the page sent.
1095pub const EVENT: &str = "x-lmjtfy-event";
1096
1097/// Keeps an event. A failure loses that one row and nothing else.
1098pub async fn event(env: &Env, event: ::archive::Event) {
1099    let _ = tell(env, Ask::Event(Box::new(event))).await;
1100}
1101
1102pub async fn live(env: &Env, request: Request) -> worker::Result<Response> {
1103    let stub = env.durable_object(BINDING)?.id_from_name(NAME)?.get_stub()?;
1104    stub.fetch_with_request(request).await
1105}
1106
1107/// `who`'s vote on Jev's answer to `input`, and the votes after it. `None`
1108/// if there is nothing to vote on or the archive cannot be asked.
1109pub async fn rate(env: &Env, input: &str, browser: &str, vote: Vote) -> Option<Rating> {
1110    match tell(env, Ask::Rate { input: input.to_owned(), who: crate::asker(browser, input), vote, browser: browser.to_owned() }).await {
1111        Ok(Told::Rating(rating)) => rating,
1112        _ => None,
1113    }
1114}
1115
1116/// The votes on Jev's answer to `input`, and `who`'s.
1117pub async fn rating(env: &Env, input: &str, who: Option<Asker>) -> Option<Rating> {
1118    match tell(env, Ask::Rating { input: input.to_owned(), who }).await {
1119        Ok(Told::Rating(rating)) => rating,
1120        _ => None,
1121    }
1122}
1123
1124/// The page of the feed after `after`. `None` if the archive cannot be
1125/// asked.
1126pub async fn older(env: &Env, after: Cursor) -> Option<Vec<Entry>> {
1127    match tell(env, Ask::Older { after }).await {
1128        Ok(Told::Older(entries)) => Some(entries),
1129        _ => None,
1130    }
1131}
1132
1133pub async fn home(env: &Env) -> Option<Home> {
1134    match tell(env, Ask::Home).await {
1135        Ok(Told::Home(home)) => Some(home),
1136        _ => None,
1137    }
1138}
1139
1140#[cfg(test)]
1141mod tests {
1142    use super::MIGRATIONS;
1143
1144    /// Every `(table, column)` the migrations make: `CREATE TABLE` columns
1145    /// and `ALTER TABLE ... ADD COLUMN`.
1146    fn columns() -> Vec<(String, String)> {
1147        let mut found = Vec::new();
1148        for statement in MIGRATIONS.iter().flat_map(|step| step.iter()) {
1149            let words: Vec<&str> = statement.split_whitespace().collect();
1150            // A table a later step drops is not in the archive any more.
1151            if words.first() == Some(&"DROP") {
1152                let dropped = words.last().copied().unwrap_or_default();
1153                found.retain(|(table, _)| table != dropped);
1154                continue;
1155            }
1156            if let Some(at) = words.iter().position(|word| *word == "TABLE") {
1157                let table = words[at..].iter().find(|word| !matches!(**word, "TABLE" | "IF" | "NOT" | "EXISTS")).copied();
1158                let Some(table) = table.map(|table| table.trim_end_matches('(')) else { continue };
1159                if words.first() == Some(&"CREATE") {
1160                    let body = statement.split_once('(').map(|(_, body)| body).unwrap_or_default();
1161                    // Columns until the table's own key clause, whose list of
1162                    // names is not columns.
1163                    for line in body.split(',') {
1164                        let column = line.split_whitespace().next().unwrap_or_default();
1165                        if column == "PRIMARY" {
1166                            break;
1167                        }
1168                        if !column.is_empty() {
1169                            found.push((table.to_owned(), column.to_owned()));
1170                        }
1171                    }
1172                } else if let Some(at) = words.iter().position(|word| *word == "COLUMN") {
1173                    found.push((table.to_owned(), words[at + 1].to_owned()));
1174                }
1175            }
1176        }
1177        found
1178    }
1179
1180    #[test]
1181    fn the_events_table_is_the_shape_of_an_event() {
1182        let made: Vec<String> = columns().into_iter().filter(|(table, _)| table == "events").map(|(_, column)| column).collect();
1183        let mut wanted = vec!["id", "at_ms"];
1184        wanted.extend(::archive::Event::COLUMNS);
1185        assert_eq!(made, wanted);
1186    }
1187
1188    #[test]
1189    fn the_readme_diagram_has_every_column_the_migrations_make() {
1190        let readme = include_str!("../README.md");
1191        let diagram = readme.split("### The archive").nth(1).and_then(|rest| rest.split("```").nth(1)).expect("the archive's erDiagram");
1192        let columns = columns();
1193        assert!(columns.len() >= 15, "{columns:?}");
1194        for (table, column) in columns {
1195            let entity = diagram.split(&format!("  {table} {{")).nth(1).and_then(|rest| rest.split('}').next());
1196            let entity = entity.unwrap_or_else(|| panic!("the diagram has no `{table}`"));
1197            assert!(
1198                entity.lines().any(|line| line.split_whitespace().nth(1) == Some(column.as_str())),
1199                "the diagram's `{table}` lacks `{column}`"
1200            );
1201        }
1202        assert!(!diagram.contains("  calls {"), "`calls` became `versions` (MIGRATIONS step 7)");
1203    }
1204}