Something went wrong. Try again.
Web frontend and supporting services for lance.blue
Something went wrong. Try again.
48 kB · 1148 lines
Rust
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149//! Match lifecycle: create a manifest, start the Fargate task, answer the//! waiting screen's polls, take arena's callbacks and uploads. The status//! vocabulary is the matches table's `status` column://!//! lobby -> deploying -> (deleted; a new row starts its own chain below)//! starting -> running -> ready -> over//! \-> failed//!//! `running` means the task reports RUNNING; `ready` means a browser sent to//! the match can play - what the waiting screen waits for.//!//! `ready` arrives two ways. arena posts it when both ports inside the//! container are open, which is the fast path and the accurate one; the HTTP//! probe in `advance` is the fallback for a container that cannot reach us, and//! only sees Suramadu's port. Whichever lands first wins, and they agree.//!//! `lobby` and `deploying` never share a row with `starting`/`running`/etc.//! `lobby.rs` opens a row directly at `lobby`, and its live edits//! (`set_scenario`, `claim_seat`, `assign_bot`) never call anything here.//! Resolving one into a real match is `routes.rs`'s `deploy_lobby`: it flips//! the row to `deploying` with `Db::claim_deploy` (an atomic//! `UPDATE ... WHERE status = 'lobby'`, guarding against two callers//! deploying the same lobby at once — the same shape `Db::claim_match_player`//! uses for a seat), then calls `Matches::start` exactly the way//! `create_match` does. `start` mints a brand-new id and inserts a brand-new//! row at `starting` — it never reuses the lobby row's own id, so the lobby//! and the match it becomes are two different rows for their whole//! overlapping moment. Once `start` succeeds, `deploy_lobby` deletes the//! lobby row and its seats (`Db::delete_lobby`) and broadcasts//! `LobbyEvent::Deployed` with the new id, so every socket still on the//! lobby's channel — not just the HTTP caller — learns where to go. A//! failed `start` leaves the row parked at `deploying` on purpose: a dead//! end for a human to notice, not seat claims silently erased for someone//! to redo from scratch.
mod aws;pub mod daily;pub mod lobby;pub mod manifest;pub mod scenarios;pub mod sizing;
use std::sync::Arc;
use base64::Engine;use base64::engine::general_purpose::URL_SAFE_NO_PAD;use rand::RngCore;
use crate::config::{Config, MatchConfig};use crate::db::{Db, MatchPlayer, MatchRow};
/// One seat as the routes hand it to the launcher: identity settled, camo/// still by reference. Staging turns this into a manifest player.pub struct NewPlayer { /// A faction name from the scenario's `Factions=` line. pub slot: String, pub bot: bool, /// Humans always have one; bots never do. pub did: Option<String>, /// The verified handle. Bots are labelled by the launcher itself. pub handle: Option<String>, /// A `blue.lance.camo` record key in the seat-holder's own repository. /// Human seats only. pub camo_rkey: Option<String>, /// A rolled picture from the browser, for a seat with no repository to /// read from. Bot seats only. pub camo_png: Option<Vec<u8>>, /// The MegaMek team this slot fights on, defaulted from the scenario. /// /// `None` for a scenario that declares none — the daily challenges carry /// their own teams in their files, and the catalog's come from the same /// place. A lobby team control would set this; until there is one, the /// file is the only thing with an opinion. pub team: Option<i32>,}
pub struct Matches { aws: aws::Aws, config: MatchConfig, public_url: String, db: Db, http: reqwest::Client, /// Per match, when ECS was last asked about it and whether its task /// timeline has been logged. Keyed by match id and swept in /// `probe_window_open`; a match that never polls leaves nothing behind /// because nothing was ever inserted. probes: std::sync::Mutex<std::collections::HashMap<String, ProbeMark>>,}
struct ProbeMark { at: std::time::Instant, timing_logged: bool,}
pub struct StartFailed;
#[derive(serde::Serialize)]pub struct MatchView { pub status: String, #[serde(skip_serializing_if = "Option::is_none")] pub url: Option<String>,}
/// The result document a finished match uploaded: MegaMek's outcome, the teams,/// and each unit's state, extracted by arena rather than written by the engine.////// Under `derived/` because it is ours: a claim we made out of what MegaMek/// wrote, not one of its own files. The prefixes say what a file *is* — what/// the match was given, what MegaMek wrote, what we made of it — rather than/// which upload key sent it. See `docs/match-artifacts.md`.pub const RESULT_KEY: &str = "derived/result.json";
/// The board as it stood at the end, which is the picture a shared match is/// drawn on.////// Under `derived/socials/` rather than `raw/`: the bytes are MegaMek's and/// unaltered, but choosing this frame out of the thirty-one it rendered is/// arena's act, and what the file is for is somebody looking at it. The card/// itself — the 1200×630 with the forces lettered over it — is drawn here per/// request by `card.rs`, and would be the next thing to live beside it.pub const BOARD_KEY: &str = "derived/socials/share-card.png";
impl Matches { /// None when MATCH_CLUSTER is unset - the routes then answer 503. pub async fn new(config: &Config, db: Db) -> Option<Matches> { let match_config = config.matches.clone()?; Some(Matches { aws: aws::Aws::new().await, config: match_config, public_url: config .public_url .clone() .unwrap_or_else(|| format!("http://127.0.0.1:{}", config.bind_addr.port())), db, http: reqwest::Client::new(), probes: std::sync::Mutex::new(std::collections::HashMap::new()), }) }
/// Settle a daily-challenge attempt whose match has ended. /// /// Two outcomes and no third: the match produced a result, in which case /// the attempt is scored; or it did not, in which case the attempt is /// given back. A match that crashed, was reaped, or ended without reaching /// victory costs the player nothing — an attempt lost to our own /// infrastructure is a bug they should not pay for, and "one attempt" only /// means anything if it is one attempt they actually got. /// /// Safe to call more than once and for matches that are not attempts at /// all: an already-scored row is left alone, and a match with no attempt /// behind it is a no-op. The result upload can land after the event that /// says the match ended, so this is called again when the card asks for /// the player's score. /// The result document a finished match uploaded, parsed. /// /// `None` for a match that has not produced one - not uploaded yet, or /// never will be. Which of those it is belongs to the caller, who has the /// match's status; nothing here guesses. /// Where this API answers, for the addresses it has to write into a /// document rather than let a caller construct - the launch manifest's /// upload targets, and the share page's own canonical URL. pub fn public_url(&self) -> &str { &self.public_url }
/// One artifact of a finished match, as bytes. /// /// Only for artifacts this service is willing to hand out; the caller /// names the key, and today that is the share card and nothing else. /// The camo this service staged for one slot, for anything that draws the /// match afterwards. `None` for a seat that had none, and for a match old /// enough that its staged files have aged out. /// /// The slot is checked before it is spelled into a key. It reaches here /// from the result document, which is uploaded by the match container - /// so however plain the catalog's own slot names are, this is a string /// from outside deciding what this service reads. Slot names are faction /// names from a scenario file: words, digits and the odd dash. Anything /// else is not a slot we staged a camo under, and is refused rather than /// resolved. /// One force's colours. `Ok(None)` is a slot with none staged, which is /// ordinary; `Err` is not knowing, which the caller must not cache as if /// it were. pub async fn camo( &self, match_id: &str, slot: &str, ) -> Result<Option<Vec<u8>>, aws::AwsFailed> { let plain = !slot.is_empty() && slot.len() <= 64 && slot .chars() .all(|c| c.is_ascii_alphanumeric() || matches!(c, ' ' | '.' | '_' | '-')); if !plain { tracing::warn!(match_id, slot, "camo: slot name refused, not read"); return Ok(None); } let found = self .aws .get_camo( &self.config.artifacts_bucket, match_id, &format!("camo-{slot}.png"), ) .await; if found.is_err() { tracing::warn!(match_id, slot, "camo: read failed"); } found }
/// Whether the match left this artifact behind. `false` covers both "no /// such file" and "could not tell"; the caller shortens its caching /// either way, because a board can still arrive from the exit pass. pub async fn has_artifact(&self, match_id: &str, name: &str) -> bool { match self .aws .has_artifact(&self.config.artifacts_bucket, match_id, name) .await { Ok(found) => found, Err(_) => { tracing::warn!(match_id, name, "artifact head failed"); false } } }
pub async fn artifact(&self, match_id: &str, name: &str) -> Option<Vec<u8>> { self.aws .get_artifact(&self.config.artifacts_bucket, match_id, name) .await .ok() }
/// The board render a shared match is drawn on. `None` for a match that /// never uploaded one, and for one old enough that its artifacts have aged /// out. pub async fn board(&self, match_id: &str) -> Option<Vec<u8>> { self.artifact(match_id, BOARD_KEY).await }
/// Whether this match has a board render. pub async fn has_board(&self, match_id: &str) -> bool { self.has_artifact(match_id, BOARD_KEY).await }
/// The result document's bytes. Callers that want it parsed use /// [`Self::result`]. async fn result_bytes(&self, match_id: &str) -> Option<Vec<u8>> { self.artifact(match_id, RESULT_KEY).await }
pub async fn result(&self, match_id: &str) -> Option<serde_json::Value> { let bytes = self.result_bytes(match_id).await?; match serde_json::from_slice(&bytes) { Ok(value) => Some(value), Err(_) => { // The artifact is there and is not JSON, which no arena // produces. Worth a line, because the alternative reading - // "the match never uploaded one" - is what the caller would // otherwise show a player. tracing::error!(match_id, "match: the result document is not JSON"); None } } }
pub async fn settle_daily(&self, match_id: &str) { let Ok(Some(attempt)) = self.db.daily_attempt_for_match(match_id).await else { return; }; if attempt.score.is_some() { return; } let bytes = match self.result_bytes(match_id).await { Some(bytes) => bytes, None => { // No result yet. Whether that is "not uploaded yet" or "never // will be" is decided by the match's own status, not here: // releasing on a read that raced the upload would hand back an // attempt that was actually played. if let Ok(Some(row)) = self.db.get_match(match_id).await && row.status == "failed" && self .db .release_daily_attempt(match_id) .await .unwrap_or(false) { tracing::info!( match_id, slug = attempt.slug, "daily: match failed, attempt released" ); } return; } };
let Ok(result) = serde_json::from_slice::<serde_json::Value>(&bytes) else { tracing::error!(match_id, "daily: the result document is not JSON"); return; }; let Some((mine, theirs)) = daily::sides_from_result(&result) else { // A result from an arena that does not report battle value. There // is nothing to score against, and inventing a denominator would // rank a player on a number nobody measured — so the attempt goes // back rather than being recorded as a defeat. tracing::warn!( match_id, slug = attempt.slug, "daily: the result carries no battle value; attempt released" ); let _ = self.db.release_daily_attempt(match_id).await; return; }; let Some(score) = daily::score(mine, theirs) else { tracing::warn!(match_id, slug = attempt.slug, "daily: unscorable result"); let _ = self.db.release_daily_attempt(match_id).await; return; }; if let Err(e) = self .db .score_daily_attempt( match_id, score, (mine.bv_start, mine.bv_left), (theirs.bv_start, theirs.bv_left), ) .await { tracing::error!(match_id, ?e, "daily: score not stored"); return; } tracing::info!( match_id, slug = attempt.slug, score, bv_left = mine.bv_left, bv_start = mine.bv_start, foe_bv_left = theirs.bv_left, foe_bv_start = theirs.bv_start, "daily: scored" ); }
/// The launch configuration, for the legacy request shape whose scenario /// and slots come from the environment rather than the catalog. pub fn match_config(&self) -> &MatchConfig { &self.config }
/// Launch a match. /// /// `scenario` is where the fight comes from and `label` is what the match /// row records it as — a library path for a catalog match, `daily:<slug>` /// for a daily challenge, whose file is in this binary rather than in any /// image. pub async fn start( &self, atproto: &crate::atproto::Atproto, owner_did: &str, scenario: manifest::ScenarioSource<'_>, label: &str, players: Vec<NewPlayer>, ) -> Result<String, StartFailed> { let id = uuid::Uuid::new_v4().to_string(); let mut token_bytes = [0u8; 24]; rand::rngs::OsRng.fill_bytes(&mut token_bytes); let token = URL_SAFE_NO_PAD.encode(token_bytes);
// Camo before the manifest, because it goes in it; after nothing, // because it is allowed to fail. A camo that cannot be fetched or // cannot be made safe is a seat in default colours, which is what a // player who never made one gets and is not worth refusing to start // over. let mut manifest_players = Vec::with_capacity(players.len()); let mut seats = Vec::with_capacity(players.len()); for player in &players { let camo = if player.bot { match &player.camo_png { Some(png) => self.stage_rolled_camo(&id, &player.slot, png).await, None => None, } } else { match (&player.did, &player.camo_rkey) { (Some(did), Some(rkey)) => { self.stage_camo(atproto, &id, &player.slot, did, rkey).await } _ => None, } }; manifest_players.push(manifest::Player { slot: player.slot.clone(), bot: player.bot, did: player.did.clone(), handle: player.handle.clone(), team: player.team, camo, }); seats.push(MatchPlayer { slot: player.slot.clone(), control: if player.bot { "bot" } else { "human" }.to_owned(), did: player.did.clone(), // Resolved by now: whatever the lobby said, or what the // scenario said when it said nothing. The seats table // records what the match was launched with, not what was // overridden to get there. team: player.team, // The seats table stores what the seat fought as, so the bot // label lands here too — the matches list reads only this. handle: if player.bot { Some(manifest::BOT_HANDLE.to_owned()) } else { player.handle.clone() }, }); }
let manifest = manifest::build(&manifest::ManifestInput { match_id: &id, token: &token, public_url: &self.public_url, scenario, players: &manifest_players, });
let key = self .aws .put_manifest(&self.config.artifacts_bucket, &id, &manifest) .await .map_err(|e| self.log_fail(&id, e))?; let manifest_url = self .aws .presign_get(&self.config.artifacts_bucket, &key) .await .map_err(|e| self.log_fail(&id, e))?;
self.db .insert_match(&MatchRow { id: id.clone(), owner_did: owner_did.to_owned(), status: "starting".into(), token, task_arn: None, task_ip: None, created_at: 0, // the insert uses the column default started_at: None, ended_at: None, scenario: Some(label.to_owned()), }) .await .map_err(|e| { tracing::error!(match_id = id, ?e, "match: row insert failed"); StartFailed })?; // Fatal, unlike camo: the seat rows are how an invited player's list // shows the match and how the proxy admits them, so a match without // them is a lobby half its players cannot reach. self.db .insert_match_players(&id, &seats) .await .map_err(|e| { tracing::error!(match_id = id, ?e, "match: seat insert failed"); StartFailed })?;
// Sized by what is in it: every human is a client JVM, every bot is // Princess load in the host, and a big scenario multiplies both. An // unknown scenario (the legacy env-configured one) counts zero units // and still sizes correctly off its seats. let humans = players.iter().filter(|p| !p.bot).count() as u32; let bots = players.len() as u32 - humans; // A catalog scenario's unit count is in the catalog; a carried one's // is in the file, which is right here. Both matter: unit count is what // the host and Princess spend CPU on, and sizing a daily challenge as // if it were empty would under-provision every one of them. let units = match scenario { manifest::ScenarioSource::Library(path) => { scenarios::find(path).map_or(0, |s| s.unit_total()) } manifest::ScenarioSource::Carried(text) => headquarters_mms::factions(text) .iter() .map(|f| f.units) .sum::<usize>() as u32, }; let size = sizing::task_size(humans, bots, units); tracing::info!( match_id = id, cpu = size.cpu, memory = size.memory, "match: task sized" );
let task_arn = self .aws .run_task(&self.config, &manifest_url, &size) .await .map_err(|e| self.log_fail(&id, e))?; if let Err(e) = self.db.set_match_task(&id, &task_arn).await { // The Fargate task is already running and billed; without the // arn stored, advance() can never find it, and this line is the // only place it survives to stop by hand. tracing::error!( match_id = id, task_arn, ?e, "match: task arn not stored, task now unreachable" ); }
tracing::info!( match_id = id, owner = owner_did, scenario = label, seats = players.len(), humans = players.iter().filter(|p| !p.bot).count(), "match: started" ); Ok(id) }
/// Read a published camo, make it safe, and put it where the match can /// fetch it. /// /// The re-encode in `stage_png` is the security boundary, and it is there /// rather than at upload because a player owns their repository: they can /// write any bytes into this collection with any client, and headquarters /// is not in that path. This *is* the path where those bytes are handed /// to other players' browsers to decode, so it is the last place they can /// be rewritten from a pixel buffer (`camo.rs` states the same rule from /// the other end, as does arena's `container/init/30-assets.sh`). /// /// Every failure returns None and is logged. None of them are worth /// failing a match: the camo is cosmetic, and the player would rather /// play in grey than not play. async fn stage_camo( &self, atproto: &crate::atproto::Atproto, match_id: &str, slot: &str, did: &str, rkey: &str, ) -> Option<manifest::Camo> { let published = match atproto.published_camo(did, rkey).await { Ok(published) => published, Err(e) => { tracing::warn!(match_id, slot, ?e, "camo: not staged, the read failed"); return None; } };
// Keyed by slot, because every seat can carry one and they must not // overwrite each other. Slot names come from the catalog, so they are // plain ASCII words (scenarios.rs holds the test). let (url, sha256) = self .stage_png(match_id, &format!("camo-{slot}.png"), &published.png) .await?; tracing::info!(match_id, slot, rkey, sha256, "camo: staged for the match"); Some(manifest::Camo { url, cid: Some(published.cid), sha256, }) }
/// A bot seat's camo, rolled by the browser for this match alone. /// /// Through the same re-encode as the player's, and for the same reason: /// these bytes are decoded by a browser during the match, and here they /// arrived from one. Nothing is published and nothing is kept, so there /// is no record and no `cid` to carry - the digest is the only thing /// arena checks anyway. async fn stage_rolled_camo( &self, match_id: &str, slot: &str, png: &[u8], ) -> Option<manifest::Camo> { let (url, sha256) = self .stage_png(match_id, &format!("camo-{slot}.png"), png) .await?; tracing::info!(match_id, slot, sha256, "camo: staged for the bot"); Some(manifest::Camo { url, cid: None, sha256, }) }
/// Re-encode, upload, presign. The digest is over what arena will /// actually download, which is the re-encoded file and not what arrived. async fn stage_png( &self, match_id: &str, name: &str, bytes: &[u8], ) -> Option<(String, String)> { let png = match crate::camo::reencode(bytes) { Ok(png) => png, Err(why) => { // Nobody is waiting on this message; it is here so a camo // that silently stops appearing has one line saying why. tracing::warn!( match_id, name, why, "camo: not staged, the image was rejected" ); return None; } };
let sha256 = { use sha2::Digest; format!("{:x}", sha2::Sha256::digest(&png)) };
let key = match self .aws .put_camo(&self.config.artifacts_bucket, match_id, name, png) .await { Ok(key) => key, Err(e) => { tracing::error!(match_id, name, stage = e.0, "camo: not staged"); return None; } }; match self .aws .presign_get(&self.config.artifacts_bucket, &key) .await { Ok(url) => Some((url, sha256)), Err(e) => { tracing::error!(match_id, name, stage = e.0, "camo: not staged"); None } } }
fn log_fail(&self, match_id: &str, e: aws::AwsFailed) -> StartFailed { tracing::error!(match_id, stage = e.0, "match: start failed"); StartFailed }
/// Drives the row toward `ready` on each poll from the waiting screen, and /// past it once the task the match was is gone. pub async fn poll(&self, row: &MatchRow) -> MatchView { let status = match row.status.as_str() { "starting" | "running" => self.advance(row).await, "ready" => self.outlived(row).await, other => other.to_owned(), }; MatchView { // A readiness marker more than an address: the routes replace it // with the caller's own client path - /match/{id}/play/<did>/ - // or clear it for a caller with no fighting seat. url: (status == "ready").then(|| format!("/match/{}", row.id)), status, } }
/// Whether a match that reached `ready` is still a match anyone can join. /// /// `ready` used to be the end of the line here: nothing polled a row once /// it got there, on the understanding that arena reports its own ending /// with `game-over` or `exiting`. It does, when it can. A container that is /// killed, that runs out of memory, or that loses the network before it /// finalises reports nothing at all, and the row stayed `ready` for good — /// so the matches list kept offering an Open link to a task that no longer /// exists, and the proxy answered it with a 502. /// /// Only the task's own existence is asked about. `advance` probes the /// container's web server, which is the right question on the way up and /// the wrong one here: a live match whose probe times out would be reported /// as no longer ready, and a player would be told a match they are sitting /// in has ended. /// /// `over` and not `failed`. The match reached ready, so it ran and someone /// could get into it; what is missing is only the event that would have /// said how it ended. `failed` in this repo means a match that never /// started, and using it here would file every dropped ending under /// failure in the player's history. /// /// What this cannot reach: ECS keeps a stopped task describable for about /// an hour and then forgets it, and a forgotten task comes back as /// `AwsFailed` rather than as a status. A row that nobody polls inside that /// window therefore stays `ready` for good, and the rows already stuck that /// way when this shipped are past it. Expiring rows by age is the other /// half, and is still open in plan/match-lifecycle.md. async fn outlived(&self, row: &MatchRow) -> String { let Some(task_arn) = row.task_arn.as_deref() else { return row.status.clone(); }; // The same one-second window `advance` uses, and for the same reason: // the matches list drives this now, and a player refreshing it must not // turn into a DescribeTasks per refresh per match. if !self.probe_window_open(&row.id) { return row.status.clone(); } // Unreachable ECS is not evidence the match ended. Leaving the row // alone means the next poll asks again. let Ok(task) = self.aws.task_status(&self.config.cluster, task_arn).await else { return row.status.clone(); }; if !task_is_gone(&task.status) { return row.status.clone(); } if let Err(e) = self.db.set_match_status(&row.id, "over").await { // Unwritten, so this re-enters on every later poll: the line is // what makes that converge on one diagnosis rather than a stream of // identical warnings. tracing::error!( match_id = row.id, status = "over", ?e, "match: status not stored" ); } // Info and not warn: a task ending is what every match does. What is // worth seeing is that we learned it from ECS rather than from arena, // which means the container's own ending event never arrived. tracing::info!( match_id = row.id, task_status = task.status, "match: over, found by polling a ready row" ); "over".to_owned() }
async fn advance(&self, row: &MatchRow) -> String { let Some(task_arn) = row.task_arn.as_deref() else { return row.status.clone(); }; // The waiting screen polls several times a second so that a `ready` it // has already been told about reaches the player quickly. That must not // turn into several DescribeTasks a second per waiting player: the ECS // API is throttled account-wide, and being throttled would slow down // every match start in the account to speed up one poll. The stored row // is what a poll inside the window answers from. if !self.probe_window_open(&row.id) { return row.status.clone(); } let Ok(task) = self.aws.task_status(&self.config.cluster, task_arn).await else { return row.status.clone(); }; let aws::TaskState { status: task_status, ip, timing, } = task;
if task_is_gone(&task_status) { // Stopped before ever becoming ready is a failure. A row that got // as far as ready is `outlived`'s, and ends as `over` instead. if let Err(e) = self.db.set_match_status(&row.id, "failed").await { // Unwritten, so this branch re-enters on every later // poll: this line is what makes that converge on one // diagnosis instead of a stream of repeated warnings. tracing::error!( match_id = row.id, status = "failed", ?e, "match: status not stored" ); } tracing::warn!(match_id = row.id, "match: task stopped before ready"); return "failed".to_owned(); } if task_status != "RUNNING" { return row.status.clone(); // PROVISIONING / PENDING }
// Logged at RUNNING rather than at ready, because arena's own callback // can carry the match to ready without this path ever running again - // and everything in this timeline happens before the container exists, // so nothing inside it can report the pull. if self.first_timing(&row.id) { tracing::info!( match_id = row.id, provision_ms = timing.provision_ms, pull_ms = timing.pull_ms, start_ms = timing.start_ms, total_ms = timing.total_ms, "match: task timeline" ); }
let Some(ip) = ip else { return row.status.clone(); }; if row.task_ip.is_none() && let Err(e) = self.db.set_match_ip(&row.id, &ip).await { tracing::error!(match_id = row.id, ip, ?e, "match: ip not stored"); } if row.status == "starting" && let Err(e) = self.db.set_match_status(&row.id, "running").await { tracing::error!( match_id = row.id, status = "running", ?e, "match: status not stored" ); }
// Ready means the web server inside actually answers. let probe = self .http .get(format!("http://{ip}:{}/", crate::proxy::ARENA_PORT)) .timeout(std::time::Duration::from_secs(3)) .send() .await; if probe.is_ok_and(|response| response.status().is_success()) { let _ = self.db.set_match_status(&row.id, "ready").await; tracing::info!(match_id = row.id, via = "probe", "match: ready"); "ready".to_owned() } else { "running".to_owned() } }
/// True at most once per second per match, which is what keeps the poll /// rate and the DescribeTasks rate separate. Entries are dropped once they /// are older than any match could still be starting. fn probe_window_open(&self, match_id: &str) -> bool { const WINDOW: std::time::Duration = std::time::Duration::from_secs(1); const KEEP: std::time::Duration = std::time::Duration::from_secs(900); let now = std::time::Instant::now(); let mut probes = self.probes.lock().expect("probe map is never poisoned"); if let Some(mark) = probes.get(match_id) && now.duration_since(mark.at) < WINDOW { return false; } probes.retain(|_, mark| now.duration_since(mark.at) < KEEP); probes .entry(match_id.to_owned()) .and_modify(|mark| mark.at = now) .or_insert(ProbeMark { at: now, timing_logged: false, }); true }
/// True once per match. The task timeline is a measurement, not a status, /// so it is worth exactly one line - and which poll happens to observe /// RUNNING first is arbitrary. fn first_timing(&self, match_id: &str) -> bool { let mut probes = self.probes.lock().expect("probe map is never poisoned"); match probes.get_mut(match_id) { Some(mark) if !mark.timing_logged => { mark.timing_logged = true; true } _ => false, } }
/// A presigned PUT for one artifact key, minted on request. pub async fn presign_artifact_put(&self, match_id: &str, name: &str) -> Option<String> { match self .aws .presign_artifact_put( &self.config.artifacts_bucket, match_id, name, std::time::Duration::from_secs(ARTIFACT_PUT_EXPIRY_SECS), ) .await { Ok(url) => Some(url), Err(e) => { tracing::error!(match_id, stage = e.0, "match: artifact presign failed"); None } } }}
/// How long a minted artifact PUT is good for.////// It has to cover one upload and its retries, not a match: the container asks/// for this immediately before sending the bytes. Generous enough for a slow/// push of a 30MB log on a task that is being torn down, short enough that the/// URL is worth little if it leaks.pub const ARTIFACT_PUT_EXPIRY_SECS: u64 = 900;
/// An artifact name becomes an S3 key segment, so it stays a plain relative/// path. This is the only check standing between a match token and a key: the/// minting route signs whatever gets past it, and the caller then writes to S3/// with us out of the loop.////// This compares literal text, so it only holds because the caller gets the/// name from axum's `Path` extractor, which percent-decodes the `{*name}`/// capture first. `%2e%2e%2f` reaches here as `../`. Move this check anywhere/// that sees the raw path instead - a layer, a middleware, a hand-rolled split/// of `uri().path()` - and every encoded traversal walks straight through it./// `encoded_traversal_names_are_rejected` in `routes.rs` is what would catch/// that.pub fn artifact_name_ok(name: &str) -> bool { !name.contains("..") && !name.starts_with('/') && !name.is_empty()}
/// Whether this match may still send artifacts, checked on the minting route/// for the same reason [`artifact_name_ok`] is: what it hands back is an S3/// write credential.////// The token in the URL is the whole credential there, and until this existed it/// was good forever. Closing the window on a status is what bounds it, and it/// has to be a status no healthy match reaches with an upload still to come:////// * `over` is reached from `game-over` or `exiting`, and arena sends `game-over`/// as the *last* line of `container/exit/finalize.sh` - after the turn/// reports, the result, the screenshots, the stats table, the profiles and the/// client logs. `exiting` follows it in `container/entrypoint.sh`, once/// `finalize_once` has returned. Nothing in the container uploads after/// either. The `game-over` that `MatchHost` writes mid-match is a file in the/// state dir, read by `container/watch/turn-reports.sh` to know when to/// drain; it is not this event and does not reach us./// * `failed` is reached from a poll seeing the ECS task `STOPPED`,/// `DEPROVISIONING` or `DELETED`, all of which are after the container's/// processes are gone, so there is nothing left to upload from.////// `ready` is deliberately not a boundary: finalize runs while the row still/// says `ready`, and that is when almost every artifact is sent.pub fn uploads_open(row: &MatchRow) -> bool { !matches!(row.status.as_str(), "over" | "failed")}
/// arena's callback: advisory, but it advances what we know. Db-only on/// purpose - it must work (and be testable) without an AWS client.////// arena emits `starting`, `running`, `ready`, `error`, `idle-timeout`,/// `game-over` and `exiting`, and that is all of them - `emit` in its/// `container/lib/common.sh` is the only caller.////// `ready` is the mechanism now - readiness reported beats readiness/// discovered - and the probe in `advance` is the fallback for a container/// that cannot reach us.////// Anything unrecognised moves nothing, which is load-bearing: arena's/// `idle-timeout` relies on it. Nobody was playing, so nothing failed, and the/// shutdown it triggers ends the match through finalize's `game-over` instead.////// Once a row is [`is_terminal`], nothing here moves it again. This route is/// a bare token with no expiry, and `emit` in `container/lib/common.sh`/// retries its POST, so any event can be delivered twice or replayed late. A/// duplicate `game-over` must not paper over a real `failed` with `over`, and/// a duplicate `error` must not turn a match that ended cleanly into one the/// player's history shows as failed.pub async fn record_event(db: &Db, row: &MatchRow, event: &str, detail: &str) { tracing::info!(match_id = row.id, event, detail, "match: event"); if is_terminal(&row.status) { return; } let new_status = match (row.status.as_str(), event) { (_, "game-over") | (_, "exiting") => Some("over"), ("starting", "running") => Some("running"), // Never backwards: a late `ready` must not resurrect a match that has // already ended. The probe can produce one at any time, and an arena // that starts emitting it would send it before the ending events. // (`is_terminal` above already excludes them; this also keeps a // mid-startup `ready` from firing twice.) ("starting" | "running", "ready") => Some("ready"), (_, "error") => Some("failed"), _ => None, }; if let Some(status) = new_status { if let Err(e) = db.set_match_status(&row.id, status).await { // The route below still answers 204 regardless, so this line is // what tells an operator the transition arena believes landed // never reached the row. tracing::error!(match_id = row.id, status, ?e, "match: status not stored"); } if status == "ready" { tracing::info!(match_id = row.id, via = "callback", detail, "match: ready"); } }}
/// `over` and `failed` are terminal: nothing an arena event sends afterwards/// moves a row out of either. This happens to line up with [`uploads_open`]/// today, but the two answer different questions - one closes an S3 write/// credential, the other closes the status column - and are kept as separate/// predicates so a future change to either does not silently change the/// other.fn is_terminal(status: &str) -> bool { matches!(status, "over" | "failed")}
/// The ECS task states that mean there is no container left to reach.////// One list, because two callers act on it and they reach opposite/// conclusions from it: on the way up a gone task is a match that never/// started, and after `ready` it is a match that ended. Which of those is/// right is each caller's decision; *whether the task is gone* is one fact,/// and letting the two keep their own copy of these three strings is how they/// would come to disagree about it.////// Anything else — `RUNNING`, `PROVISIONING`, `PENDING`, a state ECS adds/// after this was written — is not evidence of an ending. Reading an unknown/// state as gone would end live matches.fn task_is_gone(task_status: &str) -> bool { matches!(task_status, "STOPPED" | "DEPROVISIONING" | "DELETED")}
/// Constant-time token comparison - the token is the whole credential on the/// callback and upload routes.pub fn token_matches(row: &MatchRow, presented: &str) -> bool { use hmac::Mac; // Reuse HMAC as a constant-time comparator: MAC both with a throwaway // key and compare tags via the crate's constant-time verify. let key = b"token-compare"; let mut a = hmac::Hmac::<sha2::Sha256>::new_from_slice(key).expect("any key length"); a.update(row.token.as_bytes()); let tag = a.finalize().into_bytes(); let mut b = hmac::Hmac::<sha2::Sha256>::new_from_slice(key).expect("any key length"); b.update(presented.as_bytes()); b.verify_slice(&tag).is_ok()}
pub type SharedMatches = Option<Arc<Matches>>;
#[cfg(test)]mod tests { use super::*;
fn row(status: &str) -> MatchRow { MatchRow { id: "m1".into(), owner_did: "did:plc:abc".into(), status: status.into(), token: "tok".into(), task_arn: None, task_ip: None, created_at: 0, started_at: None, ended_at: None, scenario: None, } }
#[test] fn token_comparison() { assert!(token_matches(&row("ready"), "tok")); assert!(!token_matches(&row("ready"), "tok2")); assert!(!token_matches(&row("ready"), "")); }
/// The three strings that decide whether a player's match is declared /// over. Both callers of this act on it, and one of them writes a terminal /// status off it, so the set is worth pinning: a state read as gone by /// mistake ends a match somebody is playing. #[test] fn only_a_stopped_task_counts_as_gone() { assert!(task_is_gone("STOPPED")); assert!(task_is_gone("DEPROVISIONING")); assert!(task_is_gone("DELETED"));
assert!(!task_is_gone("RUNNING")); assert!(!task_is_gone("PROVISIONING")); assert!(!task_is_gone("PENDING")); assert!(!task_is_gone("ACTIVATING")); // What `aws::task_status` substitutes when the task carries no // lastStatus at all. Not knowing is not an ending. assert!(!task_is_gone("UNKNOWN")); // Whatever ECS adds after this was written. assert!(!task_is_gone("SOMETHING_NEW")); }
/// The window has to outlast finalize, which uploads under `ready` and then /// posts the `game-over` that closes it. #[test] fn the_upload_window_closes_only_at_the_end() { assert!(uploads_open(&row("starting"))); assert!(uploads_open(&row("running"))); assert!(uploads_open(&row("ready"))); assert!(!uploads_open(&row("over"))); assert!(!uploads_open(&row("failed"))); }
/// Applies one callback to a stored row and reports the status left behind. async fn after_event(from: &str, event: &str) -> String { let dir = tempfile::tempdir().unwrap(); let db = crate::db::Db::open(&dir.path().join("test.sqlite")).unwrap(); let mut stored = row(from); stored.status = "starting".into(); db.insert_match(&stored).await.unwrap(); // The row as the route would have read it, at the status under test. db.set_match_status(&stored.id, from).await.unwrap(); stored.status = from.into();
record_event(&db, &stored, event, "suramadu=9100ms host=11200ms").await; db.get_match(&stored.id).await.unwrap().unwrap().status }
/// Readiness reported rather than discovered. No arena version emits this /// yet, so the mapping is the thing under test, not a live path. #[tokio::test] async fn ready_arrives_by_callback() { assert_eq!(after_event("starting", "ready").await, "ready"); assert_eq!(after_event("running", "ready").await, "ready"); }
/// A `ready` that arrives after the match has ended - a retried callback, /// or a container that outlived its game - must not put a finished match /// back in front of a player. #[tokio::test] async fn ready_never_moves_a_match_backwards() { assert_eq!(after_event("over", "ready").await, "over"); assert_eq!(after_event("failed", "ready").await, "failed"); }
/// Same reasoning as `ready`, for every other event: `emit` in /// `container/lib/common.sh` retries its POST twice, and the callback /// route is a bare token with no expiry, so a duplicate or a replay can /// land after the real ending event already has. A late `error` must not /// turn a match that finished cleanly into one the player's history /// shows as failed. #[tokio::test] async fn error_never_moves_a_match_backwards() { assert_eq!(after_event("over", "error").await, "over"); assert_eq!(after_event("failed", "error").await, "failed"); }
/// The other direction: a duplicate `game-over` (or `exiting`) must not /// paper over a genuine failure by putting `over` back over `failed`. #[tokio::test] async fn game_over_never_moves_a_match_backwards() { assert_eq!(after_event("failed", "game-over").await, "failed"); assert_eq!(after_event("failed", "exiting").await, "failed"); }
/// The transitions that were already here, unchanged by the new one. #[tokio::test] async fn the_existing_events_still_land() { assert_eq!(after_event("starting", "running").await, "running"); assert_eq!(after_event("ready", "game-over").await, "over"); assert_eq!(after_event("running", "error").await, "failed"); assert_eq!(after_event("running", "wamble").await, "running"); }}