1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
//! One framed JSON-RPC stream, rendered as the Server-Sent Events a browser
//! reads (#6155).
//!
//! Why: `search_uds` built this bridge for trusty-search and `memory_uds` needs
//! exactly the same thing for trusty-memory — the frame-to-`data:` encoding, the
//! keep-alive comment, the cancel-safe reader task, and the response head are
//! identical for both, because none of them is about which daemon is on the
//! other end. A second copy is how one bridge starts closing a failed stream
//! silently while the other reports it. #6637 hoisted it out of trusty-console
//! for the same reason one step out: trusty-code-gui's webview bridge is a
//! third consumer of the identical shape, and it lives in a different crate.
//!
//! What is NOT here: anything a service owns. Which method is streaming, which
//! JSON-RPC code becomes which HTTP status, and what the peek before the
//! response head means all stay in the per-service module, because the two
//! daemons answer with different code tables.
//!
//! ## The contract this preserves
//!
//! One stream ITEM is exactly the JSON document one SSE `data:` line carried,
//! parsed rather than prefixed. So this re-prefixes it and adds nothing: every
//! event reaches the browser byte-identical to what the daemon's own SSE route
//! wrote before ADR-0032 retired it.
//!
//! Two things those SSE routes emitted that an RPC stream does not, and what
//! happens to them here:
//!
//! - the `: heartbeat\n\n` comment every 20 s. It exists so an idle TCP body is
//! not torn down, and the browser hop is still TCP — so this emits it, on the
//! same interval.
//! - the terminal `data:` framing of a failure. A mid-stream failure becomes one
//! `{"type":"error","message":…}` event before the body closes, because a
//! consumer reads a closed stream as a COMPLETED operation and a silent close
//! would report a broken one as finished.
//!
//! Test: `sse_data_is_one_line_per_event` below. Whole streams run through a
//! real router in trusty-console's `tests/search_uds_bridge.rs` and
//! `tests/memory_uds_bridge.rs`.
use ;
use ;
use ;
use StreamExt as _;
use ;
use mpsc;
use warn;
use UdsRpcError;
use FramedStream;
/// How often an open stream emits an SSE keep-alive comment.
///
/// The same 20 s the daemons' own SSE routes used, so an idle browser
/// connection sees the byte sequence it saw before the migration.
pub const SSE_HEARTBEAT_INTERVAL: Duration = from_secs;
/// How many stream items may buffer between the socket reader and the browser.
///
/// Matches the daemon-side producer buffer, so neither side is the first to
/// accumulate behind a slow reader.
pub const SSE_BUFFER: usize = 64;
/// Build the `200 text/event-stream` response for an already-opened stream.
///
/// Why the caller passes `first` separately: a streaming method can refuse, and
/// that refusal arrives as the stream's FIRST frame after the dial has already
/// succeeded. The caller peeks it so a refusal becomes an HTTP status rather
/// than an empty `200`; the peeked item then has to lead the body, which is what
/// this takes it for. `None` is a stream that ended with no items — a
/// well-formed empty answer, and an immediately-closed event stream.
/// What: the peeked item, then [`sse_tail`], under the three headers a browser
/// and any reverse proxy in front of the console need.
/// Test: covered end to end by trusty-console's
/// `tests/memory_uds_bridge.rs`, in `a_stream_reaches_the_browser_frame_for_frame`.
///
/// [`sse_tail`]: crate::uds::sse::sse_tail
/// The rest of an open stream, as SSE frames plus keep-alive comments.
///
/// Why the reader runs in its own task rather than inside the `select!`:
/// [`FramedStream::next_frame`] reads a line off a `BufReader`, and cancelling
/// that mid-line — which a heartbeat tick would do — discards the bytes already
/// read. Moving the read behind an `mpsc` makes both arms of the select
/// cancel-safe, since `Receiver::recv` and `Interval::tick` both are.
///
/// The task also carries the disconnect signal: when the browser goes, axum
/// drops this body and the receiver drops. The read itself selects on
/// `Sender::closed()` so the task notices immediately rather than at the next
/// frame — a status stream can be silent for minutes, and waiting for a frame
/// that will never come would hold the socket, and the daemon's producer behind
/// it, open for exactly that long. Either way the task returns, dropping the
/// [`FramedStream`] and closing the socket, which is what ends the producer.
///
/// Test: covered end to end by trusty-console's
/// `tests/search_uds_bridge.rs`, in `a_mid_stream_failure_becomes_an_error_event`
/// and `a_browser_disconnect_releases_the_daemon_socket`.
///
/// [`FramedStream`]: crate::uds::stream_client::FramedStream
/// [`FramedStream::next_frame`]: crate::uds::stream_client::FramedStream::next_frame
/// Encode one stream item as an SSE `data:` frame.
///
/// Why `to_string` on a `Value` rather than passing the raw line through: the
/// stream carries parsed JSON, and re-serialising is what puts it back on one
/// line — an embedded newline would split one event into two.
/// Test: `sse_data_is_one_line_per_event`.