From 4942db645b9e10208f82bc97d31cf7c9f68d8357 Mon Sep 17 00:00:00 2001 From: Josh Poole Date: Tue, 22 Sep 2026 11:36:38 +0100 Subject: [PATCH 1/7] 20260922 - Adopt node ingest v1.4.0, and stop it renaming our state enum The contract moved twice while we were on 1.2.2, both additively. 1.3.0 added `GET`/`PUT /v1/nodes/claim` and `POST /v1/nodes/claim/resend` with the `NodeClaimRequest` and `ClaimResponse` schemas, an `invalid_claim` slug on 400, and a 409 that answers with `ClaimResponse` rather than `Error`. 1.4.0 added `claim_state`, `claim_email` and `claim_undeliverable`, all required, to `HeartbeatResponse` and `ContactResponse`. Both landed on production on 2026-09-18. Nothing was broken meanwhile, which is why this is adoption rather than a fix: `comms/client.py` parses responses into plain dicts and never validates them against the generated response models, so the three new keys have been arriving and being ignored since that deploy. Claim codes were removed server-side in the same stack and cost us nothing, because we never implemented them. The spec file is byte-identical to `contracts/nodes-v1.openapi.yaml` in retina-server, verified by blob hash rather than by eye. `NodeConfig.beam_width_deg` was checked on the way in, as the working agreement requires: still nullable, a third revision running. The reason this is not a pure regeneration is `ClaimResponse.state`. It carries `title: State`, exactly like `HeartbeatRequest.state`, and datamodel-codegen names generated enum classes after that title. Faced with two it keeps the first and renames the second `State1`, and in 1.4.0 the loser is the node's own six-value state. That breaks quietly rather than loudly: `comms/lifecycle.py` and `wire/heartbeat.py` both do `from ...wire.models import State as WireState`, so they would go on importing a name that still exists and now means "where a claim stands". The first symptom would be `WireState.streaming` raising AttributeError while building a heartbeat, at runtime, on a node. Chasing the number would be worse than useless, since which enum gets demoted depends on declaration order in someone else's file. So `tools/normalise_spec.py` grows a second rewrite alongside the nullable one, on the same terms: applied to a temporary copy, never to the checked-in contract. Each enum in `NAMED_ENUMS` is hoisted into a component of that name and every inline occurrence replaced with a `$ref` to it. The claim enum becomes `ClaimState`, `State` keeps its meaning, and the three claim-state fields share one class instead of generating two identical ones under different names. The RootModel count is unchanged at 3. `tests/test_normalise_spec.py` is new, and covers both rewrites. The test worth keeping is `test_no_two_inline_enums_share_a_title`: it is the general form of this fault, so the next collision fails a test rather than silently renaming something. Co-Authored-By: Claude Opus 5 (1M context) --- docs/node-ingest-v1.yml | 359 +++++++++++++++++++++++++++++++- retina_telemetry/wire/models.py | 49 ++++- tests/test_normalise_spec.py | 140 +++++++++++++ tools/normalise_spec.py | 111 +++++++++- vulture_whitelist.py | 13 ++ 5 files changed, 656 insertions(+), 16 deletions(-) create mode 100644 tests/test_normalise_spec.py diff --git a/docs/node-ingest-v1.yml b/docs/node-ingest-v1.yml index 557d9f1..0ec0d02 100644 --- a/docs/node-ingest-v1.yml +++ b/docs/node-ingest-v1.yml @@ -1,7 +1,7 @@ openapi: 3.1.0 info: title: RETINA node ingest - version: 1.2.2 + version: 1.4.0 description: | The RETINA server's HTTP API. The paths under `/v1/nodes` are the RETINA node ingest contract, generated from the server and versioned as a unit; everything @@ -17,13 +17,20 @@ info: | Status | Means | |---|---| - | `400` | The body was refused. `invalid_config` for a configuration that failed validation, `invalid_contact` for contact details that failed it, `invalid_body` for a body that failed the schema | + | `400` | The body was refused. `invalid_config` for a configuration that failed validation, `invalid_contact` for contact details that failed it, `invalid_claim` for an address that failed it, `invalid_body` for a body that failed the schema | | `401` | The bearer token is bad, revoked or expired | | `403` | Registration refused, without saying why | - | `409` | The frame names a `config_version` this server never issued | + | `409` | The frame names a `config_version` this server never issued, or the node already has an owner | | `413` | The body exceeded the cap for its path | | `429` | Rate limited | + One refusal does not carry `Error`, and it is the only one: the `409` on + `PUT /v1/nodes/claim` answers with `ClaimResponse` instead. Nothing the node can + change makes that request succeed, so what it needs is not a slug naming the + field it got wrong but the address that already owns the node, which it + reconciles against without a second call. Every other refusal under this prefix + wears `Error`. + `403` is deliberately opaque: unknown device, not yet accepted by Mender, already holding a valid token and in cooldown are one response, at one latency, with one `Retry-After`. Registration is limited per `node_id` and answers that @@ -73,6 +80,8 @@ tags: description: Receiver and transmitter geometry, versioned by the server. - name: contact description: Whom to contact about the node, reported by the node itself. +- name: claim + description: The address that owns the node, and where its claim stands. paths: /v1/nodes/register: post: @@ -334,6 +343,258 @@ paths: - bearerAuth: [] x-cadence: on local change only x-max-body-bytes: 2048 + /v1/nodes/claim: + get: + tags: + - claim + summary: Read where this node's claim stands. + description: |- + Where this node's claim stands. Cheap, and safe to poll every few seconds while a setup + page is open; `HeartbeatResponse` carries the same three fields for the rest of the + node's life, so there is no reason to keep polling once the page closes. + operationId: getClaim + responses: + '200': + description: Where the claim stands. + content: + application/json: + schema: + $ref: '#/components/schemas/ClaimResponse' + '401': + description: 'Token bad, revoked or expired. Surface it locally and leave the address unsent: + nothing is lost by waiting for a credential, since an address can be offered at any point + in a node''s life. Do not re-register, for the reason detection and heartbeat give on their + own 401s.' + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: never + x-terminal: true + 5XX: + description: 'Server side, though a 502 or 504 is the edge''s own response rather than the origin''s, + and carries whatever body that gateway sends instead of this `Error`. Abandon the request + and retry, backing off with jitter if it persists: on detection and heartbeat that retry is + the node''s next frame or beat rather than a resend of this one, while registration and configuration + retry the identical body.' + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: backoff + x-terminal: false + security: + - bearerAuth: [] + x-cadence: every few seconds while a setup page is open, and not otherwise + put: + tags: + - claim + summary: Offer the address that owns this node. + description: |- + Offer the address that owns this node. The server mails it a link, and clicking that + link binds the node to the account behind the address, creating one if there is not + already one. + + Sending an address the node already holds is accepted and changes nothing, so this may + be resent on every configuration sync. It does not mail anything again: asking for the + mail again is `POST /v1/nodes/claim/resend`, because a write that mailed every time it + ran would mail on every sync. + + The response says where the claim stands, and is the same shape `GET` returns. Nothing + about the node's operation depends on any of it: an address that is never verified + grants nothing, and a node with no address at all runs unowned indefinitely and can be + claimed whenever its owner gets round to it. + + Poll the `GET` every few seconds while a setup page is open and somebody is waiting. + Afterwards stop: `HeartbeatResponse` carries the same three fields once a minute, which + is how a release performed months later reaches a node that stopped polling long ago. + operationId: putClaim + requestBody: + description: The address to claim this node with. + content: + application/json: + schema: + $ref: '#/components/schemas/NodeClaimRequest' + required: true + responses: + '200': + description: Where the claim stands after this call. + content: + application/json: + schema: + $ref: '#/components/schemas/ClaimResponse' + '400': + description: The address failed validation, as `invalid_claim`. A body that is not JSON at all + lands here too, since the remedy is the same. `detail` names the offending field and nothing + else, so retina-gui can put the message next to the input it belongs to. Retrying unchanged + will not help. + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: never + x-terminal: false + '401': + description: 'Token bad, revoked or expired. Surface it locally and leave the address unsent: + nothing is lost by waiting for a credential, since an address can be offered at any point + in a node''s life. Do not re-register, for the reason detection and heartbeat give on their + own 401s.' + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: never + x-terminal: true + '409': + description: 'The node already has an owner, so there is nothing for this call to do: a nomination + names an address that is not theirs, and a resend would mail a link for a claim already settled. + Nothing was written. The body is a `ClaimResponse` rather than an `Error`: this is the one + refusal a node reconciles from rather than corrects, so it carries the state and the address + that won. A node that believed itself unclaimed should adopt what this says and stop offering + to claim. Releasing is the owner''s to do, from the dashboard, and cannot be done from here.' + content: + application/json: + schema: + $ref: '#/components/schemas/ClaimResponse' + x-retry: never + x-terminal: false + '413': + description: 'The body exceeded this path''s cap, which is checked from `Content-Length` before + parsing. Per-field bounds in the schemas do not make this unreachable: `HeartbeatRequest.errors` + alone permits 16 KiB of strings, `RegisterRequest.config` carries no bound at all, and JSON + puts no length limit on a number''s text, so a schema-valid body can still exceed it.' + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: never + x-terminal: false + '429': + description: 'Rate limited, per token and per endpoint. The skipped frames are best dropped + rather than accumulated: a frame delivered late carries an old timestamp and is rejected by + the association gate rather than paired against whatever is current.' + headers: + Retry-After: + description: Seconds to wait before retrying. Honour it, then back off with jitter. + required: true + schema: + type: integer + maximum: 86400.0 + minimum: 0.0 + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: retry-after + x-terminal: false + 5XX: + description: 'Server side, though a 502 or 504 is the edge''s own response rather than the origin''s, + and carries whatever body that gateway sends instead of this `Error`. Abandon the request + and retry, backing off with jitter if it persists: on detection and heartbeat that retry is + the node''s next frame or beat rather than a resend of this one, while registration and configuration + retry the identical body.' + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: backoff + x-terminal: false + security: + - bearerAuth: [] + x-cadence: on local change only + x-max-body-bytes: 2048 + /v1/nodes/claim/resend: + post: + tags: + - claim + summary: Send the claim link again. + description: |- + Send the link again, to the address already on file. No body. + + Separate from `PUT` because a write that mailed every time it ran would mail on every + configuration sync. This is the explicit ask, and it is what a person who did not get + the first mail presses. + + Refused on a node that already has an owner, and answered without sending on an address + that bounced hard: a second copy to an address that does not exist earns nothing but a + second bounce, and bounces are how a sending domain loses its reputation. + operationId: resendClaim + responses: + '200': + description: Where the claim stands after this call. + content: + application/json: + schema: + $ref: '#/components/schemas/ClaimResponse' + '401': + description: 'Token bad, revoked or expired. Surface it locally and leave the address unsent: + nothing is lost by waiting for a credential, since an address can be offered at any point + in a node''s life. Do not re-register, for the reason detection and heartbeat give on their + own 401s.' + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: never + x-terminal: true + '409': + description: 'The node already has an owner, so there is nothing for this call to do: a nomination + names an address that is not theirs, and a resend would mail a link for a claim already settled. + Nothing was written. The body is a `ClaimResponse` rather than an `Error`: this is the one + refusal a node reconciles from rather than corrects, so it carries the state and the address + that won. A node that believed itself unclaimed should adopt what this says and stop offering + to claim. Releasing is the owner''s to do, from the dashboard, and cannot be done from here.' + content: + application/json: + schema: + $ref: '#/components/schemas/ClaimResponse' + x-retry: never + x-terminal: false + '413': + description: 'The body exceeded this path''s cap, which is checked from `Content-Length` before + parsing. Per-field bounds in the schemas do not make this unreachable: `HeartbeatRequest.errors` + alone permits 16 KiB of strings, `RegisterRequest.config` carries no bound at all, and JSON + puts no length limit on a number''s text, so a schema-valid body can still exceed it.' + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: never + x-terminal: false + '429': + description: 'Rate limited, per token and per endpoint. The skipped frames are best dropped + rather than accumulated: a frame delivered late carries an old timestamp and is rejected by + the association gate rather than paired against whatever is current.' + headers: + Retry-After: + description: Seconds to wait before retrying. Honour it, then back off with jitter. + required: true + schema: + type: integer + maximum: 86400.0 + minimum: 0.0 + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: retry-after + x-terminal: false + 5XX: + description: 'Server side, though a 502 or 504 is the edge''s own response rather than the origin''s, + and carries whatever body that gateway sends instead of this `Error`. Abandon the request + and retry, backing off with jitter if it persists: on detection and heartbeat that retry is + the node''s next frame or beat rather than a resend of this one, while registration and configuration + retry the identical body.' + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorBody' + x-retry: backoff + x-terminal: false + security: + - bearerAuth: [] + x-cadence: when somebody asks for the mail again, and never automatically + x-max-body-bytes: 2048 /v1/nodes/detection: post: tags: @@ -600,6 +861,36 @@ components: - publication title: Agreements description: Three separately versioned records, because they are withdrawn separately. + ClaimResponse: + properties: + state: + type: string + enum: + - unclaimed + - pending + - owned + title: State + email: + anyOf: + - type: string + - type: 'null' + title: Email + undeliverable: + type: boolean + title: Undeliverable + type: object + required: + - state + - email + - undeliverable + title: ClaimResponse + description: |- + What the node is told about its own claim. + + Carried by `GET` and `PUT /v1/nodes/claim` alike, and by the 409 a node + gets for nominating an address for a node that already has an owner: the + node has to reconcile from that answer rather than retry it, so it is told + which address won in the same response. ConfigResponse: properties: config_version: @@ -616,9 +907,27 @@ components: type: string format: date-time title: Updated At + claim_state: + type: string + enum: + - unclaimed + - pending + - owned + title: Claim State + claim_email: + anyOf: + - type: string + - type: 'null' + title: Claim Email + claim_undeliverable: + type: boolean + title: Claim Undeliverable type: object required: - updated_at + - claim_state + - claim_email + - claim_undeliverable title: ContactResponse DetectionAck: properties: @@ -784,13 +1093,57 @@ components: type: string pattern: ^(nde|sim)[0-9a-z]{12}$ title: Node Ref + claim_state: + type: string + enum: + - unclaimed + - pending + - owned + title: Claim State + claim_email: + anyOf: + - type: string + - type: 'null' + title: Claim Email + claim_undeliverable: + type: boolean + title: Claim Undeliverable type: object required: - server_time - config_stale - streaming_allowed - node_ref + - claim_state + - claim_email + - claim_undeliverable title: HeartbeatResponse + NodeClaimRequest: + type: object + title: NodeClaimRequest + description: |- + The address to claim this node with. The server mails it a link; clicking that + link binds the node to the account behind the address, creating one if there is + not already one. + + Required, and the only field. Sending an address the node already holds is + accepted and changes nothing, so this may be resent freely; it does not mail + anything again. Asking for the mail again is a separate call. + + Surrounding whitespace is stripped and the address is lower cased before it is + judged, so the bound here describes the trimmed form rather than the bytes sent. + + Nothing here verifies that the address exists. Until the link is clicked the + address grants nothing, and the node runs exactly as it would with no address + at all. + properties: + email: + type: string + maxLength: 255 + description: The owner's address, lower cased and trimmed by the server. + required: + - email + additionalProperties: false NodeConfig: type: object title: NodeConfig diff --git a/retina_telemetry/wire/models.py b/retina_telemetry/wire/models.py index f428590..9ef1d30 100644 --- a/retina_telemetry/wire/models.py +++ b/retina_telemetry/wire/models.py @@ -22,10 +22,6 @@ class ConfigResponse(BaseModel): config_version: Annotated[int, Field(ge=1, title="Config Version")] -class ContactResponse(BaseModel): - updated_at: Annotated[AwareDatetime, Field(title="Updated At")] - - class DetectionAck(BaseModel): accepted: Annotated[int, Field(ge=0, title="Accepted")] config_stale: Annotated[bool, Field(title="Config Stale")] @@ -68,11 +64,17 @@ class Error(RootModel[str]): root: Annotated[str, Field(max_length=512)] -class HeartbeatResponse(BaseModel): - server_time: Annotated[AwareDatetime, Field(title="Server Time")] - config_stale: Annotated[bool, Field(title="Config Stale")] - streaming_allowed: Annotated[bool, Field(title="Streaming Allowed")] - node_ref: Annotated[str, Field(pattern="^(nde|sim)[0-9a-z]{12}$", title="Node Ref")] +class NodeClaimRequest(BaseModel): + model_config = ConfigDict( + extra="forbid", + ) + email: Annotated[ + str, + Field( + description="The owner's address, lower cased and trimmed by the server.", + max_length=255, + ), + ] class NodeConfig(BaseModel): @@ -160,6 +162,12 @@ class RegisterResponse(BaseModel): server_time: Annotated[AwareDatetime, Field(title="Server Time")] +class ClaimState(StrEnum): + unclaimed = "unclaimed" + pending = "pending" + owned = "owned" + + class Agreements(BaseModel): model_config = ConfigDict( extra="forbid", @@ -169,6 +177,19 @@ class Agreements(BaseModel): publication: PublicationChoice +class ClaimResponse(BaseModel): + state: ClaimState + email: Annotated[str | None, Field(title="Email")] + undeliverable: Annotated[bool, Field(title="Undeliverable")] + + +class ContactResponse(BaseModel): + updated_at: Annotated[AwareDatetime, Field(title="Updated At")] + claim_state: ClaimState + claim_email: Annotated[str | None, Field(title="Claim Email")] + claim_undeliverable: Annotated[bool, Field(title="Claim Undeliverable")] + + class HeartbeatRequest(BaseModel): model_config = ConfigDict( extra="forbid", @@ -182,6 +203,16 @@ class HeartbeatRequest(BaseModel): errors: Annotated[list[Error] | None, Field(max_length=32, title="Errors")] = None +class HeartbeatResponse(BaseModel): + server_time: Annotated[AwareDatetime, Field(title="Server Time")] + config_stale: Annotated[bool, Field(title="Config Stale")] + streaming_allowed: Annotated[bool, Field(title="Streaming Allowed")] + node_ref: Annotated[str, Field(pattern="^(nde|sim)[0-9a-z]{12}$", title="Node Ref")] + claim_state: ClaimState + claim_email: Annotated[str | None, Field(title="Claim Email")] + claim_undeliverable: Annotated[bool, Field(title="Claim Undeliverable")] + + class RegisterRequest(BaseModel): model_config = ConfigDict( extra="forbid", diff --git a/tests/test_normalise_spec.py b/tests/test_normalise_spec.py new file mode 100644 index 0000000..580ec94 --- /dev/null +++ b/tests/test_normalise_spec.py @@ -0,0 +1,140 @@ +"""The two rewrites that stand between the contract and the generated models. + +Neither changes what the spec means. Both change what `datamodel-codegen` +builds from it, and in ways that reach every call site, so they are worth +holding still. +""" + +from pathlib import Path + +import yaml + +from tools.normalise_spec import name_enums, normalise + +SPEC = Path(__file__).resolve().parents[1] / "docs" / "node-ingest-v1.yml" + +CLAIM_VALUES = ["unclaimed", "pending", "owned"] +NODE_STATE_VALUES = ["starting", "streaming", "stalled", "paused", "error", "stopping"] + + +def spec() -> dict: + return yaml.safe_load(SPEC.read_text(encoding="utf-8")) + + +def inline_enums(node, found=None) -> list[tuple[str | None, tuple[str, ...]]]: + """Every enum still declared inline, as (title, values).""" + found = [] if found is None else found + if isinstance(node, list): + for item in node: + inline_enums(item, found) + elif isinstance(node, dict): + values = node.get("enum") + if isinstance(values, list): + found.append((node.get("title"), tuple(values))) + for value in node.values(): + inline_enums(value, found) + return found + + +# ── nullable fields ────────────────────────────────────────────────── + + +def test_a_nullable_anyOf_becomes_a_type_array(): + """The whole reason this module exists: `anyOf` makes the generator + manufacture a RootModel wrapper per field, a type array does not.""" + rewritten = normalise({"anyOf": [{"type": "number", "minimum": -90}, {"type": "null"}]}) + + assert rewritten == {"type": ["number", "null"], "minimum": -90} + + +def test_a_nullable_ref_is_left_alone(): + """`HeartbeatRequest.health`. No type to hoist, and it already generates + correctly as `NodeHealth | None`.""" + union = {"anyOf": [{"$ref": "#/components/schemas/NodeHealth"}, {"type": "null"}]} + + assert normalise(union) == union + + +def test_a_union_that_is_not_the_idiom_is_left_alone(): + union = {"anyOf": [{"type": "string"}, {"type": "integer"}]} + + assert normalise(union) == union + + +# ── enums that would collide ───────────────────────────────────────── + + +def test_a_named_enum_is_hoisted_to_a_component(): + document = name_enums({"properties": {"state": {"type": "string", "enum": CLAIM_VALUES}}}) + + assert document["properties"]["state"] == {"$ref": "#/components/schemas/ClaimState"} + assert document["components"]["schemas"]["ClaimState"]["enum"] == CLAIM_VALUES + + +def test_every_occurrence_refs_the_one_component(): + """Three fields spell this enum in 1.4.0. One class, not three.""" + document = name_enums( + { + "a": {"type": "string", "enum": CLAIM_VALUES, "title": "State"}, + "b": {"type": "string", "enum": CLAIM_VALUES, "title": "Claim State"}, + "c": {"type": "string", "enum": CLAIM_VALUES, "title": "Claim State"}, + } + ) + + refs = {document[key]["$ref"] for key in "abc"} + assert refs == {"#/components/schemas/ClaimState"} + assert list(document["components"]["schemas"]) == ["ClaimState"] + + +def test_the_hoisted_component_is_not_rewritten_into_a_ref_to_itself(): + document = name_enums({"components": {"schemas": {}}, "x": {"enum": CLAIM_VALUES}}) + + assert document["components"]["schemas"]["ClaimState"]["enum"] == CLAIM_VALUES + + +def test_an_enum_we_do_not_name_is_untouched(): + """The node's own state stays exactly as the contract declares it.""" + node_state = {"type": "string", "enum": NODE_STATE_VALUES, "title": "State"} + + assert name_enums({"state": node_state})["state"] == node_state + + +# ── against the real contract ──────────────────────────────────────── + + +def test_the_node_state_enum_keeps_its_name(): + """The regression this was written for. + + `ClaimResponse.state` is titled `State`, exactly like `HeartbeatRequest.state`. + The generator names enum classes after that title and renames the loser + `State1`, so without the hoist the node state becomes `State1` while + `comms/lifecycle.py` and `wire/heartbeat.py` go on importing `State` and + silently get the claim enum instead. + """ + document = name_enums(normalise(spec())) + beat = document["components"]["schemas"]["HeartbeatRequest"]["properties"]["state"] + + assert beat["title"] == "State" + assert beat["enum"] == NODE_STATE_VALUES + + +def test_no_two_inline_enums_share_a_title(): + """The general form, so the next collision fails here rather than in a + rename nobody reads. A revision that introduces one adds a line to + NAMED_ENUMS.""" + by_title: dict[str | None, set[tuple[str, ...]]] = {} + for title, values in inline_enums(name_enums(normalise(spec()))): + by_title.setdefault(title, set()).add(values) + + collisions = {title: shapes for title, shapes in by_title.items() if len(shapes) > 1} + assert collisions == {} + + +def test_the_claim_enum_is_the_only_thing_hoisted(): + """Keeps the table honest: an entry that stops matching the contract shows + up as a component nothing refs.""" + document = name_enums(normalise(spec())) + schemas = document["components"]["schemas"] + + assert schemas["ClaimState"]["enum"] == CLAIM_VALUES + assert [name for name in schemas if name.endswith("State")] == ["ClaimState"] diff --git a/tools/normalise_spec.py b/tools/normalise_spec.py index 4c9f217..9f3e241 100644 --- a/tools/normalise_spec.py +++ b/tools/normalise_spec.py @@ -1,11 +1,17 @@ #!/usr/bin/env python3 -"""Rewrite the spec's nullable spelling into the one the generator understands. +"""Rewrite two spec idioms into the spellings the generator understands. -**This does not change the contract.** It reads the spec, rewrites one JSON -Schema idiom into an exactly equivalent one, and writes the result somewhere +**This does not change the contract.** It reads the spec, rewrites two JSON +Schema idioms into exactly equivalent ones, and writes the result somewhere else. ``docs/node-ingest-v1.yml`` is never touched: it stays byte-identical to what the server author sent, which is the whole point of keeping it read-only. +Both rewrites exist for the same reason: ``datamodel-codegen`` produces a +different *shape* for two spellings that mean the same thing, and the shape we +want is not the one the server's FastAPI export happens to emit. + +# 1. Nullable fields + ## The idiom OpenAPI 3.1 has two ways to say "a number between -90 and 90, or null":: @@ -52,6 +58,48 @@ class RxLat(RootModel[float]): so it passes through untouched. Anything else (a three-way union, a ``oneOf``, a bare ``anyOf`` with no null member) is not this idiom and is left exactly as written. + +# 2. Inline enums that collide on their title + +## The idiom + +The contract declares enums inline on the property rather than as named +components, and carries the name in ``title``:: + + state: claim_state: + type: string type: string + enum: [starting, streaming, enum: [unclaimed, pending, owned] + stalled, paused, title: Claim State + error, stopping] + title: State + +## Why it matters + +``datamodel-codegen`` names a generated enum class after that title, and two +unrelated enums in 1.4.0 both answer to ``State``: the node's own six-value +state on ``HeartbeatRequest``, and where a claim stands on ``ClaimResponse``. +Faced with the collision the generator keeps the first it meets and renames the +second ``State1``. + +Which one loses is a function of declaration order in someone else's file. In +1.4.0 it is the node state that becomes ``State1``, and the damage is silent: +``comms/lifecycle.py`` and ``wire/heartbeat.py`` both do ``import State as +WireState``, so they would go on importing a name that still exists and now +means something else entirely. The first sign would be ``WireState.streaming`` +raising ``AttributeError`` while building a heartbeat. + +A positional name cannot be depended on either way, so the fix is to stop the +collision happening rather than to chase the number. + +## What is rewritten + +Each enum listed in ``NAMED_ENUMS`` is hoisted into a component schema of that +name and every inline occurrence replaced by a ``$ref`` to it. The generator +then emits one class, under a name chosen here, shared by every field that +refs it, which is also what the three claim-state fields should have been all +along, since they are the same three values in all three places. + +Nothing else is touched. An enum not in the table generates exactly as before. """ from __future__ import annotations @@ -63,6 +111,19 @@ class RxLat(RootModel[float]): NULL_BRANCH = {"type": "null"} +#: Enums to hoist out of the properties that declare them, and the name each +#: one gets. Keyed by the values because the values are what identify an enum: +#: the titles are exactly what cannot be trusted here, and the three +#: claim-state fields do not agree on one anyway (``State`` on +#: ``ClaimResponse.state``, ``Claim State`` on the other two). +#: +#: One entry, added when 1.4.0 introduced a second enum titled ``State``. +#: A revision that collides again adds a line here rather than renaming +#: whatever the generator happened to demote that time. +NAMED_ENUMS: dict[tuple[str, ...], str] = { + ("unclaimed", "pending", "owned"): "ClaimState", +} + def normalise(node: Any) -> Any: """Depth-first rewrite of the nullable-``anyOf`` idiom. @@ -102,6 +163,47 @@ def _nullable_branch(any_of: Any) -> dict[str, Any] | None: return others[0] if isinstance(others[0], dict) and "type" in others[0] else None +def name_enums(document: Any) -> Any: + """Hoist every enum in ``NAMED_ENUMS`` into a component of that name. + + Runs over the whole document and replaces each inline occurrence with a + ``$ref``, then adds the components themselves. Adding them afterwards is + what stops the hoisted copy being rewritten into a reference to itself. + """ + hoisted: dict[str, dict[str, Any]] = {} + + def rewrite(node: Any) -> Any: + if isinstance(node, list): + return [rewrite(item) for item in node] + if not isinstance(node, dict): + return node + + name = _named_enum(node) + if name is None: + return {key: rewrite(value) for key, value in node.items()} + + # Everything but the title, which is the one key the occurrences + # disagree on and the one being replaced. Anything else they carry + # (a description, say) comes along from the first occurrence seen; + # today they carry nothing else. + hoisted.setdefault(name, {**{k: v for k, v in node.items() if k != "title"}, "title": name}) + return {"$ref": f"#/components/schemas/{name}"} + + document = rewrite(document) + if hoisted: + schemas = document.setdefault("components", {}).setdefault("schemas", {}) + schemas.update(hoisted) + return document + + +def _named_enum(node: dict[str, Any]) -> str | None: + """The name this node should be hoisted under, if it is one we name.""" + values = node.get("enum") + if not isinstance(values, list) or not all(isinstance(value, str) for value in values): + return None + return NAMED_ENUMS.get(tuple(values)) + + def main() -> int: if len(sys.argv) != 3: print(f"usage: {sys.argv[0]} ", file=sys.stderr) @@ -112,7 +214,8 @@ def main() -> int: document = yaml.safe_load(handle) with open(target, "w", encoding="utf-8") as handle: - yaml.safe_dump(normalise(document), handle, sort_keys=False, allow_unicode=True) + rewritten = name_enums(normalise(document)) + yaml.safe_dump(rewritten, handle, sort_keys=False, allow_unicode=True) return 0 diff --git a/vulture_whitelist.py b/vulture_whitelist.py index 4bca2af..f3a83f6 100644 --- a/vulture_whitelist.py +++ b/vulture_whitelist.py @@ -81,6 +81,19 @@ email phone country +# The claim endpoints, added in v1.3.0, and the enum v1.4.0 put on the +# heartbeat and contact responses alongside them. Nothing here offers an +# address: a node has no source for the owner's, so `PUT /v1/nodes/claim` has +# no caller until retina-gui collects one. The state itself *is* read, in +# comms/levels.py, but never through these members. It is passed through as +# the string the server sent so that a value this node does not recognise +# still reaches an operator rather than being dropped. +NodeClaimRequest +ClaimResponse +unclaimed +pending +owned + # read by socketserver.ThreadingMixIn # tools/mock_server.py:392 _.daemon_threads From b81f9c84fc3925d3ff91e1e18b069f3788c215c4 Mon Sep 17 00:00:00 2001 From: Josh Poole Date: Tue, 22 Sep 2026 11:36:47 +0100 Subject: [PATCH 2/7] 20260922 - Tell the owner where their claim stands v1.4.0 restates `claim_state`, `claim_email` and `claim_undeliverable` on every heartbeat and contact response. This reads them and writes them to the status document, which is the whole point of the revision: the node binds no ports, so an owner waiting on a claim link, or one whose address bounced, learns it there or nowhere. retina-gui already reads that file. None of it gates anything. An unclaimed node registers, streams and beats exactly as an owned one does; ownership decides who sees the data at the far end, not whether it is sent. `Claim` is carried on `Snapshot` and applied by `State.apply_levels`, which is where it belongs: the server restates the block on every beat, so it is a level rather than an edge, and a node that missed one is told again a minute later. That is also how a claim performed months after setup reaches a node that was never watching for it. It is one argument rather than three, and read as a unit gated on `claim_state`. That is the bug this shape exists to prevent. `claim_state` is required and non-nullable wherever it appears, so its presence says the block is there; `claim_email` is required and *nullable*, so its null is a value, meaning nobody has offered an address. Read field by field, a detection ack carries neither, so its absent `claim_email` would be indistinguishable from a cleared one and the ack arriving seconds after a heartbeat would blank the address. Both directions are pinned by tests. `Claim.state` is a plain string rather than the generated `ClaimState`. A fourth value from a later server should reach an operator as itself, not take the address and the bounce flag down with it on the way to being dropped. The log line says whether an address is on file and never what it is, matching what `collect/contact.py` already does when it logs the names of unrecognised fields rather than their values. The owner sees their own address in the status document; a container log is not the place for it. It logs only on a change, because an unchanged level arrives once a minute forever. Three notes for whoever picks up the rest of the claim work: - Nothing here calls a claim endpoint, and this deliberately stops short of that. `PUT /v1/nodes/claim` needs the address that owns the node, which nothing on a node knows. It is not the contact email: that answers "whom do we ring", and reusing it would mail a link that hands the node to whoever happened to be listed. The server has how an address relates to an account open as its own question. - `pending` is not durable. A claim link declined while a second is outstanding can strand a node there for about fifteen minutes before it falls back to `unclaimed`. That is a known server-side race, tracked there. Nothing should wait on `pending` or read reaching it as progress. - The mock holds the three values rather than deriving them, and exposes them through `/_control/levels`. The claim endpoints are deliberately not implemented in it: a mock that answers calls nobody makes would be asserting a shape the node has never had to parse. Co-Authored-By: Claude Opus 5 (1M context) --- retina_telemetry/comms/levels.py | 32 ++++++++++- retina_telemetry/state.py | 67 ++++++++++++++++++++- tests/comms/test_levels.py | 99 +++++++++++++++++++++++++++++++- tests/test_mock_server.py | 12 +++- tests/test_state.py | 31 +++++++++- tests/test_status.py | 28 ++++++++- tools/mock_server.py | 31 +++++++++- 7 files changed, 293 insertions(+), 7 deletions(-) diff --git a/retina_telemetry/comms/levels.py b/retina_telemetry/comms/levels.py index e1c434a..fcbb691 100644 --- a/retina_telemetry/comms/levels.py +++ b/retina_telemetry/comms/levels.py @@ -29,7 +29,7 @@ from typing import Any from retina_telemetry.comms.client import Kind, Outcome -from retina_telemetry.state import State +from retina_telemetry.state import Claim, State log = logging.getLogger(__name__) @@ -64,12 +64,42 @@ def apply_response(outcome: Outcome, state: State) -> None: streaming_allowed=_bool_or_none(body.get("streaming_allowed")), node_ref=_str_or_none(body.get("node_ref")), server_time=server_time, + claim=_claim(body), ) if server_time is not None: _warn_on_clock_offset(state) +def _claim(body: dict[str, Any]) -> Claim | None: + """The claim block, if this response carried one. + + Read as a unit, and gated on ``claim_state`` alone, because that is the + field whose presence says the block is there at all: it is required and + non-nullable everywhere it appears, while ``claim_email`` is required and + *nullable*, so its ``null`` is a value (nobody has offered an address) + rather than an absence. + + Reading the three independently would confuse those two cases in the + direction that loses information. A detection ack carries none of them, so + an absent ``claim_email`` would look exactly like a cleared one and the + ack would wipe an address the heartbeat reported a second earlier. + + An unrecognised ``claim_state`` is passed through rather than rejected. + This is a display value, and a server that grows a fourth one should show + up in the status document as itself rather than vanish. + """ + state = body.get("claim_state") + if not isinstance(state, str) or not state: + return None + + return Claim( + state=state, + email=_str_or_none(body.get("claim_email")), + undeliverable=body.get("claim_undeliverable") is True, + ) + + def _warn_on_clock_offset(state: State) -> None: offset = state.snapshot().clock_offset_s if offset is not None and abs(offset) > CLOCK_WARN_S: diff --git a/retina_telemetry/state.py b/retina_telemetry/state.py index 46b8faa..e1a93c2 100644 --- a/retina_telemetry/state.py +++ b/retina_telemetry/state.py @@ -61,7 +61,7 @@ import os import secrets import threading -from dataclasses import dataclass +from dataclasses import asdict, dataclass from datetime import UTC, datetime from pathlib import Path from time import monotonic @@ -77,6 +77,40 @@ TOKEN_MODE = 0o600 +@dataclass(frozen=True) +class Claim: + """Where this node's claim stands, as the server last stated it. + + One object rather than three fields on :class:`Snapshot`, because the + three arrive together and the distinction that matters is whether the + server has said anything at all. ``Snapshot.claim is None`` means it has + not: no response carrying the block has landed yet. That is not the same + as :attr:`state` being ``"unclaimed"``, which is the server positively + saying nobody owns this node. + + Nothing here gates anything. An unclaimed node registers, streams and + beats exactly as an owned one does; ownership decides who sees the data on + the far end, not whether it is sent. This is carried so the status + document can show an owner where they stand. + + :attr:`state` is a plain string rather than the generated ``ClaimState``. + A value this node does not recognise is still worth showing an operator, + and dropping the whole block because one field grew a fourth value would + lose the other two as well. + + **Do not build anything that waits on ``"pending"``.** It is not durable: + a claim link that is declined while a second is outstanding can leave a + node reading ``pending`` against a link nobody can redeem, until the + challenge expires about fifteen minutes later and it reads ``unclaimed`` + again. That is a known server-side race, tracked there, and the node's + part is simply to keep reporting whatever it was last told. + """ + + state: str + email: str | None + undeliverable: bool + + @dataclass(frozen=True) class Snapshot: """A consistent view of shared state, taken under one lock acquisition.""" @@ -91,6 +125,7 @@ class Snapshot: token_rejected: bool clock_offset_s: float | None process_uptime_s: int + claim: Claim | None @property def registered(self) -> bool: @@ -138,6 +173,10 @@ def redacted(self) -> dict[str, Any]: "token_rejected": self.token_rejected, "clock_offset_s": self.clock_offset_s, "process_uptime_s": self.process_uptime_s, + # Nested, and ``null`` until a response carries it. Three flat keys + # could not say "not yet told" without a null boolean, which reads + # as a bug at the other end. + "claim": None if self.claim is None else asdict(self.claim), } @@ -171,6 +210,7 @@ def __init__(self, token_path: Path | str = DEFAULT_TOKEN_PATH) -> None: self._streaming_allowed = True self._token_rejected = False self._clock_offset_s: float | None = None + self._claim: Claim | None = None #: Set when a configuration resend is due — because the server said #: ``config_stale``, because a detection POST returned 409, or because @@ -195,6 +235,7 @@ def snapshot(self) -> Snapshot: token_rejected=self._token_rejected, clock_offset_s=self._clock_offset_s, process_uptime_s=int(monotonic() - self._started_at), + claim=self._claim, ) # ── writing ────────────────────────────────────────────────────── @@ -240,6 +281,7 @@ def apply_levels( streaming_allowed: bool | None = None, node_ref: str | None = None, server_time: datetime | None = None, + claim: Claim | None = None, ) -> None: """Adopt what a response carried, in one atomic step. @@ -251,6 +293,12 @@ def apply_levels( Only the values actually present are applied; ``None`` means the response did not carry that field, which is not the same as false. Nothing here touches the disk. + + ``claim`` arrives on the heartbeat and contact responses only, and a + detection ack carries no part of it. That is why it is one argument + rather than three: a caller cannot half-apply it, so an endpoint that + says nothing about the claim cannot blank an address another endpoint + reported a second earlier. """ with self._lock: if config_version is not None: @@ -268,6 +316,23 @@ def apply_levels( if config_stale is not None: self._config_stale = config_stale + if claim is not None and claim != self._claim: + # Worth a line each time it moves: it changes rarely, and an + # owner waiting on a link has no other way to see that the + # node heard back. Whether there is an address, never the + # address itself, which is the same rule collect/contact.py + # follows when it logs the names of fields it did not + # recognise rather than their values. The status document is + # where an owner sees their own address; a container log is + # not. + log.info( + "claim is now %s (%s)%s", + claim.state, + "address on file" if claim.email else "no address on file", + ", and the last mail bounced" if claim.undeliverable else "", + ) + self._claim = claim + if server_time is not None: self._clock_offset_s = ( datetime.now(UTC) - server_time.astimezone(UTC) diff --git a/tests/comms/test_levels.py b/tests/comms/test_levels.py index 6e0498a..0d3a431 100644 --- a/tests/comms/test_levels.py +++ b/tests/comms/test_levels.py @@ -8,7 +8,7 @@ from retina_telemetry.comms.client import Kind, Outcome from retina_telemetry.comms.levels import apply_response -from retina_telemetry.state import State +from retina_telemetry.state import Claim, State @pytest.fixture @@ -121,6 +121,103 @@ def test_junk_values_are_ignored_rather_than_adopted(state): assert snapshot.node_ref == "nd_original" +# ── the claim block ────────────────────────────────────────────────── + + +#: The rest of a heartbeat response, so each test below varies only the claim. +LEVELS = { + "server_time": "2026-09-22T10:00:00Z", + "config_stale": False, + "streaming_allowed": True, + "node_ref": "nde4f2k9xq7m3b8", +} + +CLAIM = { + "claim_state": "pending", + "claim_email": "owner@example.com", + "claim_undeliverable": False, +} + + +def test_a_heartbeat_response_carries_the_claim(state): + apply_response(outcome(**LEVELS, **CLAIM), state) + + claim = state.snapshot().claim + assert claim == Claim(state="pending", email="owner@example.com", undeliverable=False) + + +def test_a_contact_response_carries_it_too(state): + """The other endpoint that restates it. Same three fields, same handling.""" + apply_response(outcome(updated_at="2026-09-22T10:00:00Z", **CLAIM), state) + + assert state.snapshot().claim.state == "pending" + + +def test_nothing_is_known_until_a_response_says_so(state): + """Distinct from `unclaimed`, which is the server saying nobody owns it.""" + assert state.snapshot().claim is None + + apply_response(outcome(**LEVELS), state) # a heartbeat from before 1.4.0 + + assert state.snapshot().claim is None + + +def test_a_detection_ack_does_not_clear_the_claim(state): + """The failure this block is read as a unit to prevent. + + A detection ack carries no claim fields at all. Read field by field, its + absent `claim_email` would be indistinguishable from a cleared one, and the + ack that arrives seconds after a heartbeat would blank the address. + """ + apply_response(outcome(**LEVELS, **CLAIM), state) + + apply_response(outcome(accepted=1, config_stale=False, streaming_allowed=True), state) + + assert state.snapshot().claim.email == "owner@example.com" + + +def test_an_owner_who_releases_the_node_clears_the_address(state): + """The other half of the same rule: a null *inside* a block that is present + is a value, and has to be adopted.""" + apply_response(outcome(**LEVELS, **CLAIM), state) + + apply_response( + outcome( + **LEVELS, + claim_state="unclaimed", + claim_email=None, + claim_undeliverable=False, + ), + state, + ) + + assert state.snapshot().claim == Claim(state="unclaimed", email=None, undeliverable=False) + + +def test_a_bounced_address_is_carried(state): + """The one claim value that is actionable: the owner will never get the + link, and nothing but this says so.""" + apply_response(outcome(**LEVELS, **CLAIM | {"claim_undeliverable": True}), state) + + assert state.snapshot().claim.undeliverable + + +def test_an_unrecognised_state_is_passed_through(state): + """Shown to the operator as itself rather than dropped. A fourth value + would otherwise take the address and the bounce flag down with it.""" + apply_response(outcome(**LEVELS, **CLAIM | {"claim_state": "disputed"}), state) + + assert state.snapshot().claim.state == "disputed" + + +def test_a_claim_state_that_is_not_a_string_takes_nothing_with_it(state): + """The gate is `claim_state`, so junk there means no block rather than a + half-applied one.""" + apply_response(outcome(**LEVELS, **CLAIM | {"claim_state": 7}), state) + + assert state.snapshot().claim is None + + # ── clock offset ───────────────────────────────────────────────────── diff --git a/tests/test_mock_server.py b/tests/test_mock_server.py index d3fe36f..2dd5af2 100644 --- a/tests/test_mock_server.py +++ b/tests/test_mock_server.py @@ -153,7 +153,17 @@ def test_heartbeat_restates_the_levels(server): status, body, _ = post(f"{server.url}/nodes/heartbeat", beat(version), token) assert status == 200 - assert set(body) == {"server_time", "config_stale", "streaming_allowed", "node_ref"} + assert set(body) == { + "server_time", + "config_stale", + "streaming_allowed", + "node_ref", + # Required on HeartbeatResponse since 1.4.0, so their absence would be + # the mock departing from the contract rather than a lean response. + "claim_state", + "claim_email", + "claim_undeliverable", + } def test_empty_frame_is_accepted(server): diff --git a/tests/test_state.py b/tests/test_state.py index 57ab05e..fa38629 100644 --- a/tests/test_state.py +++ b/tests/test_state.py @@ -4,7 +4,7 @@ import pytest -from retina_telemetry.state import State, with_uptime_fallback +from retina_telemetry.state import Claim, State, with_uptime_fallback @pytest.fixture @@ -241,6 +241,35 @@ def test_the_token_never_appears_in_the_redacted_view(state): assert "tok_abc123" not in json.dumps(state.snapshot().redacted()) +def test_the_owner_address_never_reaches_a_log_line(state, caplog): + """Same rule `collect/contact.py` follows when it logs the names of fields + it did not recognise rather than their values. An owner sees their own + address in the status document; a container log is not the place for it. + """ + registered(state) + + with caplog.at_level("INFO"): + state.apply_levels( + claim=Claim(state="pending", email="owner@example.com", undeliverable=False) + ) + + assert "owner@example.com" not in caplog.text + assert "pending" in caplog.text # the move itself is still worth a line + + +def test_a_claim_is_logged_once_rather_than_on_every_beat(state, caplog): + """It arrives on every heartbeat response. Logging an unchanged level once + a minute would bury everything else.""" + registered(state) + claim = Claim(state="owned", email="owner@example.com", undeliverable=False) + state.apply_levels(claim=claim) + + with caplog.at_level("INFO"): + state.apply_levels(claim=claim) + + assert "claim is now" not in caplog.text + + # ── response levels ────────────────────────────────────────────────── diff --git a/tests/test_status.py b/tests/test_status.py index 343600c..9bb7e33 100644 --- a/tests/test_status.py +++ b/tests/test_status.py @@ -3,7 +3,7 @@ import pytest -from retina_telemetry.state import State +from retina_telemetry.state import Claim, State from retina_telemetry.status import SCHEMA, StatusWriter @@ -67,6 +67,32 @@ def test_node_ref_reaches_the_owner(path, state): assert written(path)["node_ref"] == "nde4f2k9xq7m3b8" +def test_the_claim_reaches_the_owner_too(path, state): + """The same argument as `node_ref`, and the only one there is. + + This service binds no ports, so an owner waiting on a claim link, or one + whose address bounced, learns it here or not at all. retina-gui reads this + file already. + """ + state.apply_levels(claim=Claim(state="pending", email="owner@example.com", undeliverable=True)) + + StatusWriter(path).write(state="streaming", snapshot=state.snapshot()) + + assert written(path)["claim"] == { + "state": "pending", + "email": "owner@example.com", + "undeliverable": True, + } + + +def test_an_unasked_claim_is_null_rather_than_absent(path, state): + """Nested and null, not three flat keys. "The server has not told us" + needs saying, and a null boolean would read as a bug at the other end.""" + StatusWriter(path).write(state="streaming", snapshot=state.snapshot()) + + assert written(path)["claim"] is None + + def test_a_missing_identity_is_the_headline(path, tmp_path): """One of the three conditions that must reach a human.""" fresh = State(tmp_path / "absent-token") diff --git a/tools/mock_server.py b/tools/mock_server.py index c3d8848..ae499e2 100644 --- a/tools/mock_server.py +++ b/tools/mock_server.py @@ -724,6 +724,25 @@ class NodeRecord: contact: dict[str, Any] = field(default_factory=dict) contact_updated_at: str | None = None + #: Where the claim stands, restated on every heartbeat and contact + #: response. The mock holds it rather than deriving it, because the claim + #: endpoints that move it are not implemented here: nothing in this + #: service calls them yet, and a mock that answers calls nobody makes + #: would be asserting a shape the node has never had to parse. These are + #: set through the control channel, which is what the node reads them as + #: anyway: levels it is told and does not negotiate. + claim_state: str = "unclaimed" + claim_email: str | None = None + claim_undeliverable: bool = False + + def claim_block(self) -> dict[str, Any]: + """The three claim fields, spelled as every response carrying them does.""" + return { + "claim_state": self.claim_state, + "claim_email": self.claim_email, + "claim_undeliverable": self.claim_undeliverable, + } + def upsert_config(self, config: dict[str, Any]) -> int: """Return the active version, minting one only if the values differ. @@ -1034,6 +1053,15 @@ def _control(self) -> None: node.status = "active" if body["streaming_allowed"] else "blocked" if "node_ref" in body: node.node_ref = str(body["node_ref"]) + # The claim knobs. Unlike streaming_allowed these are held + # rather than derived, so they are set straight through. + if "claim_state" in body: + node.claim_state = str(body["claim_state"]) + if "claim_email" in body: + raw = body["claim_email"] + node.claim_email = None if raw is None else str(raw) + if "claim_undeliverable" in body: + node.claim_undeliverable = bool(body["claim_undeliverable"]) # Likewise config_stale: the server compares the version the # node reported against the active one. Moving the active # version is how staleness actually arises, and it is the @@ -1250,6 +1278,7 @@ def _heartbeat(self, body: Any) -> None: "config_stale": beat.config_version != node.active_config_version, "streaming_allowed": node.status == "active", "node_ref": node.node_ref, + **node.claim_block(), }, ) @@ -1303,7 +1332,7 @@ def _contact(self, body: Any) -> None: with self.state.lock: node.contact = contact node.contact_updated_at = _now() - self._send(200, {"updated_at": node.contact_updated_at}) + self._send(200, {"updated_at": node.contact_updated_at, **node.claim_block()}) class MockServer: From 062b4bac4a08a6d2e09f64b059b90d944a5947b9 Mon Sep 17 00:00:00 2001 From: Josh Poole Date: Tue, 22 Sep 2026 11:36:53 +0100 Subject: [PATCH 3/7] 20260922 - Write down why the claim address is not the contact email MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both docs still said v1.2.2, and neither mentioned a claim. The fact worth recording is the one that will otherwise be re-derived wrongly, probably by someone being helpful: `telemetry-contact.json` already holds an `email`, and `PUT /v1/nodes/claim` wants an email, so the two look like the same field. They are different questions. The contact address answers "whom do we ring about this node" and is optional throughout. The claim address answers "who owns this node", and getting it wrong mails a stranger a link that hands them the node. An owner may well give the same address twice; that is their answer to two questions, not a licence to infer the second from the first. The same discipline the consent records are held to, and for the same reason: it reaches a person. `docs/data-sources.md` gains that comparison in §4, next to the contact document it will be confused with, plus what a node *is* told since 1.4.0 and the fact that `pending` is not durable. `CLAUDE.md` moves to 1.4.0, records the beam-field check on this adoption, and gains a row in the retina-gui table for the address nobody collects yet, marked "no, but" rather than blocking: a node that is never claimed still registers, streams and beats. The required-and-nullable count stays fourteen. The three new nullables are response-side, and `tests/wire/test_serialise.py` counts payload schemas, which is the only place `exclude_none` could do damage. Co-Authored-By: Claude Opus 5 (1M context) --- CLAUDE.md | 21 +++++++++++++++++---- docs/data-sources.md | 37 +++++++++++++++++++++++++++++++++++++ 2 files changed, 54 insertions(+), 4 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 2a895c1..410948d 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -4,7 +4,7 @@ The node-side telemetry uplink for the RETINA passive radar fleet. One container node, owning everything sent to the server: registration, detection streaming, heartbeat, config sync. Nothing else on the node talks to `api.retina.fm`. -**Status: built, and implementing spec v1.2.2.** Verified end to end on the Owl node +**Status: built, and implementing spec v1.4.0.** Verified end to end on the Owl node against a tunnelled mock — every endpoint, every reachable state including `stalled`, and the refusal paths. @@ -105,6 +105,7 @@ Nothing here is buildable from this repo, and the first one blocks every node: | Read `/data/retina-telemetry/status.json` | no, but | We bind no ports, so it is the only way *no identity*, *revoked token* and *rejected config* reach an operator. `telemetry_status.py` reads it and the home page shows it | | Collect `location.rx.beam_width` / `beam_azimuth` | no | Deferred indefinitely. Both are nullable, so sending two nulls is correct behaviour rather than a gap | | Collect the owner's contact details | shipped | Landed 2026-09-16. A skippable wizard step after the agreements step, plus a block under Remote support on the Configuration page, writing `/data/retina-gui/telemetry-contact.json` | +| Collect the address that **owns** the node | no, but | v1.3.0 added `PUT /nodes/claim`, and nothing on a node knows the address, so we do not call it. Not the contact email: see `docs/data-sources.md` §4. Blocked on the server's own question about how an address relates to an account, so do not start here | `owl-os` separately owes a `mender-update show-provides` snapshot so `versions.retina_node` has a source. Optional field; omitted honestly until then. @@ -141,6 +142,18 @@ Full detail and citations in `docs/data-sources.md`. The short version: Unlike a refused registration, it breaks nothing: the node registers, streams and beats exactly as before, and the only loss is a way to ring the owner. `detail` is for what stops a node working. +- **The claim is read, never offered.** v1.3.0 added the endpoints that nominate an + owner's address and v1.4.0 put `claim_state`, `claim_email` and `claim_undeliverable` + on the heartbeat and contact responses. We consume those three and write them to the + status document; we call no claim endpoint, because **nothing on a node knows the + owner's address** and the contact email is a different question with a worse failure + (mailing a stranger a link that hands them the node). None of it gates anything: an + unclaimed node registers, streams and beats normally. Read the three as one block, + gated on `claim_state`: `claim_email` is required *and nullable*, so a field-by-field + read would let a detection ack, which carries none of them, blank an address the + heartbeat reported a second earlier. **`pending` is not durable** and nothing should + wait on it: a declined link can strand a node there for about fifteen minutes before + it falls back to `unclaimed`. - **Nothing in the stack pushes to us.** No event bus, no inbound ports. Every input is a poll or a file read, including "the user changed the config". - **`wire/models.py` is generated.** Regenerate with `tools/generate-models.sh`; never @@ -165,7 +178,7 @@ Full detail and citations in `docs/data-sources.md`. The short version: absence, so dropping the key produces a payload it rejects. `to_wire` also applies `mode="json"`, which is load-bearing: without it the acceptance timestamps stay as `datetime` objects and `json.dumps` refuses the registration payload outright. - **Fourteen fields are required-and-nullable in v1.2.2**, so payloads go out through + **Fourteen fields are required-and-nullable in v1.4.0**, so payloads go out through `wire.to_wire`, never `model_dump(exclude_none=True)` directly. `tests/wire/test_serialise.py` pins the inventory by name and fails if the spec grows or loses one. @@ -233,8 +246,8 @@ that get re-litigated if the reasoning is not written down. beam fields were changed with the server author's agreement, relayed by Josh, and their next revision did not carry it, so our edit was silently reverted on adoption. **Check `NodeConfig.beam_width_deg` when adopting any revision**, and expect to reapply it. - Checked on adopting `1.2.2` (2026-09-16): it survived, nullable as agreed. Keep - checking anyway. Two revisions carrying it is not yet a habit. + Checked on adopting `1.2.2` (2026-09-16) and `1.4.0` (2026-09-22): it survived + both times, nullable as agreed. Keep checking anyway. - **The spec is the scope.** If a field is not in it, we do not collect it — however cheap or obviously useful it looks. Wanting something new means asking the server author, not a field we add unilaterally. This has already removed Pi diff --git a/docs/data-sources.md b/docs/data-sources.md index 9757f1d..88fac84 100644 --- a/docs/data-sources.md +++ b/docs/data-sources.md @@ -408,6 +408,43 @@ streams and beats exactly as before, and the only loss is a way to ring the owne Nothing is ever substituted, the same discipline as the consent records and the beam geometry. These reach a person. +### The claim address, which has no source on a node + +Spec v1.3.0 added `PUT /v1/nodes/claim`: the node offers the address that owns it, the +server mails a link, and clicking it binds the node to that account. **Nothing on a node +knows that address**, so this service does not call the endpoint and will not until +something collects one. retina-gui would have to, the way it collects the consent +records. + +**The contact email is not the claim address, however convenient that looks.** They are +different questions with different consequences: + +| | `telemetry-contact.json` `email` | the claim address | +|---|---|---| +| Asks | whom to ring about this node | who *owns* this node | +| If wrong | a support call goes astray | a stranger is mailed a link that hands them the node | +| Optional | yes, entirely | there is no claim without one | + +An owner may well give the same address for both. That is their answer to two questions, +not a licence for us to infer the second from the first, and reusing the contact email +would claim ownership on behalf of whoever happened to be listed. The same discipline as +the consent records: nothing that reaches a person is ever synthesised here. + +The server has the account side of this open as its own question, so the shape of what +retina-gui should collect is not settled yet either. + +**What a node *is* told, since v1.4.0**, is where its claim stands: `claim_state`, +`claim_email` and `claim_undeliverable`, restated on every heartbeat and contact +response. Those three are read (`comms/levels.py`) and written to the status document, +which is the only way they reach an owner. None of them gates anything: an unclaimed +node registers, streams and beats exactly as an owned one does. + +`claim_state` is `unclaimed`, `pending` or `owned`. **`pending` is not durable.** A +claim link declined while a second is outstanding can leave a node reading `pending` +against a link nobody can redeem, until the challenge expires about fifteen minutes +later and it reads `unclaimed` again. That is a known server-side race, tracked there. +Nothing here should wait on `pending` or treat reaching it as progress. + ### The agreements, and the publication choice `RegisterRequest.agreements` requires three records. Today: From 93dfb90f4238194c972680238f40fcf11a27c7e3 Mon Sep 17 00:00:00 2001 From: Josh Poole Date: Tue, 22 Sep 2026 12:01:51 +0100 Subject: [PATCH 4/7] 20260922 - Stop a claim conflict asking for a configuration resend `apply_response` turned every 409 into `request_config_resend()`. That was right while the only 409 in the contract was a frame naming a `config_version` the server never issued, and v1.3.0 added a second one that means something entirely different: the node already has an owner. Nothing about it concerns our configuration. Left alone, the first `PUT /nodes/claim` against an already-owned node would have queued a `PUT /nodes/config` as a side effect, and every retry after it another. The node would have looked, from the server, like one whose geometry kept going stale for no reason anybody could trace back to a claim. A status code cannot tell the two apart, so `Outcome` now carries the path it answered. It is the last field and defaults to the empty string, so the positional constructions in the tests still read as they did. `_classify` already had the path in hand; it was being spent on the error string and discarded. On a claim path the answer is handled separately and never asks for a resend. It also adopts the body, because these endpoints speak `ClaimResponse`, which names the three fields without the `claim_` prefix the heartbeat and contact responses use, and they speak it on a refusal too: the 409 is the one refusal in the contract that does not wear `Error`. It carries the address that won precisely so a node can reconcile from it rather than retry, which is what `_apply_claim_answer` does. A refusal that does wear `Error` (`invalid_claim`, a 429, a 5xx) carries no `state`, applies nothing, and leaves the last known claim alone. `test_a_409_anywhere_else_still_asks_for_a_config_resend` guards the rule this is an exception to, rather than a replacement for. Co-Authored-By: Claude Opus 5 (1M context) --- retina_telemetry/comms/client.py | 11 +++- retina_telemetry/comms/levels.py | 51 ++++++++++++++++++ retina_telemetry/comms/lifecycle.py | 1 + tests/comms/test_levels.py | 84 ++++++++++++++++++++++++++++- 4 files changed, 143 insertions(+), 4 deletions(-) diff --git a/retina_telemetry/comms/client.py b/retina_telemetry/comms/client.py index 431606d..9cfac55 100644 --- a/retina_telemetry/comms/client.py +++ b/retina_telemetry/comms/client.py @@ -91,6 +91,12 @@ class Outcome: body: dict[str, Any] | None retry_after_s: int | None error: str | None + #: The path that produced it. Carried because the status code alone does + #: not say what a refusal means: the `409` on a claim is "this node already + #: has an owner", which is nothing like the `409` on a frame, and answering + #: the first with a configuration resend would be pure noise. Last, and + #: defaulted, so the positional constructions in the tests still read. + path: str = "" @property def ok(self) -> bool: @@ -225,6 +231,7 @@ def request( body=None, retry_after_s=None, error=f"{method} {path} unreachable: {exc}", + path=path, ) return _classify(method, path, response) @@ -236,7 +243,7 @@ def _classify(method: str, path: str, response: Any) -> Outcome: retry_after = _retry_after(response) if 200 <= status < 300: - return Outcome(Kind.OK, status, body, retry_after, None) + return Outcome(Kind.OK, status, body, retry_after, None, path) kind = { 400: Kind.INVALID, @@ -248,7 +255,7 @@ def _classify(method: str, path: str, response: Any) -> Outcome: detail = (body or {}).get("detail") or (body or {}).get("error") or "" error = f"{method} {path} → {status}" + (f": {detail}" if detail else "") - return Outcome(kind, status, body, retry_after, error) + return Outcome(kind, status, body, retry_after, error, path) def _json_or_none(response: Any) -> dict[str, Any] | None: diff --git a/retina_telemetry/comms/levels.py b/retina_telemetry/comms/levels.py index fcbb691..c733aad 100644 --- a/retina_telemetry/comms/levels.py +++ b/retina_telemetry/comms/levels.py @@ -20,6 +20,15 @@ **A 409 is not a reason to retry the frame.** It says the server does not recognise our `config_version`, so the same request would fail identically until a configuration resend has happened. + +**Except on the claim paths, where a 409 means something else entirely.** There +it says the node already has an owner, which has nothing to do with our +configuration, and it is the one refusal in the contract that carries a +`ClaimResponse` rather than an `Error`. The node reconciles from the body +instead of correcting and retrying. That is why `Outcome` carries the path it +answered: a status code alone cannot tell these two apart, and answering the +second with a configuration resend would put a `PUT /nodes/config` on the wire +every time somebody offered an address for a node that already has an owner. """ from __future__ import annotations @@ -39,6 +48,10 @@ #: settles — but a second or two of ordinary skew is not news. CLOCK_WARN_S = 5.0 +#: Matches every claim endpoint: ``PUT``/``GET /nodes/claim`` and +#: ``POST /nodes/claim/resend``. +CLAIM_PATH = "/nodes/claim" + def apply_response(outcome: Outcome, state: State) -> None: """Adopt everything a response implies, whatever endpoint produced it. @@ -51,6 +64,10 @@ def apply_response(outcome: Outcome, state: State) -> None: state.reject_token() return + if CLAIM_PATH in outcome.path: + _apply_claim_answer(outcome, state) + return + if outcome.kind is Kind.CONFLICT: state.request_config_resend() return @@ -71,6 +88,40 @@ def apply_response(outcome: Outcome, state: State) -> None: _warn_on_clock_offset(state) +def _apply_claim_answer(outcome: Outcome, state: State) -> None: + """Adopt what a claim endpoint answered, whatever its status. + + These paths speak ``ClaimResponse``, which names the three fields without + the ``claim_`` prefix the heartbeat and contact responses use, and they + speak it on a refusal too: the 409 for a node that already has an owner is + the one refusal in the contract that does not wear ``Error``. It carries + the address that won precisely so the node can reconcile without a second + call, so it is adopted exactly as a 200 would be. + + Nothing here ever asks for a configuration resend. A 409 on this path says + nothing whatever about our ``config_version``. + + A refusal that does carry ``Error`` (``invalid_claim``, a 429, a 5xx) has + no ``state``, so it applies nothing and leaves the last known claim alone. + """ + claim = _claim_response(outcome.body or {}) + if claim is not None: + state.apply_levels(claim=claim) + + +def _claim_response(body: dict[str, Any]) -> Claim | None: + """A ``ClaimResponse`` body, gated on ``state`` for the reasons in :func:`_claim`.""" + state = body.get("state") + if not isinstance(state, str) or not state: + return None + + return Claim( + state=state, + email=_str_or_none(body.get("email")), + undeliverable=body.get("undeliverable") is True, + ) + + def _claim(body: dict[str, Any]) -> Claim | None: """The claim block, if this response carried one. diff --git a/retina_telemetry/comms/lifecycle.py b/retina_telemetry/comms/lifecycle.py index f5f7f3e..1398b22 100644 --- a/retina_telemetry/comms/lifecycle.py +++ b/retina_telemetry/comms/lifecycle.py @@ -327,4 +327,5 @@ def _malformed(outcome: Outcome, message: str) -> Outcome: body=outcome.body, retry_after_s=outcome.retry_after_s, error=message, + path=outcome.path, ) diff --git a/tests/comms/test_levels.py b/tests/comms/test_levels.py index 0d3a431..5bdbced 100644 --- a/tests/comms/test_levels.py +++ b/tests/comms/test_levels.py @@ -19,8 +19,10 @@ def state(tmp_path): return state -def outcome(kind=Kind.OK, status=200, **body): - return Outcome(kind=kind, status=status, body=body or None, retry_after_s=None, error=None) +def outcome(kind=Kind.OK, status=200, path="/nodes/heartbeat", **body): + return Outcome( + kind=kind, status=status, body=body or None, retry_after_s=None, error=None, path=path + ) # ── the two rules that matter most ─────────────────────────────────── @@ -218,6 +220,84 @@ def test_a_claim_state_that_is_not_a_string_takes_nothing_with_it(state): assert state.snapshot().claim is None +# ── the claim endpoints, where a 409 means something else ──────────── + + +CLAIM_200 = {"state": "pending", "email": "owner@example.com", "undeliverable": False} + + +def test_a_claim_answer_is_adopted(state): + """`ClaimResponse` names the three without the `claim_` prefix.""" + apply_response(outcome(path="/nodes/claim", **CLAIM_200), state) + + assert state.snapshot().claim == Claim( + state="pending", email="owner@example.com", undeliverable=False + ) + + +def test_a_claim_409_reconciles_rather_than_resending_the_config(state): + """The bug this path exists to prevent. + + A 409 here says the node already has an owner, which says nothing about our + config_version. Treated as an ordinary conflict it would put a + PUT /nodes/config on the wire every time somebody offered an address for a + node that already has one. + """ + state.config_resend.clear() + + apply_response( + outcome( + Kind.CONFLICT, + 409, + path="/nodes/claim", + state="owned", + email="someone.else@example.com", + undeliverable=False, + ), + state, + ) + + assert not state.config_resend.is_set() + assert not state.snapshot().config_stale + # It carries the address that won, so the node reconciles without asking. + assert state.snapshot().claim.email == "someone.else@example.com" + assert state.snapshot().claim.state == "owned" + + +def test_a_409_anywhere_else_still_asks_for_a_config_resend(state): + """The rule the claim path is an exception to, not a replacement for.""" + state.config_resend.clear() + + apply_response( + outcome(Kind.CONFLICT, 409, path="/nodes/detection", error="unknown_config"), state + ) + + assert state.config_resend.is_set() + + +def test_a_resend_answer_is_adopted_too(state): + apply_response(outcome(path="/nodes/claim/resend", **CLAIM_200 | {"state": "owned"}), state) + + assert state.snapshot().claim.state == "owned" + + +def test_a_refused_claim_leaves_the_last_one_alone(state): + """`invalid_claim` wears `Error`, so it carries no state to adopt.""" + apply_response(outcome(path="/nodes/claim", **CLAIM_200), state) + + apply_response(outcome(Kind.INVALID, 400, path="/nodes/claim", error="invalid_claim"), state) + + assert state.snapshot().claim.state == "pending" + assert not state.config_resend.is_set() + + +def test_a_401_on_a_claim_path_still_rejects_the_token(state): + """Checked before the path is looked at, so the claim cannot shadow it.""" + apply_response(outcome(Kind.UNAUTHORIZED, 401, path="/nodes/claim"), state) + + assert state.snapshot().token_rejected + + # ── clock offset ───────────────────────────────────────────────────── From 35c75a113cabb41e0d5cf735a98c798cd605c5b6 Mon Sep 17 00:00:00 2001 From: Josh Poole Date: Tue, 22 Sep 2026 12:02:24 +0100 Subject: [PATCH 5/7] 20260922 - Drive the claim endpoints by hand, for a live test A harness, not a feature. `PUT /v1/nodes/claim` still has no caller in the service and this does not give it one: nothing on a node knows the address that owns it, and the contact email is a different question with a worse failure. What is missing is the trigger and the source of the address, so `tools/probe_claim.py` supplies exactly those two by hand and nothing else. Every request is built and sent by the service's own code: `Client`, the generated `NodeClaimRequest`, `to_wire` and `apply_response`. Unlike every other live-*.sh, `tools/live-claim.sh` does not point the node at the mock. It talks to whichever server minted the node's token, which for a node with no `RETINA_API_URL` override is production, so `offer` makes that server mail a real person and a click on the link binds the node to the account behind the address. Releasing it is the owner's to do from the dashboard. That is why the phases are separate invocations rather than one script that runs start to finish: `read` and `watch` send no mail and are safe to repeat, while `offer` and `resend` both refuse to run without `--confirm-send`, and `offer` additionally makes the caller repeat the address in `--confirm-address`. A mistyped address here mails a stranger. It deliberately does not run the service. The service registers on startup, and the spec answers a node already holding a valid token with the opaque 403, so it would never get a token; the one path where registration does succeed revokes the token the node's live container is using, which the server alerts on. A second process with the same `node_id` posting frames would also interleave `seq` and `boot_id` with the real container's and make the node look like it was flapping. So only the claim endpoints are driven, with the token the node already holds, and /data is mounted read-only so that token cannot be written even by accident. One thing the probe reports needs explaining, because it cost a confusing rehearsal: loading a token always queues a configuration resend, since a restored token never comes with a `config_version`. That is `State._load` doing its job and has nothing to do with the claim, so it is cleared on startup. Past that point a queued resend means a response asked for one, which is exactly what the claim 409 must not do. The mock grows the three endpoints so the whole sequence can be rehearsed locally before it is pointed at a real server, which is how the two bugs above were found. It models what the contract describes: offering an address the node already holds changes nothing and mails nothing, an address is trimmed and lower cased before it is judged so the 255 bound describes the trimmed form, and a claim on an owned node answers 409 with a `ClaimResponse` rather than an `Error`, having written nothing. Co-Authored-By: Claude Opus 5 (1M context) --- tests/test_mock_server.py | 126 ++++++++++++++++++++++++ tools/live-claim.sh | 84 ++++++++++++++++ tools/mock_server.py | 117 +++++++++++++++++++++++ tools/probe_claim.py | 196 ++++++++++++++++++++++++++++++++++++++ vulture_whitelist.py | 5 + 5 files changed, 528 insertions(+) create mode 100755 tools/live-claim.sh create mode 100644 tools/probe_claim.py diff --git a/tests/test_mock_server.py b/tests/test_mock_server.py index 2dd5af2..940e4c4 100644 --- a/tests/test_mock_server.py +++ b/tests/test_mock_server.py @@ -177,6 +177,132 @@ def test_empty_frame_is_accepted(server): assert body["accepted"] == 0 +# ── the claim ──────────────────────────────────────────────────────── + + +def claim_url(server): + return f"{server.url}/nodes/claim" + + +def test_a_fresh_node_is_unclaimed(server): + token, _ = register(server) + + status, body, _ = post(claim_url(server), None, token, method="GET") + + assert status == 200 + assert body == {"state": "unclaimed", "email": None, "undeliverable": False} + + +def test_offering_an_address_makes_the_claim_pending(server): + token, _ = register(server) + + status, body, _ = post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") + + assert status == 200 + assert body == {"state": "pending", "email": "owner@example.com", "undeliverable": False} + + +def test_the_address_is_trimmed_and_lower_cased(server): + """The schema says the bound describes the trimmed form, so the server + normalises before it judges.""" + token, _ = register(server) + + _, body, _ = post(claim_url(server), {"email": " Owner@Example.COM "}, token, method="PUT") + + assert body["email"] == "owner@example.com" + + +def test_offering_the_same_address_again_changes_nothing(server): + """So it may be resent on every sync. Mailing again is a separate call.""" + token, _ = register(server) + post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") + + status, body, _ = post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") + + assert status == 200 + assert body["state"] == "pending" + + +def test_a_bad_address_is_invalid_claim(server): + """`invalid_claim`, not `invalid_config`: the slugs are per-document so + retina-gui can mark the right field.""" + token, _ = register(server) + + status, body, _ = post(claim_url(server), {"email": "not-an-address"}, token, method="PUT") + + assert status == 400 + assert body["error"] == "invalid_claim" + assert body["detail"] == "email" + + +def test_an_unknown_field_is_refused(server): + """`NodeClaimRequest` is additionalProperties: false.""" + token, _ = register(server) + + status, body, _ = post( + claim_url(server), {"email": "owner@example.com", "name": "x"}, token, method="PUT" + ) + + assert status == 400 + assert body["error"] == "invalid_claim" + + +def test_claiming_an_owned_node_answers_409_without_an_error_body(server): + """The one refusal in the contract that does not wear `Error`. + + The node reconciles from this rather than correcting and retrying, so it is + told which address won in the same response. + """ + token, _ = register(server) + post(claim_url(server), {"email": "first@example.com"}, token, method="PUT") + server.state.only_node().claim_state = "owned" + + status, body, _ = post(claim_url(server), {"email": "second@example.com"}, token, method="PUT") + + assert status == 409 + assert "error" not in body + assert body == {"state": "owned", "email": "first@example.com", "undeliverable": False} + + +def test_a_refused_claim_writes_nothing(server): + token, _ = register(server) + post(claim_url(server), {"email": "first@example.com"}, token, method="PUT") + server.state.only_node().claim_state = "owned" + + post(claim_url(server), {"email": "second@example.com"}, token, method="PUT") + + assert server.state.only_node().claim_email == "first@example.com" + + +def test_a_resend_is_refused_on_an_owned_node(server): + token, _ = register(server) + post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") + server.state.only_node().claim_state = "owned" + + status, body, _ = post(f"{server.url}/nodes/claim/resend", None, token, method="POST") + + assert status == 409 + assert body["state"] == "owned" + + +def test_a_resend_on_a_pending_claim_is_accepted(server): + token, _ = register(server) + post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") + + status, body, _ = post(f"{server.url}/nodes/claim/resend", None, token, method="POST") + + assert status == 200 + assert body["state"] == "pending" + + +def test_the_claim_needs_a_token(server): + register(server) + + status, _, _ = post(claim_url(server), {"email": "owner@example.com"}, None, method="PUT") + + assert status == 401 + + # ── auth ───────────────────────────────────────────────────────────── diff --git a/tools/live-claim.sh b/tools/live-claim.sh new file mode 100755 index 0000000..8e44500 --- /dev/null +++ b/tools/live-claim.sh @@ -0,0 +1,84 @@ +#!/usr/bin/env bash +# +# Drive the claim endpoints on a real node, against the real server. +# +# tools/live-claim.sh read +# tools/live-claim.sh offer --email X --confirm-address X --confirm-send +# tools/live-claim.sh watch [--seconds 600] +# tools/live-claim.sh offer-again --email X +# tools/live-claim.sh resend --confirm-send +# +# ## This one is not like the others +# +# Every other live-*.sh points the node at tools/mock_server.py over an SSH +# reverse tunnel. This one does not. It talks to whichever server minted the +# node's token, which for a node with no RETINA_API_URL override is production. +# `offer` makes that server send a real person an email, and a click on the +# link in it binds this node to the account behind the address. Releasing it is +# the owner's to do from the dashboard and cannot be undone from the node. +# +# So the phases are separate invocations rather than one script that runs +# start to finish. `read` and `watch` send no mail and are safe to repeat. +# +# ## Why it does not run the service +# +# It would be wrong twice over. The service registers on startup, and the spec +# answers a node that already holds a valid token with the opaque 403, so it +# would never get a token; the one path where registration does succeed revokes +# the token the node's live container is using, which the server alerts on. And +# a second process with the same node_id posting frames would interleave `seq` +# and `boot_id` with the real container's and make the node look like it was +# flapping. So this drives the claim endpoints only, with the token the node +# already holds, and the live container keeps streaming throughout. +# +# ## Safety +# +# /data is mounted read-only, so the token cannot be written even by accident, +# and nothing is copied off the device. The only writable path is a scratch +# directory holding a copy of the package, removed on exit. The retina-node +# compose project is not touched and nothing is restarted. + +set -euo pipefail + +HOST="${1:?usage: tools/live-claim.sh [args...]}" +shift +PHASE="${1:?usage: tools/live-claim.sh [args...]}" +shift + +IMAGE="${PROBE_IMAGE:-python:3.11-slim}" +API_URL="${RETINA_API_URL:-}" +REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" + +echo "→ host: $HOST" +echo "→ phase: $PHASE $*" +if [ -n "$API_URL" ]; then + echo "→ api: $API_URL (overridden)" +else + echo "→ api: the built-in default, which is production" +fi +echo + +REMOTE_DIR="/tmp/retina-claim.$$" +tar czf - -C "$REPO_ROOT" retina_telemetry tools/probe_claim.py \ + | ssh "$HOST" "mkdir -p '$REMOTE_DIR/app' && tar xzf - -C '$REMOTE_DIR/app'" + +ssh "$HOST" "REMOTE_DIR='$REMOTE_DIR' IMAGE='$IMAGE' API_URL='$API_URL' bash -s" -- "$PHASE" "$@" <<'REMOTE' +set -euo pipefail +PHASE="$1"; shift + +cleanup() { rm -rf "$REMOTE_DIR"; } +trap cleanup EXIT + +# /data read-only: the token is read and must never be written. --network host +# for parity with how the node's own telemetry container reaches the server. +docker run --rm --network host \ + --pull missing \ + -e PYTHONDONTWRITEBYTECODE=1 \ + -e PYTHONUNBUFFERED=1 \ + -v "$REMOTE_DIR/app:/app:ro" \ + -v /data:/data:ro \ + -w /app \ + "$IMAGE" \ + sh -c "pip install --quiet --no-cache-dir --timeout 60 --retries 10 requests PyYAML pydantic \ + && python -m tools.probe_claim ${API_URL:+--api-url '$API_URL'} $PHASE $*" +REMOTE diff --git a/tools/mock_server.py b/tools/mock_server.py index ae499e2..7f94c21 100644 --- a/tools/mock_server.py +++ b/tools/mock_server.py @@ -101,6 +101,10 @@ # it, so this is the server declaring how little it expects rather than a # limit a node can reach. "contact": 2 * 1024, + # Both claim paths declare x-max-body-bytes: 2048, and the resend takes no + # body at all. + "claim": 2 * 1024, + "claim_resend": 2 * 1024, } MAX_BODY_BYTES = 64 * 1024 @@ -134,6 +138,11 @@ class _RegisterEnvelope(RegisterRequest): # Read by hand for the same reason as config: a body-shaped refusal # must not be able to precede the 401. "contact": ("PUT", f"{BASE_PATH}/nodes/contact", None), + # The claim trio, added in v1.3.0. Read by hand for the same reason as + # config and contact: a body-shaped refusal must not precede the 401. + "claim_read": ("GET", f"{BASE_PATH}/nodes/claim", None), + "claim": ("PUT", f"{BASE_PATH}/nodes/claim", None), + "claim_resend": ("POST", f"{BASE_PATH}/nodes/claim/resend", None), } @@ -394,6 +403,41 @@ def validate_contact(payload: Any) -> dict[str, Any]: return out +def validate_claim(payload: Any) -> str: + """Return the normalised address, or raise `invalid_claim` naming `email`. + + The schema caps the address at 255 and the server states that it strips + surrounding whitespace and lower cases before judging, so the bound + describes the trimmed form rather than the bytes sent. Normalising here + rather than only bounding is what makes a resend of the same address with + different capitalisation read as unchanged, which is the behaviour the + contract promises. + + The shape check is deliberately thin. The real server's is too: nothing + anywhere verifies that the address exists, and until the link is clicked it + grants nothing. + """ + if not isinstance(payload, dict): + raise ConfigInvalid("claim", "not an object") + + unknown = sorted(set(payload) - {"email"}) + if unknown: + raise ConfigInvalid(unknown[0], "unknown field") + + email = payload.get("email") + if not isinstance(email, str): + raise ConfigInvalid("email", "not a string") + + email = email.strip().lower() + if not email or len(email) > 255: + raise ConfigInvalid("email") + # One `@`, something either side of it, and no whitespace inside. + local, _, domain = email.partition("@") + if not local or not domain or "@" in domain or any(c.isspace() for c in email): + raise ConfigInvalid("email") + return email + + def validate_config(payload: Any) -> dict[str, Any]: """Return the normalised configuration, or raise naming exactly one field. @@ -735,6 +779,14 @@ class NodeRecord: claim_email: str | None = None claim_undeliverable: bool = False + def claim_response(self) -> dict[str, Any]: + """``ClaimResponse``, which names the three without the ``claim_`` prefix.""" + return { + "state": self.claim_state, + "email": self.claim_email, + "undeliverable": self.claim_undeliverable, + } + def claim_block(self) -> dict[str, Any]: """The three claim fields, spelled as every response carrying them does.""" return { @@ -989,6 +1041,12 @@ def do_GET(self) -> None: # noqa: N802 - BaseHTTPRequestHandler's naming }, ) return + # Everything else that is a declared GET endpoint goes through the + # ordinary path, so it is recorded, capped and rate limited like any + # other request. `GET /nodes/claim` is the only one today. + if self._endpoint_for("GET", self.path) is not None: + self._handle("GET") + return if self.path in ("/", "/_control/live"): page = LIVE_PAGE.encode() self.send_response(200) @@ -1334,6 +1392,65 @@ def _contact(self, body: Any) -> None: node.contact_updated_at = _now() self._send(200, {"updated_at": node.contact_updated_at, **node.claim_block()}) + def _claim_read(self, body: Any) -> None: + """``GET /v1/nodes/claim``. Cheap, and safe to poll every few seconds.""" + node = self._bearer_node(taxonomy=False) + if node is None: + return + with self.state.lock: + self._send(200, node.claim_response()) + + def _claim(self, body: Any) -> None: + """``PUT /v1/nodes/claim``. Offers the address that owns this node.""" + node = self._bearer_node(taxonomy=True) + if node is None: + return + + if body is MALFORMED: + # A body that is not JSON at all lands here too, since the remedy + # is the same: resending it unchanged will not help. + self._taxonomy(400, "invalid_claim", "claim") + return + + try: + email = validate_claim(body) + except ConfigInvalid as exc: + self._taxonomy(400, "invalid_claim", exc.field) + return + + with self.state.lock: + if node.claim_state == "owned": + # The one refusal in the contract that does not wear `Error`. + # Nothing is written, and the body names the address that won + # so the node reconciles without a second call. + self._send(409, node.claim_response()) + return + + # Offering an address the node already holds is accepted and + # changes nothing, so this may be resent on every sync. It does + # not mail anything again: that is what the resend path is for. + if email != node.claim_email: + node.claim_email = email + node.claim_undeliverable = False + node.claim_state = "pending" + self._send(200, node.claim_response()) + + def _claim_resend(self, body: Any) -> None: + """``POST /v1/nodes/claim/resend``. Sends the link again. No body.""" + node = self._bearer_node(taxonomy=False) + if node is None: + return + + with self.state.lock: + if node.claim_state == "owned": + self._send(409, node.claim_response()) + return + # Answered without sending on an address that bounced hard: a + # second copy to an address that does not exist earns nothing but + # a second bounce, and bounces cost a sending domain its + # reputation. Still a 200, because nothing is wrong with the ask. + self._send(200, node.claim_response()) + class MockServer: """The mock, as a context manager. Binds port 0 unless told otherwise. diff --git a/tools/probe_claim.py b/tools/probe_claim.py new file mode 100644 index 0000000..3a174ac --- /dev/null +++ b/tools/probe_claim.py @@ -0,0 +1,196 @@ +#!/usr/bin/env python3 +"""Drive the claim endpoints by hand, against a real server. + +**This is a test harness, not a feature.** `PUT /v1/nodes/claim` has no caller +in this service and deliberately so: nothing on a node knows the address that +owns it, and the contact email is a different question with a worse failure. +See `docs/data-sources.md` §4. What is missing is the trigger and the source of +the address, so this supplies both by hand and nothing else. Every request is +built and sent by the service's own code: `Client`, the generated +`NodeClaimRequest`, `to_wire`, and `apply_response`. + +## It talks to the real server + +There is no mock here. `offer` makes the server send somebody an email, and a +click on that link binds this node to the account behind the address. Releasing +it is the owner's to do from the dashboard and cannot be undone from the node. +So the two subcommands that cause a send both refuse to run without +`--confirm-send`, and `offer` additionally prints the address and makes you +name it again. + +## Phases + + read GET /nodes/claim safe, sends no mail + offer PUT /nodes/claim MAILS THE ADDRESS + watch GET /nodes/claim, polled safe, sends no mail + offer-again PUT /nodes/claim expects the 409, mails nothing + resend POST /nodes/claim/resend MAILS THE ADDRESS AGAIN + +Nothing here writes to disk. The token is read and never printed. +""" + +from __future__ import annotations + +import argparse +import json +import sys +import time +from pathlib import Path + +from retina_telemetry.comms.client import DEFAULT_BASE_URL, Client, Outcome +from retina_telemetry.comms.levels import apply_response +from retina_telemetry.state import DEFAULT_TOKEN_PATH, State + +CLAIM = "/nodes/claim" +RESEND = "/nodes/claim/resend" + + +def _report(outcome: Outcome, state: State) -> None: + """What the server said, and what the node made of it.""" + print(f" {outcome.kind.name} {outcome.status}") + if outcome.body is not None: + print(f" body: {json.dumps(outcome.body, sort_keys=True)}") + if outcome.error: + print(f" error: {outcome.error}") + + apply_response(outcome, state) + + snapshot = state.snapshot() + print(f" node now holds: {snapshot.claim}") + # Proves the claim 409 is not mistaken for a configuration conflict. A + # PUT /nodes/config on the wire here would be pure noise. Cleared at + # startup in `_client`, so a True here was caused by this response. + queued = state.config_resend.is_set() + print(f" config resend queued by this response: {queued}{' <-- WRONG' if queued else ''}") + + +def _client(args: argparse.Namespace) -> tuple[Client, State, str]: + state = State(args.token_path) + snapshot = state.snapshot() + if not snapshot.registered: + sys.exit(f"no token at {args.token_path}: this node is not registered") + # Loading a token always queues a configuration resend, because a restored + # token never comes with a config_version. That is `State._load` doing its + # job and has nothing to do with the claim, so it is cleared here: past + # this point a queued resend means a *response* asked for one, which is + # exactly what the claim 409 must not do. + state.config_resend.clear() + + # Never the token itself. Its length is enough to say one was loaded. + print(f"→ {args.api_url}") + print(f" token loaded from {args.token_path} ({len(snapshot.token)} chars)") + print(f" node_ref {snapshot.node_ref or 'unknown until a response carries it'}") + return Client(base_url=args.api_url), state, snapshot.token + + +def cmd_read(args: argparse.Namespace) -> int: + client, state, token = _client(args) + print("→ GET /nodes/claim") + # `None` rather than `{}`: requests sends no body at all, which is what a + # GET should carry. + _report(client.request("GET", CLAIM, None, token=token), state) + return 0 + + +def cmd_offer(args: argparse.Namespace) -> int: + from retina_telemetry.wire.models import NodeClaimRequest + from retina_telemetry.wire.serialise import to_wire + + if args.confirm_address != args.email: + sys.exit( + "refusing to send: --confirm-address must repeat --email exactly.\n" + f" --email {args.email}\n" + f" --confirm-address {args.confirm_address}" + ) + + client, state, token = _client(args) + # Built and validated by the generated model, so the spec's own 255-char + # bound fires here rather than at the server. + payload = to_wire(NodeClaimRequest(email=args.email)) + print(f"→ PUT /nodes/claim {json.dumps(payload)}") + print(" this sends mail") + _report(client.request("PUT", CLAIM, payload, token=token), state) + return 0 + + +def cmd_watch(args: argparse.Namespace) -> int: + client, state, token = _client(args) + print(f"→ polling GET /nodes/claim every {args.every}s for up to {args.seconds}s") + print(" (the spec's own cadence while a setup page is open)") + + deadline = time.monotonic() + args.seconds + last: str | None = None + while time.monotonic() < deadline: + outcome = client.request("GET", CLAIM, None, token=token) + apply_response(outcome, state) + claim = state.snapshot().claim + current = repr(claim) + if current != last: + print(f" [{time.strftime('%H:%M:%S')}] {current}") + last = current + if claim is not None and claim.state == "owned": + print(" owned. The link was clicked.") + return 0 + time.sleep(args.every) + + print(" still not owned when time ran out.") + return 1 + + +def cmd_offer_again(args: argparse.Namespace) -> int: + from retina_telemetry.wire.models import NodeClaimRequest + from retina_telemetry.wire.serialise import to_wire + + client, state, token = _client(args) + payload = to_wire(NodeClaimRequest(email=args.email)) + print(f"→ PUT /nodes/claim again {json.dumps(payload)}") + print(" expecting 409 with a ClaimResponse, and no configuration resend") + _report(client.request("PUT", CLAIM, payload, token=token), state) + return 0 + + +def cmd_resend(args: argparse.Namespace) -> int: + client, state, token = _client(args) + print("→ POST /nodes/claim/resend") + print(" this sends mail, unless the node is already owned or the address bounced") + _report(client.request("POST", RESEND, None, token=token), state) + return 0 + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--api-url", default=DEFAULT_BASE_URL) + parser.add_argument("--token-path", type=Path, default=DEFAULT_TOKEN_PATH) + sub = parser.add_subparsers(dest="command", required=True) + + sub.add_parser("read").set_defaults(run=cmd_read) + + offer = sub.add_parser("offer") + offer.add_argument("--email", required=True) + offer.add_argument( + "--confirm-address", + required=True, + help="repeat --email exactly. A real person is mailed by this call.", + ) + offer.add_argument("--confirm-send", action="store_true", required=True) + offer.set_defaults(run=cmd_offer) + + watch = sub.add_parser("watch") + watch.add_argument("--seconds", type=int, default=600) + watch.add_argument("--every", type=int, default=5) + watch.set_defaults(run=cmd_watch) + + again = sub.add_parser("offer-again") + again.add_argument("--email", required=True) + again.set_defaults(run=cmd_offer_again) + + resend = sub.add_parser("resend") + resend.add_argument("--confirm-send", action="store_true", required=True) + resend.set_defaults(run=cmd_resend) + + args = parser.parse_args(argv) + return int(args.run(args)) + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/vulture_whitelist.py b/vulture_whitelist.py index f3a83f6..32bfd1f 100644 --- a/vulture_whitelist.py +++ b/vulture_whitelist.py @@ -122,3 +122,8 @@ _._detection _._heartbeat _._register +# `_claim_read` and `_claim_resend` joined them with spec v1.3.0. `_claim` +# itself is not listed because comms/levels.py has a function of that name and +# vulture matches on the bare name, so listing it would whitelist that too. +_._claim_read +_._claim_resend From 91808227c42e751b804cf58c74668d3a7dbd7d68 Mon Sep 17 00:00:00 2001 From: Josh Poole Date: Tue, 22 Sep 2026 12:16:39 +0100 Subject: [PATCH 6/7] 20260922 - Correct the mock's claim against what production actually does Found by running the claim end to end against production on jn1, which is the argument for the mock being written rather than generated: a generated one would have agreed with itself. Two things were wrong, both in the direction of the mock being stricter than the server. **The 409 turns on the address, not on ownership.** The mock refused any offer on an owned node. The spec's wording is "a nomination names an address that is not theirs", and production accepts the owner's own address with a 200 on the idempotent path. Re-offering what the node already holds is not a conflict, because there is nothing to refuse. **Offering an address the node already holds really does change nothing**, including the state, which the mock was quietly moving to `pending`. The case that proves it is a declined claim, and it is worth writing down because it is a trap for whoever wires this into retina-gui: - Declining returns the node to `unclaimed` immediately rather than stranding it on `pending` until the challenge expires. That settles the open question in the server's own ticket about what a decline covers. - The address survives the decline. The node reports `unclaimed` with an address still on file, which the very first read of an untouched node does not: that returns a null address. - So offering that same address again is "an address the node already holds". It answers 200, changes nothing, mails nothing, and leaves the node `unclaimed`. An owner who declined by accident and asked to claim again would get silence. - `POST /nodes/claim/resend` is the only call that produces another link, and it is what moved jn1 from `unclaimed` back to `pending`. The mock now models all of it, and the tests name production and the date they were checked against it, so a later revision that changes any of this shows up as a test to revisit rather than a comment nobody trusts. Co-Authored-By: Claude Opus 5 (1M context) --- tests/test_mock_server.py | 49 +++++++++++++++++++++++++++++++++++++++ tools/mock_server.py | 34 +++++++++++++++++++++------ 2 files changed, 76 insertions(+), 7 deletions(-) diff --git a/tests/test_mock_server.py b/tests/test_mock_server.py index 940e4c4..0453df7 100644 --- a/tests/test_mock_server.py +++ b/tests/test_mock_server.py @@ -274,6 +274,55 @@ def test_a_refused_claim_writes_nothing(server): assert server.state.only_node().claim_email == "first@example.com" +def test_the_owners_own_address_is_accepted_rather_than_refused(server): + """The 409 turns on the address, not on ownership alone. + + The spec's wording is "a nomination names an address that is not theirs", + so re-offering the owner's own address is the idempotent path. Checked + against production on 2026-09-22, where this mock answered 409 and the real + server answered 200. + """ + token, _ = register(server) + post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") + server.state.only_node().claim_state = "owned" + + status, body, _ = post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") + + assert status == 200 + assert body["state"] == "owned" + + +def test_a_declined_claim_keeps_the_address_and_cannot_be_reoffered(server): + """The live run's most surprising result, and a trap for retina-gui. + + Declining returns the node to `unclaimed` but leaves the address on file. + Offering that same address again is therefore "an address the node already + holds", so it changes nothing and mails nothing: the node stays `unclaimed` + and an owner who declined by accident would sit there getting silence. + """ + token, _ = register(server) + post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") + # What a decline leaves behind, as production did on 2026-09-22. + server.state.only_node().claim_state = "unclaimed" + + status, body, _ = post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") + + assert status == 200 + assert body == {"state": "unclaimed", "email": "owner@example.com", "undeliverable": False} + + +def test_a_resend_is_what_gets_a_declined_node_another_link(server): + """The only way forward from the previous test.""" + token, _ = register(server) + post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") + server.state.only_node().claim_state = "unclaimed" + + status, body, _ = post(f"{server.url}/nodes/claim/resend", None, token, method="POST") + + assert status == 200 + assert body["state"] == "pending" + + def test_a_resend_is_refused_on_an_owned_node(server): token, _ = register(server) post(claim_url(server), {"email": "owner@example.com"}, token, method="PUT") diff --git a/tools/mock_server.py b/tools/mock_server.py index 7f94c21..0ffa7ca 100644 --- a/tools/mock_server.py +++ b/tools/mock_server.py @@ -1419,19 +1419,34 @@ def _claim(self, body: Any) -> None: return with self.state.lock: - if node.claim_state == "owned": + # The 409 is for a nomination that "names an address that is not + # theirs", so it turns on the address rather than on ownership + # alone. The owner's own address falls through to the idempotent + # path below, which is what production does: checked against it on + # 2026-09-22, when this returned 409 for both and was wrong. + if node.claim_state == "owned" and email != node.claim_email: # The one refusal in the contract that does not wear `Error`. # Nothing is written, and the body names the address that won # so the node reconciles without a second call. self._send(409, node.claim_response()) return - # Offering an address the node already holds is accepted and - # changes nothing, so this may be resent on every sync. It does - # not mail anything again: that is what the resend path is for. - if email != node.claim_email: - node.claim_email = email - node.claim_undeliverable = False + if email == node.claim_email: + # "Sending an address the node already holds is accepted and + # changes nothing, so this may be resent on every configuration + # sync. It does not mail anything again." It really does change + # nothing, the state included. The case that proves it is a + # declined claim: the address survives the decline, so a node + # sitting at `unclaimed` with an address on file stays there + # however many times it offers the same one, and `resend` is + # the only call that produces another link. Checked against + # production on 2026-09-22, where this mock had wrongly moved + # the node to `pending` and implied a mail that never went. + self._send(200, node.claim_response()) + return + + node.claim_email = email + node.claim_undeliverable = False node.claim_state = "pending" self._send(200, node.claim_response()) @@ -1449,6 +1464,11 @@ def _claim_resend(self, body: Any) -> None: # second copy to an address that does not exist earns nothing but # a second bounce, and bounces cost a sending domain its # reputation. Still a 200, because nothing is wrong with the ask. + if node.claim_email is not None and not node.claim_undeliverable: + # A fresh link, so the claim is pending again. This is how a + # declined node gets back to `pending`, since offering the same + # address again does nothing at all. + node.claim_state = "pending" self._send(200, node.claim_response()) From aecebc62372aecb67a24611043e29c0fd97205b3 Mon Sep 17 00:00:00 2001 From: Josh Poole Date: Tue, 22 Sep 2026 12:52:45 +0100 Subject: [PATCH 7/7] 20260922 - Offer the owner's address, and ask again when they do The other half of retina-gui's Node claim section. It writes /data/retina-gui/telemetry-claim.json; this reads it, decides which call achieves what the owner asked for, and makes it. Until now nothing here called a claim endpoint at all, because nothing on a node knew an address to offer. `collect/claim.py` returns a `Nomination`, the spec's own word for offering an address, and deliberately not `Claim`: `state.Claim` is where the claim actually stands, which is the server's answer rather than the owner's ask, and the two disagree for as long as it takes us to notice the file. ## Two keys, two kinds of thing `email` is state, so a change in it is a local change and goes out as `PUT /nodes/claim`, which is the call that mails a link. `send_requested_at` is an event, and it has to exist separately because of something the live run against production found rather than anything the spec says outright: offering an address the node already holds is accepted, writes nothing and mails nothing, and a declined link leaves the node `unclaimed` with the address still on file. So a node can sit unclaimed with an address against it and no number of offers will ever move it. `POST /nodes/claim/resend` is the only way out. Choosing between the two lives here rather than in retina-gui, which records what the owner wants and nothing about the wire. This is the only side that knows where the claim stands and what the server does with each call, and a rule implemented on both sides is one that will eventually disagree with itself. ## The bound on acting twice Nothing durable records that an ask was acted on, and the ask stays in the file, so a restart would re-read whatever it last held and mail the owner another link, every time this container came up. `CLAIM_ASK_FRESH_FOR_S` bounds that: five minutes, because somebody is looking at a page when they press that button, so an ask this service was not running to see is one they will simply make again. The cost is at most one duplicate from a restart inside the window, against a duplicate on every restart for ever. An offer adopts the stamp stored beside it for the same reason. An owner who fills in the box and presses send again in one go has the link sent by the offer; the stamp has been answered by it and must not produce a second. The resend is sent through the client directly rather than through `send_until_delivered`. The endpoint takes no body and that wrapper cannot make a request without one, and passing an empty object instead would be a shape this has never been checked against. The retry is the owner's anyway: the spec says a resend happens "when somebody asks for the mail again, and never automatically", so a failure is recorded and left for the person who is already looking at the page. Nothing here gates anything, so every failure goes to `errors[]` and never to the status document's `detail`, which is for what stops a node working. Co-Authored-By: Claude Opus 5 (1M context) --- CLAUDE.md | 50 ++++++--- docs/data-sources.md | 40 +++++-- retina_telemetry/__main__.py | 157 ++++++++++++++++++++++++++- retina_telemetry/collect/claim.py | 142 ++++++++++++++++++++++++ retina_telemetry/settings.py | 2 + retina_telemetry/wire/claim.py | 35 ++++++ tests/collect/test_claim.py | 110 +++++++++++++++++++ tests/test_service.py | 172 +++++++++++++++++++++++++++++- 8 files changed, 680 insertions(+), 28 deletions(-) create mode 100644 retina_telemetry/collect/claim.py create mode 100644 retina_telemetry/wire/claim.py create mode 100644 tests/collect/test_claim.py diff --git a/CLAUDE.md b/CLAUDE.md index 410948d..c8067a2 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -50,12 +50,13 @@ Corollary: **all unit conversion happens in stage 2.** Stage 1 hands over source under names that say so — `delay_km`, `timestamp_ms`, `rx_alt_m` — and stage 2 emits the spec's names and units. A missing conversion is then visible at the call site. -## Three files retina-gui writes, two of which gate registration +## Four files retina-gui writes, two of which gate registration -All three live under `/data/retina-gui`, mounted read-only, and nothing in any of them +All four live under `/data/retina-gui`, mounted read-only, and nothing in any of them is ever synthesised here. **`telemetry-consent.json` and `setup-wizard-completed` gate registration**: a node missing either refuses to register and says which in its status -document. `telemetry-contact.json` gates nothing and is optional throughout. +document. `telemetry-contact.json` and `telemetry-claim.json` gate nothing and are +optional throughout. **`telemetry-consent.json`** carries the three records `RegisterRequest.agreements` needs: `licence`, `remote_management` and `publication`. `publication` is a privacy @@ -70,6 +71,13 @@ node that has none never calls the endpoint, which the spec states explicitly, so an absent file is a complete answer rather than a gap and nothing here blocks on it. See `collect/contact.py`. +**`telemetry-claim.json`** carries the address that *owns* the node, for +`PUT /nodes/claim`, plus a `send_requested_at` stamp when the owner asks for their +link again. A different question from the contact email and a worse failure: the +server mails this one a link, and opening it binds the node to the account behind +it. Both keys optional, the file optional, and an unclaimed node runs exactly as a +claimed one does. See `collect/claim.py`. Shipped in retina-gui 2026-09-22. + **`setup-wizard-completed`** proves the config is the owner's rather than the shipped default. `retina-node/config/default.yml` ships a *working* configuration (Greenwich Observatory, Crystal Palace), the merger writes it on first boot, and retina-gui records @@ -105,7 +113,7 @@ Nothing here is buildable from this repo, and the first one blocks every node: | Read `/data/retina-telemetry/status.json` | no, but | We bind no ports, so it is the only way *no identity*, *revoked token* and *rejected config* reach an operator. `telemetry_status.py` reads it and the home page shows it | | Collect `location.rx.beam_width` / `beam_azimuth` | no | Deferred indefinitely. Both are nullable, so sending two nulls is correct behaviour rather than a gap | | Collect the owner's contact details | shipped | Landed 2026-09-16. A skippable wizard step after the agreements step, plus a block under Remote support on the Configuration page, writing `/data/retina-gui/telemetry-contact.json` | -| Collect the address that **owns** the node | no, but | v1.3.0 added `PUT /nodes/claim`, and nothing on a node knows the address, so we do not call it. Not the contact email: see `docs/data-sources.md` §4. Blocked on the server's own question about how an address relates to an account, so do not start here | +| Collect the address that **owns** the node | shipped | Landed 2026-09-22. A Node claim section on the Configuration page writing `/data/retina-gui/telemetry-claim.json`, with Save for the address and Send again for another link. Not the contact email: see `docs/data-sources.md` §4 | `owl-os` separately owes a `mender-update show-provides` snapshot so `versions.retina_node` has a source. Optional field; omitted honestly until then. @@ -142,18 +150,28 @@ Full detail and citations in `docs/data-sources.md`. The short version: Unlike a refused registration, it breaks nothing: the node registers, streams and beats exactly as before, and the only loss is a way to ring the owner. `detail` is for what stops a node working. -- **The claim is read, never offered.** v1.3.0 added the endpoints that nominate an - owner's address and v1.4.0 put `claim_state`, `claim_email` and `claim_undeliverable` - on the heartbeat and contact responses. We consume those three and write them to the - status document; we call no claim endpoint, because **nothing on a node knows the - owner's address** and the contact email is a different question with a worse failure - (mailing a stranger a link that hands them the node). None of it gates anything: an - unclaimed node registers, streams and beats normally. Read the three as one block, - gated on `claim_state`: `claim_email` is required *and nullable*, so a field-by-field - read would let a detection ack, which carries none of them, blank an address the - heartbeat reported a second earlier. **`pending` is not durable** and nothing should - wait on it: a declined link can strand a node there for about fifteen minutes before - it falls back to `unclaimed`. +- **The claim's address comes from retina-gui and nowhere else.** Never the contact + email: that answers "whom do we ring" and reusing it would mail a stranger a link that + hands them the node. `telemetry-claim.json` carries the address and, when the owner + presses send again, a timestamp. **Which call to make is decided here, not there**: a + changed address is a `PUT` and a fresh timestamp is a `resend`, because this is the + only side that knows what the server does with each. +- **Offering an address the node already holds does nothing at all.** It is accepted, + writes nothing and mails nothing. A declined link leaves the node `unclaimed` *with + the address still on file*, so it can sit there and no number of offers will move it; + `POST /nodes/claim/resend` is the only way out. Checked against production on + 2026-09-22. This is why the resend trigger exists, and why a "claim this node" button + that always PUTs is wrong. +- **A stale ask for a link is ignored**, past `CLAIM_ASK_FRESH_FOR_S`. Nothing durable + records that we acted on one, so without an age bound every restart would mail the + owner another link. +- **The three claim fields are read as one block, gated on `claim_state`.** + `claim_email` is required *and nullable*, so a field-by-field read would let a + detection ack, which carries none of them, blank an address the heartbeat reported a + second earlier. None of it gates anything: an unclaimed node registers, streams and + beats normally. **`pending` is not durable** and nothing should wait on it: a declined + link can strand a node there for about fifteen minutes before it falls back to + `unclaimed`. - **Nothing in the stack pushes to us.** No event bus, no inbound ports. Every input is a poll or a file read, including "the user changed the config". - **`wire/models.py` is generated.** Regenerate with `tools/generate-models.sh`; never diff --git a/docs/data-sources.md b/docs/data-sources.md index 88fac84..d88fbc3 100644 --- a/docs/data-sources.md +++ b/docs/data-sources.md @@ -408,13 +408,26 @@ streams and beats exactly as before, and the only loss is a way to ring the owne Nothing is ever substituted, the same discipline as the consent records and the beam geometry. These reach a person. -### The claim address, which has no source on a node +### The claim address -Spec v1.3.0 added `PUT /v1/nodes/claim`: the node offers the address that owns it, the -server mails a link, and clicking it binds the node to that account. **Nothing on a node -knows that address**, so this service does not call the endpoint and will not until -something collects one. retina-gui would have to, the way it collects the consent -records. +`/data/retina-gui/telemetry-claim.json`, written by retina-gui's Node claim section and +read-only to us. Two keys, both optional, and so is the file: + +| Key | Meaning | What it makes us do | +|---|---|---| +| `email` | the address that owns this node | a **change** is offered with `PUT /nodes/claim`, which is the call that mails a link | +| `send_requested_at` | when the owner last pressed "send again" | a **fresh** one triggers `POST /nodes/claim/resend` | + +`email` is state and `send_requested_at` is an event, and keeping them apart is what +lets retina-gui stay ignorant of the wire: it records what the owner wants, and +`collect/claim.py` plus the config loop decide which of the two calls achieves it. The +rule below is why that decision cannot live on the retina-gui side. + +**A stale ask is ignored**, past `CLAIM_ASK_FRESH_FOR_S` (five minutes). Nothing durable +records that we acted on one, so without an age bound every restart of this container +would re-read whatever the file last held and mail the owner another link, for ever. +Somebody is looking at a page when they press that button, so an ask we were not running +to see is one they will make again. **The contact email is not the claim address, however convenient that looks.** They are different questions with different consequences: @@ -428,10 +441,17 @@ different questions with different consequences: An owner may well give the same address for both. That is their answer to two questions, not a licence for us to infer the second from the first, and reusing the contact email would claim ownership on behalf of whoever happened to be listed. The same discipline as -the consent records: nothing that reaches a person is ever synthesised here. - -The server has the account side of this open as its own question, so the shape of what -retina-gui should collect is not settled yet either. +the consent records: nothing that reaches a person is ever synthesised here. The two +live in separate files and neither is ever read as a fallback for the other. + +**Offering an address the node already holds does nothing at all.** Not "nothing +visible": the server accepts it, writes nothing and mails nothing, and the state it +reports is unchanged. That matters because of what a declined link leaves behind, +verified against production on 2026-09-22: declining returns the node to `unclaimed` +*but leaves the address on file*, so a node can sit unclaimed with an address against it +and no number of offers will ever move it. `POST /nodes/claim/resend` is the only way +out. This is the single most surprising thing about the claim and the reason the resend +trigger exists at all. **What a node *is* told, since v1.4.0**, is where its claim stands: `claim_state`, `claim_email` and `claim_undeliverable`, restated on every heartbeat and contact diff --git a/retina_telemetry/__main__.py b/retina_telemetry/__main__.py index 32486d5..71e68b6 100644 --- a/retina_telemetry/__main__.py +++ b/retina_telemetry/__main__.py @@ -37,22 +37,26 @@ import signal import threading from collections.abc import Callable +from datetime import UTC, datetime from typing import Any import pydantic +from retina_telemetry.collect import claim as claim_reader from retina_telemetry.collect import consent as consent_reader from retina_telemetry.collect import contact as contact_reader from retina_telemetry.collect import identity as identity_reader from retina_telemetry.collect import node_config as config_reader from retina_telemetry.collect import wizard as wizard_reader from retina_telemetry.collect.blah2 import Blah2Client +from retina_telemetry.collect.claim import Nomination from retina_telemetry.collect.consent import Consent from retina_telemetry.collect.contact import Contact from retina_telemetry.collect.host import HostReader from retina_telemetry.collect.identity import IdentityUnavailable from retina_telemetry.collect.node_config import ConfigUnavailable, NodeConfigRaw from retina_telemetry.comms.client import Client, Kind, Outcome +from retina_telemetry.comms.levels import apply_response from retina_telemetry.comms.lifecycle import NodeState, Registrar, derive_state, explain from retina_telemetry.comms.reliable import is_fatal_for_config, send_until_delivered from retina_telemetry.comms.stream import DetectionStream, Slot @@ -60,6 +64,7 @@ from retina_telemetry.settings import Settings from retina_telemetry.state import State, with_uptime_fallback from retina_telemetry.status import StatusWriter +from retina_telemetry.wire.claim import build_claim from retina_telemetry.wire.config import build_node_config from retina_telemetry.wire.contact import build_contact from retina_telemetry.wire.detection import build_detection_frame @@ -69,6 +74,19 @@ log = logging.getLogger("retina_telemetry") +#: How recent an ask for another claim link has to be before we act on it. +#: +#: The ask is an event left in a file we poll, and nothing records that we +#: acted on it anywhere durable. Without a window, every restart of this +#: container would re-read whatever the file last held and mail the owner +#: another link, for ever. +#: +#: Short on purpose. Somebody is looking at a page when they press that button, +#: so an ask this service was not running to see is one they will simply make +#: again; acting on an hour-old press would mail a link nobody is waiting for. +#: The cost of the window is at most one duplicate, from a restart inside it. +CLAIM_ASK_FRESH_FOR_S = 300.0 + def _refusal_detail(outcome: Outcome) -> str: """A sentence for the status document when registration is refused. @@ -139,6 +157,14 @@ def __init__(self, settings: Settings | None = None) -> None: #: nothing but the token is persisted, so a restart re-sends once, #: which the endpoint's wholesale replace makes idempotent. self._contact_sent: Contact | None = None + #: The address last offered to the server, and the last ask for another + #: link that was acted on. Process-local, like everything else here. + #: A restart re-offers the address once, which is harmless because + #: offering one the server already holds changes nothing and mails + #: nothing. The ask is guarded by CLAIM_ASK_FRESH_FOR_S instead, + #: because re-acting on that one *would* mail somebody. + self._claim_sent: str | None = None + self._claim_asked: datetime | None = None #: Why the server last refused to register this node. Separate from #: `_config_rejected`, which is a PUT answering about a configuration #: the node is already registered to send. @@ -242,7 +268,8 @@ def heartbeat_loop(self) -> None: self.stop.wait(self.settings.heartbeat_interval_s) def config_loop(self) -> None: - """Push what changed locally: the configuration, and the contact details. + """Push what changed locally: the configuration, the contact details + and the claim. Nothing pushes at us, so a local edit is noticed by re-reading. The server asking arrives immediately through the resend event. @@ -269,6 +296,7 @@ def config_loop(self) -> None: # whose config.yml is unreadable can still say who owns it, and # that is exactly the node somebody needs to ring. self._send_contact_if_changed() + self._send_claim_if_needed() config = self.node_config() if config is None: @@ -406,6 +434,133 @@ def payload() -> dict[str, Any] | None: # that the list is cleared once a beat is acknowledged. batch.commit() + def _send_claim_if_needed(self) -> None: + """Offer the address that owns this node, or ask for its link again. + + Two different calls, chosen here rather than by retina-gui, because + this is the only side that knows where the claim stands and what the + server does with each. The file it writes says what the owner wants, + not which request to make. See ``collect/claim.py``. + + **A changed address is offered.** That is ``PUT /nodes/claim``, and it + is the call that makes the server mail a link. + + **An unchanged address with a fresh ask is resent.** This is the case + that needs the second call to exist at all: offering an address the + node already holds is accepted, changes nothing and mails nothing, so a + node whose link was declined sits at ``unclaimed`` with the address + still on file and no ``PUT`` will ever move it. + + Nothing here gates anything. A node nobody claims registers, streams + and beats exactly as a claimed one does, so every failure below goes to + ``errors[]`` and never to the status document's ``detail``. + """ + nomination = claim_reader.read_nomination(self.settings.claim_path) + + if nomination.email is None: + # Nobody is claiming this node, or the owner cleared the box. There + # is nothing to send either way: releasing an existing claim is the + # owner's to do from the dashboard and no endpoint here can do it. + # Forgetting what we sent means putting the same address back later + # counts as a change and is offered again. + self._claim_sent = None + return + + if nomination.email != self._claim_sent: + self._offer_claim(nomination) + return + + if self._ask_is_new(nomination.send_requested_at): + self._resend_claim(nomination.send_requested_at) + + def _ask_is_new(self, asked: datetime | None) -> bool: + """Whether this is an ask for another link that we have not acted on. + + Guarded by age as well as by novelty, because acting twice on the same + ask mails somebody twice. Nothing durable records what we have acted + on, so without the window every restart would re-send whatever the file + last held. See CLAIM_ASK_FRESH_FOR_S. + """ + if asked is None or asked == self._claim_asked: + return False + age = (datetime.now(UTC) - asked).total_seconds() + if age > CLAIM_ASK_FRESH_FOR_S: + # Adopted without acting, so it is not reconsidered every tick. + self._claim_asked = asked + log.info("ignoring a %.0fs-old ask for another claim link", age) + return False + return True + + def _offer_claim(self, nomination: Nomination) -> None: + """``PUT /nodes/claim``, for an address the server has not been told.""" + try: + payload = to_wire(build_claim(nomination)) + except ValueError as exc: + # retina-gui checks the same bound at the box, so this means the + # file was hand-edited. Dropped rather than retried: the next read + # is identical and would fail identically. + self.errors.add(f"claim: {exc}") + log.warning("cannot build a claim payload: %s", exc) + self._claim_sent = nomination.email + return + + outcome = send_until_delivered( + self.client, + "PUT", + "/nodes/claim", + lambda: payload, + state=self.state, + stop=self.stop, + token=self.state.snapshot().token, + max_attempts=3, + ) + if outcome is None: + return + + if outcome.ok or outcome.kind is Kind.CONFLICT: + # A 409 says the node already has an owner, which is settled rather + # than something to keep offering; `apply_response` has already + # adopted the ClaimResponse it carried. Either way the address is + # now the server's, and any ask stored beside it has been answered + # by the link this call just sent, so it must not fire a resend. + self._claim_sent = nomination.email + self._claim_asked = nomination.send_requested_at + return + + self.errors.add(f"claim: {outcome.describe()}") + if outcome.kind is Kind.INVALID: + # `invalid_claim`. Repeating it cannot help, so it is recorded as + # sent and a corrected address is what tries again. + self._claim_sent = nomination.email + + def _resend_claim(self, asked: datetime | None) -> None: + """``POST /nodes/claim/resend``, for a link that never arrived. + + Sent directly rather than through ``send_until_delivered``, for two + reasons. The endpoint takes no body, and that wrapper has no way to + make a request without one; passing an empty object instead would be a + shape this has never been checked against. + + And the retry belongs to the owner. The spec says this happens "when + somebody asks for the mail again, and never automatically", so a + failure is reported and left: the person who pressed the button is + looking at the page and can press it again, which is a better retry + than one that might mail them while they are not. + """ + outcome = self.client.request( + "POST", "/nodes/claim/resend", None, token=self.state.snapshot().token + ) + # Not called for us, unlike the wrapper above, and it is what adopts + # the ClaimResponse this answers with. + apply_response(outcome, self.state) + + # Recorded either way. A failed ask that stayed unrecorded would be + # retried on every tick from here on, mailing the owner once it began + # working. + self._claim_asked = asked + if not outcome.ok: + self.errors.add(f"claim resend: {outcome.describe()}") + def _send_contact_if_changed(self) -> None: """``PUT /nodes/contact``, on local change and never otherwise. diff --git a/retina_telemetry/collect/claim.py b/retina_telemetry/collect/claim.py new file mode 100644 index 0000000..7c34396 --- /dev/null +++ b/retina_telemetry/collect/claim.py @@ -0,0 +1,142 @@ +"""The address that owns this node, and any ask for another link. + +.. code-block:: json + + {"email": "owner@example.com", "send_requested_at": "2026-09-22T11:30:00Z"} + +Written by retina-gui, read-only to us, at +``/data/retina-gui/telemetry-claim.json``. Both keys are optional and so is the +whole file: a node nobody has claimed runs exactly as a claimed one does, so an +absent file is a complete answer rather than a gap. + +## Not the contact email, however alike the two files look + +:mod:`retina_telemetry.collect.contact` carries an ``email`` too, and it is a +different question with a worse failure. That one answers "whom do we ring +about this node", is optional throughout and grants nobody anything. This one +answers "who owns it": the server mails it a link, and opening that link binds +the node to the account behind the address. Substituting one for the other +would mail a stranger a link that hands them somebody's node, so the two are +read from separate files and nothing here ever falls back to the other. + +## Two keys, because they mean different kinds of thing + +``email`` is **state**. A change in it is a local change, which is what +``PUT /nodes/claim`` is for, and that call is the one that mails. + +``send_requested_at`` is an **event**: the owner pressed "send again". It has +to exist separately because re-offering an address the node already holds is +accepted, changes nothing and mails nothing, so a node whose link was declined +sits at ``unclaimed`` with the address still on file and no ``PUT`` will ever +move it. ``POST /nodes/claim/resend`` is the only way out, and this timestamp +is how that ask reaches a service that binds no ports and cannot be called. +Nothing in the stack pushes to us, so an event has to be left somewhere to be +found by polling. + +Which of the two calls to make is decided in stage 3, not here. This module +reports what the file says and nothing more. +""" + +from __future__ import annotations + +import json +import logging +from dataclasses import dataclass +from datetime import datetime +from pathlib import Path +from typing import Any + +log = logging.getLogger(__name__) + +DEFAULT_CLAIM_PATH = Path("/data/retina-gui/telemetry-claim.json") + + +@dataclass(frozen=True) +class Nomination: + """What the owner asked for, as retina-gui recorded it. + + The spec's own word for offering an address, and deliberately not + ``Claim``: :class:`retina_telemetry.state.Claim` is where the claim + actually stands, which is the server's answer rather than the owner's ask, + and the two disagree for as long as it takes us to notice this file. + + Frozen, so ``==`` is the change check, the same mechanism as ``Contact`` + and ``NodeConfigRaw``. + """ + + email: str | None = None + #: When the owner last asked for another link, or ``None`` if they never + #: have. Parsed here so that a hand-edited value fails once, on the read, + #: rather than at every comparison downstream. + send_requested_at: datetime | None = None + + +def read_nomination(path: Path | str = DEFAULT_CLAIM_PATH) -> Nomination: + """Read the stored claim document. + + Never raises. An unreadable claim document is indistinguishable in + consequence from an absent one: the node goes unclaimed, which costs it + nothing operationally, and it is not a reason to stop streaming. + """ + document = _load(Path(path)) + if not document: + return Nomination() + + unknown = sorted(set(document) - {"email", "send_requested_at"}) + if unknown: + # Named rather than shown: this file is the one that reaches a person. + log.warning("ignoring unknown claim fields: %s", ", ".join(unknown)) + + return Nomination( + email=_text(document.get("email")), + send_requested_at=_timestamp(document.get("send_requested_at")), + ) + + +def _load(path: Path) -> dict[str, Any] | None: + try: + raw = path.read_text(encoding="utf-8") + except FileNotFoundError: + # The ordinary state of a node nobody has claimed. + return None + except OSError as exc: + log.warning("%s could not be read: %s", path, exc) + return None + + try: + document = json.loads(raw) + except ValueError as exc: + log.warning("%s is not valid JSON: %s", path, exc) + return None + + if not isinstance(document, dict): + log.warning("%s does not contain an object", path) + return None + return document + + +def _text(value: Any) -> str | None: + """A usable string, or nothing. + + Never coerced. ``str(12345)`` would turn a hand-edited number into an + address the owner never typed, and this one gets mailed. + """ + if not isinstance(value, str): + return None + return value.strip() or None + + +def _timestamp(value: Any) -> datetime | None: + """Parse retina-gui's RFC 3339 stamp, tolerating the ``Z`` suffix. + + An unparseable value is dropped rather than treated as "now": acting on it + would mail somebody because a file was hand-edited badly, and dropping it + only costs an owner a second press of a button they are already looking at. + """ + if not isinstance(value, str) or not value: + return None + try: + return datetime.fromisoformat(value.replace("Z", "+00:00")) + except ValueError: + log.warning("ignoring unparseable send_requested_at: %r", value) + return None diff --git a/retina_telemetry/settings.py b/retina_telemetry/settings.py index 3e594a1..4035834 100644 --- a/retina_telemetry/settings.py +++ b/retina_telemetry/settings.py @@ -13,6 +13,7 @@ from pathlib import Path from retina_telemetry.collect.blah2 import DEFAULT_BASE_URL as BLAH2_URL +from retina_telemetry.collect.claim import DEFAULT_CLAIM_PATH from retina_telemetry.collect.consent import DEFAULT_CONSENT_PATH from retina_telemetry.collect.contact import DEFAULT_CONTACT_PATH from retina_telemetry.collect.host import DEFAULT_DISK_PATH @@ -49,6 +50,7 @@ class Settings: consent_path: Path = field(default_factory=lambda: _path("CONSENT_PATH", DEFAULT_CONSENT_PATH)) #: Optional throughout, and an absent file is a complete answer. contact_path: Path = field(default_factory=lambda: _path("CONTACT_PATH", DEFAULT_CONTACT_PATH)) + claim_path: Path = field(default_factory=lambda: _path("CLAIM_PATH", DEFAULT_CLAIM_PATH)) #: retina-gui's setup-wizard-completed flag. Registration waits for it, #: so that a node cannot report the shipped Greenwich/Crystal Palace #: default as though the owner had chosen it. diff --git a/retina_telemetry/wire/claim.py b/retina_telemetry/wire/claim.py new file mode 100644 index 0000000..0752ddf --- /dev/null +++ b/retina_telemetry/wire/claim.py @@ -0,0 +1,35 @@ +"""``Nomination`` → ``NodeClaimRequest``. + +As thin as the contact builder and for the same reason: retina-gui stores the +wire's own field name, so there is nothing to convert. What this contributes is +the boundary and the one bound that matters. +""" + +from __future__ import annotations + +from retina_telemetry.collect.claim import Nomination +from retina_telemetry.wire.models import NodeClaimRequest + + +def build_claim(nomination: Nomination) -> NodeClaimRequest: + """Convert the owner's nomination into the wire payload. + + | Wire field | Source | Conversion | + |---|---|---| + | ``email`` | ``nomination.email`` | none, trimmed and lower cased by the server | + + Sent as the owner typed it. The server states that it strips surrounding + whitespace and lower cases before judging, so normalising here would only + be a second implementation of a rule that already exists on the other side, + and one that could disagree with it. + + ``NodeClaimRequest`` caps the address at 255 and forbids extra properties, + so a hand-edited file fails here rather than being refused as + ``invalid_claim`` a second later with nothing to show the owner. + + Raises: + ValueError: via pydantic, if there is no address or it exceeds the cap. + The caller checks for an absent address first, because a node with + nobody claiming it is the ordinary state rather than an error. + """ + return NodeClaimRequest(email=nomination.email) diff --git a/tests/collect/test_claim.py b/tests/collect/test_claim.py new file mode 100644 index 0000000..d42bea0 --- /dev/null +++ b/tests/collect/test_claim.py @@ -0,0 +1,110 @@ +"""Reading the claim document retina-gui writes. + +Two keys that mean different kinds of thing: an address, which is state, and a +timestamp, which is an event. Most of what matters here is that a file nobody +wrote, or wrote badly, produces a node nobody claims rather than an exception, +because an unclaimed node is a perfectly ordinary one. +""" + +import json + +from retina_telemetry.collect.claim import Nomination, read_nomination + +ADDRESS = "owner@example.com" + + +def write(tmp_path, document): + path = tmp_path / "telemetry-claim.json" + path.write_text(json.dumps(document), encoding="utf-8") + return path + + +def test_an_address_and_an_ask_are_read(tmp_path): + path = write(tmp_path, {"email": ADDRESS, "send_requested_at": "2026-09-22T11:30:00Z"}) + + nomination = read_nomination(path) + + assert nomination.email == ADDRESS + assert nomination.send_requested_at.isoformat() == "2026-09-22T11:30:00+00:00" + + +def test_an_absent_file_is_a_node_nobody_has_claimed(tmp_path): + """The ordinary state, not a gap: an unclaimed node records and reports + exactly as a claimed one does.""" + assert read_nomination(tmp_path / "nothing.json") == Nomination() + + +def test_an_address_with_no_ask_is_normal(tmp_path): + """What the file holds until somebody presses send again.""" + path = write(tmp_path, {"email": ADDRESS}) + + assert read_nomination(path) == Nomination(email=ADDRESS, send_requested_at=None) + + +def test_surrounding_whitespace_is_dropped(tmp_path): + path = write(tmp_path, {"email": f" {ADDRESS}\n"}) + + assert read_nomination(path).email == ADDRESS + + +def test_the_address_is_not_lower_cased_here(tmp_path): + """The server states that it trims and lower cases before judging, so + doing it here as well would be a second implementation of somebody else's + rule, free to disagree with it.""" + path = write(tmp_path, {"email": "Owner@Example.COM"}) + + assert read_nomination(path).email == "Owner@Example.COM" + + +def test_a_non_string_address_is_dropped_rather_than_coerced(tmp_path): + """`str(12345)` would turn a hand-edited number into an address, and this + one gets mailed.""" + path = write(tmp_path, {"email": 12345}) + + assert read_nomination(path).email is None + + +def test_an_unparseable_ask_is_dropped(tmp_path): + """Not treated as "now", which would mail somebody because a file was + edited badly. Dropping it costs one press of a button instead.""" + path = write(tmp_path, {"email": ADDRESS, "send_requested_at": "last Tuesday"}) + + nomination = read_nomination(path) + + assert nomination.email == ADDRESS + assert nomination.send_requested_at is None + + +def test_malformed_json_reads_as_unclaimed(tmp_path): + path = tmp_path / "telemetry-claim.json" + path.write_text("{not json", encoding="utf-8") + + assert read_nomination(path) == Nomination() + + +def test_a_document_that_is_not_an_object_reads_as_unclaimed(tmp_path): + path = write(tmp_path, [ADDRESS]) + + assert read_nomination(path) == Nomination() + + +def test_unknown_fields_are_ignored(tmp_path, caplog): + """`NodeClaimRequest` forbids extra properties, so passing one on would + have the server refuse the lot.""" + path = write(tmp_path, {"email": ADDRESS, "nickname": "the shed"}) + + with caplog.at_level("WARNING"): + nomination = read_nomination(path) + + assert nomination.email == ADDRESS + assert "nickname" in caplog.text + + +def test_the_address_is_never_in_a_log_line(tmp_path, caplog): + """Same rule the contact reader follows: field names, never values.""" + path = write(tmp_path, {"email": ADDRESS, "nickname": "the shed"}) + + with caplog.at_level("WARNING"): + read_nomination(path) + + assert ADDRESS not in caplog.text diff --git a/tests/test_service.py b/tests/test_service.py index 9b1d207..15a1239 100644 --- a/tests/test_service.py +++ b/tests/test_service.py @@ -4,15 +4,17 @@ them the ones that catch a payload the pieces each considered fine. """ +import contextlib import dataclasses import json import threading import time +from datetime import UTC, datetime, timedelta import pytest import yaml -from retina_telemetry.__main__ import Service +from retina_telemetry.__main__ import CLAIM_ASK_FRESH_FOR_S, Service from retina_telemetry.comms.lifecycle import NodeState from retina_telemetry.settings import Settings from tests.collect.test_node_config import DEFAULTS @@ -51,6 +53,7 @@ def settings_for(node, server, **overrides): device_type_path=node / "device_type", consent_path=node / "consent.json", contact_path=node / "contact.json", + claim_path=node / "claim.json", wizard_flag_path=node / "setup-wizard-completed", config_path=node / "config.yml", disk_path=node, @@ -526,6 +529,173 @@ def test_an_unreadable_config_does_not_stop_the_contact_details(node, server): assert server.received("contact") +# ── the claim ──────────────────────────────────────────────────────── +# +# Two calls chosen here rather than by retina-gui, because this is the only +# side that knows what the server does with each. The case that makes the +# second call exist is a declined link: the address stays on file, so offering +# it again is accepted, changes nothing and mails nothing. + +CLAIM_ADDRESS = "owner@example.com" + + +def write_claim(node, **document): + (node / "claim.json").write_text(json.dumps({"email": CLAIM_ADDRESS, **document})) + + +def just_now(): + return datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%SZ") + + +@contextlib.contextmanager +def service_running(service): + """Run the service for the body, so a test can change a file mid-run.""" + service.stop.clear() + thread = threading.Thread(target=service.run, daemon=True) + thread.start() + try: + yield + finally: + service.stop.set() + thread.join(timeout=5) + + +def wait_for(predicate, seconds=3.0): + deadline = time.monotonic() + seconds + while time.monotonic() < deadline and not predicate(): + time.sleep(0.05) + return predicate() + + +def test_a_node_nobody_claims_never_calls_the_endpoint(node, server): + """The ordinary state. An unclaimed node registers, streams and beats + exactly as a claimed one does.""" + service = Service(settings_for(node, server)) + + run_briefly(service, until=lambda: server.received("config")) + + assert not server.received("claim") + + +def test_an_address_is_offered_once_registered(node, server): + write_claim(node) + service = Service(settings_for(node, server)) + + run_briefly(service, until=lambda: server.received("claim")) + + assert server.received("claim")[0].body == {"email": CLAIM_ADDRESS} + + +def test_the_same_address_is_not_offered_twice(node, server): + """On local change only, like the contact document.""" + write_claim(node) + service = Service(settings_for(node, server)) + + run_briefly(service, seconds=1.2) + + assert len(server.received("claim")) == 1 + + +def test_a_changed_address_is_offered_again(node, server): + write_claim(node) + service = Service(settings_for(node, server)) + + with service_running(service): + assert wait_for(lambda: server.received("claim")) + write_claim(node, email="someone.else@example.com") + assert wait_for(lambda: len(server.received("claim")) == 2) + + assert server.received("claim")[-1].body == {"email": "someone.else@example.com"} + + +def test_an_ask_resends_rather_than_offering_again(node, server): + """The declined-link case, and the whole reason the second call exists. + + The address has not changed, so a PUT would be accepted, change nothing and + mail nothing. Only the resend produces another link. + """ + write_claim(node) + service = Service(settings_for(node, server)) + + with service_running(service): + assert wait_for(lambda: server.received("claim")) + write_claim(node, send_requested_at=just_now()) + assert wait_for(lambda: server.received("claim_resend")) + + assert len(server.received("claim")) == 1 # the address never changed + assert len(server.received("claim_resend")) == 1 + + +def test_an_ask_is_acted_on_once(node, server): + """It stays in the file, so acting on it every tick would mail the owner + every tick.""" + write_claim(node) + service = Service(settings_for(node, server)) + + with service_running(service): + assert wait_for(lambda: server.received("claim")) + write_claim(node, send_requested_at=just_now()) + assert wait_for(lambda: server.received("claim_resend")) + time.sleep(0.6) # several more ticks + + assert len(server.received("claim_resend")) == 1 + + +def test_a_stale_ask_is_ignored(node, server): + """Nothing durable records that we acted, so without an age bound a + restart would mail the owner another link every time it came up.""" + write_claim(node) + service = Service(settings_for(node, server)) + old = datetime.now(UTC) - timedelta(seconds=CLAIM_ASK_FRESH_FOR_S + 60) + + with service_running(service): + assert wait_for(lambda: server.received("claim")) + write_claim(node, send_requested_at=old.strftime("%Y-%m-%dT%H:%M:%SZ")) + time.sleep(0.6) + + assert not server.received("claim_resend") + + +def test_an_ask_stored_with_a_first_offer_does_not_mail_twice(node, server): + """An owner who filled the box and pressed send again in one go. + + The offer is itself the call that mails, so the timestamp beside it has + already been answered and must not produce a second link. + """ + write_claim(node, send_requested_at=just_now()) + service = Service(settings_for(node, server)) + + run_briefly(service, seconds=1.2) + + assert len(server.received("claim")) == 1 + assert not server.received("claim_resend") + + +def test_a_claim_failure_never_reaches_the_status_document(node, server): + """Nothing about the claim stops a node working, so a refusal belongs in + `errors[]` and never in `detail`, which is for what does.""" + write_claim(node, email="not-an-address") + service = Service(settings_for(node, server)) + + run_briefly(service, seconds=1.2) + + document = json.loads((node / "status.json").read_text()) + # `detail` still describes the node's own state, which here is a healthy + # node waiting on its first frame. What must not be in it is the claim. + assert "claim" not in (document["detail"] or "") + + +def test_a_refused_address_is_not_offered_again(node, server): + """`invalid_claim` means repeating it cannot help. A corrected address is + what tries again.""" + write_claim(node, email="not-an-address") + service = Service(settings_for(node, server)) + + run_briefly(service, seconds=1.2) + + assert len(server.received("claim")) == 1 + + # ── the server pushing back ──────────────────────────────────────────