From 07235c2cbe5348282c5dcb7daf4802d2a8919a8b Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Wed, 30 Sep 2026 21:16:11 +0530 Subject: [PATCH 1/6] docs(spec): a DNS-SD mapping for service descriptors A ServiceDescriptor as a DNS-SD instance under the subtype _svc._sub of the _offlineprotocol._tcp type the peer-stream chapter fixes: a digest instance name, a TXT record (txtvers, sid, ver, addr, one c. per capability) and its bounds, refused rather than truncated. An imported LAN record is unsigned: it carries source "lan", lives in an application-level registry and is never registered with the engine, because a registration made from it would go out in signed discovery responses under this node's identity. A peer browser ignores a record carrying sid, so a published service is not one more connector to the same host. The service discovery guide is corrected against the services crate: responses go to the peer the query came from and are forwarded toward the originator, the status set is closed, the version is opaque, the peer-tracking hook is on_neighbor_discovered, and the limits table gains the fanout, dedup and size constants. The crate's own doc comment said multi-hop response relay was planned while the code does it. --- .../offline-protocol-services/src/services.rs | 19 +- docs/service-discovery.md | 45 ++-- docs/spec/README.md | 1 + docs/spec/dns-sd-mapping.md | 205 ++++++++++++++++++ docs/spec/stream-framing.md | 8 + 5 files changed, 251 insertions(+), 27 deletions(-) create mode 100644 docs/spec/dns-sd-mapping.md diff --git a/crates/offline-protocol-services/src/services.rs b/crates/offline-protocol-services/src/services.rs index d35394f9e..61f035090 100644 --- a/crates/offline-protocol-services/src/services.rs +++ b/crates/offline-protocol-services/src/services.rs @@ -68,14 +68,15 @@ pub enum ServiceAction { /// and service request/response routing. All methods return **actions** (messages /// to send, events to emit) rather than performing I/O directly. /// -/// ## Multi-hop discovery limitation +/// ## Multi-hop discovery /// -/// Discovery **queries** are forwarded across multiple hops via gossip. However, -/// discovery **responses** are currently sent only to the direct sender of the -/// query (one hop back) rather than being relayed all the way to the original -/// querier. In practice this means that services more than one hop away will -/// generate responses that reach intermediate forwarders but not the node that -/// initiated the query. Multi-hop response relay is planned for a future release. +/// Discovery **queries** are forwarded across multiple hops via gossip. +/// Discovery **responses** are sent to the direct sender of the query, never +/// to the `originator` the payload names, so that a spoofed originator cannot +/// make a provider leak its service list to an arbitrary peer. A node that +/// receives a response it did not originate forwards it as a unicast toward +/// the originator (see `try_handle_discover_response`) and emits no event of +/// its own; only the originator emits `ServiceDiscovered`. pub struct MeshServices { local_services: HashMap, seen_discovery_queries: HashMap, @@ -113,8 +114,8 @@ impl MeshServices { /// Generates a discovery broadcast to known peers, capped at /// [`DISCOVERY_INITIAL_BROADCAST_MAX`] recipients. /// - /// See the [multi-hop limitation](MeshServices#multi-hop-discovery-limitation) - /// note on `MeshServices`. + /// See the [multi-hop discovery](MeshServices#multi-hop-discovery) note on + /// `MeshServices` for where the responses go. pub fn discover_services( &mut self, user_id: &str, diff --git a/docs/service-discovery.md b/docs/service-discovery.md index 080a984cc..8fa173537 100644 --- a/docs/service-discovery.md +++ b/docs/service-discovery.md @@ -8,9 +8,11 @@ Service discovery turns an Offline Protocol mesh into a decentralized service ma The system has three phases: -1. **Registration** — a node declares what services it offers locally. -2. **Discovery** — a node broadcasts a query that gossips through the mesh; providers respond directly to the originator. -3. **Request/Response** — the consumer sends a typed request to a chosen provider and receives a response. +1. **Registration**: a node declares what services it offers locally. +2. **Discovery**: a node broadcasts a query that gossips through the mesh; a provider answers the peer it heard the query from, and each hop forwards the answer toward the originator. +3. **Request/Response**: the consumer sends a typed request to a chosen provider and receives a response. + +On a LAN there is a fourth path beside the mesh: a node can publish its registrations as DNS-SD instances and read its neighbours', under the [DNS-SD mapping](spec/dns-sd-mapping.md). A LAN import is unsigned and arrives as a `service_discovered` event with `source: "lan"`; the Python binding ships the bridge (`dnssd_bridge.py`). ``` Node A (consumer) Mesh Node B (provider) @@ -37,8 +39,8 @@ Every registered service is described by a `ServiceDescriptor`: | Field | Type | Description | |-------|------|-------------| -| `service_id` | `ServiceId` (non-empty string) | Unique identifier, e.g. `"weather.v1"`, `"wiki.first-aid"` | -| `version` | `String` | Semantic version of the service, e.g. `"2.0"` | +| `service_id` | `ServiceId` | Unique identifier, e.g. `"weather.v1"`, `"wiki.first-aid"`. Validated on construction: not empty, not whitespace only, at most 256 bytes, and not starting with `__`, which is reserved for control-message prefixes | +| `version` | `String` | An opaque label the provider chooses, e.g. `"2.0"`. Nothing parses or compares it; a consumer that wants an ordering defines its own | | `capabilities` | `HashMap` | Key-value metadata advertising features, formats, limits, etc. | The `capabilities` map lets providers advertise what they support **before** any request is made. Consumers inspect these in `ServiceDiscovered` events to pick the right provider. Example capabilities: @@ -56,11 +58,11 @@ The `capabilities` map lets providers advertise what they support **before** any Discovery queries propagate through the mesh using **gossip flooding**: -1. The originator sends the query to all its known peers. -2. Each receiving node checks its local service registry for matches and responds directly to the originator. -3. Each receiving node forwards the query to all its other known peers (excluding the sender and originator). -4. A **deduplication window** (60 seconds) prevents query storms — each node tracks query IDs it has already processed. -5. A **max-hops limit** (default 10) prevents unbounded propagation — each forward decrements the remaining hop counter. +1. The originator sends the query to at most 20 of its known peers. +2. Each receiving node checks its local service registry for matches and answers the peer it heard the query from, never the `originator` the payload names: a spoofed originator field would otherwise make every provider leak its service list to an arbitrary peer. A node that receives an answer it did not originate forwards it as a unicast toward the originator and emits no event; only the originator emits `ServiceDiscovered`. +3. Each receiving node forwards the query to at most 5 of its other known peers (excluding the sender and originator), chosen deterministically from the query id. +4. A **deduplication window** (60 seconds, at most 10,000 remembered ids) prevents query storms: each node tracks query IDs it has already processed. +5. A **max-hops limit** (default 10) prevents unbounded propagation: each forward decrements the remaining hop counter. ``` A ──── B ──── D @@ -81,7 +83,7 @@ Discovery responses include a `hop_count` field derived from the message's actua ### Peer Tracking -Service discovery broadcasts to **all known peers**, not just those with established MLS encryption sessions. Peers are tracked independently of encryption state — any peer discovered via `on_neighbor_found()` is eligible for service discovery messages. This means service discovery works even when encryption is disabled or before key exchange completes. +Service discovery broadcasts to **known peers**, not just those with established MLS encryption sessions. Peers are tracked independently of encryption state: any peer the engine learns of through `on_neighbor_discovered()` or as the sender of a message is eligible for service discovery messages, until it has not been seen for the known-peer lifetime. This means service discovery works even when encryption is disabled or before key exchange completes. ### Encryption Interaction @@ -110,7 +112,7 @@ protocol.register_service(ServiceDescriptor { let was_registered: bool = protocol.unregister_service("weather.v1")?; ``` -`ServiceId::new()` validates the ID is non-empty and returns `Err(Error::InvalidServiceId)` if it is. +`ServiceId::new()` returns `Err(Error::InvalidServiceId)` for an empty or whitespace-only id, an id over 256 bytes, or one starting with `__` (the control-message prefix family). Registering an id that is already registered replaces its descriptor. ### Discovering Services @@ -155,7 +157,7 @@ let message_id: MessageId = protocol.respond_to_service_request( )?; ``` -Common status values: `"ok"`, `"error"`, `"not_found"`. The status field is application-defined — use whatever values make sense for your service protocol. +The status is one of exactly three values: `"ok"`, `"not_found"` or `"error"`. The engine refuses any other status on send (`ServiceError::InvalidStatus`) and a requester drops a response carrying any other status on receipt, so a status of your own is a response nobody receives. Put application-level outcomes in the body. ## Events @@ -219,7 +221,7 @@ Emitted on the **consumer** node when a provider responds to a request. |-------|------|-------------| | `request_id` | `String` | Matches the request ID from `send_service_request()` | | `service_id` | `String` | The service that responded | -| `status` | `String` | Application-defined status (`"ok"`, `"error"`, `"not_found"`, etc.) | +| `status` | `String` | One of `"ok"`, `"not_found"`, `"error"`; a response with any other status is dropped before this event | | `body` | `String` | Response payload | | `provider_peer_id` | `String` | Peer ID of the provider | @@ -331,7 +333,7 @@ Service messages use internal control-message prefixes to distinguish them from | Prefix | Message Type | Direction | |--------|-------------|-----------| | `__SVC_DISC_Q__` | Discovery query | Broadcast + gossip forwarded | -| `__SVC_DISC_R__` | Discovery response | Direct to originator | +| `__SVC_DISC_R__` | Discovery response | To the peer the query came from; each hop forwards it toward the originator | | `__SVC_REQ__` | Service request | Direct to provider | | `__SVC_RESP__` | Service response | Direct to requester | @@ -349,10 +351,17 @@ Service messages use internal control-message prefixes to distinguish them from | Parameter | Value | Description | |-----------|-------|-------------| | Dedup TTL | 60 seconds | How long a query ID is remembered to prevent re-processing | +| Dedup entries | 10,000 | Remembered query ids; the oldest is evicted first | | Max hops | 10 | Maximum gossip forwarding depth for discovery queries | -| ServiceId | Non-empty string | Validated on construction; empty strings are rejected | - -These values are compile-time constants. The dedup map is automatically cleaned up during the protocol's periodic `cleanup_expired_entries()` cycle. +| Initial broadcast | 20 peers | How many known peers the originator sends a query to | +| Gossip fanout | 5 peers | How many known peers each hop forwards a query to, chosen deterministically from the query id | +| Service payload | 128 KiB | A control frame over this is dropped unparsed | +| Request or response body | 64 KiB | Refused on send, dropped on receipt | +| Method name | 256 bytes | Refused on send, dropped on receipt | +| Response status | `ok`, `not_found`, `error` | Refused on send, dropped on receipt | +| ServiceId | 1 to 256 bytes | Not whitespace only, not starting with `__` | + +These values are compile-time constants in `crates/offline-protocol-services/src/payloads.rs`. The dedup map is swept by the engine's periodic cleanup. ## Architecture Integration diff --git a/docs/spec/README.md b/docs/spec/README.md index b77d0f75c..9ead3a2b8 100644 --- a/docs/spec/README.md +++ b/docs/spec/README.md @@ -23,6 +23,7 @@ document says which reading is normative for the wire. | [Leaf node provisioning](leaf-provisioning.md) | What a constrained device owes at pairing, the never-committing profile, and the provisioning-time adversary | | [Bluetooth LE framing](ble-framing.md) | The GATT contract, the fragment header, and what a receiver owes on reassembly | | [Peer-stream framing](stream-framing.md) | The preamble that proves a stream's peer, the length-prefixed message frame, and the LAN discovery hint | +| [DNS-SD mapping](dns-sd-mapping.md) | A service descriptor as a DNS-SD instance: the subtype, the instance name, the TXT record and its bounds, and what an unsigned LAN import may never become | | [Username discovery and invites](username-discovery.md) | The self-certifying invite payload, and the non-authoritative username directory | | [The gateway contract](gateway-contract.md) | What a gateway is, the five verbs it implements, the gateway-daemon wire protocol, and the backbone | | [Conformance](conformance.md) | The two profiles, what every implementation owes, and how the vectors decide it | diff --git a/docs/spec/dns-sd-mapping.md b/docs/spec/dns-sd-mapping.md new file mode 100644 index 000000000..be14c7aff --- /dev/null +++ b/docs/spec/dns-sd-mapping.md @@ -0,0 +1,205 @@ +# DNS-SD mapping for service descriptors + +## What this chapter is for + +The mesh discovers services by a signed query that gossips hop by hop and a +signed response that walks back ([service discovery](../service-discovery.md)). +On a LAN there is a second way to learn what a neighbour offers: DNS-SD +(RFC 6763) over multicast DNS (RFC 6762), which every desktop operating +system already answers. This chapter fixes how a `ServiceDescriptor` is laid +out as a DNS-SD instance, so that two implementations on one LAN publish and +read the same records, and it fixes what an implementation may and may not do +with a record it did not sign. + +The mapping is optional in both directions. An implementation that never +publishes and never browses is complete. One that does either is bound by +everything below. + +## Invariants + +1. **A LAN record proves nothing.** A mesh discovery response is a signed + control frame from the provider; a DNS-SD record is an unsigned multicast + answer from whoever is on the segment. Anything an implementation derives + from a LAN record carries `source: "lan"`, and an application MUST treat it + as a claim rather than a discovery. The failure this prevents: an + application that files a LAN import beside a mesh discovery lets any host + on the LAN name any address as the provider of any service. +2. **A LAN record is never re-advertised on the mesh.** An implementation + MUST NOT pass an imported record to `register_service` or otherwise + publish it under its own identity. Imports go to the application's own + registry and to the application's events, and nowhere else. The failure + this prevents is a forgery laundered into a signature: a registration made + on the strength of an anonymous LAN record would go out in signed discovery + responses under this node's identity, and every mesh peer would trust the + claim on the strength of that signature. +3. **Only this node's own registrations are published, once each.** One + registration is one instance, named by the pair `(address, service_id)`, + and a registration that leaves the local registry leaves the LAN. +4. **A field that does not fit is refused, never truncated.** A truncated + service id names a different service; a truncated capability value is a + different claim. An implementation that cannot publish a registration as + specified publishes nothing for it and keeps the mesh registration. +5. **A service instance is not a peer hint.** The record carries this node's + address so that an application knows whom to ask, and it may carry the + port of a peer stream, but a browser looking for peers under + [peer-stream framing](stream-framing.md#finding-a-peer-on-a-lan) MUST + ignore any record that carries a `sid` entry. Without this rule every + published service is a second, third and fourth connector to one host. + +## The service type + +The service type is `_offlineprotocol._tcp`, the same type the peer-stream +chapter fixes, and the one type this protocol registers: a service label is +at most fifteen letters, digits and hyphens (RFC 6763 section 7), and this +one is exactly fifteen. Service instances are distinguished from peer-stream +records by the subtype `_svc`, so that a service instance is published under + +``` +_svc._sub._offlineprotocol._tcp. +``` + +and a browser looking for services browses that subtype. RFC 6763 section +7.1 lists a subtyped instance under its base type as well, and a responder +that follows it answers a browse of the base type with service instances, +which is why invariant 5 exists and why a peer browser keys on the `sid` +entry rather than on the name it was browsing. A responder that lists the +instance under the subtype only is conforming too; a browser MUST NOT rely +on either behaviour. + +The domain is `local.` on a multicast LAN. Nothing here depends on it. + +## The instance name + +The instance name is `svc-` followed by the first sixteen hexadecimal digits, +lower case, of + +``` +SHA-256(address || 0x00 || service_id) +``` + +where `address` is the publisher's canonical address (`off1…`) and +`service_id` is the descriptor's id, both as UTF-8. Twenty octets, under the +63-octet bound RFC 6763 section 4.1.1 places on an instance name. The name +carries no meaning: the address and the id live in the TXT record, once each, +and a name that repeated them would be a second place to get them wrong. It is +a digest rather than a counter so that a re-publish after a restart claims the +same name and a responder's conflict resolution has nothing to resolve. + +## The SRV record + +The SRV target is a host name of the publisher's own, never the machine's +`.local.`, which the operating system's responder may already +answer for. An implementation SHOULD reuse the host name its peer-stream +record uses when it publishes one, so that both resolve to one set of +addresses. The addresses are the publisher's interface addresses; a record +with none is seen by every browser and resolved by none. + +The port is the publisher's peer-stream listening port when it listens, and +`0` when it does not. A browser MUST NOT connect to it on the strength of +this record (invariant 5); the port is there so that a host running both +records publishes one consistent SRV target, not as an endpoint. + +## The TXT record + +The record is a sequence of `key=value` strings in this order: + +| Entry | Value | Required | +|-------|-------|----------| +| `txtvers=1` | The version of this mapping. First, as RFC 6763 section 6.7 recommends, so a reader can stop at the first string | Yes | +| `sid=` | The descriptor's `service_id`, as UTF-8 | Yes | +| `ver=` | The descriptor's `version`, as UTF-8; the engine does not parse it and neither does this mapping | Yes, possibly empty | +| `addr=` | The publisher's canonical address, which is what a `service_discovered` event names as `provider_peer_id` | Yes | +| `c.=` | One string per capability, keys in byte order, so that one descriptor is one record | Zero or more | + +A reader MUST ignore a record whose `txtvers` is not `1`, and MUST ignore a +record missing `sid` or `addr`, or whose `addr` is not an address. Unknown +keys are ignored. A key that appears twice makes the record malformed and it +is ignored whole. + +### Bounds + +Every bound below is refused at the publisher, never truncated +(invariant 4). A reader applies the same bounds to what it imports, so that a +publisher that ignores them cannot make a reader hold what the reader would +never have published. + +| Bound | Value | Where it comes from | +|-------|-------|---------------------| +| Service label | 15 characters | RFC 6763 section 7; `_offlineprotocol` is at the bound | +| Instance name | 63 octets | RFC 6763 section 4.1.1; `svc-` plus sixteen hex digits is 20 | +| One TXT string (`key=value`) | 255 bytes | RFC 6763 section 6.1, the length octet | +| TXT key | Printable ASCII (0x20 to 0x7E) with no `=`, at least one character | RFC 6763 section 6.4; a capability key must satisfy it or the descriptor is not publishable | +| `sid` value | 200 bytes | The engine allows 256, which does not fit one string beside `sid=`; 200 leaves the string within its bound with margin, and an id that long is not something a LAN needs to carry | +| Whole TXT record | 1300 bytes, counting one length octet per string | RFC 6763 section 6.2, so that the record fits one 1500-byte Ethernet frame with its headers | + +A capability value is opaque bytes to DNS-SD and UTF-8 here; it is bounded +only by the string bound. + +## Importing a record + +A browser that resolves an instance under the subtype, checks the TXT record +against the table above and finds it well formed has a `ServiceDescriptor` +claimed by the address in `addr`. What it does with it is bounded by +invariants 1 and 2: + +- It delivers a `service_discovered` event to the application with the fields + the mesh event has (`query_id`, `service_id`, `version`, + `provider_peer_id`, `capabilities`, `hop_count`) and one more, + `source: "lan"`. `query_id` is empty, because no query was sent, and + `hop_count` is `0`, because the record came from the segment. The mesh event + is unchanged and carries no `source`; a reader that wants to distinguish the + two checks for the field, as the bridge rules require of every additive + field ([C3](../bridges/README.md#c3-events-cross-as-opaque-json)). +- It keeps the import in an application-level registry that the engine never + reads. The engine's registry holds what this node offers; the import is + what a neighbour claims to offer. +- It never registers the import with the engine (invariant 2). + +An import that names this node's own address is ignored: it is this node's +own record coming back. + +### Lifetime + +A DNS-SD record has the lifetime its publisher gave it, and a responder that +dies without a goodbye leaves its records to age out in every browser's cache, +which for the PTR is 75 minutes by RFC 6762 section 10. An importer therefore +owes the application a shorter liveness rule of its own: an entry is kept +while the instance still resolves, re-resolved at half the importer's time to +live, and removed from the application-level registry when it has not resolved +within that time to live or when the browser reports the instance gone, +whichever is first. Removal is the shadow registry's `unregister`; it never +reaches the engine, because the engine never held the entry. + +Nothing on the mesh side has a lifetime: a mesh registration stands until it +is unregistered, and a mesh discovery response is a point-in-time answer with +no removal signal. The importer's time to live is the one lifetime in the +system, and it belongs to the LAN side only. + +## Publishing + +A publisher publishes each registration in its local registry as one instance +and withdraws it when the registration is removed. A registration whose +descriptor cannot be laid out within the bounds is not published, and the +publisher says so in its log; the mesh registration is unaffected. Publishing +is best effort: an application that needs to know whether a service is +reachable over the LAN asks the registry, not the mesh. + +What a device publishes on a LAN is visible to every device on it. That is +the same exposure as the mesh discovery response, which any querying peer +receives, and the same exposure as the peer-stream record: the address is +public by design, and the service list is what the node chose to advertise. +An application that does not want a service visible on the LAN does not +publish it there; the mapping is per registration, not all or nothing. + +## Position among the discovery paths + +| Path | Signed | Reach | Lifetime | Event field | +|------|--------|-------|----------|-------------| +| Mesh discovery response | Yes, by the provider | Every peer the gossip reaches, up to ten hops | None; a response is an answer, a registration stands until unregistered | No `source` | +| DNS-SD instance | No | The multicast segment | The importer's time to live, re-resolved at half of it | `source: "lan"` | + +A LAN import can be confirmed by a mesh discovery: an application that +receives a LAN import and then a mesh response for the same `(address, +service_id)` has a signed answer, and may treat the LAN entry as confirmed +from then on. This chapter does not require an implementation to do that +join; it requires only that the two are never mistaken for each other. diff --git a/docs/spec/stream-framing.md b/docs/spec/stream-framing.md index 3075e32bd..4802aad5a 100644 --- a/docs/spec/stream-framing.md +++ b/docs/spec/stream-framing.md @@ -259,6 +259,14 @@ and carries no meaning; an implementation SHOULD NOT put the address there, since the TXT entry already carries it and one copy is one place to get it wrong. +A record that carries a `sid` entry is not a peer hint. It is a service +instance under the [DNS-SD mapping](dns-sd-mapping.md), published under the +subtype `_svc._sub` of this type and listed under the base type as well by +responders that follow RFC 6763 section 7.1. It names the same host and the +same address as the peer record, so a browser that took it as one would open +a second connector to one host per service published there. A browser +looking for peers MUST ignore any record that carries `sid`. + What a device advertises on a LAN is visible to every device on it. That is the same exposure as a Bluetooth LE advertisement, and the same answer: the address is public by design, and what it does not reveal (who the human is, From 7ef3e82c2b41e99a37be4404cfd4f45998100aed Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Wed, 30 Sep 2026 21:16:11 +0530 Subject: [PATCH 2/6] feat(bindings): Python service wrappers and a DNS-SD bridge for services Services wraps the generated MeshServices with the copy of this node's registrations the engine cannot enumerate, changed only after the engine accepted, and refuses a response status outside the engine's closed set with the reason. DnsSdBridge publishes those registrations on the LAN and imports the LAN's into a registry of its own under the DNS-SD mapping chapter, over the existing optional lan extra imported at start(): an import is delivered as service_discovered with source "lan" and never reaches register_service; a descriptor that does not fit the record is kept on the mesh and not published; an import is re-resolved at half its time to live and dropped at the whole. The peer-stream record reader ignores a record carrying sid. The chapter's bounds, the subtype and the status set are pinned as literals in the tests (C5). The suite drives the bridge through a fake of the five responder calls; no real mDNS is exercised by it. --- CHANGELOG.md | 23 + bindings/python/README.md | 32 + .../offline_protocol_sdk/dnssd_bridge.py | 546 +++++++++++++ .../peer_stream_manager.py | 7 + .../python/offline_protocol_sdk/services.py | 195 +++++ bindings/python/tests/test_dnssd_bridge.py | 725 ++++++++++++++++++ .../python/tests/test_peer_stream_manager.py | 13 + bindings/python/tests/test_services.py | 199 +++++ docs/bridges/python.md | 13 + 9 files changed, 1753 insertions(+) create mode 100644 bindings/python/offline_protocol_sdk/dnssd_bridge.py create mode 100644 bindings/python/offline_protocol_sdk/services.py create mode 100644 bindings/python/tests/test_dnssd_bridge.py create mode 100644 bindings/python/tests/test_services.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 551396f87..3b2da696d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -171,6 +171,29 @@ archived by series under [docs/changelog/](docs/changelog/); see the package's job. Two mesh controller tests were failing: they registered two peers in a mesh with room for four, so the eviction they assert was never weighed. They now fill the mesh, as their Kotlin twins have since #120. +- **Services on the LAN, and the first Python service wrappers.** A new + specification chapter, `docs/spec/dns-sd-mapping.md`, lays a + `ServiceDescriptor` out as a DNS-SD instance under the subtype + `_svc._sub._offlineprotocol._tcp`: a digest instance name, a TXT record + (`txtvers`, `sid`, `ver`, `addr`, one `c.` per capability) and its + bounds (a service id over 200 bytes, a record over 1300 bytes or a + capability key DNS-SD cannot carry is refused, never truncated). An + imported LAN record is unsigned: it arrives as a `service_discovered` + event with `source: "lan"`, is kept in an application-level registry, and + is never registered with the engine, because a registration made from it + would go out in signed discovery responses under this node's identity. + A peer-stream browser now ignores a record carrying `sid`, so a published + service is not one more connector to the same host. In Python, + `services.Services` wraps the generated `MeshServices` with the copy of + this node's registrations the engine cannot enumerate, and refuses a + response status outside the engine's closed set with the reason; + `dnssd_bridge.DnsSdBridge` publishes those registrations and imports the + LAN's, over the existing optional `lan` extra, re-resolving an import at + half its time to live and dropping it at the whole. The service discovery + guide is corrected where it disagreed with the engine: discovery responses + go to the peer the query came from and are forwarded toward the + originator, the response status is one of exactly three values, the + version is opaque, and the peer-tracking hook is `on_neighbor_discovered`. ### Fixed diff --git a/bindings/python/README.md b/bindings/python/README.md index 90fecc561..cbf394864 100644 --- a/bindings/python/README.md +++ b/bindings/python/README.md @@ -98,6 +98,36 @@ await pm.peer_stream.start() A peer is announced only under the address its preamble proves, and `pm.stop()` stops the stream layer with everything else. +### Services, and services on the LAN + +`Services` wraps the generated `MeshServices` and keeps the one thing the +engine cannot give back: the list of what this node registered. `respond` +refuses a status the engine would refuse (`ok`, `not_found` and `error` are +the whole set) with the reason instead of an opaque core error. + +```python +from offline_protocol_sdk.services import Services +from offline_protocol_sdk.dnssd_bridge import DnsSdBridge + +services = Services(pm.protocol) +services.register("weather.v1", "2.0", {"format": "json"}) # before or after pm.start() + +# Publish this node's registrations on the LAN and import the neighbours'. +# Needs: pip install 'offline-protocol-sdk[lan]' +bridge = DnsSdBridge(services, on_event=handle_event) +await bridge.start(address=pm.local_address, port=pm.peer_stream.listen_port) +... +await bridge.stop() # before pm.stop() +``` + +A LAN import arrives at `handle_event` as a `service_discovered` event with +`source: "lan"` and is kept in `bridge.lan_services()`; it is an unsigned +claim by whoever answered on the segment, never a discovery, and the bridge +never registers it with the engine. Registrations that do not fit the record +(a service id over 200 bytes, a capability key DNS-SD cannot carry, a record +over 1300 bytes) are kept on the mesh and not published, with a warning. The +mapping is [docs/spec/dns-sd-mapping.md](../../docs/spec/dns-sd-mapping.md). + ## Architecture ``` @@ -106,6 +136,8 @@ offline_protocol_sdk/ ├── protocol_manager.py # High-level wrapper (processing loop, lifecycle) ├── internet_manager.py # WebSocket transport (websockets library) ├── peer_stream_manager.py # TCP peer streams + DNS-SD (the wifi_direct slot) +├── services.py # Service registry wrappers with the copy the engine lacks +├── dnssd_bridge.py # Services on the LAN: publish own, import neighbours' (optional extra `lan`) ├── ble_manager.py # BLE transport (bleak library) ├── secure_storage.py # MLS key storage (keyring library) ├── state_storage.py # Restartable protocol state (application data) diff --git a/bindings/python/offline_protocol_sdk/dnssd_bridge.py b/bindings/python/offline_protocol_sdk/dnssd_bridge.py new file mode 100644 index 000000000..ab9927f8c --- /dev/null +++ b/bindings/python/offline_protocol_sdk/dnssd_bridge.py @@ -0,0 +1,546 @@ +"""DNS-SD for services: publish this node's registrations on the LAN, and +import what the LAN's neighbours publish, under ``docs/spec/dns-sd-mapping.md``. + +Three rules from the chapter shape everything here, each with the failure it +prevents: + +1. **A LAN record proves nothing.** A mesh discovery response is signed by + the provider; a DNS-SD record is an unsigned multicast answer from whoever + is on the segment. Every import leaves this module as a + ``service_discovered`` event with ``source: "lan"``, and as a + :class:`~.services.ServiceRecord` with ``source == "lan"``, so that an + application cannot mistake it for a discovery without ignoring the field. +2. **A LAN record is never re-advertised on the mesh.** Imports go to this + bridge's own registry and to the application's event handler, and never + to :meth:`Services.register`. A registration made on the strength of an + anonymous LAN record would go out in signed discovery responses under this + node's identity, and every mesh peer would trust the forged claim on the + strength of that signature. +3. **A field that does not fit is refused, never truncated.** A truncated + service id names a different service. A registration whose descriptor + does not fit the record is not published, and the mesh registration is + untouched; :meth:`DnsSdBridge.publish` returns ``False`` and the log says + why. + +The bridge needs the optional ``zeroconf`` dependency +(``pip install 'offline-protocol-sdk[lan]'``), imported at :meth:`start` so +the base install never imports it. Records are published under the subtype +``_svc._sub._offlineprotocol._tcp`` of the type the peer-stream transport +already uses; a peer-stream browser ignores any record carrying ``sid``, and +this browser only ever browses the subtype. + +The bounds are the chapter's, mirrored here as literals and pinned as +literals in ``test_dnssd_bridge.py`` for the reason the bridge rules give +(C5): a test that read them from this module would agree with any edit. +""" + +from __future__ import annotations + +import asyncio +import hashlib +import logging +import time +from dataclasses import dataclass +from typing import Any, Callable, Mapping + +from .peer_stream_manager import SERVICE_TYPE, advertised_addresses, service_host_name +from .services import SOURCE_LAN, SOURCE_LOCAL, ServiceRecord, Services + +logger = logging.getLogger(__name__) + +#: Service instances live under this subtype of the peer-stream type. +SUBTYPE = "_svc._sub." + SERVICE_TYPE +#: The first TXT entry, and the only mapping version a reader accepts. +TXT_VERSION = "1" +#: The instance name is this prefix and sixteen hex digits of a digest. +INSTANCE_PREFIX = "svc-" +#: RFC 6763 section 4.1.1: an instance name is at most 63 octets. +MAX_INSTANCE_NAME_OCTETS = 63 +#: RFC 6763 section 6.1: one ``key=value`` string is at most 255 bytes. +MAX_TXT_STRING_BYTES = 255 +#: The chapter's bound on ``sid``: the engine's 256 does not fit one string. +MAX_SID_BYTES = 200 +#: RFC 6763 section 6.2: the whole record, one length octet per string. +MAX_TXT_RECORD_BYTES = 1300 +#: What a canonical address starts with; a record's ``addr`` must too. +ADDRESS_PREFIX = "off1" +#: How long an import is kept without resolving again, in seconds. Records +#: age out of a browser's cache after 75 minutes; a headless neighbour that +#: died without a goodbye should not be offered for that long. +DEFAULT_TTL = 300.0 +#: How long one resolution may take. +RESOLVE_TIMEOUT_MS = 3000 + +_TXT_KEYS_FIXED = ("txtvers", "sid", "ver", "addr") +_CAPABILITY_PREFIX = "c." + + +class RecordRefused(ValueError): + """A descriptor that does not fit the record, and the reason.""" + + +# -- the mapping, as pure functions --------------------------------------------- + + +def instance_label(address: str, service_id: str) -> str: + """``svc-`` plus the first sixteen hex digits of + ``SHA-256(address || 0x00 || service_id)``.""" + digest = hashlib.sha256(address.encode("utf-8") + b"\x00" + service_id.encode("utf-8")) + return INSTANCE_PREFIX + digest.hexdigest()[:16] + + +def service_instance_name(address: str, service_id: str) -> str: + """The fully qualified instance name under the peer-stream type. The + subtype is where it is browsed, not part of its name.""" + return f"{instance_label(address, service_id)}.{SERVICE_TYPE}" + + +def _check_key(key: str) -> None: + if not key: + raise RecordRefused("a TXT key cannot be empty") + for ch in key: + code = ord(ch) + if code < 0x20 or code > 0x7E or ch == "=": + raise RecordRefused( + f"TXT key {key!r} is not printable ASCII without '=' (RFC 6763 section 6.4)" + ) + + +def txt_record(record: ServiceRecord, address: str) -> dict[str, str]: + """Lay a descriptor out as the chapter's TXT record, in the chapter's + order, or raise :class:`RecordRefused`. Nothing is ever shortened.""" + if not address.startswith(ADDRESS_PREFIX): + raise RecordRefused(f"{address!r} is not a canonical address") + sid = record.service_id.encode("utf-8") + if len(sid) > MAX_SID_BYTES: + raise RecordRefused( + f"service id is {len(sid)} bytes; the record carries at most {MAX_SID_BYTES}" + ) + entries: dict[str, str] = { + "txtvers": TXT_VERSION, + "sid": record.service_id, + "ver": record.version, + "addr": address, + } + for key in sorted(record.capabilities, key=lambda k: k.encode("utf-8")): + _check_key(key) + entries[_CAPABILITY_PREFIX + key] = record.capabilities[key] + for key, value in entries.items(): + string = len(key.encode("utf-8")) + 1 + len(value.encode("utf-8")) + if string > MAX_TXT_STRING_BYTES: + raise RecordRefused( + f"TXT entry {key!r} is {string} bytes; one string carries at most " + f"{MAX_TXT_STRING_BYTES} (RFC 6763 section 6.1)" + ) + size = txt_record_size(entries) + if size > MAX_TXT_RECORD_BYTES: + raise RecordRefused( + f"the TXT record is {size} bytes; the chapter allows {MAX_TXT_RECORD_BYTES}" + ) + return entries + + +def txt_record_size(entries: Mapping[str, str]) -> int: + """The record's size on the wire: one length octet per ``key=value``.""" + return sum(1 + len(k.encode("utf-8")) + 1 + len(v.encode("utf-8")) for k, v in entries.items()) + + +def parse_txt(raw: bytes) -> list[tuple[bytes, bytes | None]] | None: + """The strings of a TXT record as ``(key, value)`` pairs, in order, or + ``None`` when a length octet runs past the end. A string with no ``=`` + is a key with no value (RFC 6763 section 6.4); an empty string is + skipped, as the RFC says a reader must.""" + pairs: list[tuple[bytes, bytes | None]] = [] + offset = 0 + while offset < len(raw): + length = raw[offset] + offset += 1 + if offset + length > len(raw): + return None + string = raw[offset : offset + length] + offset += length + if not string: + continue + key, sep, value = string.partition(b"=") + pairs.append((key, value if sep else None)) + return pairs + + +def record_from_txt(raw: bytes) -> ServiceRecord | None: + """A LAN neighbour's claim, from the raw TXT record of a resolved + instance, or ``None`` for a record the chapter says to ignore: another + ``txtvers``, a missing ``sid`` or ``addr``, an ``addr`` that is not an + address, a key twice, a value that is not UTF-8, or anything over the + bounds a publisher would have been refused on.""" + pairs = parse_txt(raw) + if pairs is None or len(raw) > MAX_TXT_RECORD_BYTES: + return None + seen: set[bytes] = set() + fields: dict[str, str] = {} + capabilities: dict[str, str] = {} + for key, value in pairs: + if key in seen: + return None + seen.add(key) + try: + name = key.decode("ascii") + _check_key(name) + text = "" if value is None else value.decode("utf-8") + except (UnicodeDecodeError, RecordRefused): + return None + if name in _TXT_KEYS_FIXED: + fields[name] = text + elif name.startswith(_CAPABILITY_PREFIX) and len(name) > len(_CAPABILITY_PREFIX): + capabilities[name[len(_CAPABILITY_PREFIX) :]] = text + if fields.get("txtvers") != TXT_VERSION: + return None + service_id = fields.get("sid") + address = fields.get("addr") + if not service_id or len(service_id.encode("utf-8")) > MAX_SID_BYTES: + return None + if not address or not address.startswith(ADDRESS_PREFIX): + return None + return ServiceRecord( + service_id=service_id, + version=fields.get("ver", ""), + capabilities=capabilities, + provider=address, + source=SOURCE_LAN, + ) + + +def lan_service_event(record: ServiceRecord) -> dict[str, Any]: + """The event a LAN import produces: the mesh ``service_discovered`` + fields, an empty ``query_id`` (no query was sent), ``hop_count`` 0 (the + record came from the segment) and ``source: "lan"``.""" + return { + "type": "service_discovered", + "source": SOURCE_LAN, + "query_id": "", + "service_id": record.service_id, + "version": record.version, + "provider_peer_id": record.provider, + "capabilities": dict(record.capabilities), + "hop_count": 0, + } + + +# -- the responder behind the bridge ------------------------------------------ + + +class ZeroconfBackend: + """The optional ``zeroconf`` package, behind the five calls the bridge + makes. Tests substitute a fake with the same five.""" + + def __init__(self) -> None: + try: + from zeroconf import IPVersion, ServiceInfo, ServiceStateChange + from zeroconf.asyncio import AsyncServiceBrowser, AsyncServiceInfo, AsyncZeroconf + except ImportError as exc: + raise ImportError( + "DNS-SD for services needs the optional dependency: " + "pip install 'offline-protocol-sdk[lan]'" + ) from exc + self._ServiceInfo = ServiceInfo + self._ServiceStateChange = ServiceStateChange + self._AsyncServiceBrowser = AsyncServiceBrowser + self._AsyncServiceInfo = AsyncServiceInfo + self._zc = AsyncZeroconf(ip_version=IPVersion.All) + self._browser: Any = None + + async def register( + self, *, name: str, port: int, txt: Mapping[str, str], server: str, addresses: list[str] + ) -> Any: + info = self._ServiceInfo( + SUBTYPE, + name, + port=port, + properties=dict(txt), + server=server, + parsed_addresses=list(addresses), + ) + await self._zc.async_register_service(info) + return info + + async def unregister(self, handle: Any) -> None: + await self._zc.async_unregister_service(handle) + + def browse(self, on_change: Callable[[str, bool], None]) -> None: + removed = self._ServiceStateChange.Removed + + def handler(zeroconf: Any, service_type: str, name: str, state_change: Any) -> None: + on_change(name, state_change is removed) + + self._browser = self._AsyncServiceBrowser(self._zc.zeroconf, SUBTYPE, handlers=[handler]) + + async def resolve(self, name: str, timeout_ms: int) -> bytes | None: + info = self._AsyncServiceInfo(SUBTYPE, name) + if await info.async_request(self._zc.zeroconf, timeout_ms): + return bytes(info.text) + return None + + async def close(self) -> None: + if self._browser is not None: + await self._browser.async_cancel() + self._browser = None + await self._zc.async_close() + + +@dataclass +class _LanEntry: + record: ServiceRecord + seen_at: float + + +class DnsSdBridge: + """Publishes a :class:`Services` registry on the LAN and imports the + LAN's into a registry of its own. + + Start it after ``ProtocolManager.start()``, with the address the engine + derived (``pm.local_address``); stop it before the manager. Registrations + made while it runs are published as they happen, through the + :class:`~.services.ServicesListener` it installs; registrations made + before are published at :meth:`start`. + + ``on_event`` receives each import as the event + :func:`lan_service_event` builds, on the event loop's thread. Removal is + not an event: the mesh has no removal signal either, and + :meth:`lan_services` is the current view. + """ + + def __init__( + self, + services: Services, + *, + on_event: Callable[[dict[str, Any]], None] | None = None, + publish: bool = True, + browse: bool = True, + ttl: float = DEFAULT_TTL, + resolve_timeout_ms: int = RESOLVE_TIMEOUT_MS, + backend: Any | None = None, + ) -> None: + if ttl <= 0: + raise ValueError("ttl must be positive") + self._services = services + self._on_event = on_event + self._publish = publish + self._browse = browse + self._ttl = float(ttl) + self._resolve_timeout_ms = int(resolve_timeout_ms) + self._backend_given = backend + self._backend: Any = None + self._loop: asyncio.AbstractEventLoop | None = None + self._address: str | None = None + self._port = 0 + self._addresses: list[str] = [] + self._published: dict[str, Any] = {} + self._lan: dict[str, _LanEntry] = {} + self._tasks: set[asyncio.Task[Any]] = set() + self._sweeper: asyncio.Task[None] | None = None + self._running = False + + # -- lifecycle ------------------------------------------------------------ + + @property + def running(self) -> bool: + return self._running + + @property + def address(self) -> str | None: + return self._address + + async def start( + self, + *, + address: str, + port: int | None = None, + listen_host: str = "0.0.0.0", + addresses: list[str] | None = None, + ) -> None: + """Publish and browse. + + ``address`` is this node's canonical address, ``port`` its + peer-stream listening port when it has one (the record's SRV port, + 0 otherwise), and ``addresses`` the interface addresses to publish + (every non-loopback address of every interface, by default, as the + peer-stream record does). + """ + if self._running: + return + if not address.startswith(ADDRESS_PREFIX): + raise ValueError(f"{address!r} is not a canonical address") + self._loop = asyncio.get_running_loop() + self._address = address + self._port = int(port or 0) + self._backend = self._backend_given if self._backend_given is not None else ZeroconfBackend() + if self._publish: + self._addresses = list(addresses) if addresses is not None else advertised_addresses(listen_host) + if not self._addresses: + await self._backend.close() + self._backend = None + raise RuntimeError("publishing found no interface address to put in the record") + self._running = True + self._services.add_listener(self) + try: + if self._publish: + for record in self._services.registered(): + await self.publish(record) + if self._browse: + self._backend.browse(self._on_change) + self._sweeper = self._loop.create_task(self._sweep_forever()) + except BaseException: + # A responder that refuses a registration (a name conflict, a + # closed socket) must not leave a bridge that is listening to + # the registry and holds an open responder while reporting + # itself stopped. + await self.stop() + raise + + async def stop(self) -> None: + if not self._running: + return + self._running = False + self._services.remove_listener(self) + if self._sweeper is not None: + self._sweeper.cancel() + self._sweeper = None + for task in list(self._tasks): + task.cancel() + self._tasks.clear() + for service_id in list(self._published): + handle = self._published.pop(service_id) + try: + await self._backend.unregister(handle) + except Exception: + logger.exception("withdrawing %r from the LAN failed", service_id) + self._lan.clear() + try: + await self._backend.close() + finally: + self._backend = None + + # -- publishing ----------------------------------------------------------- + + def published(self) -> list[str]: + """The service ids currently on the LAN.""" + return list(self._published) + + async def publish(self, record: ServiceRecord) -> bool: + """Publish one of this node's registrations. Returns ``False`` when + the descriptor does not fit the record (rule 3); the mesh + registration stands either way. A record that is not this node's + own is refused outright (rule 2).""" + if record.source != SOURCE_LOCAL: + raise ValueError("only this node's own registrations are published (dns-sd-mapping invariant 2)") + if not self._running or not self._publish or self._address is None: + return False + try: + txt = txt_record(record, self._address) + except RecordRefused as exc: + logger.warning("service %r not published on the LAN: %s", record.service_id, exc) + return False + previous = self._published.pop(record.service_id, None) + if previous is not None: + await self._backend.unregister(previous) + handle = await self._backend.register( + name=service_instance_name(self._address, record.service_id), + port=self._port, + txt=txt, + server=service_host_name(self._address), + addresses=self._addresses, + ) + self._published[record.service_id] = handle + return True + + async def withdraw(self, record: ServiceRecord) -> None: + handle = self._published.pop(record.service_id, None) + if handle is not None and self._backend is not None: + await self._backend.unregister(handle) + + # ServicesListener: called on whatever thread changed the registry. + def on_registered(self, record: ServiceRecord) -> None: + self._later(lambda: self.publish(record)) + + def on_unregistered(self, record: ServiceRecord) -> None: + self._later(lambda: self.withdraw(record)) + + def _later(self, make: Callable[[], Any]) -> None: + loop = self._loop + if loop is None or not self._running: + return + loop.call_soon_threadsafe(self._spawn, make) + + def _spawn(self, make: Callable[[], Any]) -> None: + if not self._running or self._loop is None: + return + task = self._loop.create_task(make()) + self._tasks.add(task) + task.add_done_callback(self._tasks.discard) + + # -- importing ------------------------------------------------------------ + + def lan_services(self) -> list[ServiceRecord]: + """What the LAN's neighbours currently claim to offer.""" + return [entry.record for entry in self._lan.values()] + + def _on_change(self, name: str, removed: bool) -> None: + if removed: + self._lan.pop(name, None) + return + self._spawn(lambda: self._resolve_and_import(name)) + + async def _resolve_and_import(self, name: str) -> None: + if self._backend is None: + return + raw = await self._backend.resolve(name, self._resolve_timeout_ms) + if raw is not None: + self.import_txt(name, raw) + + def import_txt(self, name: str, raw: bytes, now: float | None = None) -> bool: + """Take a resolved instance's TXT record into the LAN registry. + Returns whether an event was delivered: a new instance, or one whose + record changed. A record the chapter says to ignore removes any + earlier import under that name; this node's own record is ignored.""" + record = record_from_txt(raw) + if record is None: + self._lan.pop(name, None) + return False + if record.provider == self._address: + return False + seen_at = time.monotonic() if now is None else now + previous = self._lan.get(name) + self._lan[name] = _LanEntry(record, seen_at) + if previous is not None and previous.record == record: + return False + self._emit(lan_service_event(record)) + return True + + async def sweep(self, now: float | None = None) -> None: + """Re-resolve every import older than half the time to live, and + drop every one older than the whole of it. The periodic task calls + this; a test calls it with a clock of its own.""" + current = time.monotonic() if now is None else now + for name, entry in list(self._lan.items()): + age = current - entry.seen_at + if age >= self._ttl: + self._lan.pop(name, None) + elif age >= self._ttl / 2 and self._backend is not None: + raw = await self._backend.resolve(name, self._resolve_timeout_ms) + if raw is not None: + self.import_txt(name, raw, now=current) + + async def _sweep_forever(self) -> None: + while self._running: + await asyncio.sleep(self._ttl / 2) + try: + await self.sweep() + except Exception: + logger.exception("the LAN sweep failed; the next one runs on schedule") + + def _emit(self, event: dict[str, Any]) -> None: + if self._on_event is None: + return + try: + self._on_event(event) + except Exception: + logger.exception("event handler failed on %s", event.get("type")) diff --git a/bindings/python/offline_protocol_sdk/peer_stream_manager.py b/bindings/python/offline_protocol_sdk/peer_stream_manager.py index e336f801f..e7db4f413 100644 --- a/bindings/python/offline_protocol_sdk/peer_stream_manager.py +++ b/bindings/python/offline_protocol_sdk/peer_stream_manager.py @@ -1266,6 +1266,13 @@ def peers_from_record(info: Any) -> list[PeerEntry]: unscoped one cannot be connected to. """ properties = getattr(info, "properties", None) or {} + # A record carrying `sid` is a service instance under the DNS-SD mapping + # chapter, published under a subtype of this type. It names the same + # host and address as the peer record and would become a second + # connector to it per service; the chapter's invariant 5 has a peer + # browser ignore it. + if b"sid" in properties: + return [] raw = properties.get(b"addr") if raw is None: return [] diff --git a/bindings/python/offline_protocol_sdk/services.py b/bindings/python/offline_protocol_sdk/services.py new file mode 100644 index 000000000..89b86b59f --- /dev/null +++ b/bindings/python/offline_protocol_sdk/services.py @@ -0,0 +1,195 @@ +"""Application-level service wrappers over the generated ``MeshServices``. + +The engine keeps the registry of what this node offers, and it cannot list +it: the FFI has register, unregister, discover, request and respond, and no +enumeration. Anything that needs the list (a DNS-SD publisher, a status page, +a restart that re-announces) keeps a copy. This module keeps that copy once, +beside the calls that change it, so that the copy and the engine's registry +move together: the engine call runs first, and the copy changes only when +the engine accepted. + +What is here is the mesh side only. A LAN import under +``docs/spec/dns-sd-mapping.md`` is a claim by an unsigned record, and it is +never registered with the engine through this module or any other: that +would sign, under this node's identity, a claim made by whoever answered on +the segment (invariant 2 of the chapter). :class:`ServiceRecord` carries a +``source`` so that a reader can tell the two apart, and only records with +``source == "local"`` ever reach :meth:`Services.register`. + +The response status set is closed. The engine accepts ``ok``, ``not_found`` +and ``error``, refuses any other status on send, and drops a response with +any other status on receipt; a status the guide once called +"application-defined" is a response nobody receives. The set is mirrored +here as a literal, pinned by ``test_services.py`` the way the bridge rules +pin every hand-mirrored constant (C5), so that a respond that would be +refused fails here with the reason instead of as an opaque core error. +""" + +from __future__ import annotations + +import logging +import threading +from dataclasses import dataclass, field +from typing import Any, Callable, Mapping, Protocol + +from .offline_protocol import MeshServices, OfflineProtocol + +logger = logging.getLogger(__name__) + +#: The response statuses the engine accepts (``VALID_SERVICE_STATUSES`` in +#: the services crate). Mirrored as a literal; see the module docstring. +VALID_STATUSES: tuple[str, ...] = ("ok", "not_found", "error") + +#: A registration this node made. +SOURCE_LOCAL = "local" +#: An import from an unsigned LAN record (``docs/spec/dns-sd-mapping.md``). +SOURCE_LAN = "lan" + + +@dataclass(frozen=True) +class ServiceRecord: + """One service, as this node offers it or as a LAN neighbour claims to. + + ``provider`` is ``None`` for this node's own registrations and the + claimed address for a LAN import. ``capabilities`` is a copy: a caller + mutating its own map afterwards does not change what was registered. + """ + + service_id: str + version: str = "" + capabilities: Mapping[str, str] = field(default_factory=dict) + provider: str | None = None + source: str = SOURCE_LOCAL + + def __post_init__(self) -> None: + object.__setattr__(self, "capabilities", dict(self.capabilities)) + + @property + def is_local(self) -> bool: + return self.source == SOURCE_LOCAL + + +class ServicesListener(Protocol): + """What a publisher (the DNS-SD bridge) hears from the registry. + + Both are called outside the registry's lock, on the thread that made the + change, after the engine accepted it. A listener that raises is logged + and does not undo the registration: the mesh registration is the primary + and a LAN copy that could not be published is the LAN's loss only. + """ + + def on_registered(self, record: ServiceRecord) -> None: ... + + def on_unregistered(self, record: ServiceRecord) -> None: ... + + +class Services: + """Register, unregister, discover, request and respond, with the copy of + this node's registrations the engine cannot provide. + + Construct it over a live :class:`OfflineProtocol` (``pm.protocol``); the + generated :class:`MeshServices` is built here. The engine takes a + registration before ``start()`` as well as after, so a host may register + everything it offers before it starts. + """ + + def __init__(self, protocol: OfflineProtocol, *, mesh_services: Any | None = None) -> None: + self._mesh: Any = mesh_services if mesh_services is not None else MeshServices(protocol) + self._lock = threading.Lock() + self._local: dict[str, ServiceRecord] = {} + self._listeners: list[ServicesListener] = [] + + # -- the registry --------------------------------------------------------- + + def register( + self, + service_id: str, + version: str = "", + capabilities: Mapping[str, str] | None = None, + ) -> ServiceRecord: + """Register a service this node offers. + + Registering an id again replaces the descriptor, as the engine does; + listeners hear the old record unregistered and the new one registered. + Raises the core's error when the engine refuses the id (empty, + whitespace, over 256 bytes, or a reserved ``__`` prefix). + """ + record = ServiceRecord(service_id, version, capabilities or {}) + self._mesh.register_service(service_id, version, dict(record.capabilities)) + with self._lock: + previous = self._local.get(service_id) + self._local[service_id] = record + listeners = list(self._listeners) + if previous is not None: + self._notify(listeners, "on_unregistered", previous) + self._notify(listeners, "on_registered", record) + return record + + def unregister(self, service_id: str) -> bool: + """Unregister a service. Returns what the engine returns: whether it + was registered. The copy is kept in step with the engine's answer, + so an id the engine did not hold is not in the copy afterwards + either, whatever this side believed.""" + found = bool(self._mesh.unregister_service(service_id)) + with self._lock: + previous = self._local.pop(service_id, None) + listeners = list(self._listeners) + if previous is not None: + self._notify(listeners, "on_unregistered", previous) + return found + + def registered(self) -> list[ServiceRecord]: + """This node's registrations, in registration order.""" + with self._lock: + return list(self._local.values()) + + def get(self, service_id: str) -> ServiceRecord | None: + with self._lock: + return self._local.get(service_id) + + # -- the calls that carry no state ----------------------------------------- + + def discover(self, service_id: str | None = None) -> str: + """Broadcast a discovery query; returns the query id the + ``service_discovered`` events will carry.""" + return str(self._mesh.discover_services(service_id)) + + def request(self, provider: str, service_id: str, method: str, body: str) -> str: + """Send a request to a provider; returns the request id the + ``service_response_received`` event will carry.""" + return str(self._mesh.send_service_request(provider, service_id, method, body)) + + def respond(self, request_id: str, requester: str, service_id: str, status: str, body: str) -> str: + """Answer a ``service_request_received`` event. + + ``status`` must be one of :data:`VALID_STATUSES`; the engine refuses + any other and the requester drops any other, so it is refused here + with the reason. + """ + if status not in VALID_STATUSES: + raise ValueError( + f"status {status!r} is not one the engine sends or a requester keeps; " + f"use one of {', '.join(VALID_STATUSES)}" + ) + return str(self._mesh.respond_to_service_request(request_id, requester, service_id, status, body)) + + # -- listeners ------------------------------------------------------------ + + def add_listener(self, listener: ServicesListener) -> None: + with self._lock: + if listener not in self._listeners: + self._listeners.append(listener) + + def remove_listener(self, listener: ServicesListener) -> None: + with self._lock: + if listener in self._listeners: + self._listeners.remove(listener) + + @staticmethod + def _notify(listeners: list[ServicesListener], method: str, record: ServiceRecord) -> None: + for listener in listeners: + callback: Callable[[ServiceRecord], None] = getattr(listener, method) + try: + callback(record) + except Exception: + logger.exception("service listener %r failed in %s", listener, method) diff --git a/bindings/python/tests/test_dnssd_bridge.py b/bindings/python/tests/test_dnssd_bridge.py new file mode 100644 index 000000000..e3a1d8c3e --- /dev/null +++ b/bindings/python/tests/test_dnssd_bridge.py @@ -0,0 +1,725 @@ +"""Tests for the DNS-SD bridge under ``docs/spec/dns-sd-mapping.md``. + +No network and no responder: the mapping is tested as pure functions on +bytes, and the bridge is driven through a fake of the five calls it makes on +the responder. The chapter's bounds are asserted as literals, not read from +the module, for the reason the bridge rules give (C5). +""" + +from __future__ import annotations + +import asyncio +import logging +from typing import Any +from unittest.mock import MagicMock + +import pytest + +from offline_protocol_sdk.dnssd_bridge import ( + ADDRESS_PREFIX, + DEFAULT_TTL, + INSTANCE_PREFIX, + MAX_INSTANCE_NAME_OCTETS, + MAX_SID_BYTES, + MAX_TXT_RECORD_BYTES, + MAX_TXT_STRING_BYTES, + SUBTYPE, + TXT_VERSION, + DnsSdBridge, + RecordRefused, + ZeroconfBackend, + instance_label, + lan_service_event, + parse_txt, + record_from_txt, + service_instance_name, + txt_record, + txt_record_size, +) +from offline_protocol_sdk.peer_stream_manager import SERVICE_TYPE, service_host_name +from offline_protocol_sdk.services import SOURCE_LAN, ServiceRecord, Services + +OUR_ADDRESS = "off1qyulwy7s5ezz20cy222zrw04rwds39uapqv9j8r0" +PEER_ADDRESS = "off1qysluvwl5922yctzd0u9gpr06gn3k7ldfvgtwgvn" + + +def encode_txt(entries: dict[str, str]) -> bytes: + """The wire form of a TXT record: one length octet per string.""" + out = bytearray() + for key, value in entries.items(): + string = f"{key}={value}".encode("utf-8") + out.append(len(string)) + out += string + return bytes(out) + + +def raw_strings(strings: list[bytes]) -> bytes: + out = bytearray() + for string in strings: + out.append(len(string)) + out += string + return bytes(out) + + +class FakeBackend: + """The responder's five calls, remembered.""" + + def __init__(self) -> None: + self.registered: dict[str, tuple[object, dict[str, Any]]] = {} + self.unregistered: list[object] = [] + self.on_change: Any = None + self.txt_by_name: dict[str, bytes] = {} + self.resolve_calls: list[str] = [] + self.closed = False + + async def register(self, **kw: Any) -> object: + handle = object() + self.registered[kw["name"]] = (handle, kw) + return handle + + async def unregister(self, handle: object) -> None: + for name, (held, _) in list(self.registered.items()): + if held is handle: + del self.registered[name] + self.unregistered.append(handle) + + def browse(self, on_change: Any) -> None: + self.on_change = on_change + + async def resolve(self, name: str, timeout_ms: int) -> bytes | None: + self.resolve_calls.append(name) + return self.txt_by_name.get(name) + + async def close(self) -> None: + self.closed = True + + +@pytest.fixture +def mesh() -> MagicMock: + m = MagicMock() + m.unregister_service = MagicMock(return_value=True) + return m + + +@pytest.fixture +def services(mesh: MagicMock) -> Services: + return Services(MagicMock(), mesh_services=mesh) + + +@pytest.fixture +def backend() -> FakeBackend: + return FakeBackend() + + +@pytest.fixture +def events() -> list[dict[str, Any]]: + return [] + + +@pytest.fixture +async def bridge(services, backend, events): + b = DnsSdBridge(services, on_event=events.append, backend=backend, ttl=100.0) + await b.start(address=OUR_ADDRESS, port=7878, addresses=["192.168.1.9"]) + try: + yield b + finally: + await b.stop() + + +async def settle() -> None: + """Let the tasks the listener scheduled run.""" + for _ in range(3): + await asyncio.sleep(0) + + +# --------------------------------------------------------------------------- +# The chapter's literals +# --------------------------------------------------------------------------- + + +class TestChapterLiterals: + def test_the_type_the_subtype_and_the_version(self): + assert SERVICE_TYPE == "_offlineprotocol._tcp.local." + assert len("_offlineprotocol") == 16 # underscore plus fifteen, the RFC 6763 bound + assert SUBTYPE == "_svc._sub._offlineprotocol._tcp.local." + assert TXT_VERSION == "1" + assert INSTANCE_PREFIX == "svc-" + assert ADDRESS_PREFIX == "off1" + + def test_the_bounds(self): + assert MAX_INSTANCE_NAME_OCTETS == 63 + assert MAX_TXT_STRING_BYTES == 255 + assert MAX_SID_BYTES == 200 + assert MAX_TXT_RECORD_BYTES == 1300 + assert DEFAULT_TTL == 300.0 + + +# --------------------------------------------------------------------------- +# The instance name +# --------------------------------------------------------------------------- + + +class TestInstanceName: + def test_is_a_digest_of_address_and_id(self): + label = instance_label(OUR_ADDRESS, "weather.v1") + assert label.startswith("svc-") and len(label) == 20 + assert all(c in "0123456789abcdef" for c in label[4:]) + assert len(label.encode("utf-8")) <= 63 + assert instance_label(OUR_ADDRESS, "weather.v1") == label + assert instance_label(PEER_ADDRESS, "weather.v1") != label + assert instance_label(OUR_ADDRESS, "weather.v2") != label + + def test_the_separator_keeps_address_and_id_apart(self): + # Without the 0x00 between them, "ab" + "c" and "a" + "bc" collide. + assert instance_label("off1ab", "c") != instance_label("off1a", "bc") + + def test_the_full_name_is_under_the_type_and_carries_neither_input(self): + name = service_instance_name(OUR_ADDRESS, "weather.v1") + assert name.endswith("." + SERVICE_TYPE) + assert OUR_ADDRESS not in name and "weather" not in name + + +# --------------------------------------------------------------------------- +# The TXT record, publisher side +# --------------------------------------------------------------------------- + + +class TestTxtRecord: + def test_the_entries_and_their_order(self): + record = ServiceRecord("weather.v1", "2.0", {"format": "json", "coverage": "us,eu"}) + entries = txt_record(record, OUR_ADDRESS) + assert list(entries.items()) == [ + ("txtvers", "1"), + ("sid", "weather.v1"), + ("ver", "2.0"), + ("addr", OUR_ADDRESS), + ("c.coverage", "us,eu"), + ("c.format", "json"), + ] + + def test_capability_keys_are_in_byte_order(self): + entries = txt_record(ServiceRecord("s", "", {"b": "1", "B": "2", "a": "3"}), OUR_ADDRESS) + assert [k for k in entries if k.startswith("c.")] == ["c.B", "c.a", "c.b"] + + def test_the_size_counts_a_length_octet_per_string(self): + entries = {"txtvers": "1", "sid": "s"} + assert txt_record_size(entries) == (1 + len("txtvers=1")) + (1 + len("sid=s")) + assert txt_record_size(entries) == len(encode_txt(entries)) + + def test_a_service_id_at_the_bound_fits_and_one_over_is_refused(self): + assert txt_record(ServiceRecord("a" * 200), OUR_ADDRESS)["sid"] == "a" * 200 + with pytest.raises(RecordRefused, match="201 bytes"): + txt_record(ServiceRecord("a" * 201), OUR_ADDRESS) + + def test_the_id_bound_is_in_bytes_not_characters(self): + with pytest.raises(RecordRefused): + txt_record(ServiceRecord("é" * 101), OUR_ADDRESS) # 202 bytes + + def test_a_record_at_the_bound_fits_and_one_over_is_refused(self): + fixed = txt_record_size(txt_record(ServiceRecord("s"), OUR_ADDRESS)) + # Fill with 250-byte strings ("c.kN=" plus value) up to the bound. + caps: dict[str, str] = {} + remaining = 1300 - fixed + n = 0 + while remaining >= 1 + 250: + caps[f"k{n:02d}"] = "v" * (250 - len(f"c.k{n:02d}=") ) + remaining -= 1 + 250 + n += 1 + caps["last"] = "v" * (remaining - 1 - len("c.last=")) + entries = txt_record(ServiceRecord("s", "", caps), OUR_ADDRESS) + assert txt_record_size(entries) == 1300 + caps["last"] += "v" + with pytest.raises(RecordRefused, match="1301 bytes"): + txt_record(ServiceRecord("s", "", caps), OUR_ADDRESS) + + def test_a_string_over_255_bytes_is_refused_not_shortened(self): + record = ServiceRecord("s", "", {"k": "v" * (255 - len("c.k=") + 1)}) + with pytest.raises(RecordRefused, match="256 bytes"): + txt_record(record, OUR_ADDRESS) + assert txt_record(ServiceRecord("s", "", {"k": "v" * (255 - len("c.k="))}), OUR_ADDRESS) + + def test_a_version_over_the_string_bound_is_refused(self): + with pytest.raises(RecordRefused, match="'ver'"): + txt_record(ServiceRecord("s", "9" * 252), OUR_ADDRESS) + + @pytest.mark.parametrize("key", ["", "a=b", "tab\tkey", "é", "\x7f"]) + def test_a_capability_key_dns_sd_cannot_carry_is_refused(self, key): + with pytest.raises(RecordRefused): + txt_record(ServiceRecord("s", "", {key: "v"}), OUR_ADDRESS) + + def test_a_capability_value_may_be_any_utf8(self): + entries = txt_record(ServiceRecord("s", "", {"lang": "français=oui"}), OUR_ADDRESS) + assert entries["c.lang"] == "français=oui" + + def test_an_address_that_is_not_one_is_refused(self): + with pytest.raises(RecordRefused, match="canonical address"): + txt_record(ServiceRecord("s"), "alice") + + +# --------------------------------------------------------------------------- +# The TXT record, reader side +# --------------------------------------------------------------------------- + + +class TestParseTxt: + def test_round_trip(self): + entries = {"txtvers": "1", "sid": "weather.v1", "ver": "", "addr": PEER_ADDRESS, "c.f": "json"} + assert parse_txt(encode_txt(entries)) == [(k.encode(), v.encode()) for k, v in entries.items()] + + def test_a_key_without_a_value_and_an_empty_string(self): + assert parse_txt(raw_strings([b"flag", b"", b"k="])) == [(b"flag", None), (b"k", b"")] + + def test_a_length_past_the_end_is_malformed(self): + assert parse_txt(b"\x05ab") is None + assert parse_txt(b"") == [] + + +class TestRecordFromTxt: + def good(self, **overrides: Any) -> dict[str, str]: + entries = {"txtvers": "1", "sid": "weather.v1", "ver": "2.0", "addr": PEER_ADDRESS, "c.format": "json"} + entries.update(overrides) + return entries + + def test_a_well_formed_record_is_a_lan_claim(self): + record = record_from_txt(encode_txt(self.good())) + assert record == ServiceRecord( + "weather.v1", "2.0", {"format": "json"}, provider=PEER_ADDRESS, source=SOURCE_LAN + ) + assert not record.is_local + + def test_the_version_defaults_to_empty(self): + entries = self.good() + del entries["ver"] + assert record_from_txt(encode_txt(entries)).version == "" + + @pytest.mark.parametrize("txtvers", ["0", "2", ""]) + def test_another_mapping_version_is_ignored(self, txtvers): + assert record_from_txt(encode_txt(self.good(txtvers=txtvers))) is None + + def test_txtvers_need_not_be_first_to_be_read(self): + entries = {"sid": "weather.v1", "addr": PEER_ADDRESS, "txtvers": "1"} + assert record_from_txt(encode_txt(entries)) is not None + + def test_missing_sid_or_addr_is_ignored(self): + for key in ("sid", "addr"): + entries = self.good() + del entries[key] + assert record_from_txt(encode_txt(entries)) is None + assert record_from_txt(encode_txt(self.good(sid=""))) is None + + def test_an_addr_that_is_not_an_address_is_ignored(self): + assert record_from_txt(encode_txt(self.good(addr="alice"))) is None + + def test_a_key_twice_is_ignored_whole(self): + raw = raw_strings([b"txtvers=1", b"sid=a", b"sid=b", b"addr=" + PEER_ADDRESS.encode()]) + assert record_from_txt(raw) is None + + def test_a_value_that_is_not_utf8_is_ignored(self): + raw = raw_strings([b"txtvers=1", b"sid=\xff\xfe", b"addr=" + PEER_ADDRESS.encode()]) + assert record_from_txt(raw) is None + + def test_a_key_dns_sd_forbids_is_ignored(self): + raw = raw_strings([b"txtvers=1", b"sid=a", b"addr=" + PEER_ADDRESS.encode(), b"\x01=x"]) + assert record_from_txt(raw) is None + + def test_unknown_keys_and_an_empty_capability_key_are_skipped(self): + record = record_from_txt(encode_txt(self.good(**{"other": "x", "c.": "y"}))) + assert record.capabilities == {"format": "json"} + + def test_the_reader_holds_the_publishers_bounds(self): + assert record_from_txt(encode_txt(self.good(sid="a" * 201))) is None + assert record_from_txt(encode_txt(self.good(sid="a" * 200))) is not None + big = self.good(**{f"c.k{n}": "v" * 240 for n in range(6)}) + assert len(encode_txt(big)) > 1300 + assert record_from_txt(encode_txt(big)) is None + + def test_a_malformed_record_is_ignored(self): + assert record_from_txt(b"\x09txtvers=1\x0fsid") is None + + +class TestLanEvent: + def test_the_shape(self): + record = ServiceRecord("weather.v1", "2.0", {"f": "j"}, provider=PEER_ADDRESS, source=SOURCE_LAN) + assert lan_service_event(record) == { + "type": "service_discovered", + "source": "lan", + "query_id": "", + "service_id": "weather.v1", + "version": "2.0", + "provider_peer_id": PEER_ADDRESS, + "capabilities": {"f": "j"}, + "hop_count": 0, + } + + +# --------------------------------------------------------------------------- +# The bridge: publishing +# --------------------------------------------------------------------------- + + +class TestPublishing: + async def test_registrations_made_before_start_are_published_at_start(self, services, backend, events): + record = services.register("weather.v1", "2.0", {"format": "json"}) + bridge = DnsSdBridge(services, on_event=events.append, backend=backend) + await bridge.start(address=OUR_ADDRESS, port=7878, addresses=["192.168.1.9", "fe80::1"]) + try: + name = service_instance_name(OUR_ADDRESS, "weather.v1") + assert bridge.published() == ["weather.v1"] + _, kw = backend.registered[name] + assert kw == { + "name": name, + "port": 7878, + "txt": txt_record(record, OUR_ADDRESS), + "server": service_host_name(OUR_ADDRESS), + "addresses": ["192.168.1.9", "fe80::1"], + } + finally: + await bridge.stop() + + async def test_a_registration_made_while_running_is_published(self, bridge, services, backend): + services.register("wiki.first-aid", "1.0") + await settle() + assert service_instance_name(OUR_ADDRESS, "wiki.first-aid") in backend.registered + + async def test_an_unregistration_withdraws(self, bridge, services, backend): + services.register("wiki.first-aid") + await settle() + handle, _ = backend.registered[service_instance_name(OUR_ADDRESS, "wiki.first-aid")] + services.unregister("wiki.first-aid") + await settle() + assert backend.registered == {} + assert backend.unregistered == [handle] + assert bridge.published() == [] + + async def test_registering_again_replaces_the_instance(self, bridge, services, backend): + services.register("wiki", "1.0") + await settle() + name = service_instance_name(OUR_ADDRESS, "wiki") + first, _ = backend.registered[name] + services.register("wiki", "2.0") + await settle() + second, kw = backend.registered[name] + assert second is not first + assert first in backend.unregistered + assert kw["txt"]["ver"] == "2.0" + assert bridge.published() == ["wiki"] + + async def test_publishing_the_same_id_twice_directly_keeps_one_instance(self, bridge, backend): + # The listener path withdraws before it publishes; a direct caller + # (start() over an id something already published) must not leave + # two handles for one name. + assert await bridge.publish(ServiceRecord("wiki", "1.0")) is True + name = service_instance_name(OUR_ADDRESS, "wiki") + first, _ = backend.registered[name] + assert await bridge.publish(ServiceRecord("wiki", "2.0")) is True + second, kw = backend.registered[name] + assert second is not first + assert backend.unregistered == [first] + assert kw["txt"]["ver"] == "2.0" + assert bridge.published() == ["wiki"] + + async def test_a_port_of_none_publishes_zero(self, services, backend): + services.register("s") + bridge = DnsSdBridge(services, backend=backend) + await bridge.start(address=OUR_ADDRESS, addresses=["10.0.0.1"]) + try: + _, kw = backend.registered[service_instance_name(OUR_ADDRESS, "s")] + assert kw["port"] == 0 + finally: + await bridge.stop() + + async def test_a_descriptor_that_does_not_fit_is_not_published_and_the_mesh_keeps_it( + self, bridge, services, backend, mesh, caplog + ): + with caplog.at_level(logging.WARNING): + services.register("a" * 201) + await settle() + assert backend.registered == {} + assert bridge.published() == [] + assert services.registered()[0].service_id == "a" * 201 + mesh.register_service.assert_called_once() + assert "not published on the LAN" in caplog.text and "201 bytes" in caplog.text + + async def test_publish_refuses_anything_but_this_nodes_own(self, bridge): + # Invariant 2: a LAN claim never goes out under this node's identity. + lan = ServiceRecord("s", provider=PEER_ADDRESS, source=SOURCE_LAN) + with pytest.raises(ValueError, match="invariant 2"): + await bridge.publish(lan) + + async def test_publish_before_start_does_nothing(self, services, backend): + bridge = DnsSdBridge(services, backend=backend) + assert await bridge.publish(ServiceRecord("s")) is False + assert backend.registered == {} + + async def test_a_listener_call_before_start_or_after_stop_is_dropped(self, services, backend): + bridge = DnsSdBridge(services, backend=backend) + bridge.on_registered(ServiceRecord("s")) # no loop yet: ignored + await bridge.start(address=OUR_ADDRESS, addresses=["10.0.0.1"]) + await bridge.stop() + bridge.on_registered(ServiceRecord("s")) + await settle() + assert backend.registered == {} + + async def test_browse_only_publishes_nothing(self, services, backend): + services.register("s") + bridge = DnsSdBridge(services, backend=backend, publish=False) + await bridge.start(address=OUR_ADDRESS) + try: + assert backend.registered == {} + assert backend.on_change is not None + finally: + await bridge.stop() + + async def test_publish_only_does_not_browse(self, services, backend): + bridge = DnsSdBridge(services, backend=backend, browse=False) + await bridge.start(address=OUR_ADDRESS, addresses=["10.0.0.1"]) + try: + assert backend.on_change is None + finally: + await bridge.stop() + + +# --------------------------------------------------------------------------- +# The bridge: importing +# --------------------------------------------------------------------------- + + +def peer_txt(**overrides: Any) -> bytes: + entries = {"txtvers": "1", "sid": "weather.v1", "ver": "2.0", "addr": PEER_ADDRESS, "c.format": "json"} + entries.update(overrides) + return encode_txt(entries) + + +PEER_NAME = service_instance_name(PEER_ADDRESS, "weather.v1") + + +class TestImporting: + async def test_a_browsed_instance_is_resolved_imported_and_announced(self, bridge, backend, events): + backend.txt_by_name[PEER_NAME] = peer_txt() + backend.on_change(PEER_NAME, False) + await settle() + assert backend.resolve_calls == [PEER_NAME] + assert events == [ + { + "type": "service_discovered", + "source": "lan", + "query_id": "", + "service_id": "weather.v1", + "version": "2.0", + "provider_peer_id": PEER_ADDRESS, + "capabilities": {"format": "json"}, + "hop_count": 0, + } + ] + assert bridge.lan_services() == [ + ServiceRecord("weather.v1", "2.0", {"format": "json"}, provider=PEER_ADDRESS, source=SOURCE_LAN) + ] + + async def test_an_import_never_reaches_the_engine(self, bridge, backend, mesh): + # Invariant 2, at the seam that matters: nothing on the import path + # calls the engine, whatever the record says. + backend.txt_by_name[PEER_NAME] = peer_txt() + backend.on_change(PEER_NAME, False) + await settle() + assert bridge.lan_services() != [] + mesh.register_service.assert_not_called() + mesh.unregister_service.assert_not_called() + + async def test_our_own_record_is_ignored(self, bridge, backend, events): + name = service_instance_name(OUR_ADDRESS, "weather.v1") + backend.txt_by_name[name] = peer_txt(addr=OUR_ADDRESS) + backend.on_change(name, False) + await settle() + assert events == [] and bridge.lan_services() == [] + + async def test_an_unresolvable_instance_is_not_imported(self, bridge, backend, events): + backend.on_change(PEER_NAME, False) + await settle() + assert events == [] and bridge.lan_services() == [] + + async def test_a_record_the_chapter_ignores_removes_an_earlier_import(self, bridge, backend, events): + backend.txt_by_name[PEER_NAME] = peer_txt() + backend.on_change(PEER_NAME, False) + await settle() + assert len(bridge.lan_services()) == 1 + backend.txt_by_name[PEER_NAME] = peer_txt(txtvers="2") + backend.on_change(PEER_NAME, False) + await settle() + assert bridge.lan_services() == [] + assert len(events) == 1 + + async def test_a_changed_record_is_announced_again_and_an_unchanged_one_is_not(self, bridge, backend, events): + backend.txt_by_name[PEER_NAME] = peer_txt() + backend.on_change(PEER_NAME, False) + await settle() + backend.on_change(PEER_NAME, False) # Updated with the same content + await settle() + assert len(events) == 1 + backend.txt_by_name[PEER_NAME] = peer_txt(ver="3.0") + backend.on_change(PEER_NAME, False) + await settle() + assert len(events) == 2 and events[1]["version"] == "3.0" + assert [r.version for r in bridge.lan_services()] == ["3.0"] + + async def test_a_removed_instance_leaves_the_registry(self, bridge, backend): + backend.txt_by_name[PEER_NAME] = peer_txt() + backend.on_change(PEER_NAME, False) + await settle() + backend.on_change(PEER_NAME, True) + assert bridge.lan_services() == [] + + async def test_two_services_from_one_neighbour_are_two_imports(self, bridge, backend): + other = service_instance_name(PEER_ADDRESS, "wiki") + backend.txt_by_name[PEER_NAME] = peer_txt() + backend.txt_by_name[other] = peer_txt(sid="wiki") + backend.on_change(PEER_NAME, False) + backend.on_change(other, False) + await settle() + assert sorted(r.service_id for r in bridge.lan_services()) == ["weather.v1", "wiki"] + + async def test_a_failing_event_handler_is_contained(self, services, backend, caplog): + handler = MagicMock(side_effect=RuntimeError("boom")) + bridge = DnsSdBridge(services, on_event=handler, backend=backend, publish=False) + await bridge.start(address=OUR_ADDRESS) + try: + backend.txt_by_name[PEER_NAME] = peer_txt() + backend.on_change(PEER_NAME, False) + await settle() + assert len(bridge.lan_services()) == 1 + assert "event handler failed" in caplog.text + finally: + await bridge.stop() + + +class TestLifetime: + async def test_an_import_is_re_resolved_at_half_the_ttl_and_dropped_at_the_ttl(self, bridge, backend): + backend.txt_by_name[PEER_NAME] = peer_txt() + bridge.import_txt(PEER_NAME, peer_txt(), now=1000.0) + backend.resolve_calls.clear() + + await bridge.sweep(now=1000.0 + 49.0) # under half of 100: untouched + assert backend.resolve_calls == [] + await bridge.sweep(now=1000.0 + 50.0) # half: re-resolved, and it answers + assert backend.resolve_calls == [PEER_NAME] + assert len(bridge.lan_services()) == 1 + + del backend.txt_by_name[PEER_NAME] # the neighbour is gone without a goodbye + await bridge.sweep(now=1050.0 + 50.0) # half again since the refresh: asked, no answer + assert backend.resolve_calls == [PEER_NAME, PEER_NAME] + assert len(bridge.lan_services()) == 1 # still within the ttl of its last answer + await bridge.sweep(now=1050.0 + 99.0) + assert len(bridge.lan_services()) == 1 + await bridge.sweep(now=1050.0 + 100.0) # the ttl since the last answer + assert bridge.lan_services() == [] + + async def test_a_successful_re_resolution_refreshes_the_clock_silently(self, bridge, backend, events): + bridge.import_txt(PEER_NAME, peer_txt(), now=0.0) + backend.txt_by_name[PEER_NAME] = peer_txt() + await bridge.sweep(now=60.0) + await bridge.sweep(now=120.0) # would have expired at 100 without the refresh at 60 + assert len(bridge.lan_services()) == 1 + assert len(events) == 1 + + async def test_the_sweeper_runs_on_its_own(self, services, backend): + bridge = DnsSdBridge(services, backend=backend, publish=False, ttl=0.04) + await bridge.start(address=OUR_ADDRESS) + try: + bridge.import_txt(PEER_NAME, peer_txt()) + await asyncio.sleep(0.15) + assert bridge.lan_services() == [] + assert PEER_NAME in backend.resolve_calls + finally: + await bridge.stop() + + def test_a_ttl_must_be_positive(self, services, backend): + with pytest.raises(ValueError): + DnsSdBridge(services, backend=backend, ttl=0) + + +# --------------------------------------------------------------------------- +# The bridge: lifecycle +# --------------------------------------------------------------------------- + + +class TestLifecycle: + async def test_stop_withdraws_everything_and_closes_the_responder(self, services, backend): + services.register("a") + services.register("b") + bridge = DnsSdBridge(services, backend=backend) + await bridge.start(address=OUR_ADDRESS, addresses=["10.0.0.1"]) + backend.txt_by_name[PEER_NAME] = peer_txt() + backend.on_change(PEER_NAME, False) + await settle() + assert len(backend.registered) == 2 and len(bridge.lan_services()) == 1 + await bridge.stop() + assert backend.registered == {} and len(backend.unregistered) == 2 + assert backend.closed and not bridge.running + assert bridge.lan_services() == [] and bridge.published() == [] + await bridge.stop() # idempotent + + async def test_start_is_idempotent(self, bridge, backend): + await bridge.start(address=PEER_ADDRESS, addresses=["10.0.0.2"]) + assert bridge.address == OUR_ADDRESS + + async def test_start_refuses_a_non_address(self, services, backend): + bridge = DnsSdBridge(services, backend=backend) + with pytest.raises(ValueError, match="canonical address"): + await bridge.start(address="alice", addresses=["10.0.0.1"]) + assert not bridge.running + + async def test_start_with_nothing_to_publish_from_fails_and_closes(self, services, backend): + bridge = DnsSdBridge(services, backend=backend) + with pytest.raises(RuntimeError, match="no interface address"): + await bridge.start(address=OUR_ADDRESS, addresses=[]) + assert backend.closed and not bridge.running + + async def test_a_responder_refusal_at_start_unwinds(self, services, backend): + # A name conflict or a closed socket at the first publish: the bridge + # must not stay half-started, listening to the registry with an open + # responder while reporting itself stopped. + services.register("a") + + async def refuse(**kw: Any) -> object: + raise RuntimeError("name conflict") + + backend.register = refuse # type: ignore[assignment] + bridge = DnsSdBridge(services, backend=backend) + with pytest.raises(RuntimeError, match="name conflict"): + await bridge.start(address=OUR_ADDRESS, addresses=["10.0.0.1"]) + assert not bridge.running + assert bridge not in services._listeners + assert backend.closed + assert backend.on_change is None + + async def test_the_listener_is_removed_at_stop(self, services, backend): + bridge = DnsSdBridge(services, backend=backend) + await bridge.start(address=OUR_ADDRESS, addresses=["10.0.0.1"]) + assert bridge in services._listeners + await bridge.stop() + assert bridge not in services._listeners + + def test_the_real_responder_registers_under_the_subtype(self): + # The one line of real zeroconf exercised here: the constructor's + # name check accepts an instance under the base type registered + # with the subtype, and refuses a name that is not under the type. + zeroconf = pytest.importorskip("zeroconf") + info = zeroconf.ServiceInfo( + SUBTYPE, + service_instance_name(OUR_ADDRESS, "s"), + port=0, + properties=txt_record(ServiceRecord("s"), OUR_ADDRESS), + server=service_host_name(OUR_ADDRESS), + parsed_addresses=["10.0.0.1"], + ) + assert info.type == SUBTYPE + assert record_from_txt(bytes(info.text)) == ServiceRecord( + "s", provider=OUR_ADDRESS, source=SOURCE_LAN + ) + with pytest.raises(zeroconf.BadTypeInNameException): + zeroconf.ServiceInfo(SUBTYPE, "s._other._tcp.local.", port=0) + + def test_the_real_responder_is_built_lazily(self): + pytest.importorskip("zeroconf") + backend = ZeroconfBackend() + assert backend is not None + asyncio.run(backend.close()) diff --git a/bindings/python/tests/test_peer_stream_manager.py b/bindings/python/tests/test_peer_stream_manager.py index b65e46830..90a568d4e 100644 --- a/bindings/python/tests/test_peer_stream_manager.py +++ b/bindings/python/tests/test_peer_stream_manager.py @@ -224,6 +224,19 @@ def test_record_without_addr_yields_nothing(self): assert peers_from_record(self._info(properties={b"addr": b"not-an-address"})) == [] assert peers_from_record(self._info(port=0)) == [] + def test_a_service_instance_is_not_a_peer_hint(self): + # dns-sd-mapping.md invariant 5: a record carrying `sid` is a service + # instance under the subtype, listed under this type as well by some + # responders, and would otherwise be one more connector per service. + assert ( + peers_from_record( + self._info( + properties={b"txtvers": b"1", b"sid": b"weather.v1", b"addr": PEER_ADDRESS.encode()} + ) + ) + == [] + ) + def test_record_falls_back_to_the_server_name(self): assert peers_from_record(self._info(addresses=[])) == [ PeerEntry("host.local", 7878, PEER_ADDRESS) diff --git a/bindings/python/tests/test_services.py b/bindings/python/tests/test_services.py new file mode 100644 index 000000000..1041194f9 --- /dev/null +++ b/bindings/python/tests/test_services.py @@ -0,0 +1,199 @@ +"""Tests for the application-level service wrappers. + +The generated ``MeshServices`` is a mock: what is under test is the copy of +the registry this side keeps, the order in which the engine and the copy +change, and the closed status set. +""" + +from __future__ import annotations + +import subprocess +import sys +from unittest.mock import MagicMock + +import pytest + +from offline_protocol_sdk.services import ( + SOURCE_LAN, + SOURCE_LOCAL, + VALID_STATUSES, + ServiceRecord, + Services, +) + + +class Recorder: + """A listener that writes down what it hears, in order.""" + + def __init__(self) -> None: + self.calls: list[tuple[str, ServiceRecord]] = [] + + def on_registered(self, record: ServiceRecord) -> None: + self.calls.append(("registered", record)) + + def on_unregistered(self, record: ServiceRecord) -> None: + self.calls.append(("unregistered", record)) + + +@pytest.fixture +def mesh() -> MagicMock: + m = MagicMock() + m.register_service = MagicMock(return_value=None) + m.unregister_service = MagicMock(return_value=True) + m.discover_services = MagicMock(return_value="query-1") + m.send_service_request = MagicMock(return_value="request-1") + m.respond_to_service_request = MagicMock(return_value="message-1") + return m + + +@pytest.fixture +def services(mesh: MagicMock) -> Services: + return Services(MagicMock(), mesh_services=mesh) + + +class TestLiterals: + def test_the_status_set_is_the_engines_three(self): + # As literals, not read from the module: a test that read them would + # agree with any edit. The engine's VALID_SERVICE_STATUSES. + assert VALID_STATUSES == ("ok", "not_found", "error") + + def test_the_sources(self): + assert SOURCE_LOCAL == "local" + assert SOURCE_LAN == "lan" + + +class TestRecord: + def test_capabilities_are_copied(self): + caps = {"format": "json"} + record = ServiceRecord("weather.v1", "2.0", caps) + caps["format"] = "xml" + assert record.capabilities == {"format": "json"} + assert record.is_local + assert record.provider is None + + def test_a_lan_record_is_not_local(self): + record = ServiceRecord("weather.v1", provider="off1abc", source=SOURCE_LAN) + assert not record.is_local + + +class TestRegistry: + def test_register_reaches_the_engine_first_and_then_the_copy(self, services, mesh): + listener = Recorder() + services.add_listener(listener) + record = services.register("weather.v1", "2.0", {"format": "json"}) + mesh.register_service.assert_called_once_with("weather.v1", "2.0", {"format": "json"}) + assert services.registered() == [record] + assert services.get("weather.v1") == record + assert listener.calls == [("registered", record)] + + def test_an_engine_refusal_leaves_the_copy_untouched(self, services, mesh): + listener = Recorder() + services.add_listener(listener) + mesh.register_service.side_effect = RuntimeError("refused") + with pytest.raises(RuntimeError): + services.register("__reserved") + assert services.registered() == [] + assert listener.calls == [] + + def test_registering_again_replaces_and_says_so(self, services): + listener = Recorder() + services.add_listener(listener) + first = services.register("weather.v1", "1.0") + second = services.register("weather.v1", "2.0") + assert services.registered() == [second] + assert listener.calls == [ + ("registered", first), + ("unregistered", first), + ("registered", second), + ] + + def test_unregister_returns_the_engines_answer_and_drops_the_copy(self, services, mesh): + listener = Recorder() + record = services.register("weather.v1") + services.add_listener(listener) + assert services.unregister("weather.v1") is True + mesh.unregister_service.assert_called_once_with("weather.v1") + assert services.registered() == [] + assert listener.calls == [("unregistered", record)] + + def test_unregister_of_an_id_the_engine_lost_still_drops_the_copy(self, services, mesh): + services.register("weather.v1") + mesh.unregister_service.return_value = False + assert services.unregister("weather.v1") is False + assert services.registered() == [] + + def test_unregister_of_an_unknown_id_notifies_nobody(self, services, mesh): + listener = Recorder() + services.add_listener(listener) + mesh.unregister_service.return_value = False + assert services.unregister("nothing") is False + assert listener.calls == [] + + def test_registration_order_is_kept(self, services): + a = services.register("a") + b = services.register("b") + c = services.register("c") + assert services.registered() == [a, b, c] + + def test_a_failing_listener_does_not_undo_the_registration(self, services, caplog): + broken = MagicMock() + broken.on_registered.side_effect = RuntimeError("boom") + services.add_listener(broken) + record = services.register("weather.v1") + assert services.registered() == [record] + assert "failed in on_registered" in caplog.text + + def test_a_listener_is_added_once_and_can_leave(self, services): + listener = Recorder() + services.add_listener(listener) + services.add_listener(listener) + services.register("a") + assert len(listener.calls) == 1 + services.remove_listener(listener) + services.remove_listener(listener) + services.register("b") + assert len(listener.calls) == 1 + + +class TestCalls: + def test_discover_and_request_delegate(self, services, mesh): + assert services.discover() == "query-1" + mesh.discover_services.assert_called_once_with(None) + assert services.discover("weather.v1") == "query-1" + assert services.request("off1peer", "weather.v1", "get", "{}") == "request-1" + mesh.send_service_request.assert_called_once_with("off1peer", "weather.v1", "get", "{}") + + @pytest.mark.parametrize("status", ["ok", "not_found", "error"]) + def test_a_valid_status_crosses(self, services, mesh, status): + assert services.respond("r-1", "off1peer", "weather.v1", status, "{}") == "message-1" + mesh.respond_to_service_request.assert_called_once_with("r-1", "off1peer", "weather.v1", status, "{}") + + @pytest.mark.parametrize("status", ["OK", "success", "", "not-found"]) + def test_a_status_the_engine_would_refuse_is_refused_here_with_the_reason(self, services, mesh, status): + with pytest.raises(ValueError, match="ok, not_found, error"): + services.respond("r-1", "off1peer", "weather.v1", status, "{}") + mesh.respond_to_service_request.assert_not_called() + + +class TestBaseInstall: + def test_the_wrappers_import_without_zeroconf(self): + # The base install has no zeroconf; both modules must import, and + # only the responder's constructor may say what is missing. In a + # fresh interpreter, because blocking the import here would reload + # the modules under every other test. + script = ( + "import sys\n" + "sys.modules['zeroconf'] = None\n" + "sys.modules['zeroconf.asyncio'] = None\n" + "from offline_protocol_sdk import services, dnssd_bridge\n" + "try:\n" + " dnssd_bridge.ZeroconfBackend()\n" + "except ImportError as exc:\n" + " assert 'offline-protocol-sdk[lan]' in str(exc), exc\n" + " print('refused')\n" + "else:\n" + " raise SystemExit('the responder was built without zeroconf')\n" + ) + done = subprocess.run([sys.executable, "-c", script], capture_output=True, text=True, timeout=120) + assert done.returncode == 0, done.stderr + assert done.stdout.strip() == "refused" diff --git a/docs/bridges/python.md b/docs/bridges/python.md index 0b6fd2c12..20bac7cd2 100644 --- a/docs/bridges/python.md +++ b/docs/bridges/python.md @@ -225,6 +225,19 @@ bounds as literals (`4`, `1_048_576`, `96`) for the C5 reason. DNS-SD needs the optional extra (`pip install 'offline-protocol-sdk[lan]'`); the base install carries no LGPL dependency. +The same type carries service instances under the subtype `_svc._sub` +([DNS-SD mapping](../spec/dns-sd-mapping.md)), published and read by +`dnssd_bridge.py` over the same optional extra. Two rules of that chapter are +the binding's to hold, because nothing in the core can see a LAN record: an +import is delivered with `source: "lan"` and never reaches +`register_service` (a registration made from an unsigned LAN record would go +out in signed discovery responses under this node's identity), and a peer +browser ignores any record carrying `sid` (each published service would +otherwise be one more connector to the same host). `test_dnssd_bridge.py` +asserts the chapter's bounds (`200`, `255`, `1300`, `63`) and the subtype as +literals for the C5 reason, and `services.py` mirrors the engine's closed +status set (`ok`, `not_found`, `error`) as a literal pinned the same way. + ## P10. Freeing a core object never blocks, wherever the interpreter frees it The generated bindings keep every callback object in one handle map behind From 41cdc4bc3127b54ec1619b404c0ee3097f56d49e Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Wed, 30 Sep 2026 21:47:29 +0530 Subject: [PATCH 3/6] docs(spec): TXT key case, the browsed set and what removal is not Keys are case-insensitive (RFC 6763 section 6.4): a reader folds a key before it looks it up, a publisher refuses a descriptor whose capability keys collide under folding, and the chapter says how its whole-record refusal on a duplicate relates to the RFC's keep-the-first rule. The importer keeps a browsed set apart from the listed set, because a browser reports an instance once and again only after its cache has forgotten it, so an importer that stopped re-resolving a name when it left the listed set would lose a neighbour that missed one window for up to 75 minutes; delivery is at least once. Removal is the importer dropping its own entry and never unregister_service, which the earlier sentence named. A subtype-only listing is not called conforming: section 7.1 defines a subtype as an additional PTR to an instance listed under the parent type. --- docs/spec/dns-sd-mapping.md | 68 +++++++++++++++++++++++++++++-------- 1 file changed, 53 insertions(+), 15 deletions(-) diff --git a/docs/spec/dns-sd-mapping.md b/docs/spec/dns-sd-mapping.md index be14c7aff..dd21af7e8 100644 --- a/docs/spec/dns-sd-mapping.md +++ b/docs/spec/dns-sd-mapping.md @@ -59,12 +59,13 @@ _svc._sub._offlineprotocol._tcp. ``` and a browser looking for services browses that subtype. RFC 6763 section -7.1 lists a subtyped instance under its base type as well, and a responder -that follows it answers a browse of the base type with service instances, -which is why invariant 5 exists and why a peer browser keys on the `sid` -entry rather than on the name it was browsing. A responder that lists the -instance under the subtype only is conforming too; a browser MUST NOT rely -on either behaviour. +7.1 defines a subtype as an additional PTR to an instance that is listed +under its parent type, so a responder that follows it answers a browse of +the base type with service instances as well, which is why invariant 5 +exists and why a peer browser keys on the `sid` entry rather than on the +name it was browsing. Some responders answer only the subtype browse +(python-zeroconf indexes an instance under the one type it was registered +with); a browser MUST NOT rely on either behaviour. The domain is `local.` on a multicast LAN. Nothing here depends on it. @@ -109,12 +110,28 @@ The record is a sequence of `key=value` strings in this order: | `sid=` | The descriptor's `service_id`, as UTF-8 | Yes | | `ver=` | The descriptor's `version`, as UTF-8; the engine does not parse it and neither does this mapping | Yes, possibly empty | | `addr=` | The publisher's canonical address, which is what a `service_discovered` event names as `provider_peer_id` | Yes | -| `c.=` | One string per capability, keys in byte order, so that one descriptor is one record | Zero or more | +| `c.=` | One string per capability, keys in the byte order of their lower-case form, so that one descriptor is one record | Zero or more | + +Keys are case-insensitive, as RFC 6763 section 6.4 makes every TXT key: +`SID` is `sid`, and `c.Foo` and `c.foo` are one key. A reader folds a key +to lower case before it looks the key up, so an imported capability key is +the publisher's key in lower case. A publisher whose descriptor holds two +capability keys that are one key under case folding MUST refuse the +descriptor rather than publish either or both (invariant 4): a record with +one of them dropped is a different claim, and a record with both is the +duplicate below. A reader MUST ignore a record whose `txtvers` is not `1`, and MUST ignore a record missing `sid` or `addr`, or whose `addr` is not an address. Unknown -keys are ignored. A key that appears twice makes the record malformed and it -is ignored whole. +keys are ignored. A key that appears twice, under case folding, makes the +record malformed and it is ignored whole. RFC 6763 section 6.4 has a client +keep the first of a duplicated key and ignore the rest; this mapping is +stricter about what it accepts as its own record, in the way it is about +`txtvers`: a publisher under this chapter never emits a duplicate, and two +readers that kept different occurrences would disagree on which address a +record claims, so a duplicate marks a record that is not this mapping's. A +reader that applies the RFC's rule instead keeps the first occurrence and +never the last. ### Bounds @@ -163,12 +180,33 @@ own record coming back. A DNS-SD record has the lifetime its publisher gave it, and a responder that dies without a goodbye leaves its records to age out in every browser's cache, which for the PTR is 75 minutes by RFC 6762 section 10. An importer therefore -owes the application a shorter liveness rule of its own: an entry is kept -while the instance still resolves, re-resolved at half the importer's time to -live, and removed from the application-level registry when it has not resolved -within that time to live or when the browser reports the instance gone, -whichever is first. Removal is the shadow registry's `unregister`; it never -reaches the engine, because the engine never held the entry. +owes the application a shorter liveness rule of its own, and it keeps two +sets to hold it: + +- **The browsed set** is every instance name the browser has reported and + not yet reported gone. A name enters it when the browser reports the + instance and leaves it only when the browser reports the instance gone. + The importer re-resolves every browsed name at half its time to live, + whether the name is listed or not. +- **The listed set** is what the application sees: every browsed name whose + record resolved and was well formed within the importer's time to live. + An entry leaves it when it has not resolved within that time, or when the + browser reports the instance gone, whichever is first. + +The two sets are kept apart because a browser reports an instance once: it +reports the instance again only when its own cache has forgotten the record, +which for the PTR is the 75 minutes above. An importer that stopped +re-resolving a name when the name left the listed set would lose a +neighbour that missed one resolve window until that cache turned over. + +Delivery to the application is therefore at least once: an entry that +expires from the listed set and resolves again is announced again, with the +same event. An application that keys its own state on `(address, +service_id)` sees the second announcement as a refresh. + +Removal is the importer dropping the entry from its own registry; it never +calls `unregister_service`: the engine never held the entry, and the id may +be one this node offers itself. Nothing on the mesh side has a lifetime: a mesh registration stands until it is unregistered, and a mesh discovery response is a point-in-time answer with From cb08fe72c6cb45483711b2ea8a0ecfa6cf4a165c Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Wed, 30 Sep 2026 21:47:29 +0530 Subject: [PATCH 4/6] fix(bindings): re-resolve every browsed name, fold TXT keys, log a refused publish A dropped import never came back: the sweep re-resolved listed entries only, and python-zeroconf reports a name again only once its cache has forgotten the PTR. The bridge keeps the browsed set apart from the listed set, re-resolves every browsed name at half the time to live whether listed or not, gates listing on the time to live, and drops a name from the browsed set only on the browser's Removed; a resolve that lands after Removed lists nothing. TXT keys are folded to lower case before the fixed-key match and the duplicate check, and a descriptor whose capability keys collide under folding is refused. A publish or withdraw scheduled from the registry listener runs under a guard that logs the service id and the reason, where a responder refusal used to surface as an unretrieved task exception at collection. The seam test asserts the mocked engine saw no call at all. The Python README names python-zeroconf and ifaddr as the lan extra's runtime dependencies under the License section. --- bindings/python/README.md | 12 ++ .../offline_protocol_sdk/dnssd_bridge.py | 120 +++++++++++++---- bindings/python/tests/test_dnssd_bridge.py | 125 +++++++++++++++++- 3 files changed, 223 insertions(+), 34 deletions(-) diff --git a/bindings/python/README.md b/bindings/python/README.md index cbf394864..31ed023e8 100644 --- a/bindings/python/README.md +++ b/bindings/python/README.md @@ -413,3 +413,15 @@ Both license texts, along with `THIRD-PARTY-NOTICES.md`, are also installed with package under `offline_protocol_sdk-.dist-info/licenses/`. The links above are absolute because this file is the PyPI long description, and PyPI does not resolve repository-relative links. + +### Optional dependencies + +`THIRD-PARTY-NOTICES.md` covers the Rust crates compiled into the native +library. The `lan` extra (`pip install 'offline-protocol-sdk[lan]'`) adds two +runtime dependencies that pip installs from PyPI and that are never +redistributed in this wheel: [python-zeroconf](https://pypi.org/project/zeroconf/) +(LGPL-2.1-or-later), used by `peer_stream_manager.py` and `dnssd_bridge.py` +for DNS-SD, and [ifaddr](https://pypi.org/project/ifaddr/) (MIT), used to +list the interface addresses a record publishes. Both are imported only when a +manager or bridge is asked to advertise or discover; the base install imports +neither. diff --git a/bindings/python/offline_protocol_sdk/dnssd_bridge.py b/bindings/python/offline_protocol_sdk/dnssd_bridge.py index ab9927f8c..3a793f0df 100644 --- a/bindings/python/offline_protocol_sdk/dnssd_bridge.py +++ b/bindings/python/offline_protocol_sdk/dnssd_bridge.py @@ -122,8 +122,20 @@ def txt_record(record: ServiceRecord, address: str) -> dict[str, str]: "ver": record.version, "addr": address, } - for key in sorted(record.capabilities, key=lambda k: k.encode("utf-8")): + # TXT keys are case-insensitive (RFC 6763 section 6.4), so two + # capability keys that fold to one are one key: refused, never one of + # them dropped (a different claim) or both published (a duplicate the + # reader ignores the record for). + folded: dict[str, str] = {} + for key in record.capabilities: _check_key(key) + other = folded.setdefault(key.lower(), key) + if other != key: + raise RecordRefused( + f"capability keys {other!r} and {key!r} are one key under case folding " + "(RFC 6763 section 6.4)" + ) + for key in sorted(record.capabilities, key=lambda k: k.lower().encode("utf-8")): entries[_CAPABILITY_PREFIX + key] = record.capabilities[key] for key, value in entries.items(): string = len(key.encode("utf-8")) + 1 + len(value.encode("utf-8")) @@ -175,19 +187,22 @@ def record_from_txt(raw: bytes) -> ServiceRecord | None: pairs = parse_txt(raw) if pairs is None or len(raw) > MAX_TXT_RECORD_BYTES: return None - seen: set[bytes] = set() + seen: set[str] = set() fields: dict[str, str] = {} capabilities: dict[str, str] = {} for key, value in pairs: - if key in seen: - return None - seen.add(key) try: - name = key.decode("ascii") + # Keys are case-insensitive (RFC 6763 section 6.4): folded + # before the duplicate check and the fixed-key match, so + # `SID=` is `sid=` and `c.Foo` is `c.foo`. + name = key.decode("ascii").lower() _check_key(name) text = "" if value is None else value.decode("utf-8") except (UnicodeDecodeError, RecordRefused): return None + if name in seen: + return None + seen.add(name) if name in _TXT_KEYS_FIXED: fields[name] = text elif name.startswith(_CAPABILITY_PREFIX) and len(name) > len(_CAPABILITY_PREFIX): @@ -334,7 +349,16 @@ def __init__( self._port = 0 self._addresses: list[str] = [] self._published: dict[str, Any] = {} + # The listed set: what the application sees. self._lan: dict[str, _LanEntry] = {} + # The browsed set: every name the browser reported and has not + # reported gone, by the time of its last resolve attempt. Kept apart + # from the listed set because a browser reports a name once, and + # again only when its own cache has forgotten it (75 minutes for the + # PTR): a bridge that stopped re-resolving a name when the name left + # the listed set would lose a neighbour that missed one resolve + # window for that long. + self._browsed: dict[str, float] = {} self._tasks: set[asyncio.Task[Any]] = set() self._sweeper: asyncio.Task[None] | None = None self._running = False @@ -414,6 +438,7 @@ async def stop(self) -> None: except Exception: logger.exception("withdrawing %r from the LAN failed", service_id) self._lan.clear() + self._browsed.clear() try: await self._backend.close() finally: @@ -459,55 +484,91 @@ async def withdraw(self, record: ServiceRecord) -> None: # ServicesListener: called on whatever thread changed the registry. def on_registered(self, record: ServiceRecord) -> None: - self._later(lambda: self.publish(record)) + self._later( + lambda: self.publish(record), + f"service {record.service_id!r} not published on the LAN", + ) def on_unregistered(self, record: ServiceRecord) -> None: - self._later(lambda: self.withdraw(record)) + self._later( + lambda: self.withdraw(record), + f"service {record.service_id!r} not withdrawn from the LAN", + ) - def _later(self, make: Callable[[], Any]) -> None: + def _later(self, make: Callable[[], Any], what: str) -> None: loop = self._loop if loop is None or not self._running: return - loop.call_soon_threadsafe(self._spawn, make) + loop.call_soon_threadsafe(self._spawn, make, what) - def _spawn(self, make: Callable[[], Any]) -> None: + def _spawn(self, make: Callable[[], Any], what: str) -> None: if not self._running or self._loop is None: return - task = self._loop.create_task(make()) + task = self._loop.create_task(self._guarded(make, what)) self._tasks.add(task) task.add_done_callback(self._tasks.discard) + @staticmethod + async def _guarded(make: Callable[[], Any], what: str) -> None: + """Run one scheduled step and log what it was when it fails. A task + nobody awaits otherwise reports "exception was never retrieved" at + collection, without the service id, or nothing at all.""" + try: + await make() + except asyncio.CancelledError: + raise + except Exception as exc: + logger.warning("%s: %s", what, exc, exc_info=True) + # -- importing ------------------------------------------------------------ def lan_services(self) -> list[ServiceRecord]: - """What the LAN's neighbours currently claim to offer.""" + """What the LAN's neighbours currently claim to offer: the listed + set.""" return [entry.record for entry in self._lan.values()] + def browsed_names(self) -> list[str]: + """Every instance name the browser has reported and not yet reported + gone, listed or not: the browsed set, which the sweep re-resolves.""" + return list(self._browsed) + def _on_change(self, name: str, removed: bool) -> None: if removed: + self._browsed.pop(name, None) self._lan.pop(name, None) return - self._spawn(lambda: self._resolve_and_import(name)) + self._browsed[name] = time.monotonic() + self._spawn( + lambda: self._resolve_and_import(name), + f"instance {name!r} not resolved", + ) - async def _resolve_and_import(self, name: str) -> None: + async def _resolve_and_import(self, name: str, now: float | None = None) -> None: if self._backend is None: return raw = await self._backend.resolve(name, self._resolve_timeout_ms) + if name not in self._browsed: + # Reported gone while the resolve was in flight: an import now + # would list an instance the browser has already withdrawn. + return if raw is not None: - self.import_txt(name, raw) + self.import_txt(name, raw, now=now) def import_txt(self, name: str, raw: bytes, now: float | None = None) -> bool: """Take a resolved instance's TXT record into the LAN registry. - Returns whether an event was delivered: a new instance, or one whose - record changed. A record the chapter says to ignore removes any - earlier import under that name; this node's own record is ignored.""" + Returns whether an event was delivered: a new instance, one whose + record changed, or one that had expired from the listed set and + resolved again (delivery is at least once). A record the chapter + says to ignore removes any earlier import under that name; this + node's own record is ignored. A resolved name is a browsed name.""" + seen_at = time.monotonic() if now is None else now + self._browsed[name] = seen_at record = record_from_txt(raw) if record is None: self._lan.pop(name, None) return False if record.provider == self._address: return False - seen_at = time.monotonic() if now is None else now previous = self._lan.get(name) self._lan[name] = _LanEntry(record, seen_at) if previous is not None and previous.record == record: @@ -516,18 +577,19 @@ def import_txt(self, name: str, raw: bytes, now: float | None = None) -> bool: return True async def sweep(self, now: float | None = None) -> None: - """Re-resolve every import older than half the time to live, and - drop every one older than the whole of it. The periodic task calls - this; a test calls it with a clock of its own.""" + """Drop every listed entry older than the time to live, then + re-resolve every browsed name whose last resolve attempt is older + than half of it, listed or not. The periodic task calls this; a + test calls it with a clock of its own.""" current = time.monotonic() if now is None else now for name, entry in list(self._lan.items()): - age = current - entry.seen_at - if age >= self._ttl: + if current - entry.seen_at >= self._ttl: self._lan.pop(name, None) - elif age >= self._ttl / 2 and self._backend is not None: - raw = await self._backend.resolve(name, self._resolve_timeout_ms) - if raw is not None: - self.import_txt(name, raw, now=current) + for name, attempted_at in list(self._browsed.items()): + if current - attempted_at < self._ttl / 2 or name not in self._browsed: + continue + self._browsed[name] = current + await self._resolve_and_import(name, now=current) async def _sweep_forever(self) -> None: while self._running: diff --git a/bindings/python/tests/test_dnssd_bridge.py b/bindings/python/tests/test_dnssd_bridge.py index e3a1d8c3e..582a1ce12 100644 --- a/bindings/python/tests/test_dnssd_bridge.py +++ b/bindings/python/tests/test_dnssd_bridge.py @@ -10,6 +10,7 @@ import asyncio import logging +import time from typing import Any from unittest.mock import MagicMock @@ -197,9 +198,17 @@ def test_the_entries_and_their_order(self): ("c.format", "json"), ] - def test_capability_keys_are_in_byte_order(self): - entries = txt_record(ServiceRecord("s", "", {"b": "1", "B": "2", "a": "3"}), OUR_ADDRESS) - assert [k for k in entries if k.startswith("c.")] == ["c.B", "c.a", "c.b"] + def test_capability_keys_are_in_the_byte_order_of_their_lower_case_form(self): + # Raw byte order would put `B` (0x42) before `a` (0x61). + entries = txt_record(ServiceRecord("s", "", {"B": "1", "a": "2", "c": "3"}), OUR_ADDRESS) + assert [k for k in entries if k.startswith("c.")] == ["c.a", "c.B", "c.c"] + + def test_capability_keys_that_fold_to_one_key_are_refused_not_dropped(self): + # RFC 6763 section 6.4: keys are case-insensitive. Publishing both + # is a duplicate the reader ignores the record for; publishing one + # is a different claim. + with pytest.raises(RecordRefused, match="'Foo' and 'foo' are one key"): + txt_record(ServiceRecord("s", "", {"Foo": "1", "foo": "2"}), OUR_ADDRESS) def test_the_size_counts_a_length_octet_per_string(self): entries = {"txtvers": "1", "sid": "s"} @@ -314,6 +323,22 @@ def test_a_key_twice_is_ignored_whole(self): raw = raw_strings([b"txtvers=1", b"sid=a", b"sid=b", b"addr=" + PEER_ADDRESS.encode()]) assert record_from_txt(raw) is None + def test_keys_are_read_case_insensitively(self): + # RFC 6763 section 6.4: `SID=` is `sid=`, and `c.Foo` is `c.foo`. + raw = raw_strings( + [b"TXTVERS=1", b"SID=weather.v1", b"Addr=" + PEER_ADDRESS.encode(), b"c.Foo=1"] + ) + record = record_from_txt(raw) + assert record is not None + assert record.service_id == "weather.v1" + assert record.capabilities == {"foo": "1"} + + def test_a_key_twice_under_case_folding_is_ignored_whole(self): + raw = raw_strings([b"txtvers=1", b"sid=a", b"SID=b", b"addr=" + PEER_ADDRESS.encode()]) + assert record_from_txt(raw) is None + raw = raw_strings([b"txtvers=1", b"sid=a", b"addr=" + PEER_ADDRESS.encode(), b"c.Foo=1", b"c.foo=2"]) + assert record_from_txt(raw) is None + def test_a_value_that_is_not_utf8_is_ignored(self): raw = raw_strings([b"txtvers=1", b"sid=\xff\xfe", b"addr=" + PEER_ADDRESS.encode()]) assert record_from_txt(raw) is None @@ -440,6 +465,24 @@ async def test_a_descriptor_that_does_not_fit_is_not_published_and_the_mesh_keep mesh.register_service.assert_called_once() assert "not published on the LAN" in caplog.text and "201 bytes" in caplog.text + async def test_a_responder_refusal_while_running_is_logged_with_the_service_id( + self, bridge, services, backend, caplog + ): + # The listener path runs in a task nobody awaits. Without the guard + # a refusal surfaces as "Task exception was never retrieved" at + # collection, without the service id, or not at all. + async def refuse(**kw: Any) -> object: + raise RuntimeError("name conflict") + + backend.register = refuse # type: ignore[assignment] + with caplog.at_level(logging.WARNING): + services.register("wiki") + await settle() + assert bridge.published() == [] + assert bridge._tasks == set() + assert "service 'wiki' not published on the LAN: name conflict" in caplog.text + assert "never retrieved" not in caplog.text + async def test_publish_refuses_anything_but_this_nodes_own(self, bridge): # Invariant 2: a LAN claim never goes out under this node's identity. lan = ServiceRecord("s", provider=PEER_ADDRESS, source=SOURCE_LAN) @@ -522,8 +565,9 @@ async def test_an_import_never_reaches_the_engine(self, bridge, backend, mesh): backend.on_change(PEER_NAME, False) await settle() assert bridge.lan_services() != [] - mesh.register_service.assert_not_called() - mesh.unregister_service.assert_not_called() + # The whole engine surface, not two methods of it: a discover or a + # request on the import path would be as wrong as a registration. + assert mesh.mock_calls == [] async def test_our_own_record_is_ignored(self, bridge, backend, events): name = service_instance_name(OUR_ADDRESS, "weather.v1") @@ -631,6 +675,77 @@ async def test_the_sweeper_runs_on_its_own(self, services, backend): finally: await bridge.stop() + async def test_a_dropped_import_comes_back_when_it_resolves_again(self, bridge, backend, events): + # The browser reports a name once, and again only when its own + # cache has forgotten it (75 minutes for the PTR). A bridge that + # stopped re-resolving a name when it left the listed set would + # lose a neighbour that missed one resolve window for that long. + base = time.monotonic() + backend.txt_by_name[PEER_NAME] = peer_txt() + backend.on_change(PEER_NAME, False) + await settle() + assert len(events) == 1 and len(bridge.lan_services()) == 1 + + # `base` was read before the browser reported the name, so every + # clock below is a second past the boundary it exercises. + del backend.txt_by_name[PEER_NAME] # one missed window + await bridge.sweep(now=base + 101.0) + assert bridge.lan_services() == [] + assert bridge.browsed_names() == [PEER_NAME] + + backend.txt_by_name[PEER_NAME] = peer_txt() # the neighbour is back + await bridge.sweep(now=base + 152.0) + assert len(bridge.lan_services()) == 1 + assert len(events) == 2, "delivery is at least once: an expired entry is announced again" + assert events[1] == events[0] + + backend.on_change(PEER_NAME, True) # only Removed leaves the browsed set + assert bridge.browsed_names() == [] and bridge.lan_services() == [] + backend.resolve_calls.clear() + await bridge.sweep(now=base + 300.0) + assert backend.resolve_calls == [] + + async def test_a_name_that_never_resolved_is_re_resolved_until_it_does(self, bridge, backend, events): + base = time.monotonic() + backend.on_change(PEER_NAME, False) # reported, but the resolve times out + await settle() + assert bridge.lan_services() == [] and bridge.browsed_names() == [PEER_NAME] + assert backend.resolve_calls == [PEER_NAME] + await bridge.sweep(now=base + 49.0) + assert backend.resolve_calls == [PEER_NAME] + await bridge.sweep(now=base + 51.0) # `base` predates the report by a hair + assert backend.resolve_calls == [PEER_NAME, PEER_NAME] + await bridge.sweep(now=base + 52.0) # a failed attempt still counts as one + assert backend.resolve_calls == [PEER_NAME, PEER_NAME] + backend.txt_by_name[PEER_NAME] = peer_txt() + await bridge.sweep(now=base + 102.0) + assert len(bridge.lan_services()) == 1 and len(events) == 1 + + async def test_a_resolve_that_lands_after_the_instance_is_gone_lists_nothing(self, services, events): + class SlowBackend(FakeBackend): + def __init__(self) -> None: + super().__init__() + self.release = asyncio.Event() + + async def resolve(self, name: str, timeout_ms: int) -> bytes | None: + await self.release.wait() + return await super().resolve(name, timeout_ms) + + backend = SlowBackend() + backend.txt_by_name[PEER_NAME] = peer_txt() + bridge = DnsSdBridge(services, on_event=events.append, backend=backend, publish=False) + await bridge.start(address=OUR_ADDRESS) + try: + backend.on_change(PEER_NAME, False) + await settle() + backend.on_change(PEER_NAME, True) # gone while the resolve is in flight + backend.release.set() + await settle() + assert bridge.lan_services() == [] and bridge.browsed_names() == [] + assert events == [] + finally: + await bridge.stop() + def test_a_ttl_must_be_positive(self, services, backend): with pytest.raises(ValueError): DnsSdBridge(services, backend=backend, ttl=0) From 87da5cb3798af3f3d499cd0b66021cdad615ec26 Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Wed, 30 Sep 2026 23:02:29 +0530 Subject: [PATCH 5/6] docs(spec): an importer re-resolves before it expires With a sweep at half the time to live, one missed window puts the next attempt on the drop boundary. Resolving first refreshes an entry whose resolve on the boundary answers; dropping first would show the application an empty list during the resolve and a duplicate announcement after it. --- docs/spec/dns-sd-mapping.md | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/docs/spec/dns-sd-mapping.md b/docs/spec/dns-sd-mapping.md index dd21af7e8..9999cf2e4 100644 --- a/docs/spec/dns-sd-mapping.md +++ b/docs/spec/dns-sd-mapping.md @@ -202,7 +202,12 @@ neighbour that missed one resolve window until that cache turned over. Delivery to the application is therefore at least once: an entry that expires from the listed set and resolves again is announced again, with the same event. An application that keys its own state on `(address, -service_id)` sees the second announcement as a refresh. +service_id)` sees the second announcement as a refresh. At a sweep the +importer re-resolves before it expires, so an entry whose resolve on the +boundary answers is refreshed rather than dropped and announced again: with +a sweep at half the time to live, one missed window puts the next attempt on +that boundary, and the other order would show the application an empty list +during the resolve and a duplicate announcement after it. Removal is the importer dropping the entry from its own registry; it never calls `unregister_service`: the engine never held the entry, and the id may From 7a6e0aa9f2c830ae71260f1eddc7d49ffaf503c3 Mon Sep 17 00:00:00 2001 From: bahdotsh Date: Wed, 30 Sep 2026 23:02:29 +0530 Subject: [PATCH 6/6] fix(bindings): resolve every browsed name before dropping listed entries The sweep dropped listed entries older than the time to live before it re-resolved browsed names older than half of it. With the sweep at half the time to live, one failed resolve put the next attempt on the drop boundary, so a single missed window whose next attempt succeeded was observed as an empty listed set during the resolve and a duplicate service_discovered after it. The two loops are swapped: an answer on the boundary now refreshes the entry and emits nothing. A test with an explicit clock observes the listed set from inside the resolve and pins the ordering; swapping the loops back fails it. --- .../offline_protocol_sdk/dnssd_bridge.py | 20 +++++++---- bindings/python/tests/test_dnssd_bridge.py | 35 +++++++++++++++++++ 2 files changed, 48 insertions(+), 7 deletions(-) diff --git a/bindings/python/offline_protocol_sdk/dnssd_bridge.py b/bindings/python/offline_protocol_sdk/dnssd_bridge.py index 3a793f0df..69cb704a2 100644 --- a/bindings/python/offline_protocol_sdk/dnssd_bridge.py +++ b/bindings/python/offline_protocol_sdk/dnssd_bridge.py @@ -577,19 +577,25 @@ def import_txt(self, name: str, raw: bytes, now: float | None = None) -> bool: return True async def sweep(self, now: float | None = None) -> None: - """Drop every listed entry older than the time to live, then - re-resolve every browsed name whose last resolve attempt is older - than half of it, listed or not. The periodic task calls this; a - test calls it with a clock of its own.""" + """Re-resolve every browsed name whose last resolve attempt is older + than half the time to live, listed or not, then drop every listed + entry still older than the whole of it. The periodic task calls + this; a test calls it with a clock of its own. + + Resolve first. The sweep runs at half the time to live, so one + missed window puts the next attempt on the drop boundary; dropping + first would show the application an empty list during that resolve + and announce the same record a second time when it answers, where + an answer refreshes the entry and shows nothing.""" current = time.monotonic() if now is None else now - for name, entry in list(self._lan.items()): - if current - entry.seen_at >= self._ttl: - self._lan.pop(name, None) for name, attempted_at in list(self._browsed.items()): if current - attempted_at < self._ttl / 2 or name not in self._browsed: continue self._browsed[name] = current await self._resolve_and_import(name, now=current) + for name, entry in list(self._lan.items()): + if current - entry.seen_at >= self._ttl: + self._lan.pop(name, None) async def _sweep_forever(self) -> None: while self._running: diff --git a/bindings/python/tests/test_dnssd_bridge.py b/bindings/python/tests/test_dnssd_bridge.py index 582a1ce12..945e181bb 100644 --- a/bindings/python/tests/test_dnssd_bridge.py +++ b/bindings/python/tests/test_dnssd_bridge.py @@ -721,6 +721,41 @@ async def test_a_name_that_never_resolved_is_re_resolved_until_it_does(self, bri await bridge.sweep(now=base + 102.0) assert len(bridge.lan_services()) == 1 and len(events) == 1 + async def test_a_resolve_that_answers_on_the_drop_boundary_refreshes_rather_than_re_announces( + self, services, events + ): + # The sweep runs at half the ttl, so one missed window puts the next + # attempt on the drop boundary. The sweep resolves before it drops: + # the application never sees an empty list during that resolve, and + # the answer is a refresh, not a second announcement. + class Observing(FakeBackend): + def __init__(self, bridge_ref: list[DnsSdBridge]) -> None: + super().__init__() + self.listed_during_resolve: list[int] = [] + self.bridge_ref = bridge_ref + + async def resolve(self, name: str, timeout_ms: int) -> bytes | None: + self.listed_during_resolve.append(len(self.bridge_ref[0].lan_services())) + return await super().resolve(name, timeout_ms) + + holder: list[DnsSdBridge] = [] + backend = Observing(holder) + bridge = DnsSdBridge(services, on_event=events.append, backend=backend, publish=False, ttl=100.0) + holder.append(bridge) + await bridge.start(address=OUR_ADDRESS) + try: + bridge.import_txt(PEER_NAME, peer_txt(), now=0.0) + assert len(events) == 1 + await bridge.sweep(now=50.0) # one missed window: no answer + assert len(bridge.lan_services()) == 1 + backend.txt_by_name[PEER_NAME] = peer_txt() # the neighbour answers again + await bridge.sweep(now=100.0) # the drop boundary + assert backend.listed_during_resolve == [1, 1], "listed throughout, never emptied first" + assert len(bridge.lan_services()) == 1 + assert len(events) == 1, "an answer on the boundary is a refresh, not a second announcement" + finally: + await bridge.stop() + async def test_a_resolve_that_lands_after_the_instance_is_gone_lists_nothing(self, services, events): class SlowBackend(FakeBackend): def __init__(self) -> None: