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
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
//! The UDS ingest listener producers dial (issue #6848, DOC-73 §4.2).
//!
//! Why: DOC-73 §4.2 makes console "the only ingester" — every harness pushes
//! `HarnessEvent` frames over one socket rather than console polling each of
//! them. This module is that socket's server half: bind it hardened (the same
//! `0700` directory / `0600` socket / peer-uid convention every other socket
//! in the workspace uses), accept connections, and read each one as
//! newline-delimited `HarnessEvent` JSON for as long as the producer keeps the
//! connection open — a producer's `control_bus::PushClient` (§4.2, a parallel
//! slice) is expected to hold one long-lived connection and stream frames on
//! it rather than reconnect per event. The wire format is exactly the
//! newline-terminated JSON [`trusty_common::uds::send_framed_notification`]
//! writes — `PushClient::flush` (issue #6847) dials, sends one frame, and
//! half-closes per call today, but nothing here assumes a connection carries
//! only one frame, so a later `PushClient` revision that holds the connection
//! open and streams needs no ingest-side change.
//! What: [`ingest_socket_path`] resolves the socket
//! (`daemon_socket_path("trusty-console")`, the same cross-crate convention
//! [`trusty_common::daemon_socket_path`] documents); [`bind_ingest`] binds it
//! through [`trusty_common::uds::bind_singleton_hardened`], which reclaims a
//! stale socket file left by an unclean shutdown rather than refusing to bind
//! forever; [`serve_ingest`] accepts until shutdown, bounding concurrent
//! connections at [`MAX_CONCURRENT_CONNECTIONS`] and spawning one task per
//! connection. A connection that sends a line that is not valid `HarnessEvent`
//! JSON logs a warning and keeps reading — one malformed frame must not cost
//! every frame after it on the same connection, still less every other
//! producer's connection. A line over [`MAX_LINE_BYTES`] ends that connection
//! only, and the allocation for that line is capped at read time (a fresh
//! [`tokio::io::AsyncReadExt::take`] budget per line) rather than checked only
//! after an unbounded `read_until` returns — an unterminated line otherwise
//! grows the buffer without limit before the size check ever runs. A
//! connection that goes [`READ_IDLE_TIMEOUT`] without producing a full line is
//! dropped the same way.
//! Test: `super::tests` — `ingest_over_uds_socket_reaches_the_bus`,
//! `malformed_line_does_not_kill_the_listener`,
//! `oversized_line_ends_only_its_own_connection`,
//! `unterminated_line_never_grows_past_the_line_cap`,
//! `stale_socket_file_is_reclaimed_on_bind`,
//! `idle_connection_is_dropped_after_the_read_timeout`,
//! `connections_beyond_the_limit_wait_for_a_free_slot`,
//! `shutdown_is_observed_while_the_connection_pool_is_saturated`.
use ;
use Arc;
use Duration;
use ;
use ;
use Semaphore;
use HarnessEvent;
use ;
use EventBus;
/// Largest single newline-delimited frame this listener accepts.
///
/// Why not [`trusty_common::uds::MAX_FRAME_BYTES`] (8 MiB, sized for a request
/// on the request/response RPC transport that module also defines): a
/// `HarnessEvent` is a small, flat envelope around one domain-tagged payload —
/// the largest arm is `HarnessPayload::Hook`'s open-ended `serde_json::Value`,
/// still nowhere near an RPC response carrying whole search results. 1 MiB is
/// generous headroom over any hook payload observed in the workspace today
/// while still bounding a misbehaving or malicious producer's memory cost per
/// connection to a fixed budget rather than the connection's lifetime.
/// Test: `super::tests::oversized_line_ends_only_its_own_connection`,
/// `super::tests::unterminated_line_never_grows_past_the_line_cap`.
pub const MAX_LINE_BYTES: usize = 1024 * 1024;
/// How long one connection may go without completing a line before this
/// listener gives up on it.
///
/// Why 60 s: a `PushClient` that holds a connection open and streams (the
/// intended shape once §4.2's revision lands) sends far more often than this
/// under any real workload; a producer that goes a full minute mid-frame is
/// indistinguishable from a hung peer holding a slot another producer could
/// use. Mirrors the idle-timeout shape `webhook_relay::serve::serve_until`
/// applies per connection, scaled up because that listener serves one
/// request/response pair per connection while this one serves a long-lived
/// stream.
/// Test: `super::tests::idle_connection_is_dropped_after_the_read_timeout`.
const READ_IDLE_TIMEOUT: Duration = from_secs;
/// Largest number of producer connections this listener serves at once.
///
/// Why a cap at all: an unbounded accept loop hands every accepted socket its
/// own task and `BufReader`, so a burst of connections (malicious or merely a
/// thundering herd of restarted harnesses) costs memory proportional to
/// connection count with no ceiling. 256 is comfortably above the number of
/// harnesses this workspace runs per host today (`trusty-mpm`, `trusty-code`,
/// `trusty-agents`, `trusty-analyze` — DOC-73 §4.1) with headroom for several
/// concurrent reconnect storms, while still bounding worst-case memory to a
/// fixed multiple of [`MAX_LINE_BYTES`] rather than the connection count the
/// kernel's accept queue happens to deliver.
/// Test: `super::tests::connections_beyond_the_limit_wait_for_a_free_slot`.
const MAX_CONCURRENT_CONNECTIONS: usize = 256;
/// Resolve the console event-bus ingest socket path (DOC-73 §4.2).
///
/// # Errors
///
/// Whatever [`trusty_common::daemon_socket_path`] returns — the console data
/// directory could not be resolved or created.
pub
/// Why the ingest socket could not be bound.
pub
/// Bind the ingest socket, reclaiming a stale socket file per
/// [`trusty_common::uds::bind_singleton_hardened`].
///
/// Why `bind_singleton_hardened` rather than `bind_hardened`: console is
/// supervised and restarted like every other daemon in this workspace, and
/// `bind_hardened` refuses an occupied path outright — a predecessor that
/// exits uncleanly (a SIGKILL, a crash) never unlinks its socket, so every
/// bind after that would fail forever with no operator-visible cause. This
/// mirrors `trusty-memory`, `trusty-analyze` and `trusty-review`'s own
/// listener binds, all of which made the same switch for the same reason.
///
/// # Errors
///
/// [`IngestError::Bind`] when the path cannot be bound — including when the
/// probe proves another process is still serving it, in which case this
/// process correctly does not take over.
///
/// Test: `super::tests::ingest_over_uds_socket_reaches_the_bus`,
/// `super::tests::stale_socket_file_is_reclaimed_on_bind`.
pub async
/// Accept and serve ingest connections until `shutdown` resolves.
///
/// Each connection is handed to its own `tokio::spawn` rather than served
/// inline, matching the accept-loop shape every other in-process UDS listener
/// in this workspace uses (`webhook_relay::serve_until`): one slow or stalled
/// producer must never delay another producer's connection from being
/// accepted. Concurrency is bounded by [`MAX_CONCURRENT_CONNECTIONS`]: a
/// [`tokio::sync::Semaphore`] permit is acquired before a connection is
/// spawned, so once the limit is reached, further accepted connections wait
/// for a slot rather than piling up an unbounded number of tasks.
///
/// Test: `super::tests::ingest_over_uds_socket_reaches_the_bus`,
/// `super::tests::malformed_line_does_not_kill_the_listener`,
/// `super::tests::oversized_line_ends_only_its_own_connection`,
/// `super::tests::connections_beyond_the_limit_wait_for_a_free_slot`,
/// `super::tests::shutdown_is_observed_while_the_connection_pool_is_saturated`.
pub async
/// [`serve_ingest`] with the concurrency cap taken explicitly rather than
/// read from [`MAX_CONCURRENT_CONNECTIONS`], so a test can prove the gating
/// behavior against a limit of 1 instead of opening 257 real connections.
///
/// Test: `super::tests::connections_beyond_the_limit_wait_for_a_free_slot`,
/// `super::tests::shutdown_is_observed_while_the_connection_pool_is_saturated`.
pub async
/// Serve one accepted connection: verify the peer, then read newline-delimited
/// `HarnessEvent` JSON until EOF, a read error, or [`READ_IDLE_TIMEOUT`]
/// elapses with no line completed.
///
/// Why a peer check at all: the same ADR-0034 §3 trust boundary every other
/// hardened socket in this crate enforces — the `0600` mode is a documented
/// intention until something actually reads the peer's uid off the accepted
/// connection.
async
/// [`handle_connection`] with the idle timeout taken explicitly rather than
/// read from [`READ_IDLE_TIMEOUT`], so a test can prove the drop behavior in
/// milliseconds instead of the real 60 s.
///
/// Test: `super::tests::idle_connection_is_dropped_after_the_read_timeout`.
pub async
/// Read up to and including the next newline, with the allocation for THIS
/// line capped at [`MAX_LINE_BYTES`] before any byte of it is read.
///
/// Why: the prior shape read with a plain `BufReader::read_until` and checked
/// `line.len() > MAX_LINE_BYTES` only after it returned — an unterminated line
/// from a peer that never sends `\n` grows `line` without limit in the
/// meantime, because `read_until` has no bound of its own. Re-taking a fresh
/// `.take(MAX_LINE_BYTES)` budget on `&mut reader` for every line — the same
/// pattern `trusty_common::uds::stream_client::UdsStreamClient::read_line` and
/// `uds::rpc::read_one_frame` use — caps the allocation up front instead.
/// `take` borrows the reader rather than consuming it, so bytes already
/// buffered past the cap for an oversized line stay in `reader` and are
/// visible to the caller's own oversized-line check on `line.len()`.
///
/// Test: `super::tests::unterminated_line_never_grows_past_the_line_cap`.
pub async