From 93411966d7e8dfe52b9fe12e29a53a8012c45b1a Mon Sep 17 00:00:00 2001 From: wtclaude Date: Fri, 4 Sep 2026 19:31:31 -0500 Subject: [PATCH] feat(sidecar): carry the event plane (Phase 11b) Protocol 6 amended in place, so PROTOCOL_VERSION is unchanged. Six routes and a fourth responder; no store migration and no new machinery. `event_call` is `admin_call` without the required `actor`: an event verb's author is a RUN, which the body carries as `runId`, and demanding a human name for something no human is doing would have the runner inventing one. `respond_event` exists for two mappings the generic responder gets wrong. A drifted lease is a 200 -- the shard was asked to compare and set, it compared, and it declined to overwrite somebody's deliberate change, which is the mechanism working -- and deliberately not the 409 the version gate owns, for the same reason 425 is not. And the event plane being switched off is a 403 rather than a reason-sniffed 400: it is an operator's deliberate refusal, and a 400 would send an administrator hunting a bug in a step that is written correctly. `participation.snapshot` is a POST for a read, because it carries the caller's idempotency key and the shard may refuse it as a repeat in flight. Co-Authored-By: Claude --- sidecar/src/web.rs | 269 +++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 269 insertions(+) diff --git a/sidecar/src/web.rs b/sidecar/src/web.rs index 9909f1a..fa609a2 100644 --- a/sidecar/src/web.rs +++ b/sidecar/src/web.rs @@ -72,6 +72,22 @@ pub async fn serve(addr: &str, state: AppState) -> anyhow::Result<()> { .route("/admin/ban", post(admin_ban)) .route("/admin/unban", post(admin_unban)) .route("/admin/broadcast", post(admin_broadcast)) + // The event plane (protocol 6, EVENTS_PLAN.md Phase 11b). Leases are a live config value + // the website holds for a bounded time; the shard restores baseline when the deadline + // passes whether or not anyone asks it to. GET lists the whole catalog with current values, + // which is the one read both `read()` and `inForce()` on the website's side are served by. + .route("/lease", get(lease_list).post(lease_apply)) + .route("/lease/release", post(lease_release)) + // The run-scoped participation ledger. `snapshot` is a POST despite being a read: it + // carries the caller's `idempotencyKey`, and on a well-attended run the shard walks its + // members across ticks rather than in one inbound call -- so a repeat arriving mid-walk is + // answered `bridge.busy`, and a read that can be refused as a repeat is not a GET. + .route("/participation", post(participation_open)) + .route( + "/participation/:run_id/snapshot", + post(participation_snapshot), + ) + .route("/participation/:run_id/close", post(participation_close)) // Help-page (support) queue: snapshot the open queue, respond to / close a page. .route("/pages", get(pages_list)) .route("/pages/:id/respond", post(page_respond)) @@ -337,6 +353,56 @@ fn respond_admin(result: Result) -> (StatusCode, Json) { } } +/// Like `respond`, but for the event plane: leases and the participation ledger. +/// +/// Two mappings are the point of it existing rather than reusing `respond`. +/// +/// **`lease.drifted` is a 200.** The shard was asked to compare and set, it compared, and it +/// refused to overwrite somebody's deliberate change -- that is the mechanism working, not a +/// failure, and `cleanup.js` on the website treats `drifted` as a distinct successful outcome +/// rather than an error. It is also why this is not a 409: 409 is the protocol-version gate's, and +/// a version mismatch and a drifted lease want opposite dispositions from a caller. The same +/// argument protocol 6 made for `bridge.busy` being a 425. +/// +/// **The event plane being switched off is a 403**, not the 400 the generic responder's +/// reason-sniffing would produce. `Bridge.EventsEnabled` is an operator's deliberate refusal to let +/// the website change the world on a schedule, and telling the website it sent a bad request would +/// send an administrator hunting a bug in a step that is written correctly. +fn respond_event(result: Result) -> (StatusCode, Json) { + match result { + Ok(value) => { + let kind = value.get("kind").and_then(|k| k.as_str()).unwrap_or(""); + if kind == BUSY_KIND { + (BUSY_STATUS, Json(value)) + } else if kind.ends_with(".error") { + let reason = value + .get("reason") + .and_then(|r| r.as_str()) + .unwrap_or("request rejected"); + let code = if reason.contains("disabled") { + StatusCode::FORBIDDEN + } else if reason.contains("no lease is offered") || reason.contains("not counting") + { + StatusCode::NOT_FOUND + } else { + StatusCode::BAD_REQUEST + }; + (code, Json(value)) + } else { + (StatusCode::OK, Json(value)) + } + } + Err(RpcError::NoShard) => ( + StatusCode::SERVICE_UNAVAILABLE, + Json(json!({"error": "shard not connected"})), + ), + Err(RpcError::Timeout) => ( + StatusCode::GATEWAY_TIMEOUT, + Json(json!({"error": "shard did not reply in time"})), + ), + } +} + /// Like `respond`, but for the account-provisioning plane. Maps an `account.error` reply to a /// status by its reason: a name clash is a 409, the per-IP cap is a 429, a disabled/protected/ /// refused action is a 403, an unknown target or "not linked" is a 404, anything else a 400. @@ -536,6 +602,123 @@ async fn admin_broadcast(State(st): State, Json(body): Json) -> admin_call(&st, "admin.broadcast", body).await } +// ---- event plane handlers (protocol 6, Phase 11b) ---- + +/// Forwards an event-plane command to the shard, correlated on a fresh reqId. +/// +/// Deliberately NOT `admin_call`: that one requires an `actor`, because every verb behind it is a +/// staff member pressing a button and the shard's audit trail has to name them. An event verb's +/// author is a RUN, which the body already carries as `runId` -- and demanding an actor here would +/// have the runner inventing a human name for something no human is doing. +/// +/// Everything else about it is the same, and the `idempotencyKey` passthrough matters for the same +/// reason it does there: the key is one of the body's remaining fields, and a refactor that +/// narrowed this to a known field list would silently make every retried lease a possible duplicate. +async fn event_call(st: &AppState, kind: &str, body: Value) -> (StatusCode, Json) { + let mut obj = match body { + Value::Object(m) => m, + Value::Null => serde_json::Map::new(), + _ => { + return ( + StatusCode::BAD_REQUEST, + Json(json!({"error": "body must be a JSON object"})), + ) + } + }; + + let req_id = st.rpc.next_req_id(); + obj.insert("kind".to_string(), json!(kind)); + obj.insert("reqId".to_string(), json!(req_id)); + + respond_event(st.rpc.call(&st.shard, Value::Object(obj), &req_id).await) +} + +/// Every lease this shard offers, with what each is worth right now and what is holding it. +/// +/// One read answers both questions the website asks about a lease: `read()` wants the current value +/// before it applies anything, and `inForce()` wants to know whether the shard still has a record +/// of the hold. Splitting them would be two round trips for one key. +/// +/// **`held` means "the shard still has a record of this lease", not "the value is still +/// overridden".** A lease whose deadline has already fired stays listed, with `expired: true`, +/// until teardown collects its verdict -- otherwise a reconcile in that window would report it gone +/// and the website would write off a correctly-working backstop as an orphaned resource. +async fn lease_list(State(st): State) -> impl IntoResponse { + event_call(&st, "lease.list", Value::Null).await +} + +/// Body: {"key":"...","value":"...","holdMs":,"untilMs":,"runId":,"idempotencyKey":}. +/// +/// **`holdMs` is authoritative and `untilMs` is carried for display.** An absolute deadline computed +/// on the website and honoured on the shard is a deadline measured against two clocks, and a shard +/// running ten minutes fast would restore a ten-minute lease the moment it took it. A duration is +/// immune to that; the absolute time is still worth sending so a console can say when the hold ends. +/// +/// Values cross as TEXT whatever the lease's declared type, because JSON would otherwise decide for +/// us: `1200` and `1200.0` are one number to a parser and two strings to a compare-and-set. +async fn lease_apply(State(st): State, Json(body): Json) -> impl IntoResponse { + event_call(&st, "lease.apply", body).await +} + +/// Body: {"key":"...","expected":"...","baseline":"...","idempotencyKey":}. +/// +/// `expected` is what the event applied and `baseline` is what to put back, both out of the +/// website's ledger rather than the shard's memory -- so a release still works after a reconnect, +/// and a shard that has forgotten the lease entirely (a restart, which reverts every lease anyway) +/// can answer honestly instead of refusing. +/// +/// A mismatch comes back `lease.drifted` with a **200**: see `respond_event`. +async fn lease_release(State(st): State, Json(body): Json) -> impl IntoResponse { + event_call(&st, "lease.release", body).await +} + +/// Body: {"runId":"...","map":"Felucca","x":N,"y":N,"radius":N,"holdMs":}. +/// +/// Declares where a run happens and starts counting who is there. The area is a map, a point and a +/// radius rather than a region name, because protocol 6's own live walk established that the most +/// specific region containing an event is routinely anonymous. +async fn participation_open( + State(st): State, + Json(body): Json, +) -> impl IntoResponse { + event_call(&st, "participation.open", body).await +} + +/// Body: {"idempotencyKey":}. Answers the run's tally, best-effort resolved to accounts. +/// +/// A POST for a read, and the reason is worth keeping: on a well-attended run the shard walks its +/// members in chunks across Core ticks rather than handing the whole resolve to one inbound call, +/// so the handler completes after its call returned and a repeat arriving in between is answered +/// `bridge.busy`. A read that can legitimately be refused as a repeat in flight is not a GET. +async fn participation_snapshot( + State(st): State, + Path(run_id): Path, + Json(body): Json, +) -> impl IntoResponse { + let mut obj = match body { + Value::Object(m) => m, + _ => serde_json::Map::new(), + }; + obj.insert("runId".to_string(), json!(run_id)); + event_call(&st, "participation.snapshot", Value::Object(obj)).await +} + +/// Body: {"idempotencyKey":}. Stops counting; the tally stays readable through the shard's +/// grace window, because closing an event and collecting its results are two steps and either can +/// be retried. +async fn participation_close( + State(st): State, + Path(run_id): Path, + Json(body): Json, +) -> impl IntoResponse { + let mut obj = match body { + Value::Object(m) => m, + _ => serde_json::Map::new(), + }; + obj.insert("runId".to_string(), json!(run_id)); + event_call(&st, "participation.close", Value::Object(obj)).await +} + // ---- help-page queue handlers ---- /// The open help-page queue, correlated on reqId. Returns a pages.list. @@ -1031,6 +1214,92 @@ mod tests { respond_account(reply("bridge.busy")).0, StatusCode::TOO_EARLY ); + assert_eq!(respond_event(reply("bridge.busy")).0, StatusCode::TOO_EARLY); + } + + /// The event plane is the FIRST place `bridge.busy` is reachable on a live shard rather than + /// only in a unit test: `participation.snapshot` walks a well-attended run's members across + /// Core ticks, so it completes after its inbound call returned and a repeat can genuinely land + /// mid-flight. 11a built the door and had nothing to walk through it. + #[test] + fn a_drifted_lease_is_a_200_not_a_409() { + let value = + json!({"kind": "lease.drifted", "key": "PlayerCaps.SkillCap", "current": "1300"}); + let (status, body) = respond_event(Ok(value)); + + // The shard was asked to compare and set, it compared, and it declined to overwrite + // somebody's deliberate change. That is the mechanism working; the website records + // `drifted` as a distinct successful outcome rather than an error. + assert_eq!(status, StatusCode::OK); + assert_eq!(body.0.get("current").and_then(|v| v.as_str()), Some("1300")); + + // And explicitly not the version gate's status, for the reason 425 is not either: a + // mismatched deployment and a moved value want opposite dispositions from a caller. + assert_ne!(status, StatusCode::CONFLICT); + } + + /// The event plane being switched off is an operator's refusal, not a malformed request. A 400 + /// would send an administrator hunting a bug in a step that is written correctly. + #[test] + fn the_event_gate_being_off_is_a_403() { + assert_eq!( + respond_event(Ok(json!({ + "kind": "lease.error", + "reason": "the event plane is disabled on this shard (Bridge.EventsEnabled)" + }))) + .0, + StatusCode::FORBIDDEN + ); + assert_eq!( + respond_event(Ok(json!({ + "kind": "participation.error", + "reason": "the event plane is disabled on this shard (Bridge.EventsEnabled)" + }))) + .0, + StatusCode::FORBIDDEN + ); + } + + /// An unknown lease key and an unknown run are not-founds; anything else the shard refuses is a + /// bad request. The catalog is short and a typo in a step is the likely cause of both. + #[test] + fn unknown_lease_and_run_are_404s() { + assert_eq!( + respond_event(Ok(json!({ + "kind": "lease.error", + "reason": "no lease is offered for key 'Loot.MaxProps'" + }))) + .0, + StatusCode::NOT_FOUND + ); + assert_eq!( + respond_event(Ok(json!({ + "kind": "participation.error", + "reason": "this shard is not counting run '42'" + }))) + .0, + StatusCode::NOT_FOUND + ); + assert_eq!( + respond_event(Ok(json!({ + "kind": "lease.error", + "reason": "a lease needs a positive holdMs" + }))) + .0, + StatusCode::BAD_REQUEST + ); + } + + /// A lease taken, a tally answered: an ordinary success carries straight through. + #[test] + fn event_successes_are_200s() { + assert_eq!(respond_event(reply("lease.ok")).0, StatusCode::OK); + assert_eq!(respond_event(reply("lease.list.ok")).0, StatusCode::OK); + assert_eq!(respond_event(reply("participation.ok")).0, StatusCode::OK); + assert_eq!( + respond_event(reply("participation.snapshot.ok")).0, + StatusCode::OK + ); } /// 425 must not collide with the protocol-version gate's 409: a mismatch is a deployment fault -- 2.49.1