|
|
|
@@ -72,6 +72,22 @@ pub async fn serve(addr: &str, state: AppState) -> anyhow::Result<()> {
|
|
|
|
.route("/admin/ban", post(admin_ban))
|
|
|
|
.route("/admin/ban", post(admin_ban))
|
|
|
|
.route("/admin/unban", post(admin_unban))
|
|
|
|
.route("/admin/unban", post(admin_unban))
|
|
|
|
.route("/admin/broadcast", post(admin_broadcast))
|
|
|
|
.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.
|
|
|
|
// Help-page (support) queue: snapshot the open queue, respond to / close a page.
|
|
|
|
.route("/pages", get(pages_list))
|
|
|
|
.route("/pages", get(pages_list))
|
|
|
|
.route("/pages/:id/respond", post(page_respond))
|
|
|
|
.route("/pages/:id/respond", post(page_respond))
|
|
|
|
@@ -246,13 +262,31 @@ fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
|
|
|
|
|
|
|
|
|
|
|
|
// ---- shared reply handling ----
|
|
|
|
// ---- 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;
|
|
|
|
/// 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>) {
|
|
|
|
fn respond(result: Result<Value, RpcError>) -> (StatusCode, Json<Value>) {
|
|
|
|
match result {
|
|
|
|
match result {
|
|
|
|
Ok(value) => {
|
|
|
|
Ok(value) => {
|
|
|
|
let kind = value.get("kind").and_then(|k| k.as_str()).unwrap_or("");
|
|
|
|
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
|
|
|
|
let reason = value
|
|
|
|
.get("reason")
|
|
|
|
.get("reason")
|
|
|
|
.and_then(|r| r.as_str())
|
|
|
|
.and_then(|r| r.as_str())
|
|
|
|
@@ -286,7 +320,9 @@ fn respond_admin(result: Result<Value, RpcError>) -> (StatusCode, Json<Value>) {
|
|
|
|
match result {
|
|
|
|
match result {
|
|
|
|
Ok(value) => {
|
|
|
|
Ok(value) => {
|
|
|
|
let kind = value.get("kind").and_then(|k| k.as_str()).unwrap_or("");
|
|
|
|
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
|
|
|
|
let reason = value
|
|
|
|
.get("reason")
|
|
|
|
.get("reason")
|
|
|
|
.and_then(|r| r.as_str())
|
|
|
|
.and_then(|r| r.as_str())
|
|
|
|
@@ -317,6 +353,56 @@ 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")
|
|
|
|
|
|
|
|
{
|
|
|
|
|
|
|
|
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
|
|
|
|
/// 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/
|
|
|
|
/// 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.
|
|
|
|
/// refused action is a 403, an unknown target or "not linked" is a 404, anything else a 400.
|
|
|
|
@@ -324,7 +410,9 @@ fn respond_account(result: Result<Value, RpcError>) -> (StatusCode, Json<Value>)
|
|
|
|
match result {
|
|
|
|
match result {
|
|
|
|
Ok(value) => {
|
|
|
|
Ok(value) => {
|
|
|
|
let kind = value.get("kind").and_then(|k| k.as_str()).unwrap_or("");
|
|
|
|
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
|
|
|
|
let reason = value
|
|
|
|
.get("reason")
|
|
|
|
.get("reason")
|
|
|
|
.and_then(|r| r.as_str())
|
|
|
|
.and_then(|r| r.as_str())
|
|
|
|
@@ -454,6 +542,16 @@ async fn link_delete(
|
|
|
|
/// Forwards a staff moderation command to the shard, correlated on a fresh reqId. Injects `kind`
|
|
|
|
/// 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
|
|
|
|
/// 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.
|
|
|
|
/// 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>) {
|
|
|
|
async fn admin_call(st: &AppState, kind: &str, body: Value) -> (StatusCode, Json<Value>) {
|
|
|
|
let mut obj = match body {
|
|
|
|
let mut obj = match body {
|
|
|
|
Value::Object(m) => m,
|
|
|
|
Value::Object(m) => m,
|
|
|
|
@@ -504,6 +602,123 @@ async fn admin_broadcast(State(st): State<AppState>, Json(body): Json<Value>) ->
|
|
|
|
admin_call(&st, "admin.broadcast", body).await
|
|
|
|
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.
|
|
|
|
|
|
|
|
async fn lease_list(State(st): State<AppState>) -> impl IntoResponse {
|
|
|
|
|
|
|
|
event_call(&st, "lease.list", Value::Null).await
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
/// 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
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// ---- help-page queue handlers ----
|
|
|
|
// ---- help-page queue handlers ----
|
|
|
|
|
|
|
|
|
|
|
|
/// The open help-page queue, correlated on reqId. Returns a pages.list.
|
|
|
|
/// The open help-page queue, correlated on reqId. Returns a pages.list.
|
|
|
|
@@ -978,3 +1193,153 @@ async fn ws_client(mut socket: WebSocket, state: AppState) {
|
|
|
|
|
|
|
|
|
|
|
|
info!("ws client disconnected");
|
|
|
|
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
|
|
|
|
|
|
|
|
);
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
/// 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
|
|
|
|
|
|
|
|
/// 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
|
|
|
|
|
|
|
|
);
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
}
|
|
|
|
|