//! 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, /// The verified handle. Bots are labelled by the launcher itself. pub handle: Option, /// A `blue.lance.camo` record key in the seat-holder's own repository. /// Human seats only. pub camo_rkey: Option, /// A rolled picture from the browser, for a seat with no repository to /// read from. Bot seats only. pub camo_png: Option>, /// 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, } 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>, } 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, } /// 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 { 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>, 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> { 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> { 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> { self.artifact(match_id, RESULT_KEY).await } pub async fn result(&self, match_id: &str) -> Option { 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::(&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:` /// 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, ) -> Result { 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::() 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 { 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 { 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// - // 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 { 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::::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::::new_from_slice(key).expect("any key length"); b.update(presented.as_bytes()); b.verify_slice(&tag).is_ok() } pub type SharedMatches = Option>; #[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"); } }