lmjtfy.git / apps / lmjtfy / src / archive.rs

The archive, kept where every isolate sees the same rows: one Durable Object with a SQL table of every request that was answered, and one of the questions people asked. The messages are archive's.

The same request is never sent to the same model twice (the user, 2026-10-02), unless a visitor asks for it to be (Pick::Fresh, the page's ↻), and then every response it got is kept, numbered. Three things make that hold:

  • versions has (sent_to, request, version) as its primary key: the exact body, where it went, and which time. No hash stands in for the body, so two different requests cannot be taken for one.
  • A call is looked up before anything is spent, and its response is kept before anyone is told.
  • The archive makes the call itself. Identical asks that arrive while it is on the wire wait on that one call (running) and are told what it was told. The call is handed to wait_until, so it finishes and is kept even if every visitor waiting on it has left.

A call with no response (an error, a timeout) keeps nothing: there is nothing to return next time, so the next ask sends it.

23use std::cell::RefCell;
24use std::collections::HashMap;
25use std::rc::Rc;
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};

The binding's name in wrangler.toml.

42const BINDING: &str = "ARCHIVE";

One object for the whole Worker: a request anyone had answered is answered for everyone.

45const NAME: &str = ::archive::OBJECT;

The schema, as the steps that made it, in order. A step runs once: the object records how many it has run (migrated) and runs the rest when it starts. Never edit a step that has shipped; add one.

calls.sent_to is Jev's endpoint or a Workers AI model id. (A Jev body 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];

The step (from 1) that gave askers and ratings their browser, after which the rows already there are named from events, once.

203const ASKERS_NAMED: usize = 10;

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}
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}

The repository whose clones and pulls the pages show. jevcrates is counted too, but a clone of lmjtfy fetches it, so it is not shown twice.

248const CODE: &str = "lmjtfy.git";

The names in counts.

251const SENT: &str = "sent";
252const KEPT: &str = "kept";

A call on the wire, which every identical ask waits on.

255type Running = Shared<LocalBoxFuture<'static, Called>>;

What a call needs after the request that started it has gone: shared, so the call can outlive that request.

259struct Shelf {
260    env: Env,
261    sql: SqlStorage,
262    running: RefCell<HashMap<(String, String), Running>>,

What the home page shows, kept from the last time it was read until something it is made of changes (changed). Every home page asks for it, and its numbers are a pass over every question.

266    home: RefCell<Option<Home>>,
267}
269impl Shelf {

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    }

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    }

Keeps a response as the request's newest version, and says which it is. A plain INSERT: two calls numbering the same version would mean the archive sent twice what it meant to send once, and that should 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    }
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    }

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    }

A page that connected (the event in row came) has gone: the same 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    }

Whether the feed may show input: the owner's say if they have one, 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    }

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    }

Fills in whose each old askers and ratings row is, where an event says which browser asked or voted on that question: the browser's id and the question hash to the row's who. A row from 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    }
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    }

Adds one to a tally. A failure loses one from a number on the home 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    }
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    }

Records what Jev said to input, and whether this is a new asking of it: by a browser that has not asked it before, or by one that gave no id. Only a new asking counts and moves it up the feed; Jev's answer 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    }

Which answer to input a vote is on: the SHA-256 of the answers as 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    }

who's vote on the answer to input as kept now; the vote it already 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    }
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    }

What Jev said to input, listed or not: it is for whoever holds the 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    }

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    }
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    }

What the home page is made of has changed: read it again next time.

620    fn changed(&self) {
621        self.home.replace(None);
622    }
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}

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}

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}
751impl Archive {

The response to request, sent at most once: the kept one, the one a call already on the wire is about to get, or the one call gets now.

From the lookup to the insert into running nothing awaits, so no other ask can run in between and start the same call.

pick says which kept response is wanted, and whether to send it 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}
814impl Archive {

Sends live to every open page. A page that cannot be told has gone, and its close is on its way.

817    fn broadcast(&self, live: &Live) {
818        self.send(&self.state.get_websockets(), live);
819    }
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    }

Sends every open page the activity feeds as they now stand: asked at the top of the feed, if it was a question the feed may show, and 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    }

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    }

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    }

Everyone hears that pages came or went once, ONLINE_EVERY_MS after the first change, however many changed in between. A deploy closes every socket and they all reconnect within a second or two; telling 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}

What is kept on a page's socket, if anything.

890fn seen(socket: &WebSocket) -> Option<Seen> {
891    socket.deserialize_attachment::<Seen>().ok().flatten()
892}

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}
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    }

A page says nothing the archive listens to.

1005    async fn websocket_message(&self, _socket: WebSocket, _message: WebSocketIncomingMessage) -> worker::Result<()> {
1006        Ok(())
1007    }
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    }

Pages came or went since the last time everyone was told. A page that came through a Worker that passes no place (during a deploy) is asked to reconnect once the deploy has settled, so it gets one; until then 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}
1048async fn tell(env: &Env, ask: Ask) -> Result<Told, String> {
1049    crate::object::tell(env, BINDING, NAME, "the archive", &ask).await
1050}

A call's response: kept, or sent now. Every Jev and LLM call the Worker 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}

Records that input was asked and what Jev said. A failure loses one 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}

What Jev said to input the last time it was asked, if it ever was. 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}

The feed and the numbers, for the home page. None if the archive cannot be asked. 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}

/live: the request, handed to the archive, which keeps the socket. How long after pages come or go every page is told, at most once.

1086const ONLINE_EVERY_MS: u64 = 1000;

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";

Set by a Worker that passes places on, with or without one to pass.

1092pub const PLACED: &str = "x-lmjtfy-placed";

The request as an archive::Event, JSON and percent-encoded: who connected. Set by the Worker, over any the page sent.

1095pub const EVENT: &str = "x-lmjtfy-event";

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}
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}

who's vote on Jev's answer to input, and the votes after it. None 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}

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}

The page of the feed after after. None if the archive cannot be 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}
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;

Every (table, column) the migrations make: CREATE TABLE columns 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    }
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}