diff --git a/README.md b/README.md index e136ce2..69d2399 100644 --- a/README.md +++ b/README.md @@ -353,9 +353,10 @@ What a move cannot do is reach a machine that has already synced the host, which everything else about revocation has. Keys, passwords and buckets take theirs from a standing "new items go to" picker on the Keychain screen and cannot be moved yet. -**A group can be moved too, and it takes its contents with it** — MOVE beside EDIT and DELETE over the group -cards, or "Move to another vault…" on the card's own menu. That is the desktop only, because the phone draws -groups as headings in the host list and has never had a way to delete or move one. It is the same re-seal +**A group can be moved too, and it takes its contents with it** — "Move to another vault…" on the group +card's right-click menu, beside Open, Edit and Delete, which is the whole of what can be done to a group on +the desktop. That is the desktop only, because the phone draws groups as headings in the host list and has +never had a way to delete or move one. It is the same re-seal and tombstone underneath, applied to every item involved: the group, the groups nested inside it, and every host filed under any of them, each taking a new id in the destination. Moving less than that was never coherent — the machines and the child groups are items of the vault the group is leaving, so a group that @@ -374,6 +375,42 @@ Both default to your personal vault and neither moves on its own, because an ite visible to everybody holding that vault's key. Choosing a vault in the host editor also decides which groups it can be filed under: a group is an item like any other and lives in exactly one vault. +### Changes that do not wait + +A client holds a WebSocket open to the server — `GET /api/v1/events`, subprotocol +`dodossh.events.v1` — and the server sends a line down it whenever something you can read has moved. The +client's answer is the same delta pull it would have run on its timer, only now rather than in up to a +minute. Two things you can see: an edit somebody else makes appears while you are looking at the list, and +a vault shared with you turns up as soon as they share it. + +**What is on that socket is a notice, not your data.** A frame says which vault changed and how far its +change log has got, and nothing else: no item, no ciphertext, not even which item it was. That is the +decision the rest of this rests on, and it is deliberate twice over — the server has nothing else it +*could* send, and keeping it that way means there is still exactly one path that applies a change to your +keychain, so the socket can be wrong or absent without anything being applied incorrectly. + +**Polling is still there and is still what guarantees a pass.** The minute timer is unchanged. A network +that eats WebSockets, a server with `Events:Enabled` off, an older server, a proxy that will not upgrade, a +notice dropped because your machine was too slow to read it — every one of those leaves you with exactly +what this product did before the socket existed. Nothing is reachable only this way, and nothing is +supposed to become so. + +Three limits are worth knowing rather than discovering: + +- **One node.** Fan-out is in-process, so a deployment running more than one API replica only pushes for + writes that its own replica handled. The rest arrive on the timer. The seam for a PostgreSQL + `LISTEN`/`NOTIFY` backplane is in place and is not implemented, because an untested backplane would be + worse than a documented gap. +- **The socket does not outlive your access token.** It is closed at the token's expiry and the client + reconnects with a fresh one, which is a gap you will not see. That, plus re-reading your vault list every + few minutes, is what bounds how long a withdrawn grant can keep producing notices — and what it bounds is + *metadata*, because reading a vault needs a key the server has never held. +- **You are told about your own writes.** Your client pushed, so it has already pulled; the extra pass finds + nothing. Notices are coalesced over a quarter of a second so that a burst is one pass rather than a dozen. + +The reasoning, including why this is a WebSocket rather than server-sent events and where a shared terminal +session will attach to it, is in [ADR 0012](docs/adr/0012-realtime-push.md). + ### The Android head `src/DodoSSH.Client.Android` is a phone-first head that shares every view model with the desktop one — the @@ -652,6 +689,20 @@ keychain plus a terminal — and the spike that gates all of it. directories, an interrupted **upload** starts again rather than resuming (an object cannot be written from the middle), and a rename is a copy and a delete rather than one atomic operation. Downloads do resume — a ranged GET is part of the protocol, which is the one place a bucket beats SFTP. + + *Realtime done:* a WebSocket the client holds open, over which the server says which vault has moved so a + pull happens now rather than within the minute. What crosses it is a notice and never an item, which is + what keeps one code path applying changes and makes a dropped socket cost latency and nothing else — the + timer is unchanged and is still the guarantee. Two limits are stated rather than implied: fan-out is + in-process, so a multi-replica deployment falls back to the timer for writes another replica handled, and + a socket is closed at its access token's expiry rather than outliving the credential that authorised it. + See [Changes that do not wait](#changes-that-do-not-wait) and + [ADR 0012](docs/adr/0012-realtime-push.md). + + It is also the transport a **shared terminal session** will use — one person's shell, watched or driven by + somebody else. Nothing of that exists yet, and ADR 0012 records the one decision made early so it need not + be renegotiated: session data will be binary frames on this same socket, because base64 in a JSON envelope + is the wrong shape for the one payload here that is continuous rather than occasional. - **M3 — shared vaults**, sharing, ACLs. *Done.* Membership with roles, a public-key directory, the append-only key log served for clients to verify against, shared vaults, and vault key grants wrapped by a client and stored opaquely by the server. `VaultAccessService` now resolves team @@ -675,10 +726,12 @@ keychain plus a terminal — and the spike that gates all of it. the rotation is re-sealed as it is pushed, so nothing reaches the server under a superseded key at all. See [ADR 0010](docs/adr/0010-vault-key-rotation.md). - **A vault shared with you arrives on the next synchronisation pass**, within the minute, with no sign-in - and nothing to press. There is no push channel, so each pass asks the server which vaults this account can - reach before syncing the ones it already knows — which is also how a vault that has been deleted, or one - whose grant was withdrawn, stops being listed. + **A vault shared with you arrives at once**, with no sign-in and nothing to press. Each pass asks the + server which vaults this account can reach before syncing the ones it already knows — which is also how a + vault that has been deleted, or one whose grant was withdrawn, stops being listed — and the server now + says so the moment somebody wraps a key to you rather than leaving it for the next pass. Without a + reachable socket that becomes "within the minute", which is what it always was; see + [Changes that do not wait](#changes-that-do-not-wait). **Ownership transfer is here, and it is one write rather than two.** The member you name becomes owner and you become an admin, in a single transaction — because ownership is sole, so promoting first leaves diff --git a/docs/adr/0003-sync-protocol.md b/docs/adr/0003-sync-protocol.md index cdfbda5..05443d8 100644 --- a/docs/adr/0003-sync-protocol.md +++ b/docs/adr/0003-sync-protocol.md @@ -58,7 +58,11 @@ the Npgsql connection string** — the default; do not enable multiplexing. - One place enforces revision, change-log and ACL invariants. That halves both the endpoint count and the authorization surface, which is the main reason for the single write path. - Delta pull makes frequent polling cheap, so multi-device feels live; push notification over - SSE or the existing WebSocket can layer on with polling as the fallback. + SSE or the existing WebSocket can layer on with polling as the fallback. **That has since been + built — see [ADR 0012](0012-realtime-push.md)** — and nothing in this ADR changed to accommodate + it. The socket carries a notice naming a vault and a sequence, whose answer is the delta pull + above, so there is still exactly one path that applies a change; and polling is still what + guarantees a pass rather than a legacy route kept for old clients. - Conflict resolution is entirely client-side. The client retains a `BaseCiphertext` common ancestor and performs a field-level three-way merge for structured items, or creates a visible conflicted copy for opaque ones. **It must never silently drop a key or a host.** diff --git a/docs/adr/0012-realtime-push.md b/docs/adr/0012-realtime-push.md new file mode 100644 index 0000000..2bb7b61 --- /dev/null +++ b/docs/adr/0012-realtime-push.md @@ -0,0 +1,167 @@ +# ADR 0012 — A WebSocket that carries notices, not data + +- Status: accepted +- Date: 2026-08-04 + +## Context + +[ADR 0003](0003-sync-protocol.md) built a delta pull that is cheap enough to run on a timer, and +the client does: one pass a minute. That is the difference between a colleague's change appearing +"soon" and appearing *now*, and it shows up in three places that are not equally forgivable. + +- **A vault shared with you** arrives on the next pass. `AdmitNewVaultsAsync` says so in its own + remarks — "the recipient is handed nothing — there is no push channel" — and the README repeats + it. Sharing works and looks broken. +- **Two people editing one keychain** see each other up to a minute late, which is long enough to + make the same edit twice and produce a conflict that nobody needed to have. +- **A revoked grant** keeps serving a client that has not noticed yet, for up to a pass. + +Shortening the interval is the obvious answer and the wrong one: it costs a request per client per +interval whether or not anything happened, and it does not converge on *immediate* — it converges +on a busier server that is still late. + +There is also a second thing coming that this decision has to not preclude. The intended feature is +a **shared terminal session** — one person's shell, watched or driven by another, TeamViewer-shaped. +That is bidirectional, continuous, and latency-sensitive in a way a keychain notice is not. + +## Decision + +### One WebSocket per signed-in client, at `GET /api/v1/events` + +Subprotocol `dodossh.events.v1`. The client opens it after unlock and keeps it open; the server +sends a notice whenever something the client can read has changed. + +**Not SSE.** Server-sent events would carry today's notices perfectly well and would be less code. +It is one-directional, so the shared-session feature would need a second mechanism next to it, and +then two transports would need reconnection, authorization and lifetime rules that agree. The cost +of a WebSocket over SSE is small; the cost of two transports is not. + +**Not SignalR.** It brings hub protocol negotiation, its own serialisation and transport fallbacks, +none of which are wanted here: `DodoSSH.Contracts` and its source-generated serialiser are "the +actual contract between the two sides", and a second wire format alongside it is exactly the silent +drift `Setup/Json.cs` records having already cost this project once. + +### The notice carries no ciphertext + +A `vault.changed` frame is `{ kind, vaultId, sequence }` and nothing else. The client's answer to it +is the pull it would have done on the timer anyway. + +This is the load-bearing decision, and it is worth being explicit about why the tempting alternative +is refused. Pushing the changed items themselves would save a round trip and would fork the code +path that applies a change into two — one that arrives by pull and one that arrives by socket — with +the cursor, the merge and the tombstone rules duplicated across both. ADR 0003 put every mutation +through one write path for exactly that reason; this keeps every *read* on one path for the same +one. The socket decides *when* to sync. It never decides *what* a vault contains. + +It also means a dropped notice is harmless, which is what lets everything below be simple. + +### Polling stays, and is the fallback rather than a legacy path + +The one-minute pass is unchanged. The socket makes it *early*; it does not make it *necessary*. A +client on a network that eats WebSockets, an older client, a server that has the feature off, a +notice dropped under backpressure, a second API replica that did not see the write — every one of +those degrades to what the product does today, which is correct and up to a minute late. + +Nothing may be reachable only by socket. That is a rule about future features, not an observation +about this one. + +### Authorization: the bearer token on the upgrade, not a ticket + +[ADR 0004](0004-relay-authorization.md) gives the relay a two-step ticket so its WebSocket carries +no API authority. This one goes the other way and takes the ordinary bearer JWT on the upgrade +request, which is an ordinary authenticated HTTP request. The difference is not inconsistency: + +- The relay's socket is a **byte pipe to a third party**, and its whole authorization decision — + which host, which IPs, which port — is made *before* the socket opens and never revisited. It is + also the extraction seam for a standalone relay process that must not hold ACL code. +- This socket is a **view of the caller's own vault list**, and it has to keep answering "what may + this account read" for as long as it is open. It needs the full ACL context, in-process, for the + life of the connection. A ticket would carry that context in a token instead, and it would be + wrong the moment the account's access changed. + +A long-lived connection authorised by a short-lived token is the problem this creates, and it is met +head-on rather than ignored: + +1. **The socket does not outlive the token.** The `exp` claim is read at accept, and the connection + is closed with `4401` when it passes. The client reconnects with a fresh token; that is a + sub-second gap in a channel whose failure mode is already "poll instead". +2. **The vault set is re-resolved periodically** (`Events:AccessRefreshInterval`, default five + minutes) as well as on the changes that are known to affect it. A withdrawn grant therefore stops + producing notices within that window at the latest, and immediately in the ordinary case. + +Both are bounds on **metadata** — the fact that a vault changed and roughly when — because that is +all a notice contains. Nobody's ciphertext is behind this socket, and a client that stayed subscribed +one interval too long could still not read a byte of it: reading requires a vault key grant, which +this server has never held. + +### The frames + +Text frames, JSON, `DodoSshJsonContext`. Server to client: + +| kind | meaning | +| --- | --- | +| `hello` | accepted; carries the heartbeat interval and the vault count subscribed | +| `vault.changed` | `vaultId` moved to `sequence`; pull it | +| `vaults.changed` | the set of vaults this account can reach is different; re-read it | +| `ping` | heartbeat; the client answers `pong` | + +Client to server: `ping`, answered with `pong`. Nothing else — subscription is decided by the server +from the caller's access, not asked for by the client, because a client that could ask to subscribe +to a vault id is a client that can probe for vault ids. + +`kind` is a **string**, not an enum, and that is deliberate. `UseStringEnumConverter` throws on a +value it does not know, so a newer server sending a kind an older client has never heard of would +not add an unknown frame — it would break that client's socket entirely. A string is ignored +instead, which is what makes the table above extensible. `ProblemCodes` is the same shape for the +same reason. + +### Where the shared session will attach + +The socket is the seam, and one thing about it is chosen now so that it need not be renegotiated +later: **session data will be binary frames on this same connection, not JSON on the table above.** +Terminal output base64'd into a JSON envelope would cost a third of the bandwidth for nothing, on +the one payload here that is continuous rather than occasional. Control — offer, accept, resize, +end — is JSON like everything else. + +That is as far as this ADR goes. Two questions are open and are not being answered by implication: +whether a shared session's bytes go through the API at all or peer-to-peer past it, and what +end-to-end encryption means when the second party is watching a stream rather than holding a key. +Both are ADR 0001 questions and deserve their own decision. What this one buys is that they will not +also be transport questions. + +## Consequences + +- **Fan-out is in-process, and the deployment is therefore single-node for this feature.** Every + connection is held by the node that accepted it; a write handled by another node produces no + notice on this one. `IVaultEventPublisher` is the seam a backplane implements — PostgreSQL + `LISTEN`/`NOTIFY` needs no infrastructure this stack does not already run — and it is deliberately + **not implemented**, because an untested backplane is worse than a documented gap. Multiple API + replicas do not break: they degrade to polling, which is the state before this ADR. `/api/v1/meta` + advertises `events` so a client knows which it is getting. +- **Per-connection queues are bounded and drop the oldest.** A notice is "pull vault X, which is at + least at sequence N", so the newest is strictly more useful than the one it displaces and the + client's answer is identical either way. A slow reader costs itself latency, never the publisher's + progress — the publish path never blocks and never awaits a socket. +- **The publish happens after the transaction commits**, outside the advisory lock ADR 0003 takes. + A notice sent from inside it would name a sequence a reader cannot yet see, and would hold the + per-vault write lock across a socket write. +- **A client is notified of its own writes.** It pushed, so it already pulled; the extra pass finds + nothing. The client coalesces notices over a short window rather than the server suppressing an + echo, because suppressing it correctly needs a per-*device* identity on the socket and the same + user's other machines must still be told. +- **Connections are capped** per user and per node (`Events:MaxConnectionsPerUser`, + `Events:MaxConnectionsTotal`). A socket is cheap but not free, and an unbounded count of them is a + denial of service that authenticates first. +- The feature can be turned off entirely (`Events:Enabled`). A deployment behind a proxy that will + not upgrade should say so rather than have every client discover it by failing. + +### Rejected + +- **Shorter polling.** Cheaper to build, converges on a busier server that is still late. +- **Long polling.** No new transport and genuinely immediate, but it holds a request thread and a + connection per client for the same money as a WebSocket while offering none of the bidirectionality + the shared session needs. +- **Pushing the changed items down the socket.** Saves a round trip; forks the apply path in two. See + above. +- **Client-chosen subscriptions.** A `subscribe(vaultId)` frame is an existence oracle for vault ids, + which is the disclosure `SyncPullEndpoint` answers 404 rather than 403 to avoid. diff --git a/docs/design-import-gaps.md b/docs/design-import-gaps.md index 0292ce6..131b6c3 100644 --- a/docs/design-import-gaps.md +++ b/docs/design-import-gaps.md @@ -157,7 +157,7 @@ the chrome, hosts and terminals, file transfer, the vault, teams, and preference > | **Add Telnet**, and **Serial** in the toolbar | Omitted. `ISshConnection` is the only transport there is. This is also why the card subtitle's `ssh` is a constant today rather than a reading — it is stated in `HostRowViewModel.Summary`, which is the one place in this interface where a constant is printed on purpose. | > | **+ SSH ID, Certificate, FIDO2** | Omitted. `IDENTITIES` and `CERTIFICATES` have been on this document's list since the first import — neither is even a reserved `SyncEntityType` — and there is no security-key path anywhere in the SSH layer. One control offering three item types that do not exist. | > | The **Backspace / Default** row | Omitted. It is a terminal setting, and the client has no preferences store and no frame to carry one to the renderer — see the Preferences section. It would be a control whose value could not survive the window closing. | -> | The **chevron beside the vault name** | The name alone, and the move behind the pane's ⋯ menu instead. A host *can* now be moved between vaults, so the gap is no longer that there is nothing to offer — it is that a chevron on a subtitle implies an edit, and this is not one: the two vaults are encrypted under different keys, so it is a re-seal into one and a tombstone in the other, the host takes a new id, and its group and tags stay behind. A control that implied "just change this field" would be describing something else. Where a *new* host goes is still asked in the host editor, as a picker beside the name. A group moves too, from MOVE over the group cards, and takes its nested groups and every host filed under them; keys, passwords and buckets take theirs from the keychain screen's standing picker and cannot be moved yet. | +> | The **chevron beside the vault name** | The name alone, and the move behind the pane's ⋯ menu instead. A host *can* now be moved between vaults, so the gap is no longer that there is nothing to offer — it is that a chevron on a subtitle implies an edit, and this is not one: the two vaults are encrypted under different keys, so it is a re-seal into one and a tombstone in the other, the host takes a new id, and its group and tags stay behind. A control that implied "just change this field" would be describing something else. Where a *new* host goes is still asked in the host editor, as a picker beside the name. A group moves too, from its card's right-click menu, and takes its nested groups and every host filed under them; keys, passwords and buckets take theirs from the keychain screen's standing picker and cannot be moved yet. | > | **Show more ⌄** | Not drawn as a disclosure. What it would hide — notes, the relay switch, forgetting the host key — is in the editor, one press away, and a second fold inside a pane that already scrolls is a second place for a field to be missing from. | > | **Port Forwarding** in the sidebar | Nothing, for the third time in this document. | > | The host grid's toolbar avatar, share and tag-filter controls | Omitted, as in v3 and for the same reasons. | diff --git a/docs/manual-checks.md b/docs/manual-checks.md index 460c146..b9cef04 100644 --- a/docs/manual-checks.md +++ b/docs/manual-checks.md @@ -324,7 +324,8 @@ phone's, whose list has no room for a row of group cards and draws the whole tre host filed, the grid says so in a sentence rather than sitting empty. **Then press a group card once.** It is marked as chosen and **nothing else happens** — the grid is still the -level it was, and EDIT and DELETE now aim at that group. **Then double-press it.** The group opens: its hosts +level it was, and no buttons appear beside the GROUPS heading: editing and deleting a group are on the card's +own right-click menu, which is 7.9. **Then double-press it.** The group opens: its hosts are the grid, the trail above the cards reads `ALL HOSTS › ›`, each card carrying the group's name as an accent chip, and the card grid shows what is *inside* that group rather than every group in the keychain. Pressing ALL HOSTS goes back to the outermost level. @@ -347,9 +348,9 @@ Make two groups and file one under the other with the parent picker in the group **Pass:** only the outer group has a card to start with. Double-press it and the inner one is the only card shown, with the trail reading `ALL HOSTS › ›`. Double-press that, and the cards disappear entirely — -it has nothing inside it — while the trail, EDIT and DELETE stay: with no card selected the two buttons act -on the group the trail ends with, so a group with nothing in it can still be renamed after being opened. -Pressing the **middle** crumb goes back one level rather than all the way out. +it has nothing inside it — while the trail stays. Pressing the **middle** crumb goes back one level rather +than all the way out, which is also how a group with nothing inside it is renamed: back out to the level +where it has a card, and right-click that. **Failure means:** cards for groups that are not at this level is `VisibleGroups` having been bound past — the flat `Groups` is the phone's and the lookups'. A group that cannot be reached at all is worse and is the @@ -370,13 +371,35 @@ takes only the group that is open. ### 3.3 Deleting a group with hosts in it -Select a group with hosts and press DELETE. +Right-click a group with hosts in it and choose **Delete…**. Do it twice: once leaving the tick alone, and +once — on another group — ticking it. -**Pass:** the question names how many hosts are filed under it and says they stay. Agreeing removes the -group; the hosts lose their chip and are otherwise unchanged. +**Pass:** the question names how many hosts are filed under it, says they stay and move to UNGROUPED, and +offers a tick that would delete them as well. The tick starts clear, and it starts clear again on the next +group even if it was set on the last one. Left clear, agreeing removes the group and the hosts stay, without +a chip and otherwise unchanged. Ticked, the hosts go with it — and only the hosts that were filed under that +group. A group with nothing under it is asked no second question and shows no tick. -**Failure means:** if the hosts vanish, the delete is rewriting host payloads, which it must not — see -`HostGroupRepository`. +**Failure means:** a tick that carries from one question to the next is the reset in `OnPendingDeletionChanged` +having gone, and it deletes machines on the strength of a decision about a different group. Hosts that keep +the chip after an unticked delete are the unfiling not happening: they still name a group that is gone, which +is what this used to do on purpose and no longer should. + +### 3.3a Moving a group to another vault · **needs a second vault** + +Build `outer › inner` with a host in `inner`, all in your personal vault, then right-click **outer** and +choose **Move to another vault…**. Pick the shared vault and press MOVE. + +**Pass:** the panel says what travels and what does not before you press anything. Afterwards all three items +carry the destination's badge, `inner` is still inside `outer` and the host is still inside `inner` — every +one of them under an id it did not have a moment ago. The sentence names the vault, the counts, and the fact +that the group now sits at the top level if it was nested. Nothing is left behind in the vault it came from. + +**Failure means:** a host under UNGROUPED in the destination is the group id having been carried across +rather than remapped — the ids are the destination's making, so every reference has to be rewritten as its +target lands. Anything still in the source vault is a partial move, which is survivable by design but should +not happen with the network up: the groups are written top-down and the hosts last, so an interruption leaves +hosts behind and never a shelf with nothing on it. ### 3.4 A group deleted on another machine · **needs two machines** @@ -802,8 +825,9 @@ Delete. card, rather than back out to ALL HOSTS. Right-clicking the space around the group cards opens no menu. **Failure means:** the menu is reading `GroupTarget`'s fallback, which is the group whose contents are on -screen. That fallback is right for the EDIT and DELETE buttons beside the heading and wrong for a menu that -opened on a card. +screen rather than the card the pointer is on. This menu is the only way to edit or delete a group on the +desktop — there are no buttons beside the GROUPS heading any more — so a menu aimed wrongly is the whole of +the mistake. ### 7.10 Clicking a host in the palette connects @@ -1500,3 +1524,83 @@ delivery survives the failure because it is held against the transfer rather tha **Failure means:** a retry that succeeds but leaves the destination empty is the delivery having been dropped on the failure. An error saying the staged file is missing is the copy having been deleted at the stop, which is what `QueueDeliveredDownload` documents it does not do. + +## Phase 15 — Changes that arrive without a timer + +The socket is covered by tests on both sides: the endpoint suite opens a real one against a real +`TestServer` and proves a push produces a notice, that another account's push does not, and that a frame +carries no ciphertext; the shell suite proves a notice wakes the synchronisation loop long before the +minute. What none of that can reach is **the network in between**, and that is where this feature is most +likely to fail: a reverse proxy that will not upgrade, one that drops an idle socket without telling either +end, a corporate middlebox, a phone moving between Wi-Fi and mobile data. Every one of those looks the same +from inside a test host, which has no proxy and no radio. + +The pass condition throughout is *two* things, and the second matters as much as the first: it arrives +quickly, **and** it still arrives when the socket is gone. A build where the timer had stopped working would +pass every "it was fast" check here and fail nobody until somebody's proxy changed. + +### 15.1 A colleague's edit appears while you are looking at it + +Two accounts sharing a vault, both unlocked, both on the Hosts screen. On the first machine, rename a host +in the shared vault and save. + +**Pass:** the second machine's list shows the new name within a second or two, with nothing pressed and no +screen flicker — the row updates, the selection does not move, and the status line is not repainted with a +sync report. + +**Failure means:** nothing within a minute, then the new name, is the socket not being established at all — +that is the timer doing its job, which is the correct fallback and not the feature. Check `/api/v1/meta` +lists `events`, then whether the proxy in front of the API forwards `Upgrade` and `Connection`. A list that +never updates at all is a synchronisation failure and has nothing to do with this phase. + +### 15.2 A vault shared with you turns up as it is shared + +The second account signed in and unlocked, sitting on the VAULTS screen. From the first, add them to a team +and press SHARE KEY. + +**Pass:** the vault appears in their list within a second or two of the key being wrapped, and reads as +waiting for a key until the share, then as readable. + +**Failure means:** the vault appearing only on the minute is the `vaults.changed` notice not being published +or not being followed. Both the membership add and the grant publish one; if the membership arrives promptly +and the key does not, the grant path is the one to look at. + +### 15.3 It still works with the socket taken away + +On the second machine, block the WebSocket — the simplest way is a proxy rule rejecting the upgrade, or +setting `Events:Enabled` to `false` on the server and restarting it. + +**Pass:** everything above still happens, within the minute rather than within seconds. Nothing on the +screen says anything is wrong, because nothing is: no error, no OFFLINE badge, no repeated status message. +The Sync button still works and still reports. + +**Failure means:** an error message, a titlebar claiming to be offline, or a status line that repaints with +a socket failure is the client treating an absent push channel as a fault. It is not one — the timer is the +guarantee and the socket is the optimisation, and a user with a strict proxy must never be told their +keychain is broken. + +### 15.4 A laptop that slept comes back on its own + +With the second machine idle and connected, close the lid for a few minutes — or disable Wi-Fi for two +minutes and re-enable it. Then make a change on the first machine. + +**Pass:** the change arrives quickly again, without the vault having been locked or the application +restarted. The reconnection is invisible. + +**Failure means:** changes that arrive only on the timer from then on are the stream having given up after +its socket died — the reconnection loop is what should make that impossible, and a client that reconnects +once and not twice is the specific defect its tests exist to catch. Changes that never arrive again, timer +included, are a different and worse bug in the synchronisation loop rather than in the socket. + +### 15.5 An expiring token does not end the push + +This one needs a short access-token lifetime in the identity provider — the dev realm's Keycloak client can +be set to a couple of minutes. Leave a machine unlocked and idle for longer than that, then make a change +elsewhere. + +**Pass:** the change still arrives quickly. The socket is closed by the server at the token's expiry and the +client reconnects with a fresh one, which should be invisible. + +**Failure means:** notices stopping at roughly the token's lifetime is the reconnection not asking for a new +token — it would be dialling with the spent one and being closed again immediately. A burst of reconnection +attempts in the server log is the same defect seen from the other end. diff --git a/src/DodoSSH.Api/Features/Events/EventsEndpoint.cs b/src/DodoSSH.Api/Features/Events/EventsEndpoint.cs new file mode 100644 index 0000000..46246b2 --- /dev/null +++ b/src/DodoSSH.Api/Features/Events/EventsEndpoint.cs @@ -0,0 +1,533 @@ +using System.Collections.Frozen; +using System.Globalization; +using System.Net.WebSockets; +using System.Security.Claims; +using System.Text.Json; +using DodoSSH.Api.Authorization; +using DodoSSH.Api.Setup; +using DodoSSH.Contracts; +using DodoSSH.Domain.Authorization; +using FastEndpoints; +using Microsoft.Extensions.Options; + +namespace DodoSSH.Api.Features.Events; + +/// +/// The socket that says "pull now" so a client does not have to wait for its timer. +/// +/// +/// +/// Everything this endpoint sends is a notice. It never carries an item, a payload or a +/// cursor: the client's answer to a notice is the delta pull it would have run on its own anyway, so +/// there is exactly one code path that applies a change and this is not it. See ADR 0012 for why +/// pushing the items themselves is refused. +/// +/// +/// The bearer token authorises the upgrade, unlike the relay's ticket in ADR 0004. The relay's socket +/// is a byte pipe whose whole authorization decision is made before it opens; this one is a view of +/// the caller's own vault list and has to keep answering "what may this account read" for as long as +/// it is held. Its two bounds on that — the token's own expiry, and a periodic re-resolve — are in +/// . +/// +/// +internal sealed class VaultEventsEndpoint( + IServiceScopeFactory scopes, + VaultEventHub hub, + IOptions options, + TimeProvider clock, + IHostApplicationLifetime lifetime, + ILogger logger) + : EndpointWithoutRequest +{ + /// + /// The largest message this endpoint will read from a client. + /// + /// + /// A client sends nothing but ping, so the cap is three orders of magnitude of headroom and + /// still small enough that a hostile client cannot make the server buffer anything worth having. + /// + private const int MaxInboundFrameBytes = 4 * 1024; + + /// How long to wait for the close handshake before dropping the socket. + private static readonly TimeSpan CloseTimeout = TimeSpan.FromSeconds(5); + + /// + public override void Configure() + { + Get(VaultEvents.Path); + + // Enrolled, matching sync. A caller with no identity key holds no vault key either, so every + // notice this socket could send is about ciphertext they cannot read. + Policies(Auth.EnrolledPolicy); + + Description(b => b + .WithName("VaultEvents") + .WithSummary("Pushes a notice when a vault the caller can read has changed.") + .WithTags("Events")); + } + + /// + public override async Task HandleAsync(CancellationToken ct) + { + if (await RefusedAsync().ConfigureAwait(false)) + { + return; + } + + var (userId, vaults) = await ResolveAccessAsync(ct).ConfigureAwait(false); + + // Admitted before the upgrade so a refusal costs nothing, but answered *through* the socket + // rather than as an HTTP status: a constrained WebSocket client cannot read the status of a + // failed upgrade, and "you have too many open" is precisely the case where the client needs + // to know to back off rather than retry. Same reasoning as ADR 0004 on request headers. + var connection = hub.TryAdmit(userId, vaults); + + try + { + using var socket = await HttpContext.WebSockets + .AcceptWebSocketAsync(new WebSocketAcceptContext { SubProtocol = VaultEvents.SubProtocol }) + .ConfigureAwait(false); + + if (connection is null) + { + await CloseAsync( + socket, + new Closure( + VaultEvents.TooManyConnectionsCloseCode, "Too many open event sockets.")) + .ConfigureAwait(false); + + return; + } + + await PumpAsync(socket, connection, ct).ConfigureAwait(false); + } + finally + { + // Inside a try that starts *before* the upgrade, because an accept that throws — a client + // that abandoned the handshake — would otherwise leave an admitted connection in the hub + // for the life of the process, counting against this account's cap and taking a slot from + // the sockets that did open. + if (connection is not null) + { + hub.Remove(connection); + } + } + } + + /// + /// Answers the requests that are not an event socket at all, as ordinary HTTP. + /// + /// Whether a response was sent and the handler should stop. + /// + /// All three answers are problem documents rather than bare statuses, because each one has a + /// different remedy and a client that cannot tell them apart would retry the two that will never + /// succeed. Answered before the upgrade, so a caller that got the handshake wrong reads why in a + /// body rather than inferring it from a socket that closed. + /// + private async Task RefusedAsync() + { + if (!options.Value.Enabled) + { + // 404 rather than 501: the feature is absent from this deployment, and /api/v1/meta does + // not advertise it. A client that dialled anyway keeps polling, which is correct. + await Send.ResultAsync(Problems.Coded( + StatusCodes.Status404NotFound, + ProblemCodes.EventsUnavailable, + "This server does not push vault changes. Synchronise on a timer instead; " + + "GET /api/v1/meta lists the features it does offer.")) + .ConfigureAwait(false); + + return true; + } + + if (!HttpContext.WebSockets.IsWebSocketRequest) + { + await Send.ResultAsync(Problems.Coded( + StatusCodes.Status400BadRequest, + ProblemCodes.MalformedRequest, + "This endpoint is a WebSocket. Send an upgrade request offering the " + + $"'{VaultEvents.SubProtocol}' subprotocol.")) + .ConfigureAwait(false); + + return true; + } + + // The subprotocol is this API's version negotiation for the socket, so an upgrade that does + // not offer it is refused rather than accepted and answered in a dialect the caller may not + // read. See VaultEvents.SubProtocol. + if (!HttpContext.WebSockets.WebSocketRequestedProtocols + .Contains(VaultEvents.SubProtocol, StringComparer.Ordinal)) + { + await Send.ResultAsync(Problems.Coded( + StatusCodes.Status400BadRequest, + ProblemCodes.MalformedRequest, + $"This server speaks '{VaultEvents.SubProtocol}', which the upgrade request did " + + "not offer.")) + .ConfigureAwait(false); + + return true; + } + + return false; + } + + /// + /// Reads who the caller is and which vaults they may follow, in a scope of its own. + /// + /// + /// A fresh scope, disposed at once, rather than services injected into this endpoint — which is + /// the lesson ADR 0004 records paying for on the relay. This handler runs for as long as the + /// socket is open, so anything scoped it held would be a DbContext alive for hours, and a + /// few hundred of those exhaust the connection pool. The database is touched here and in + /// , briefly, and nowhere else. + /// + private async Task<(Guid UserId, FrozenSet Vaults)> ResolveAccessAsync( + CancellationToken cancellationToken) + { + var scope = scopes.CreateAsyncScope(); + await using var _ = scope.ConfigureAwait(false); + + var currentUser = scope.ServiceProvider.GetRequiredService(); + var vaultAccess = scope.ServiceProvider.GetRequiredService(); + + var user = await currentUser.GetOrProvisionAsync(cancellationToken).ConfigureAwait(false); + + return (user.Id, await ReachableAsync(vaultAccess, user.Id, cancellationToken).ConfigureAwait(false)); + } + + private static async Task> ReachableAsync( + IVaultAccessService vaultAccess, + Guid userId, + CancellationToken cancellationToken) + { + var accessible = await vaultAccess.ListAsync(userId, cancellationToken).ConfigureAwait(false); + + return accessible + .Where(access => access.Vault is not null + && access.Permissions.HasFlag(PermissionFlags.Read)) + .Select(access => access.Vault!.Id) + .ToFrozenSet(); + } + + /// Runs the socket until something ends it, then closes it politely. + /// + /// Three loops rather than one: reading a socket and writing to it are independent waits, and the + /// clock is a third. They are joined by and then all + /// awaited before the close is written, because a close frame racing a notice frame is a protocol + /// violation that presents as a client dropping its connection for no visible reason. + /// + private async Task PumpAsync( + WebSocket socket, + VaultEventConnection connection, + CancellationToken requestAborted) + { + var settings = options.Value; + + using var pump = CancellationTokenSource.CreateLinkedTokenSource( + requestAborted, lifetime.ApplicationStopping); + + var closure = new Closure( + (int)WebSocketCloseStatus.NormalClosure, string.Empty); + + connection.TryEnqueue(new VaultEvent( + VaultEventKinds.Hello, + ServerTime: clock.GetUtcNow(), + HeartbeatSeconds: (int)settings.HeartbeatInterval.TotalSeconds, + VaultCount: connection.VaultCount)); + + var sending = SendAsync(socket, connection, pump.Token); + var receiving = ReceiveAsync(socket, connection, pump.Token); + var minding = MindAsync(connection, closure, pump.Token); + + await Task.WhenAny(sending, receiving, minding).ConfigureAwait(false); + + await pump.CancelAsync().ConfigureAwait(false); + + // Nothing may still be mid-send when the close frame goes out. + await Task.WhenAll(Settled(sending), Settled(receiving), Settled(minding)).ConfigureAwait(false); + + if (lifetime.ApplicationStopping.IsCancellationRequested) + { + // 1001 "going away", so a client knows to reconnect immediately rather than treating a + // rolling deployment as a server that has broken. + closure.Set((int)WebSocketCloseStatus.EndpointUnavailable, "The server is shutting down."); + } + + EventsLog.ClosingConnection(logger, connection.UserId, closure.Reason); + + await CloseAsync(socket, closure).ConfigureAwait(false); + } + + /// + /// Writes queued frames to the socket, one at a time. + /// + /// + /// The only writer, which is what makes concurrent sends impossible without a lock: the + /// heartbeat and the pong both go into the same queue rather than to the socket. A WebSocket + /// permits one send at a time and faults permanently on a second, so this is not a tidiness + /// preference. + /// + private async Task SendAsync( + WebSocket socket, + VaultEventConnection connection, + CancellationToken cancellationToken) + { + await foreach (var frame in connection.Outbound + .ReadAllAsync(cancellationToken) + .ConfigureAwait(false)) + { + // Re-resolved before the notice is forwarded, not after: the client's answer to this frame + // is to re-read its vault list, and the point of a newly shared vault is that the *next* + // change to it produces a notice too. A socket that forwarded first would not follow the + // new vault until its next periodic refresh. + if (string.Equals(frame.Kind, VaultEventKinds.VaultsChanged, StringComparison.Ordinal)) + { + await RefreshAsync(connection, cancellationToken).ConfigureAwait(false); + } + + var bytes = JsonSerializer.SerializeToUtf8Bytes(frame, DodoSshJsonContext.Default.VaultEvent); + + await socket + .SendAsync(bytes, WebSocketMessageType.Text, endOfMessage: true, cancellationToken) + .ConfigureAwait(false); + } + } + + /// + /// Reads what the client sends, which in this version is heartbeats and a close. + /// + /// + /// A client cannot ask to follow a vault, and that is deliberate rather than unfinished: a + /// subscribe(vaultId) frame is an existence oracle for vault ids, which is the disclosure + /// SyncPullEndpoint answers 404 rather than 403 to avoid. Subscription is decided from the + /// caller's access and nothing else. + /// + private static async Task ReceiveAsync( + WebSocket socket, + VaultEventConnection connection, + CancellationToken cancellationToken) + { + var buffer = new byte[MaxInboundFrameBytes]; + + while (!cancellationToken.IsCancellationRequested) + { + var received = await socket.ReceiveAsync(buffer, cancellationToken).ConfigureAwait(false); + + if (received.MessageType == WebSocketMessageType.Close) + { + return; + } + + // Oversized, or split across frames. Nothing this protocol sends is either, so the client + // is broken or probing; ending the socket is cheaper than reassembling for it. + if (!received.EndOfMessage) + { + return; + } + + // Binary is unused in v1 and skipped rather than refused, because ADR 0012 reserves it for + // shared-session data — an older server meeting a newer client must ignore those, not + // close on them. + if (received.MessageType != WebSocketMessageType.Text) + { + continue; + } + + if (Kind(buffer.AsSpan(0, received.Count)) is VaultEventKinds.Ping) + { + connection.TryEnqueue(new VaultEvent(VaultEventKinds.Pong)); + } + } + } + + /// + /// Reads a frame's kind, or null if it is not one this server understands. + /// + /// + /// A frame that will not parse is skipped rather than closing the socket. This is a control + /// channel whose failure mode is "the client polls instead", so tolerating a frame from a newer + /// client costs nothing and refusing one costs that client its push for the whole session. + /// + private static string? Kind(ReadOnlySpan utf8) + { + try + { + return JsonSerializer.Deserialize(utf8, DodoSshJsonContext.Default.VaultEvent)?.Kind; + } + catch (JsonException) + { + return null; + } + } + + /// + /// Keeps the heartbeat going, the vault set current, and the socket inside its token's lifetime. + /// + /// + /// + /// The deadline is the earlier of the access token's exp and a hard cap on how long any one + /// socket may live. Closing on expiry is what keeps a long-lived connection from outliving the + /// short-lived credential that authorised it; the client answers by reconnecting with a fresh + /// token, which is a sub-second gap in a channel that degrades to polling anyway. + /// + /// + /// The wait is the shorter of the heartbeat and the time left, so the deadline is met to within a + /// tick rather than to within a heartbeat. + /// + /// + private async Task MindAsync( + VaultEventConnection connection, + Closure closure, + CancellationToken cancellationToken) + { + var settings = options.Value; + var started = clock.GetUtcNow(); + + var deadline = TokenExpiry() is { } expiry && expiry < started + settings.MaxConnectionDuration + ? (Expiry: expiry, ForToken: true) + : (Expiry: started + settings.MaxConnectionDuration, ForToken: false); + + var refreshed = started; + + while (!cancellationToken.IsCancellationRequested) + { + var now = clock.GetUtcNow(); + var remaining = deadline.Expiry - now; + + if (remaining <= TimeSpan.Zero) + { + closure.Set( + deadline.ForToken + ? VaultEvents.TokenExpiredCloseCode + : (int)WebSocketCloseStatus.NormalClosure, + deadline.ForToken + ? "The access token has expired. Reconnect with a fresh one." + : "This connection reached its maximum lifetime."); + + return; + } + + var wait = settings.HeartbeatInterval < remaining ? settings.HeartbeatInterval : remaining; + + await Task.Delay(wait, clock, cancellationToken).ConfigureAwait(false); + + now = clock.GetUtcNow(); + + if (now - refreshed >= settings.AccessRefreshInterval) + { + // The backstop for a grant withdrawn while this socket was open. What it bounds is + // metadata — that a vault changed — because that is all a notice carries and reading + // the vault still needs a key this server has never held. See ADR 0012. + await RefreshAsync(connection, cancellationToken).ConfigureAwait(false); + refreshed = now; + } + + connection.TryEnqueue(new VaultEvent(VaultEventKinds.Ping, ServerTime: now)); + } + } + + /// Re-reads which vaults this socket may follow. + /// + /// A failure is logged and swallowed. The alternative is dropping a working socket because one + /// database call timed out, which would trade an occasionally stale vault set for an outage. + /// + private async Task RefreshAsync(VaultEventConnection connection, CancellationToken cancellationToken) + { + try + { + var scope = scopes.CreateAsyncScope(); + await using var _ = scope.ConfigureAwait(false); + + var vaultAccess = scope.ServiceProvider.GetRequiredService(); + + connection.Resubscribe( + await ReachableAsync(vaultAccess, connection.UserId, cancellationToken) + .ConfigureAwait(false)); + } + catch (OperationCanceledException) + { + // The socket is closing. + } + catch (Exception exception) when (exception is not OutOfMemoryException) + { + EventsLog.AccessRefreshFailed(logger, connection.UserId, exception); + } + } + + /// + /// MapInboundClaims is off — see — so the claim is spelled as the + /// provider issued it rather than as a WS-Federation URI. Null is treated as "no bound from the + /// token", which the bearer handler's RequireExpirationTime should make unreachable; the + /// lifetime cap covers it either way. + /// + private DateTimeOffset? TokenExpiry() => + long.TryParse( + HttpContext.User.FindFirstValue("exp"), + NumberStyles.Integer, + CultureInfo.InvariantCulture, + out var seconds) + ? DateTimeOffset.FromUnixTimeSeconds(seconds) + : null; + + private static async Task CloseAsync(WebSocket socket, Closure closure) + { + if (socket.State is not (WebSocketState.Open or WebSocketState.CloseReceived)) + { + return; + } + + using var timeout = new CancellationTokenSource(CloseTimeout); + + try + { + await socket + .CloseOutputAsync((WebSocketCloseStatus)closure.Code, closure.Reason, timeout.Token) + .ConfigureAwait(false); + } + catch (Exception exception) + when (exception is OperationCanceledException or WebSocketException or ObjectDisposedException) + { + // The peer is already gone. There is nothing to tell it and nothing to recover. + } + } + + /// + /// Awaits a pump loop, treating its cancellation and its socket faults as the ordinary end. + /// + /// + /// Every one of these loops ends by being cancelled or by the socket going away, so an exception + /// here is the expected shape of "this connection is over" rather than a fault to propagate — and + /// propagating it would skip the close frame the other side is waiting for. + /// + private static async Task Settled(Task loop) + { + try + { + await loop.ConfigureAwait(false); + } + catch (Exception exception) + when (exception is OperationCanceledException or WebSocketException or ObjectDisposedException) + { + // Expected. + } + } + + /// Why the socket is being closed, decided by whichever loop ended first. + /// + /// Mutable and shared, and safe without a lock for one specific reason: it is written by the pump + /// loops and read only after over all of them, which is a + /// memory barrier. Writes race only with each other, and any of them is a true answer. + /// + private sealed class Closure(int code, string reason) + { + internal int Code { get; private set; } = code; + + internal string Reason { get; private set; } = reason; + + internal void Set(int code, string reason) + { + Code = code; + Reason = reason; + } + } +} diff --git a/src/DodoSSH.Api/Features/Events/EventsLog.cs b/src/DodoSSH.Api/Features/Events/EventsLog.cs new file mode 100644 index 0000000..b912b04 --- /dev/null +++ b/src/DodoSSH.Api/Features/Events/EventsLog.cs @@ -0,0 +1,68 @@ +namespace DodoSSH.Api.Features.Events; + +/// Source-generated log messages for the event socket. +/// +/// Ids and counts only, as everywhere else. A notice carries no ciphertext to leak, but which vault +/// changed and when is still the metadata ADR 0001 asks be kept to what is diagnostically useful. +/// +internal static partial class EventsLog +{ + [LoggerMessage( + EventId = 2201, + Level = LogLevel.Debug, + Message = "Event socket opened for user {UserId} following {VaultCount} vault(s); " + + "{ConnectionCount} open on this node.")] + internal static partial void ConnectionOpened( + ILogger logger, + Guid userId, + int vaultCount, + int connectionCount); + + [LoggerMessage( + EventId = 2202, + Level = LogLevel.Debug, + Message = "Event socket closed for user {UserId}; {ConnectionCount} open on this node.")] + internal static partial void ConnectionClosed(ILogger logger, Guid userId, int connectionCount); + + /// + /// Information rather than Debug: a refused socket is a client that will poll for the rest of its + /// session, and an operator seeing these has a cap to raise. + /// + [LoggerMessage( + EventId = 2203, + Level = LogLevel.Information, + Message = "Refused an event socket for user {UserId}: {Limit} is already reached.")] + internal static partial void ConnectionRefused(ILogger logger, Guid userId, string limit); + + [LoggerMessage( + EventId = 2204, + Level = LogLevel.Debug, + Message = "Announced vault {VaultId} at sequence {Sequence} to {ConnectionCount} socket(s).")] + internal static partial void VaultChangePublished( + ILogger logger, + Guid vaultId, + long sequence, + int connectionCount); + + [LoggerMessage( + EventId = 2205, + Level = LogLevel.Debug, + Message = "Announced a vault access change to {ConnectionCount} socket(s) of user {UserId}.")] + internal static partial void AccessChangePublished(ILogger logger, Guid userId, int connectionCount); + + /// + /// Warning, and it is worth being loud: the socket is still open and still delivering, but it is + /// delivering about a vault set that may be stale. Everything else here is routine. + /// + [LoggerMessage( + EventId = 2206, + Level = LogLevel.Warning, + Message = "Could not re-resolve which vaults user {UserId}'s event socket may follow.")] + internal static partial void AccessRefreshFailed(ILogger logger, Guid userId, Exception exception); + + [LoggerMessage( + EventId = 2207, + Level = LogLevel.Debug, + Message = "Closing user {UserId}'s event socket: {Reason}.")] + internal static partial void ClosingConnection(ILogger logger, Guid userId, string reason); +} diff --git a/src/DodoSSH.Api/Features/Events/VaultEventHub.cs b/src/DodoSSH.Api/Features/Events/VaultEventHub.cs new file mode 100644 index 0000000..dbd918e --- /dev/null +++ b/src/DodoSSH.Api/Features/Events/VaultEventHub.cs @@ -0,0 +1,242 @@ +using System.Collections.Concurrent; +using System.Collections.Frozen; +using System.Threading.Channels; +using DodoSSH.Api.Setup; +using DodoSSH.Contracts; +using Microsoft.Extensions.Options; + +namespace DodoSSH.Api.Features.Events; + +/// +/// Tells connected clients that something they can read has moved. +/// +/// +/// +/// Every method is void and returns having queued, never having sent. That is the contract, not +/// an implementation detail: the callers are write paths that have just committed a transaction, and a +/// publish that could block on a slow socket would make one client's bad network everybody else's +/// latency. A notice that cannot be queued is dropped, which is safe because the client polls anyway. +/// See ADR 0012. +/// +/// +/// An interface because this is the seam a multi-node backplane implements — PostgreSQL +/// LISTEN/NOTIFY is the obvious one and needs no infrastructure this stack does not +/// already run. It is deliberately not implemented: fan-out today is in-process, so a deployment with +/// more than one API replica notices writes handled by other replicas on the polling interval rather +/// than at once. That is the behaviour before this feature existed, which is why it degrades rather +/// than breaks. +/// +/// +public interface IVaultEventPublisher +{ + /// Announces that a vault's change log has reached . + /// + /// Call after the transaction commits, and outside the per-vault advisory lock ADR 0003 + /// takes. A notice sent from inside names a sequence no reader can see yet, and holds the vault's + /// write lock across a socket write. + /// + void VaultChanged(Guid vaultId, long sequence); + + /// Announces that the set of vaults an account can reach is no longer what it was. + /// + /// Takes the recipient, not the actor. Sharing is something one account does to another's + /// list, and it is the other account that has to re-read. + /// + void VaultAccessChanged(Guid userId); +} + +/// +/// Every event socket this node is holding. +/// +/// +/// +/// Publishing walks the whole connection list and asks each one whether it cares, rather than keeping +/// an index from vault to subscribers. With a per-node connection cap in the hundreds and an event rate +/// bounded by how often people edit keychains, the walk is not measurable — and the index is not free: +/// a connection's vault set is re-resolved while it is live, so every re-subscription would have to +/// move it between buckets under a lock that publishing also takes. The simpler shape is the one whose +/// races are obvious. +/// +/// +/// A singleton, holding no scoped service and no database context. Connections outlive requests by +/// design and anything request-scoped they captured would outlive its scope with them. +/// +/// +internal sealed class VaultEventHub( + IOptions options, + TimeProvider clock, + ILogger logger) : IVaultEventPublisher +{ + private readonly ConcurrentDictionary connections = new(); + + /// + /// Serialises admission so the caps are caps rather than approximations. + /// + /// + /// Counting and inserting under one lock, because the two done separately let N simultaneous + /// connects all read the same under-cap count and all insert. Contended only by connects, which + /// happen once per client per session; publishing never takes it. + /// + private readonly Lock admission = new(); + + /// How many sockets this node is holding. Diagnostics and tests. + internal int Count => connections.Count; + + /// + /// Admits a socket, or refuses it because a cap is already met. + /// + /// The connection, or null when a limit refused it. + internal VaultEventConnection? TryAdmit(Guid userId, FrozenSet vaults) + { + var limits = options.Value; + + lock (admission) + { + if (connections.Count >= limits.MaxConnectionsTotal) + { + EventsLog.ConnectionRefused(logger, userId, "the node limit"); + return null; + } + + var held = 0; + foreach (var existing in connections.Values) + { + if (existing.UserId == userId && ++held >= limits.MaxConnectionsPerUser) + { + EventsLog.ConnectionRefused(logger, userId, "the per-account limit"); + return null; + } + } + + var connection = new VaultEventConnection(userId, vaults, limits.OutboundQueueDepth); + + // Cannot collide: the id is fresh and this is the only insert. + connections[connection.Id] = connection; + + EventsLog.ConnectionOpened(logger, userId, vaults.Count, connections.Count); + + return connection; + } + } + + /// Forgets a socket that has closed. + internal void Remove(VaultEventConnection connection) + { + ArgumentNullException.ThrowIfNull(connection); + + connections.TryRemove(connection.Id, out _); + connection.Complete(); + + EventsLog.ConnectionClosed(logger, connection.UserId, connections.Count); + } + + /// + public void VaultChanged(Guid vaultId, long sequence) + { + var notice = new VaultEvent( + VaultEventKinds.VaultChanged, + VaultId: vaultId, + Sequence: sequence, + ServerTime: clock.GetUtcNow()); + + var delivered = 0; + + foreach (var connection in connections.Values) + { + if (connection.IsSubscribedTo(vaultId) && connection.TryEnqueue(notice)) + { + delivered++; + } + } + + if (delivered > 0) + { + EventsLog.VaultChangePublished(logger, vaultId, sequence, delivered); + } + } + + /// + public void VaultAccessChanged(Guid userId) + { + var notice = new VaultEvent(VaultEventKinds.VaultsChanged, ServerTime: clock.GetUtcNow()); + + var delivered = 0; + + foreach (var connection in connections.Values) + { + if (connection.UserId == userId && connection.TryEnqueue(notice)) + { + delivered++; + } + } + + if (delivered > 0) + { + EventsLog.AccessChangePublished(logger, userId, delivered); + } + } +} + +/// +/// One open socket, as the hub sees it. +/// +/// +/// Deliberately knows nothing about WebSockets. The hub queues frames here and the endpoint's pump +/// takes them away, which is what keeps a publish from ever touching a socket — and what lets the +/// whole fan-out be tested without one. +/// +internal sealed class VaultEventConnection +{ + private readonly Channel outbound; + + private FrozenSet vaults; + + internal VaultEventConnection(Guid userId, FrozenSet vaults, int queueDepth) + { + Id = Guid.CreateVersion7(); + UserId = userId; + this.vaults = vaults; + + // DropOldest, and the choice is what makes a slow reader harmless. A notice says "vault X has + // moved to at least sequence N", so a newer one subsumes the one it displaces and the client's + // answer — pull that vault — is identical either way. The writer therefore never waits and + // TryWrite never fails, which is what lets the publish path be non-blocking and void. + outbound = Channel.CreateBounded(new BoundedChannelOptions(queueDepth) + { + FullMode = BoundedChannelFullMode.DropOldest, + SingleReader = true, + SingleWriter = false, + }); + } + + /// Identifies this connection within the hub. Never sent to a client. + internal Guid Id { get; } + + /// The account that opened it. + internal Guid UserId { get; } + + /// Frames waiting to be written to the socket. + internal ChannelReader Outbound => outbound.Reader; + + /// How many vaults this socket currently follows. + internal int VaultCount => Volatile.Read(ref vaults).Count; + + /// Whether a change to this vault concerns this socket. + internal bool IsSubscribedTo(Guid vaultId) => Volatile.Read(ref vaults).Contains(vaultId); + + /// + /// Replaces what this socket follows, after its account's access was re-resolved. + /// + /// + /// A whole-set swap of an immutable set rather than a mutation, so a publish walking the list + /// concurrently reads either the old set or the new one and never a half-built one. No lock: the + /// only writer is this connection's own pump. + /// + internal void Resubscribe(FrozenSet replacement) => Volatile.Write(ref vaults, replacement); + + /// Queues a frame. Never blocks, and never fails — see the channel's full mode. + internal bool TryEnqueue(VaultEvent frame) => outbound.Writer.TryWrite(frame); + + /// Signals that nothing more will be queued, which ends the pump's drain loop. + internal void Complete() => outbound.Writer.TryComplete(); +} diff --git a/src/DodoSSH.Api/Features/Meta/MetaEndpoints.cs b/src/DodoSSH.Api/Features/Meta/MetaEndpoints.cs index 3533372..e9a3b7a 100644 --- a/src/DodoSSH.Api/Features/Meta/MetaEndpoints.cs +++ b/src/DodoSSH.Api/Features/Meta/MetaEndpoints.cs @@ -19,6 +19,7 @@ namespace DodoSSH.Api.Features.Meta; internal sealed class GetMetaEndpoint( IOptions sync, IOptions relay, + IOptions events, IOptions server) : EndpointWithoutRequest> { @@ -52,6 +53,14 @@ internal sealed class GetMetaEndpoint( features.Add(RelayFeature); } + // Advertised so a client knows whether to hold a socket open or rely on its timer. Absence is + // not an error — synchronising on a timer is the supported behaviour and the socket only makes + // it early — which is why this is a feature flag rather than a version bump. See ADR 0012. + if (events.Value.Enabled) + { + features.Add(VaultEvents.Feature); + } + return Task.FromResult(TypedResults.Ok(new MetaResponse( ServerVersion: ServerVersion, ApiVersions: [1], diff --git a/src/DodoSSH.Api/Features/Sync/SyncEndpoints.cs b/src/DodoSSH.Api/Features/Sync/SyncEndpoints.cs index bf402da..0f8d1d8 100644 --- a/src/DodoSSH.Api/Features/Sync/SyncEndpoints.cs +++ b/src/DodoSSH.Api/Features/Sync/SyncEndpoints.cs @@ -1,4 +1,5 @@ using DodoSSH.Api.Authorization; +using DodoSSH.Api.Features.Events; using DodoSSH.Api.Setup; using DodoSSH.Contracts; using DodoSSH.Domain.Authorization; @@ -81,6 +82,7 @@ internal sealed class SyncPullEndpoint( internal sealed class SyncPushEndpoint( ICurrentUserContext currentUser, IVaultAccessService vaultAccess, + IVaultEventPublisher events, SyncService sync) : Endpoint, NotFound, ProblemHttpResult>> { @@ -127,6 +129,8 @@ internal sealed class SyncPushEndpoint( // single stale item cannot block everything else a client queued while offline. var response = await sync.PushAsync(access.Vault!, user.Id, req, ct).ConfigureAwait(false); + Announce(access.Vault!.Id, response); + return TypedResults.Ok(response); } catch (PushBatchTooLargeException exception) @@ -142,4 +146,41 @@ internal sealed class SyncPushEndpoint( StatusCodes.Status400BadRequest, ProblemCodes.PushBatchTooLarge, exception.Message); } } + + /// + /// Tells every socket following this vault that it has moved. + /// + /// + /// + /// Here rather than inside , and that placement is the point: + /// the push has committed and released the per-vault advisory lock by the time this runs. Announced + /// from inside, it would name a sequence no reader could see yet and would hold the lock that + /// serialises writers across a fan-out. See ADR 0003 and ADR 0012. + /// + /// + /// The highest applied sequence, ignoring duplicates: a duplicate means an earlier push of + /// that operation already landed, and it was announced then. Nothing applied means nothing to say — + /// a batch of pure conflicts moved no vault, and announcing one anyway would have every client pull + /// for a change that is not there. + /// + /// + private void Announce(Guid vaultId, SyncPushResponse response) + { + var highest = 0L; + + foreach (var result in response.Results) + { + if (result.Status == SyncOperationStatus.Applied + && result.ChangeSequence is { } sequence + && sequence > highest) + { + highest = sequence; + } + } + + if (highest > 0) + { + events.VaultChanged(vaultId, highest); + } + } } diff --git a/src/DodoSSH.Api/Features/Teams/TeamService.cs b/src/DodoSSH.Api/Features/Teams/TeamService.cs index 737d786..ff681f3 100644 --- a/src/DodoSSH.Api/Features/Teams/TeamService.cs +++ b/src/DodoSSH.Api/Features/Teams/TeamService.cs @@ -1,4 +1,5 @@ using System.Globalization; +using DodoSSH.Api.Features.Events; using DodoSSH.Contracts; using DodoSSH.Domain; using DodoSSH.Infrastructure; @@ -62,6 +63,7 @@ internal readonly record struct TeamAccess(Team? Team, TeamRole Role) internal sealed class TeamService( DodoDbContext database, TimeProvider clock, + IVaultEventPublisher events, ILogger logger) { /// Longest acceptable slug. Matches the column. @@ -548,6 +550,11 @@ internal sealed class TeamService( TeamLog.MemberAdded(logger, teamId, target.Id, role, actor.Id); + // Membership is what the server will serve, so every vault this team owns has just appeared in + // the new member's list — before anybody wraps a key to them, which is a separate act and its + // own notice. Told at once rather than on their next pass. See ADR 0012. + events.VaultAccessChanged(target.Id); + return await DescribeAsync(target, membership, cancellationToken).ConfigureAwait(false); } @@ -744,6 +751,11 @@ internal sealed class TeamService( }).ConfigureAwait(false); TeamLog.MemberRemoved(logger, teamId, memberId, actor.Id, revoked); + + // After the commit, so their client re-reads a list the server has already stopped serving + // those vaults from. Their open socket re-resolves as it forwards this, which is what stops it + // announcing changes to vaults they have just lost. + events.VaultAccessChanged(memberId); } /// Revokes one user's grants on every vault a team owns, and flags each for rekey. diff --git a/src/DodoSSH.Api/Features/Teams/VaultGrantService.cs b/src/DodoSSH.Api/Features/Teams/VaultGrantService.cs index 277d8b1..3506c01 100644 --- a/src/DodoSSH.Api/Features/Teams/VaultGrantService.cs +++ b/src/DodoSSH.Api/Features/Teams/VaultGrantService.cs @@ -1,4 +1,5 @@ using System.Security.Cryptography; +using DodoSSH.Api.Features.Events; using DodoSSH.Contracts; using DodoSSH.Domain; using DodoSSH.Infrastructure; @@ -27,6 +28,7 @@ namespace DodoSSH.Api.Features.Teams; internal sealed class VaultGrantService( DodoDbContext database, TimeProvider clock, + IVaultEventPublisher events, ILogger logger) { /// @@ -414,6 +416,11 @@ internal sealed class VaultGrantService( TeamLog.GrantIssued( logger, vault.Id, generation, request.RecipientUserId, actor.Id); + + // The recipient, never the actor. This is the whole of what makes a shared vault arrive at + // once rather than on the recipient's next pass — and it is the case the README has had to + // apologise for since sharing shipped. See ADR 0012. + events.VaultAccessChanged(request.RecipientUserId); } /// @@ -701,6 +708,11 @@ internal sealed class VaultGrantService( TeamLog.GrantRevoked(logger, vault.Id, recipientUserId, actor.Id); + // Told so their client stops showing a vault it can no longer open, rather than leaving it + // listed until the next pass. It does not reach what they already pulled — nothing can, see + // ADR 0001 — and the server-side effect is immediate regardless of whether this arrives. + events.VaultAccessChanged(recipientUserId); + return true; } diff --git a/src/DodoSSH.Api/Program.cs b/src/DodoSSH.Api/Program.cs index abc8cdb..d5702d6 100644 --- a/src/DodoSSH.Api/Program.cs +++ b/src/DodoSSH.Api/Program.cs @@ -1,4 +1,5 @@ using DodoSSH.Api.Authorization; +using DodoSSH.Api.Features.Events; using DodoSSH.Api.Features.Identity; using DodoSSH.Api.Features.Sync; using DodoSSH.Api.Features.Teams; @@ -47,6 +48,14 @@ builder.Services.AddScoped(); builder.Services.AddScoped(); builder.Services.AddSingleton(); +// A singleton, because the sockets it holds outlive the requests that opened them. Registered twice +// resolving to the same instance, for the reason the invitation claim above is: the endpoint needs the +// whole hub — admit, remove, count — while the write paths that announce a change need only the two +// methods that announce one, and should not gain a reference to connection management to get them. +builder.Services.AddSingleton(); +builder.Services.AddSingleton( + provider => provider.GetRequiredService()); + // Scoped rather than the AddAuthorization default of singleton: the handler reads the request's // DbContext, and a singleton would capture one for the lifetime of the process. builder.Services.AddScoped(); @@ -65,6 +74,13 @@ var app = builder.Build(); app.BlockFastEndpointsRouteTable(); +// Before the authentication middleware, because the upgrade handshake has to survive it: the events +// endpoint answers an ordinary authenticated request that happens to become a socket, and without +// this the upgrade is never offered and the handler sees a plain GET. No allow-list of origins is +// configured, deliberately — every client here is a native application sending a bearer token, so +// there is no browser origin to trust and nothing a cross-site request could reach without one. +app.UseWebSockets(); + app.UseAuthentication(); app.UseAuthorization(); diff --git a/src/DodoSSH.Api/Setup/Configuration.cs b/src/DodoSSH.Api/Setup/Configuration.cs index 6ed888f..1deb6ef 100644 --- a/src/DodoSSH.Api/Setup/Configuration.cs +++ b/src/DodoSSH.Api/Setup/Configuration.cs @@ -38,6 +38,22 @@ internal static class Configuration "Sync:DefaultPullLimit must not exceed Sync:MaxPullLimit.") .ValidateOnStart(); + services.AddOptions() + .BindConfiguration(EventsOptions.SectionName) + .ValidateDataAnnotations() + .Validate( + options => options.HeartbeatInterval > TimeSpan.Zero, + "Events:HeartbeatInterval must be greater than zero.") + .Validate( + options => options.AccessRefreshInterval > TimeSpan.Zero, + "Events:AccessRefreshInterval must be greater than zero.") + .Validate( + options => options.MaxConnectionDuration > options.AccessRefreshInterval, + "Events:MaxConnectionDuration must exceed Events:AccessRefreshInterval; a connection " + + "that never lives long enough to re-read its own access has none of the bound that " + + "setting exists to provide.") + .ValidateOnStart(); + services.AddOptions() .BindConfiguration(RelayOptions.SectionName) .ValidateDataAnnotations() diff --git a/src/DodoSSH.Api/Setup/DodoOptions.cs b/src/DodoSSH.Api/Setup/DodoOptions.cs index d51d5c1..31ed25d 100644 --- a/src/DodoSSH.Api/Setup/DodoOptions.cs +++ b/src/DodoSSH.Api/Setup/DodoOptions.cs @@ -159,6 +159,82 @@ public sealed class RelayOptions public TimeSpan DrainTimeout { get; set; } = TimeSpan.FromSeconds(30); } +/// Realtime push settings. See ADR 0012. +/// +/// Every one of these bounds a socket rather than a feature: with the whole thing off, or every cap +/// met, clients synchronise on their timer exactly as they did before this existed. That is what +/// makes it safe for an operator to turn any of them down. +/// +public sealed class EventsOptions +{ + /// Configuration section name. + public const string SectionName = "Events"; + + /// + /// Whether this deployment pushes vault changes at all. + /// + /// + /// On by default, unlike the relay: this needs no outbound network, no target resolution and no + /// new trust, and a deployment behind a proxy that will not upgrade should say so here rather than + /// have every client discover it by failing. + /// + public bool Enabled { get; set; } = true; + + /// Maximum concurrent sockets per node. + [Range(1, 100_000)] + public int MaxConnectionsTotal { get; set; } = 500; + + /// + /// Maximum concurrent sockets per account. + /// + /// + /// Per account rather than per device, because the server cannot see a device here. Eight is a + /// laptop, a desktop, a phone and room to reconnect before the old socket has been reaped. + /// + [Range(1, 1000)] + public int MaxConnectionsPerUser { get; set; } = 8; + + /// + /// How many notices may be queued for one socket before the oldest are dropped. + /// + /// + /// A notice names a vault and a position, so a newer one subsumes the one it replaces. The depth + /// therefore buys smoothness over a brief stall and nothing else — losing the tail of a burst + /// costs a client nothing, because the newest notice still says to pull. + /// + [Range(1, 10_000)] + public int OutboundQueueDepth { get; set; } = 64; + + /// + /// How often the server pings an idle socket. + /// + /// + /// Below the sixty seconds most reverse proxies idle out at, because a silent socket that a proxy + /// has quietly dropped is indistinguishable from a quiet one until something is sent down it. + /// + public TimeSpan HeartbeatInterval { get; set; } = TimeSpan.FromSeconds(30); + + /// + /// How often an open socket re-reads which vaults its account may follow. + /// + /// + /// The backstop for a grant withdrawn mid-connection. Grants and membership changes publish + /// immediately, so this is what covers the paths that do not — and what bounds the window if one + /// is ever added without remembering to. + /// + public TimeSpan AccessRefreshInterval { get; set; } = TimeSpan.FromMinutes(5); + + /// + /// The longest any one socket may live, regardless of its token. + /// + /// + /// A socket normally ends at its access token's expiry, which is far shorter. This is the bound + /// for a provider that issues long-lived tokens, and it is what makes "no connection is older than + /// this" a property of the server rather than of the identity provider's configuration. + /// + public TimeSpan MaxConnectionDuration { get; set; } = TimeSpan.FromHours(12); +} + /// Sync protocol limits. public sealed class SyncOptions { diff --git a/src/DodoSSH.Api/Setup/EndpointRegistration.cs b/src/DodoSSH.Api/Setup/EndpointRegistration.cs index bb9d545..5260569 100644 --- a/src/DodoSSH.Api/Setup/EndpointRegistration.cs +++ b/src/DodoSSH.Api/Setup/EndpointRegistration.cs @@ -1,3 +1,4 @@ +using DodoSSH.Api.Features.Events; using DodoSSH.Api.Features.Identity; using DodoSSH.Api.Features.Meta; using DodoSSH.Api.Features.Sync; @@ -43,6 +44,7 @@ internal static class EndpointRegistration typeof(ReadKeyLogEndpoint), typeof(SyncPullEndpoint), typeof(SyncPushEndpoint), + typeof(VaultEventsEndpoint), typeof(CreateTeamEndpoint), typeof(ListTeamsEndpoint), typeof(UpdateTeamEndpoint), diff --git a/src/DodoSSH.Api/appsettings.json b/src/DodoSSH.Api/appsettings.json index 6df91bf..fde102b 100644 --- a/src/DodoSSH.Api/appsettings.json +++ b/src/DodoSSH.Api/appsettings.json @@ -26,6 +26,11 @@ "MaxConcurrentSessionsPerUser": 10, "MaxConcurrentSessionsTotal": 200 }, + "Events": { + "Enabled": true, + "MaxConnectionsTotal": 500, + "MaxConnectionsPerUser": 8 + }, "Sync": { "MaxOperationsPerPush": 500, "MaxPayloadBytes": 8388608, diff --git a/src/DodoSSH.Client.Api/VaultEventStream.cs b/src/DodoSSH.Client.Api/VaultEventStream.cs new file mode 100644 index 0000000..df1cb0b --- /dev/null +++ b/src/DodoSSH.Client.Api/VaultEventStream.cs @@ -0,0 +1,552 @@ +using System.Net.WebSockets; +using System.Security.Cryptography; +using System.Text.Json; +using System.Threading.Channels; +using DodoSSH.Contracts; + +namespace DodoSSH.Client.Api; + +/// +/// A server's "pull now" notices, as everything above the transport needs them. +/// +/// +/// +/// A queue to read from rather than an event to subscribe to, and that shape is the point: the one +/// consumer is a synchronisation loop that already waits on a timer, so it can wait on this the same +/// way and keep every continuation on the thread it started from. An event would deliver on whichever +/// thread the socket happened to complete on, which in a user interface is the difference between +/// working and an intermittent rendering fault nobody can reproduce. +/// +/// +/// Reading this is never how a change is applied. A notice says which vault moved and nothing +/// else; the answer to it is the ordinary delta pull. See ADR 0012. +/// +/// +public interface IVaultEventStream : IDisposable +{ + /// Whether a socket is currently established. + /// + /// For the interface to say whether it is live, not for a caller to branch on before reading: + /// synchronising is correct whether or not this is true, because the timer is the fallback. + /// + bool IsConnected { get; } + + /// + /// Waits for the next notice. + /// + /// + /// Connects on the first call and reconnects for as long as it is read, so a caller neither starts + /// nor restarts anything. A server that cannot be reached is not an error here — it is a wait that + /// has not finished — because the caller's alternative is the timer it is already running. + /// + ValueTask ReadAsync(CancellationToken cancellationToken); + + /// + /// Takes a notice if one is already waiting, without blocking. + /// + /// Whether there was one. + /// + /// How a caller coalesces a burst. Five people saving at once produces five notices whose answer + /// is a single synchronisation pass, so the loop reads one, waits a moment, and swallows the rest + /// rather than running the same pull five times. + /// + bool TryRead(out VaultEvent notice); +} + +/// +/// A stream that never delivers anything. +/// +/// +/// For a server that does not advertise the events feature, and for tests. Deliberately waits +/// for ever rather than completing: a caller selecting between this and a timer must fall through to +/// the timer, and a read that returned immediately would spin that loop as fast as the machine allows. +/// +public sealed class IdleVaultEventStream : IVaultEventStream +{ + /// The one instance. It holds nothing. + public static IdleVaultEventStream Instance { get; } = new(); + + /// + public bool IsConnected => false; + + /// + public async ValueTask ReadAsync(CancellationToken cancellationToken) + { + await Task.Delay(System.Threading.Timeout.Infinite, cancellationToken).ConfigureAwait(false); + + // Unreachable: the delay above only ever ends by throwing. + return new VaultEvent(VaultEventKinds.Ping); + } + + /// + public bool TryRead(out VaultEvent notice) + { + notice = new VaultEvent(VaultEventKinds.Ping); + return false; + } + + /// + public void Dispose() + { + // Nothing is held. + } +} + +/// Tuning for . +/// +/// Every value here bounds a reconnection rather than a feature. With the socket permanently +/// unavailable the client synchronises on its timer, so the cost of getting these wrong is latency, +/// never correctness. +/// +public sealed record VaultEventStreamOptions +{ + /// The defaults. + public static VaultEventStreamOptions Default { get; } = new(); + + /// How long to wait before the first reconnection attempt. + public TimeSpan InitialBackoff { get; init; } = TimeSpan.FromSeconds(1); + + /// The longest the backoff may grow to. + /// + /// A minute, which is the polling interval: past that point reconnecting sooner buys nothing, + /// because the timer has already done the work the socket would have prompted. + /// + public TimeSpan MaxBackoff { get; init; } = TimeSpan.FromMinutes(1); + + /// + /// How long a socket may be silent before it is presumed dead. + /// + /// + /// The server pings on an interval it states in its hello, so silence past a multiple of + /// that means the connection is gone rather than idle — which is otherwise indistinguishable, and + /// is exactly what a reverse proxy that quietly drops idle sockets produces. Used only until a + /// hello arrives; after that the server's own figure is trusted. + /// + public TimeSpan InitialSilenceTimeout { get; init; } = TimeSpan.FromSeconds(90); + + /// How many notices may be waiting before the oldest are dropped. + /// + /// Small on purpose. A notice means "pull that vault", so a newer one subsumes the one it + /// displaces; a backlog would only make the loop pull repeatedly for work it has already done. + /// + public int QueueDepth { get; init; } = 32; +} + +/// +/// Holds a socket to one server open, and hands over what it says. +/// +/// +/// +/// The whole class is a reconnection policy. A dropped socket is the ordinary case — laptops sleep, +/// proxies time out, tokens expire, servers are redeployed — so nothing here treats a failure as +/// exceptional: it backs off and dials again, for as long as somebody is reading. +/// +/// +/// It is safe to have no server at all. Every failure path ends in "wait, then try again", and the +/// caller's synchronisation timer runs regardless, which is what makes it correct for this class to +/// stay silent about problems rather than surface them. +/// +/// +public sealed class VaultEventStream : IVaultEventStream, IAsyncDisposable +{ + private readonly Uri endpoint; + private readonly IAccessTokenProvider tokens; + private readonly TimeProvider clock; + private readonly VaultEventStreamOptions options; + private readonly Func> connect; + private readonly Channel notices; + private readonly CancellationTokenSource closing = new(); + private readonly Lock starting = new(); + + private Task? pump; + private bool disposed; + + /// Creates a stream against one server. + /// The server's base URL, as an ordinary http or https address. + /// Supplies a bearer token, refreshing it when it is due. + /// Time source, for the backoff and the silence timeout. + /// Tuning, or null for the defaults. + public VaultEventStream( + Uri serverUrl, + IAccessTokenProvider tokens, + TimeProvider clock, + VaultEventStreamOptions? options = null) + : this(serverUrl, tokens, clock, DialAsync, options) + { + } + + /// + /// The connector is injected so the suite can drive this against a test host's in-memory socket. + /// Reconnection is the entire behaviour of this class, and testing it against a real network would + /// mean testing it against the one thing that cannot be made to fail on demand. + /// + internal VaultEventStream( + Uri serverUrl, + IAccessTokenProvider tokens, + TimeProvider clock, + Func> connect, + VaultEventStreamOptions? options = null) + { + ArgumentNullException.ThrowIfNull(serverUrl); + ArgumentNullException.ThrowIfNull(tokens); + ArgumentNullException.ThrowIfNull(clock); + ArgumentNullException.ThrowIfNull(connect); + + endpoint = EventsUrl(serverUrl); + this.tokens = tokens; + this.clock = clock; + this.connect = connect; + this.options = options ?? VaultEventStreamOptions.Default; + + notices = Channel.CreateBounded(new BoundedChannelOptions(this.options.QueueDepth) + { + FullMode = BoundedChannelFullMode.DropOldest, + SingleReader = true, + SingleWriter = true, + }); + } + + /// + public bool IsConnected { get; private set; } + + /// + public ValueTask ReadAsync(CancellationToken cancellationToken) + { + ObjectDisposedException.ThrowIf(disposed, this); + + Start(); + + return notices.Reader.ReadAsync(cancellationToken); + } + + /// + /// + /// Does not start the connection, unlike : a caller draining a burst has + /// already read one notice, and "is there another right now" is not a reason to dial a server. + /// + public bool TryRead(out VaultEvent notice) => notices.Reader.TryRead(out notice!); + + /// + /// Ends the connection, without waiting for it to unwind. + /// + /// + /// What a shell calls when a connection is dropped, from a synchronous path that must not block — + /// IVaultServer is , and blocking on a socket teardown from the + /// user-interface thread is exactly the sync-over-async this repository bans. Cancelling is enough: + /// every loop reads the token, and the pump has nothing to flush. + /// + public void Dispose() + { + if (disposed) + { + return; + } + + disposed = true; + closing.Cancel(); + + // The source is deliberately left undisposed. The pump may still be inside a linked token + // source derived from this one, and disposing a parent out from under a live child is how a + // clean shutdown becomes an ObjectDisposedException on a background thread. It holds no timer + // and no handle once cancelled; DisposeAsync is the path that cleans it up properly. + } + + /// Ends the connection and waits for it to unwind. + /// The deterministic form, for a caller that can await one — tests, mostly. + public async ValueTask DisposeAsync() + { + if (disposed) + { + return; + } + + disposed = true; + + await closing.CancelAsync().ConfigureAwait(false); + + if (pump is not null) + { + try + { + await pump.ConfigureAwait(false); + } + catch (OperationCanceledException) + { + // The point of the cancel above. + } + } + + closing.Dispose(); + } + + /// + /// Turns a server's base URL into its event socket's. + /// + /// + /// The path is replaced rather than appended, matching every other call in this client: request + /// paths here are absolute — /api/v1/… — so a deployment behind a path prefix is already + /// unsupported, and pretending otherwise in this one place would be a difference nobody could act + /// on. + /// + private static Uri EventsUrl(Uri serverUrl) => + new UriBuilder(serverUrl) + { + Scheme = string.Equals(serverUrl.Scheme, Uri.UriSchemeHttps, StringComparison.OrdinalIgnoreCase) + ? "wss" + : "ws", + Path = VaultEvents.Path, + Query = string.Empty, + Fragment = string.Empty, + }.Uri; + + private static async Task DialAsync(Uri url, string token, CancellationToken cancellationToken) + { + var socket = new ClientWebSocket(); + + try + { + socket.Options.AddSubProtocol(VaultEvents.SubProtocol); + + // A header rather than the Sec-WebSocket-Protocol smuggling ADR 0004 needs for the relay: + // this client is a native application and can set one, and the token here is the ordinary + // bearer credential rather than a ticket. + socket.Options.SetRequestHeader("Authorization", $"Bearer {token}"); + + await socket.ConnectAsync(url, cancellationToken).ConfigureAwait(false); + + return socket; + } + catch + { + socket.Dispose(); + throw; + } + } + + private void Start() + { + if (pump is not null) + { + return; + } + + lock (starting) + { + pump ??= Task.Run(() => RunAsync(closing.Token), closing.Token); + } + } + + /// Connects, reads until it cannot, waits, and does it again. + private async Task RunAsync(CancellationToken cancellationToken) + { + var backoff = options.InitialBackoff; + + while (!cancellationToken.IsCancellationRequested) + { + var outcome = await AttemptAsync(cancellationToken).ConfigureAwait(false); + + // A socket that lived long enough to say hello proves the server is there and willing, so + // the next failure starts from the bottom again rather than inheriting the backoff that + // got us here. Without this a laptop that woke, connected, and then lost its network an + // hour later would wait a full minute before trying, having already proved it need not. + if (outcome == Outcome.Established) + { + backoff = options.InitialBackoff; + } + + // The server said this token is spent, which the token provider can fix without waiting. + // Reconnecting at once is the whole reason that close code is distinct. + var wait = outcome == Outcome.TokenExpired ? TimeSpan.Zero : Jitter(backoff); + + if (wait > TimeSpan.Zero) + { + try + { + await Task.Delay(wait, clock, cancellationToken).ConfigureAwait(false); + } + catch (OperationCanceledException) + { + return; + } + + backoff = backoff < options.MaxBackoff + ? Shorter(backoff * 2, options.MaxBackoff) + : options.MaxBackoff; + } + } + } + + /// One connection, from dial to close. + private async Task AttemptAsync(CancellationToken cancellationToken) + { + WebSocket? socket = null; + + try + { + var token = await tokens.GetAccessTokenAsync(cancellationToken).ConfigureAwait(false); + + socket = await connect(endpoint, token, cancellationToken).ConfigureAwait(false); + + IsConnected = true; + + return await PumpAsync(socket, cancellationToken).ConfigureAwait(false); + } + catch (OperationCanceledException) + { + return Outcome.Cancelled; + } + catch (Exception exception) when (exception is not OutOfMemoryException) + { + // Every failure this can meet — no network, a refused upgrade, a server that has not been + // deployed with this feature, a token that cannot be refreshed — has the same remedy, and + // none of them is worth telling a user about. The synchronisation timer is still running. + return Outcome.Failed; + } + finally + { + IsConnected = false; + socket?.Dispose(); + } + } + + /// Reads frames until the socket ends or goes quiet. + private async Task PumpAsync(WebSocket socket, CancellationToken cancellationToken) + { + var buffer = new byte[8 * 1024]; + var silence = options.InitialSilenceTimeout; + var established = false; + + while (socket.State == WebSocketState.Open && !cancellationToken.IsCancellationRequested) + { + // Rebuilt per frame rather than reset, because a linked source cannot be un-cancelled and + // the deadline is what detects a socket that has silently gone away. + using var deadline = new CancellationTokenSource(silence, clock); + using var quiet = CancellationTokenSource.CreateLinkedTokenSource( + cancellationToken, deadline.Token); + + WebSocketReceiveResult received; + + try + { + received = await socket.ReceiveAsync(buffer, quiet.Token).ConfigureAwait(false); + } + catch (OperationCanceledException) when (!cancellationToken.IsCancellationRequested) + { + // Silent for longer than the server said it would be. The socket is gone in a way that + // only reconnecting can discover, which is what a proxy dropping an idle connection + // looks like from this end. + return Ended(established); + } + + if (received.MessageType == WebSocketMessageType.Close) + { + return (int?)received.CloseStatus == VaultEvents.TokenExpiredCloseCode + ? Outcome.TokenExpired + : Ended(established); + } + + // Binary is reserved by ADR 0012 for shared-session data, and text that arrived in pieces + // is longer than anything this protocol defines. Skipped rather than fatal, so a newer + // server does not cost this client its push for the whole session. + if (received.MessageType != WebSocketMessageType.Text || !received.EndOfMessage) + { + continue; + } + + if (Parse(buffer.AsSpan(0, received.Count)) is not { } frame) + { + continue; + } + + established = true; + silence = await AbsorbAsync(socket, frame, silence, cancellationToken).ConfigureAwait(false); + } + + return Ended(established); + } + + /// Deals with one frame, and says how long the socket may now stay quiet. + private async Task AbsorbAsync( + WebSocket socket, + VaultEvent frame, + TimeSpan silence, + CancellationToken cancellationToken) + { + if (frame.HeartbeatSeconds is > 0 and var seconds) + { + // Three missed heartbeats. Two is within one paused thread of a false positive, and a + // false positive here costs a reconnection rather than anything a user sees. + silence = TimeSpan.FromSeconds(seconds * 3); + } + + if (string.Equals(frame.Kind, VaultEventKinds.Ping, StringComparison.Ordinal)) + { + await socket.SendAsync( + JsonSerializer.SerializeToUtf8Bytes( + new VaultEvent(VaultEventKinds.Pong), DodoSshJsonContext.Default.VaultEvent), + WebSocketMessageType.Text, + endOfMessage: true, + cancellationToken) + .ConfigureAwait(false); + + return silence; + } + + // Everything else, including a kind this build has never heard of, goes to the reader — which + // is what makes the frame table extensible. An unrecognised kind is one the caller ignores; + // refusing it here would be this class deciding what a newer server may say. + notices.Writer.TryWrite(frame); + + return silence; + } + + private static Outcome Ended(bool established) => + established ? Outcome.Established : Outcome.Failed; + + private static VaultEvent? Parse(ReadOnlySpan utf8) + { + try + { + return JsonSerializer.Deserialize(utf8, DodoSshJsonContext.Default.VaultEvent); + } + catch (JsonException) + { + return null; + } + } + + /// + /// Spreads reconnections out, so a server that restarts is not met by every client at once. + /// + /// + /// RandomNumberGenerator because System.Random is banned repo-wide. Nothing here is + /// security-relevant — the ban exists so that nothing key-, token- or nonce-adjacent can reach for + /// the weak one by habit, and paying a few microseconds to keep that rule absolute is the cheaper + /// side of the trade. + /// + private static TimeSpan Jitter(TimeSpan delay) + { + var milliseconds = (int)Math.Clamp(delay.TotalMilliseconds, 1, int.MaxValue / 2); + + return TimeSpan.FromMilliseconds( + milliseconds + RandomNumberGenerator.GetInt32(0, Math.Max(1, milliseconds / 2))); + } + + private static TimeSpan Shorter(TimeSpan left, TimeSpan right) => left < right ? left : right; + + /// How one connection attempt ended. + private enum Outcome + { + /// Never got as far as a frame. Back off. + Failed, + + /// Ran, and then ended. Back off, but from the bottom. + Established, + + /// The server closed it because the token expired. Reconnect at once with a new one. + TokenExpired, + + /// The stream is being disposed. + Cancelled, + } +} diff --git a/src/DodoSSH.Client.App/Views/HostDrawer.axaml b/src/DodoSSH.Client.App/Views/HostDrawer.axaml index 08dffe6..5dc6afd 100644 --- a/src/DodoSSH.Client.App/Views/HostDrawer.axaml +++ b/src/DodoSSH.Client.App/Views/HostDrawer.axaml @@ -569,8 +569,8 @@ - -