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}