Merge pull request 'feat(sidecar): protocol 7 — the Event System's command plane (Phase 16b cutover, 1 of 6)' (#40) from edge into main
Some checks failed
sync-project-tree / sync (push) Successful in 21s
SonarQube / analysis (push) Successful in 1m26s
Release sidecar / release (push) Failing after 8m39s

Reviewed-on: #40
Reviewed-by: Colby Whitlock <whitlocktech@gmail.com>
This commit is contained in:
2026-09-09 19:54:28 +00:00
2 changed files with 659 additions and 5 deletions

View File

@@ -72,7 +72,26 @@ use tracing_subscriber::EnvFilter;
/// index only the columns they already had, so the new fields ride inside the stored JSON and the
/// new kind lands in `events` like any other. That is the dumb-forwarder property doing its job:
/// the sidecar defines no schema for a frame's contents and so needs no change when they grow.
pub const PROTOCOL_VERSION: u32 = 5;
///
/// v6 (Protocol 6): the first bump that is about a GUARANTEE rather than about data, and the first
/// the sidecar mostly gets for free. Two things:
///
/// * **`idempotencyKey` on inbound commands.** A command that carries one is executed by the shard
/// at most once; a repeat is answered with the original reply rather than re-run. That is what
/// makes a world-writing verb retryable at all — until now a lost acknowledgement was
/// indistinguishable from a command that never applied, so the website had to declare every
/// write un-retryable and accept losing one rather than risk doubling it. The sidecar's part is
/// to CARRY the key (it rides in the command body, which every write endpoint already passes
/// through verbatim) and to understand the one new answer the shard can now give: `bridge.busy`,
/// meaning a command under that key is still in flight. See `web::respond`.
/// * **`champ.boss.killed` is a new kind**: a champion's defeat, with the damage table only the
/// shard ever sees. It was previously inferable from `champ.update` going `bossUp` true then
/// false alongside a nearby `mob.killed`, which is fragile and says nothing about who did the
/// work. It lands in `events` and on the feed like any other kind, with no code here at all —
/// the dumb-forwarder property again.
///
/// **No store migration.** Nothing gains a column; the new kind is persisted whole like every other.
pub const PROTOCOL_VERSION: u32 = 7;
// Not `#[tokio::main]`: on Windows the SCM dispatcher takes over this thread and starts the runtime
// itself, on its own thread, once the service actually begins. The runtime is built by whichever

View File

@@ -72,6 +72,41 @@ 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))
// The world verbs (protocol 7, EVENTS_PLAN.md Phase 12a). Five things an event author
// can place -- creatures, an enhanced "boss", an oracle NPC, a temporary gate,
// decoration -- and ONE command family, because each of them ends in "an object exists
// and this run owns it". POST places, GET says what the run still owns, POST .../despawn
// gives it back. Ownership is held on the shard, so despawn cannot be pointed at a serial
// the run did not create.
.route("/world", post(world_spawn))
.route("/world/:run_id", get(world_owned))
.route("/world/:run_id/despawn", post(world_despawn))
// The one-shots (protocol 7 part b, EVENTS_PLAN.md Phase 12b). Neither owned nor
// borrowed: an item put into somebody's hands, and a world save. Both are `done is
// done`, which is why they are not in the world family -- there is nothing to give
// back and no ledger row core would come back for.
//
// `GET /items` is the shard's own grant allowlist, so the website's dropdown offers
// what this shard will actually build rather than what a module guessed.
.route("/items", get(item_catalog))
.route("/items/grant", post(item_grant))
.route("/world/save", post(world_save))
// 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))
@@ -246,13 +281,31 @@ fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
// ---- shared reply handling ----
/// Protocol 6. `bridge.busy` says a command carrying this `idempotencyKey` is already in flight on
/// the shard: nothing was run, and the caller should come back.
///
/// It maps to **425 Too Early**, which is what that status is for — a server unwilling to risk
/// processing a request that might be a replay. The obvious alternative, 409, is already the
/// protocol-version gate's answer, and those two want opposite dispositions from a client: a version
/// mismatch is a deployment fault nobody should retry, and a busy shard is a retry that should
/// succeed on its own. Sharing a status would have made the difference readable only by inspecting
/// the body, which is exactly how a retry loop ends up hiding a mismatched deployment.
///
/// It is checked BEFORE the `.error` suffix test in each responder below, and it is deliberately not
/// spelled `bridge.busy.error`: nothing is wrong. The work is happening.
const BUSY_KIND: &str = "bridge.busy";
const BUSY_STATUS: StatusCode = StatusCode::TOO_EARLY;
/// Turns an RPC result into an HTTP response. A `bridge.error` reply from the shard becomes a 4xx;
/// a real reply is returned as-is; transport failures map to 503/504.
/// a `bridge.busy` reply becomes a 425; a real reply is returned as-is; transport failures map to
/// 503/504.
fn respond(result: Result<Value, RpcError>) -> (StatusCode, Json<Value>) {
match result {
Ok(value) => {
let kind = value.get("kind").and_then(|k| k.as_str()).unwrap_or("");
if kind == "bridge.error" || kind.ends_with(".error") {
if kind == BUSY_KIND {
(BUSY_STATUS, Json(value))
} else if kind == "bridge.error" || kind.ends_with(".error") {
let reason = value
.get("reason")
.and_then(|r| r.as_str())
@@ -286,7 +339,9 @@ fn respond_admin(result: Result<Value, RpcError>) -> (StatusCode, Json<Value>) {
match result {
Ok(value) => {
let kind = value.get("kind").and_then(|k| k.as_str()).unwrap_or("");
if kind == "admin.error" {
if kind == BUSY_KIND {
(BUSY_STATUS, Json(value))
} else if kind == "admin.error" {
let reason = value
.get("reason")
.and_then(|r| r.as_str())
@@ -317,6 +372,70 @@ fn respond_admin(result: Result<Value, RpcError>) -> (StatusCode, Json<Value>) {
}
}
/// 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<Value, RpcError>) -> (StatusCode, Json<Value>) {
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")
// Phase 12b. A grant against a run this shard has never been told to count
// is the same shape as an unknown lease key: the caller named something that
// does not exist here, which is a 404 and never a retry. It is deliberately
// NOT the same as a run whose ledger is open and empty -- that is a 200 with
// `granted: 0`, because "nobody came" is a result rather than a mistake.
|| reason.contains("no participation ledger")
{
StatusCode::NOT_FOUND
} else if reason.contains("saves at most every") {
// A save refused because one just happened is the shard's rate limit, and it
// is TRANSIENT in a way nothing else on this plane is: the same request will
// succeed once the interval passes. 429 says exactly that, and keeps it out of
// the module's permanent-status set so a phase boundary is retried rather than
// abandoned.
StatusCode::TOO_MANY_REQUESTS
} 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.
@@ -324,7 +443,9 @@ fn respond_account(result: Result<Value, RpcError>) -> (StatusCode, Json<Value>)
match result {
Ok(value) => {
let kind = value.get("kind").and_then(|k| k.as_str()).unwrap_or("");
if kind == "account.error" {
if kind == BUSY_KIND {
(BUSY_STATUS, Json(value))
} else if kind == "account.error" {
let reason = value
.get("reason")
.and_then(|r| r.as_str())
@@ -454,6 +575,16 @@ async fn link_delete(
/// Forwards a staff moderation command to the shard, correlated on a fresh reqId. Injects `kind`
/// and `reqId`, requiring the caller-supplied `actor` up front (the shard enforces it too). The
/// body's remaining fields (account/serial/durationSec/reason/text/hue) pass straight through.
///
/// **Protocol 6: `idempotencyKey` is one of those remaining fields**, and passing it through is the
/// whole of the sidecar's part in the guarantee. It is worth stating rather than leaving to the
/// word "remaining", because a later refactor that narrowed this to a known field list would quietly
/// turn every retried world write back into a possible duplicate, and nothing here would fail.
///
/// The key belongs to the CALLER's unit of work — the website's event step — so the sidecar neither
/// generates one nor validates it. Note also that `reqId` is regenerated on every call: a retry
/// carries the same idempotency key under a NEW correlation id, which is exactly why the shard
/// re-stamps a replayed reply rather than echoing the id the first attempt used.
async fn admin_call(st: &AppState, kind: &str, body: Value) -> (StatusCode, Json<Value>) {
let mut obj = match body {
Value::Object(m) => m,
@@ -504,6 +635,241 @@ async fn admin_broadcast(State(st): State<AppState>, Json(body): Json<Value>) ->
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<Value>) {
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.
/// **`?key=` and `?target=` narrow it to one row, and a targeted key needs them** (protocol 7 part
/// b). `Spawner.MaxCount` is one capability over thousands of spawners, so it has no single
/// "current" and the catalog walk cannot fill one in -- while the website's `read()` needs exactly
/// one value for exactly one target before it applies anything. Naming both answers that.
///
/// The frame also carries `holds`: every lease this shard is actually holding, whatever key or
/// target it is on. A catalog walk can enumerate the KEYS but never the holds on a targeted one --
/// there is no list of spawners to walk -- so without it a reconcile after an outage would have no
/// way to ask "what are you still holding?".
async fn lease_list(State(st): State<AppState>, Query(q): Query<LeaseQuery>) -> impl IntoResponse {
let mut body = serde_json::Map::new();
if let Some(key) = q.key {
body.insert("key".to_string(), json!(key));
}
if let Some(target) = q.target {
body.insert("target".to_string(), json!(target));
}
let arg = if body.is_empty() {
Value::Null
} else {
Value::Object(body)
};
event_call(&st, "lease.list", arg).await
}
/// Narrowing for `GET /lease`. Both optional: absent means the whole catalog, as before.
#[derive(Debug, Deserialize)]
struct LeaseQuery {
key: Option<String>,
target: Option<String>,
}
/// Body: {"key":"...","value":"...","holdMs":<ms>,"untilMs":<opt>,"runId":<opt>,"idempotencyKey":<opt>}.
///
/// **`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<AppState>, Json(body): Json<Value>) -> impl IntoResponse {
event_call(&st, "lease.apply", body).await
}
/// Body: {"key":"...","expected":"...","baseline":"...","idempotencyKey":<opt>}.
///
/// `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<AppState>, Json(body): Json<Value>) -> impl IntoResponse {
event_call(&st, "lease.release", body).await
}
/// Body: {"runId":"...","map":"Felucca","x":N,"y":N,"radius":N,"holdMs":<opt>}.
///
/// 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<AppState>,
Json(body): Json<Value>,
) -> impl IntoResponse {
event_call(&st, "participation.open", body).await
}
/// Body: {"idempotencyKey":<opt>}. 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<AppState>,
Path(run_id): Path<String>,
Json(body): Json<Value>,
) -> 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":<opt>}. 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<AppState>,
Path(run_id): Path<String>,
Json(body): Json<Value>,
) -> 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
}
/// Body: {"runId":"...","what":"creature|boss|npc|gate|decor","map":"...","x":N,"y":N,...}.
///
/// One route for five author-facing verbs. The `what` discriminator is a wire detail: the
/// differences between them -- a boss's multipliers, an oracle's lines, a gate's destination and
/// `holdMs` -- are fields on one command rather than five commands, so there is one ledger shape,
/// one teardown path and one reconcile instead of five near-identical ones in three repos.
///
/// The shard registers every serial it places against the run and PERSISTS that registry beside
/// the world save, which is what makes `world_despawn` below safe: a spawned creature survives a
/// restart, so an in-memory registry would leave the website holding serials the shard would not
/// vouch for.
async fn world_spawn(State(st): State<AppState>, Json(body): Json<Value>) -> impl IntoResponse {
event_call(&st, "world.spawn", body).await
}
/// What the run still owns, and the answer the website's `reconcile()` is built on.
///
/// A GET, unlike `participation_snapshot`: it carries no idempotency key and the shard answers it
/// in one pass, pruning rows whose object the world has already lost as it walks. Anything not
/// listed is gone -- which is the shape core wants, because it takes a row out of its ledger only
/// on an explicit reply and this is that reply.
async fn world_owned(State(st): State<AppState>, Path(run_id): Path<String>) -> impl IntoResponse {
event_call(&st, "world.owned", json!({ "runId": run_id })).await
}
/// Body: {"serials":[...]} -- or no serials at all, which means everything the run owns and is the
/// call teardown actually makes.
///
/// Three answers, and the split is why the shard keeps a registry at all. `removed` was found and
/// deleted; `gone` was owned but already absent, which is what happens when a player kills an event
/// creature and is a SUCCESS; `refused` was never this run's to delete, and is the only answer here
/// that means somebody asked for something they should not have.
async fn world_despawn(
State(st): State<AppState>,
Path(run_id): Path<String>,
Json(body): Json<Value>,
) -> 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, "world.despawn", Value::Object(obj)).await
}
// ---- the one-shots (protocol 7 part b) ----
/// What this shard is willing to grant, and the bounds it will grant within.
///
/// A read, so the website's option source offers what this shard will actually build. The module
/// holds the same list, which is two copies of a short allowlist on purpose and exactly how the
/// lease bounds are already carried: the module's copy is what makes a bad value a refusal on a
/// form, and this one is what is true when the website is wrong.
async fn item_catalog(State(st): State<AppState>) -> impl IntoResponse {
event_call(&st, "item.catalog", Value::Null).await
}
/// Body: {"runId":"...","item":"gold","amount":N,"hue":<opt>,"name":<opt>,"where":<opt>,"idempotencyKey":<opt>}.
///
/// **The recipients are not in the body, and that is the design.** The shard already holds the
/// run's participation ledger (protocol 6 part b), keyed by the same character serials the
/// website's `member_key` holds, so the grant names a run and the shard resolves who was there.
/// Sending a list would mean the same list crossing the wire twice with a window in which the two
/// disagree -- and it would have needed a core surface handing a module core's own participants.
///
/// A run with no ledger open is a 404, not an empty success: "nobody came" and "you never told me
/// to count" are different facts, and only the first is a result a run should record.
///
/// **Retryable, and protocol 6 is why.** `EVENTS.md` §G called a grant un-retryable because a lost
/// acknowledgement and a grant that never applied looked the same -- exactly the argument that made
/// `uo.broadcast` answer `retry: false` in Phase 9. An `idempotencyKey` closes that: a repeat is
/// answered by the original reply, so a retried grant cannot be one winner receiving two.
async fn item_grant(State(st): State<AppState>, Json(body): Json<Value>) -> impl IntoResponse {
event_call(&st, "item.grant", body).await
}
/// Body: {"idempotencyKey":<opt>}. Starts a world save, useful as a phase boundary.
///
/// The reply says the save was STARTED and nothing more. What actually happened rides
/// `world.save.before` / `world.save.after`, which have been on the event stream since protocol 2 --
/// so this route asserts nothing it cannot know, and a caller that needs the completion watches the
/// stream it is already connected to.
///
/// **A save too soon after the last one is refused, not queued**, and the shard counts ServUO's own
/// autosave as the last one. A save stops the world; a queued one would land at a moment nobody
/// chose, in the middle of whatever the next step is doing.
async fn world_save(State(st): State<AppState>, Json(body): Json<Value>) -> impl IntoResponse {
event_call(&st, "world.save", body).await
}
// ---- help-page queue handlers ----
/// The open help-page queue, correlated on reqId. Returns a pages.list.
@@ -978,3 +1344,272 @@ async fn ws_client(mut socket: WebSocket, state: AppState) {
info!("ws client disconnected");
}
#[cfg(test)]
mod tests {
use super::*;
fn reply(kind: &str) -> Result<Value, RpcError> {
Ok(json!({"t": 1, "kind": kind, "reqId": "r-9"}))
}
/// Protocol 6. Every responder must recognise `bridge.busy`, because every write plane can be
/// retried: the staff plane, the account plane and the plain command plane all reach handlers
/// that a keyed retry can arrive at. A responder that missed it would return 200 with a body
/// saying nothing happened, which is the worst of the three possible answers.
#[test]
fn busy_maps_to_425_on_every_plane() {
assert_eq!(respond(reply("bridge.busy")).0, StatusCode::TOO_EARLY);
assert_eq!(respond_admin(reply("bridge.busy")).0, StatusCode::TOO_EARLY);
assert_eq!(
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
);
}
/// Protocol 7's world verbs go through the same responder, and this pins the two mappings
/// they depend on rather than trusting that the reason-sniffing above keeps covering a kind
/// it was written before.
///
/// A CEILING refusal is a 400 on purpose. It is permanent -- retrying "you asked for 80
/// creatures and this shard places 30" gets the same answer forever -- and it is the module's
/// `PERMANENT_STATUSES` that has to see it as such, so classifying it as anything retryable
/// would put a run in a loop against a limit that will never move.
#[test]
fn a_world_refusal_is_a_400_and_the_gate_is_still_a_403() {
assert_eq!(
respond_event(Ok(json!({
"kind": "world.error",
"action": "spawn",
"reason": "this shard places 1 to 30 of 'creature' at a time, and 80 was asked for"
})))
.0,
StatusCode::BAD_REQUEST
);
assert_eq!(
respond_event(Ok(json!({
"kind": "world.error",
"action": "spawn",
"reason": "the event plane is disabled on this shard (Bridge.EventsEnabled)"
})))
.0,
StatusCode::FORBIDDEN
);
}
/// A run the shard has no registry rows for answers with an EMPTY hand, not a 404, and the
/// distinction is load-bearing for reconcile.
///
/// "This run owns nothing" and "I have never heard of this run" are the same fact once the
/// registry is the only record of ownership, and they stay the same fact across a restart:
/// the registry is written by `EventSink.WorldSave`, so it and the objects it describes are
/// saved and lost together. A 404 here would make the website treat a run that legitimately
/// owns nothing as a shard it could not reach.
#[test]
fn a_run_owning_nothing_is_an_empty_list_not_a_404() {
let (status, body) = respond_event(Ok(json!({
"kind": "world.owned.ok",
"runId": "77",
"owned": [],
"pruned": 0
})));
assert_eq!(status, StatusCode::OK);
assert_eq!(
body.0
.get("owned")
.and_then(|v| v.as_array())
.map(|a| a.len()),
Some(0)
);
}
/// 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 run this shard was never told to count is a 404; a run that WAS counted and had no
/// attendees is a 200. Protocol 7 part b.
#[test]
fn an_uncounted_run_is_a_404_and_an_empty_one_is_not() {
assert_eq!(
respond_event(Ok(json!({
"kind": "oneshot.error",
"reason": "run 42 has no participation ledger open on this shard"
})))
.0,
StatusCode::NOT_FOUND
);
// The distinction the 404 exists to preserve. "Nobody came" is a RESULT -- an event
// nobody attended still happened -- and answering it as a failure would have the module
// retry a grant against a ledger that will be just as empty next time.
assert_eq!(
respond_event(Ok(json!({
"kind": "item.grant.ok",
"runId": "42",
"granted": 0,
"missed": []
})))
.0,
StatusCode::OK
);
}
/// The save rate limit is the one refusal on this plane that the same request will get past
/// by waiting, so it is a 429 rather than the 400 every other refusal is.
#[test]
fn a_save_refused_for_coming_too_soon_is_a_429() {
assert_eq!(
respond_event(Ok(json!({
"kind": "oneshot.error",
"reason": "this shard saves at most every 300 seconds, and the last save was 12 seconds ago"
})))
.0,
StatusCode::TOO_MANY_REQUESTS
);
// And an ordinary refusal on the same plane is still a 400, so the 429 is not swallowing
// the class it sits beside: a grant this shard does not offer will never succeed, however
// long the caller waits.
assert_eq!(
respond_event(Ok(json!({
"kind": "oneshot.error",
"reason": "this shard does not grant 'castle'"
})))
.0,
StatusCode::BAD_REQUEST
);
// The event gate being off stays a 403 on this plane too -- it is an operator's deliberate
// refusal, not a bad request.
assert_eq!(
respond_event(Ok(json!({
"kind": "oneshot.error",
"reason": "the event plane is disabled on this shard (Bridge.EventsEnabled)"
})))
.0,
StatusCode::FORBIDDEN
);
}
/// 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
/// nobody should retry, a busy shard is a retry that will succeed. Same-status would make the
/// two readable only by inspecting the body.
#[test]
fn busy_is_not_the_version_gates_status() {
assert_ne!(BUSY_STATUS, StatusCode::CONFLICT);
}
/// A replayed reply is an ordinary success. The shard marks it `replayed: true` for the log, and
/// the caller must be able to treat it exactly as it would have treated the answer it lost.
#[test]
fn a_replayed_reply_is_still_a_200() {
let value = json!({"t": 1, "kind": "admin.ok", "reqId": "r-9", "replayed": true});
let (status, body) = respond_admin(Ok(value));
assert_eq!(status, StatusCode::OK);
assert_eq!(body.0.get("replayed").and_then(|v| v.as_bool()), Some(true));
}
/// The error mapping the busy arm is threaded in front of must be untouched by it.
#[test]
fn errors_still_map_as_before() {
assert_eq!(
respond(Ok(
json!({"kind": "bridge.error", "reason": "unknown account"})
))
.0,
StatusCode::NOT_FOUND
);
assert_eq!(
respond_admin(Ok(json!({"kind": "admin.error", "reason": "protected"}))).0,
StatusCode::FORBIDDEN
);
assert_eq!(
respond_account(Ok(
json!({"kind": "account.error", "reason": "already exists"})
))
.0,
StatusCode::CONFLICT
);
}
}