From e194f0280224d77dd08f4c885b27bc0b01a888fb Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Tue, 19 May 2026 14:20:07 +0200 Subject: [PATCH] =?UTF-8?q?T-129:=20event=20streaming=20over=20the=20socke?= =?UTF-8?q?t=20=E2=80=94=20`clide=20tail=20--events`?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Sixth slice of T-99. Long-lived event subscription path, the second half of D-6. Wire shape: - Client sends `{cmd:"tail", args:{flags:{events:true, filter:X}}}`. - Server responds with `{ok:true, data:{streaming:true, filter:X}}`. - Server pushes `{type:"event", subsystem, kind, ts, data}` lines until the client closes. Server (lib/src/ipc/server.dart): - Takes a DaemonBus, subscribes to DaemonEvent on start. - Per-subsystem ring buffer (replayDepth=16 per D-6) populated on every emit. - `tail --events` connection: send ack, replay matching events from ring, register the client for future fanout. - _argv envelope now unwrapped at the server layer so the streaming check sees the inner `tail` cmd (not just `_argv`). - Broken subscriber writes drop the subscriber cleanly; the bus doesn't block on a stalled client. Client (native/clide-cli/clide.c): - Sniffs `data.streaming:true` in the ack. If set, loops reading JSON-line events to stdout (with fflush per line) until EOF. Tests: - test/ipc/server_streaming_test.dart — 8 cases covering ack shape, filter, replay buffer (size + ordering), multi-subscriber fanout, broken-subscriber cleanup. - test/cli/clide_cli_e2e_test.dart gets a tail --events test that spawns the C client, emits two events on the bus, asserts they print on stdout. T-99 children remaining: T-130 (MCP), T-131 (wrap-up). Co-Authored-By: Claude --- .pql/changelog/ticket_history/2026-05.sql | 2 + .pql/changelog/tickets/2026-05.sql | 11 ++ CHANGELOG.md | 5 + lib/main.dart | 7 +- lib/src/cli/argv_dispatch.dart | 51 ++--- lib/src/ipc/server.dart | 163 +++++++++++++++- native/clide-cli/clide.c | 28 ++- test/cli/clide_cli_e2e_test.dart | 37 ++++ test/ipc/server_streaming_test.dart | 223 ++++++++++++++++++++++ 9 files changed, 501 insertions(+), 26 deletions(-) create mode 100644 test/ipc/server_streaming_test.dart diff --git a/.pql/changelog/ticket_history/2026-05.sql b/.pql/changelog/ticket_history/2026-05.sql index e811a744..f73305ed 100644 --- a/.pql/changelog/ticket_history/2026-05.sql +++ b/.pql/changelog/ticket_history/2026-05.sql @@ -1862,3 +1862,5 @@ INSERT INTO ticket_history (ticket_id, field, old_value, new_value, changed_by, INSERT INTO ticket_history (ticket_id, field, old_value, new_value, changed_by, changed_at, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('T-127', 'status', 'in_progress', 'done', NULL, '2026-05-19 10:03:26', '2026-05-19 10:03:26', '2026-05-19 10:03:26', NULL, 'b8bf877d2057f7baa5acf6af7a13d529', 1) ON CONFLICT(hash) DO NOTHING; INSERT INTO ticket_history (ticket_id, field, old_value, new_value, changed_by, changed_at, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('T-128', 'status', 'backlog', 'in_progress', NULL, '2026-05-19 10:05:09', '2026-05-19 10:05:09', '2026-05-19 10:05:09', NULL, '763a80ec25b5fe712e5ad44b75dbb22f', 1) ON CONFLICT(hash) DO NOTHING; INSERT INTO ticket_history (ticket_id, field, old_value, new_value, changed_by, changed_at, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('T-128', 'status', 'in_progress', 'done', NULL, '2026-05-19 10:06:46', '2026-05-19 10:06:46', '2026-05-19 10:06:46', NULL, 'cc53b46764f28515d735aa77c8eeef2c', 1) ON CONFLICT(hash) DO NOTHING; +INSERT INTO ticket_history (ticket_id, field, old_value, new_value, changed_by, changed_at, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('T-129', 'status', 'backlog', 'in_progress', NULL, '2026-05-19 12:06:39', '2026-05-19 12:06:39', '2026-05-19 12:06:39', NULL, '575584f46d21b51bab3e69da1b940718', 1) ON CONFLICT(hash) DO NOTHING; +INSERT INTO ticket_history (ticket_id, field, old_value, new_value, changed_by, changed_at, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('T-129', 'status', 'in_progress', 'done', NULL, '2026-05-19 12:19:51', '2026-05-19 12:19:51', '2026-05-19 12:19:51', NULL, 'e9bf50652357521c7ee74bd6577e8108', 1) ON CONFLICT(hash) DO NOTHING; diff --git a/.pql/changelog/tickets/2026-05.sql b/.pql/changelog/tickets/2026-05.sql index b85b966c..c2fb64f9 100644 --- a/.pql/changelog/tickets/2026-05.sql +++ b/.pql/changelog/tickets/2026-05.sql @@ -2156,3 +2156,14 @@ Acceptance: 3. flutter analyze + full test suite green. Source: T-99 sketch. Depends on T-127.', 'done', 'medium', NULL, NULL, NULL, '2026-05-18 11:59:06', '2026-05-19 10:06:46', NULL, '481c51b17b027b5a04279a0ef6a3a029', 1) ON CONFLICT(id) DO UPDATE SET type=excluded.type, parent_id=excluded.parent_id, title=excluded.title, description=excluded.description, status=excluded.status, priority=excluded.priority, assigned_to=excluded.assigned_to, team=excluded.team, decision_ref=excluded.decision_ref, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at > tickets.updated_at OR (excluded.updated_at = tickets.updated_at AND excluded.hash > tickets.hash); +INSERT INTO tickets (id, type, parent_id, title, description, status, priority, assigned_to, team, decision_ref, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('T-129', 'task', 'T-99', 'event streaming over the socket — `clide tail --events`', 'Sixth slice of T-99(a). Long-lived subscription channel — the second half of D-6 (`clide tail --events [--filter [:]]`). + +Client opens a connection, sends {"subscribe": "|*"}, server pushes JSON-line events until the client closes. Per D-6: replay buffer per subsystem (default depth 16) so a late subscriber still sees recent effects. + +Acceptance: +1. `clide tail --events --filter git` streams git.* events from the running app. +2. Replay buffer per subsystem; new subscribers receive the last 16 events. +3. Server doesn''t block writes on a slow client (back-pressure handling per Q-2 — drop with a warning or apply flow control; resolve in this ticket). +4. End-to-end smoke: launch app, run `clide tail --events --filter pane` in another shell, perform a pane action in the UI, observe the event. + +Source: T-99 sketch. Depends on T-124 + T-126.', 'done', 'medium', NULL, NULL, NULL, '2026-05-18 11:59:11', '2026-05-19 12:19:51', NULL, '44c9141fea8635a9c03053afa92f3999', 1) ON CONFLICT(id) DO UPDATE SET type=excluded.type, parent_id=excluded.parent_id, title=excluded.title, description=excluded.description, status=excluded.status, priority=excluded.priority, assigned_to=excluded.assigned_to, team=excluded.team, decision_ref=excluded.decision_ref, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at > tickets.updated_at OR (excluded.updated_at = tickets.updated_at AND excluded.hash > tickets.hash); diff --git a/CHANGELOG.md b/CHANGELOG.md index 10e7f708..5e2ecc41 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -32,6 +32,11 @@ heading, and (b) bumping `pubspec.yaml` `version:` in the same commit. dumb pipe (T-99 / T-125). - `DaemonClient.reconnectAt(newPath)` — swap an active client onto a different socket without restart (project switch in T-127). +- Event streaming over the IPC socket (T-99 / T-129) — `clide tail + --events [--filter X]` opens a long-lived subscription, replays up + to 16 recent matching events per subsystem (D-6), and streams new + ones as JSON lines. C client loops on `data.streaming` ack. Slow / + broken subscribers drop themselves without blocking the bus. - C `clide` shell client at `native/clide-cli/clide.c`. Walks CWD up to the git root, hashes to the per-workspace socket (D-70), ships argv. `make clide-cli` builds it; on PATH, `clide status` works diff --git a/lib/main.dart b/lib/main.dart index 956cf064..81e985b4 100644 --- a/lib/main.dart +++ b/lib/main.dart @@ -96,7 +96,12 @@ Future main() async { ipcLog.warn('ipc', 'stop failed during swap: $e'); ipcLog.debug('ipc', '$st'); } - final server = IpcServer(dispatcher: dispatcher, workspaceRoot: workRoot.path, log: ipcLog); + final server = IpcServer( + dispatcher: dispatcher, + workspaceRoot: workRoot.path, + log: ipcLog, + events: daemonBus, + ); ipcServer = server; try { await server.start(); diff --git a/lib/src/cli/argv_dispatch.dart b/lib/src/cli/argv_dispatch.dart index e12500be..e8a2a38d 100644 --- a/lib/src/cli/argv_dispatch.dart +++ b/lib/src/cli/argv_dispatch.dart @@ -20,30 +20,37 @@ import 'package:clide/src/ipc/schema_v1.dart'; /// goes through the normal dispatcher path unchanged. const String argvSentinelCmd = '_argv'; -/// Wire the `_argv` sentinel handler onto [dispatcher]. The handler: -/// 1. Extracts `args.argv` as a List<String>. -/// 2. Calls [parseArgv]. -/// 3. If parsed → re-dispatches the inner request through the -/// *same* dispatcher (so per-handler logic runs once). -/// 4. If error → returns the pre-built [IpcResponse] verbatim, -/// patched with the outer request id so the client correlates. +/// Unwrap an `_argv` IpcRequest into the inner parsed request, or +/// return an error response if the envelope is malformed or the +/// argv doesn't parse. Pure function — no dispatch. Used by both +/// the IPC server (which needs the unwrapped cmd to decide whether +/// to enter streaming mode for `tail --events`, per T-129) and the +/// dispatcher-side handler below. +ArgvParseResult unwrapArgvRequest(IpcRequest outer) { + final raw = outer.args['argv']; + if (raw is! List) { + return ArgvError(IpcResponse.err( + id: outer.id, + error: IpcError( + code: IpcExitCode.userError, + kind: IpcErrorKind.userError, + message: '_argv requires args.argv to be a JSON array', + ), + )); + } + return parseArgv(raw.cast(), requestId: outer.id); +} + +/// Wire the `_argv` sentinel handler onto [dispatcher]. The handler +/// unwraps the inner argv via [unwrapArgvRequest], dispatches the +/// resulting request through the same dispatcher, and otherwise +/// returns the pre-built error response. Kept registered for the +/// non-streaming path; the IPC server intercepts before dispatch +/// for `tail --events` (T-129). void registerArgvUnwrap(DaemonDispatcher dispatcher) { dispatcher.register(argvSentinelCmd, (outer) async { - final raw = outer.args['argv']; - if (raw is! List) { - return IpcResponse.err( - id: outer.id, - error: IpcError( - code: IpcExitCode.userError, - kind: IpcErrorKind.userError, - message: '_argv requires args.argv to be a JSON array', - ), - ); - } - final argv = raw.cast(); - final result = parseArgv(argv, requestId: outer.id); - return switch (result) { - ArgvParsed(:final request) => dispatcher.dispatch(request), + return switch (unwrapArgvRequest(outer)) { + ArgvParsed(:final request) => await dispatcher.dispatch(request), ArgvError(:final response) => response, }; }); diff --git a/lib/src/ipc/server.dart b/lib/src/ipc/server.dart index f5165ed4..31443d4a 100644 --- a/lib/src/ipc/server.dart +++ b/lib/src/ipc/server.dart @@ -1,8 +1,13 @@ import 'dart:async'; +import 'dart:collection'; import 'dart:convert'; import 'dart:io'; +import 'package:clide/kernel/src/events/bus.dart'; +import 'package:clide/kernel/src/events/types.dart'; import 'package:clide/kernel/src/log.dart'; +import 'package:clide/src/cli/argv_dispatch.dart'; +import 'package:clide/src/cli/argv_to_request.dart'; import 'package:clide/src/daemon/dispatcher.dart'; import 'package:clide/src/ipc/envelope.dart'; import 'package:clide/src/ipc/paths.dart'; @@ -18,17 +23,45 @@ import 'package:clide/src/ipc/schema_v1.dart'; /// D-72. Per-handler isolate offload is the dispatcher / handler's /// concern, not this layer's. class IpcServer { - IpcServer({required this.dispatcher, required this.workspaceRoot, required this.log}); + IpcServer({ + required this.dispatcher, + required this.workspaceRoot, + required this.log, + this.events, + this.replayDepth = 16, + }); final DaemonDispatcher dispatcher; final String workspaceRoot; final Logger log; + /// Bus the server subscribes to for events forwarded to + /// `clide tail --events` subscribers. Optional — when null, the + /// tail handler still accepts subscriptions but never gets events + /// (useful in tests that don't need the full kernel wiring). + final DaemonBus? events; + + /// Per-subsystem replay-buffer depth (D-6: default 16). New + /// subscribers receive up to this many recent matching events on + /// connect so they don't miss effects emitted just before they + /// subscribed. + final int replayDepth; + ServerSocket? _socket; String? _socketPath; final List _clients = []; StreamSubscription? _accepts; + // Event streaming (T-129). + StreamSubscription? _busSub; + + /// Subscribers: client socket → filter (`*` or a subsystem name). + /// A connection enters this map after it sends `tail --events`. + final Map _subscribers = {}; + + /// Per-subsystem ring buffer of recent events for replay. + final Map> _replay = {}; + String get socketPath => _socketPath ?? workspaceSocketPath(workspaceRoot); bool get isRunning => _socket != null; @@ -62,6 +95,13 @@ class IpcServer { _accepts = socket.listen(_onClient, onError: (Object e, StackTrace st) { log.error('ipc', 'accept loop error', error: e, stackTrace: st); }); + // Subscribe to the bus so we can populate the replay ring AND + // fan out to live `tail --events` subscribers. Idempotent — + // we only attach when a bus is supplied. + final bus = events; + if (bus != null) { + _busSub = bus.on().listen(_onBusEvent); + } log.info('ipc', 'IPC server listening at $path'); } @@ -73,6 +113,10 @@ class IpcServer { if (s == null) return; _socket = null; _socketPath = null; + await _busSub?.cancel(); + _busSub = null; + _subscribers.clear(); + _replay.clear(); await _accepts?.cancel(); _accepts = null; for (final c in List.from(_clients)) { @@ -116,6 +160,7 @@ class IpcServer { }, onDone: () { _clients.remove(client); + _subscribers.remove(client); sub.cancel(); }, cancelOnError: true, @@ -138,7 +183,33 @@ class IpcServer { ), ); } else { - response = await dispatcher.dispatch(msg); + // Peel off the `_argv` envelope at the server layer so the + // streaming check sees the unwrapped command (T-129). Plain + // typed requests skip this path. + var req = msg; + if (req.cmd == argvSentinelCmd) { + final result = unwrapArgvRequest(req); + if (result is ArgvError) { + response = result.response; + // Fall through to write below. + try { + client.write('${response.encode()}\n'); + await client.flush(); + } catch (e) { + log.warn('ipc', 'client write failed: $e'); + } + return; + } + req = (result as ArgvParsed).request; + } + if (_isTailSubscribe(req)) { + // Long-lived subscription branch (T-129). Send the streaming + // ack, replay matching ring buffer entries, register the + // client. The connection stays open until the client closes. + await _enterStreamingMode(client, req); + return; + } + response = await dispatcher.dispatch(req); } } on FormatException catch (e) { response = IpcResponse.err( @@ -199,6 +270,94 @@ class IpcServer { } } + // -- event streaming (T-129) ---------------------------------------------- + + /// Recognise the `tail --events [--filter X]` subscription + /// request that the argv translator (T-125) produces. + bool _isTailSubscribe(IpcRequest req) { + if (req.cmd != 'tail') return false; + final flags = req.args['flags']; + return flags is Map && flags['events'] == true; + } + + Future _enterStreamingMode(Socket client, IpcRequest req) async { + final flags = req.args['flags'] as Map?; + final filter = (flags?['filter'] as String?) ?? '*'; + // Streaming ack — `data.streaming: true` tells the C client to + // loop-read instead of exiting after one response. + final ack = IpcResponse.ok(id: req.id, data: {'streaming': true, 'filter': filter}); + try { + client.write('${ack.encode()}\n'); + await client.flush(); + } catch (e) { + log.warn('ipc', 'streaming ack write failed: $e'); + return; + } + // Replay matching events from the ring. + final replay = _replayFor(filter); + for (final ev in replay) { + if (!_sendEvent(client, ev)) return; + } + _subscribers[client] = filter; + } + + Iterable _replayFor(String filter) { + if (filter == '*') { + // Flatten everything in arrival order. Per-subsystem rings + // preserve order within a subsystem; across subsystems the + // ordering is best-effort (interleaved-by-subsystem). Good + // enough for "what just happened". + return _replay.values.expand((q) => q); + } + return _replay[filter] ?? const []; + } + + void _onBusEvent(DaemonEvent e) { + final ev = IpcEvent( + subsystem: e.subsystem, + kind: e.kind, + data: e.data, + timestamp: e.ts, + ); + // Push to replay ring. + final ring = _replay.putIfAbsent(e.subsystem, () => Queue()); + ring.addLast(ev); + while (ring.length > replayDepth) { + ring.removeFirst(); + } + // Fan out to live subscribers whose filter matches. + final stale = []; + for (final entry in _subscribers.entries) { + final filter = entry.value; + if (filter != '*' && filter != e.subsystem) continue; + if (!_sendEvent(entry.key, ev)) { + stale.add(entry.key); + } + } + for (final s in stale) { + _subscribers.remove(s); + } + } + + /// Write an event line to [client]. Returns false on failure, which + /// the caller uses to drop the subscriber. We deliberately don't + /// await `flush` here — back-pressure handling per D-72: if the + /// socket's write buffer is full, dart:io's Socket.write enqueues + /// in-memory, and the kernel pushes through as it can. If the + /// client is genuinely gone the write throws or onDone fires and + /// the subscriber gets removed via _onClient's onDone. + bool _sendEvent(Socket client, IpcEvent ev) { + try { + client.write('${ev.encode()}\n'); + return true; + } catch (e) { + log.warn('ipc', 'subscriber write failed (dropping): $e'); + return false; + } + } + + // -- internals ------------------------------------------------------------ + /// `chmod` via `chmod(1)` because dart:io doesn't expose the /// syscall on unix. Cheap; only runs at start/stop. Future _chmod(String path, int modeBits) async { diff --git a/native/clide-cli/clide.c b/native/clide-cli/clide.c index 14225c7f..1ccd5c40 100644 --- a/native/clide-cli/clide.c +++ b/native/clide-cli/clide.c @@ -291,7 +291,6 @@ int main(int argc, char **argv) { close(fd); return EX_OSERR; } - close(fd); /* Pull out `ok`, `data`/`error` from the response. */ size_t ok_len = 0, data_len = 0, code_len = 0, msg_len = 0; @@ -302,11 +301,38 @@ int main(int argc, char **argv) { if (data) { fwrite(data, 1, data_len, stdout); fputc('\n', stdout); + fflush(stdout); } else { fputs("{}\n", stdout); + fflush(stdout); } + /* If the server flagged this as a streaming response + * (`tail --events` per T-129), loop-read event JSON-lines + * until the connection closes. Detection: look for the + * literal `"streaming":true` inside the data blob. */ + if (data && data_len > 0) { + char data_copy[16384]; + size_t copy_len = data_len < sizeof(data_copy) - 1 ? data_len : sizeof(data_copy) - 1; + memcpy(data_copy, data, copy_len); + data_copy[copy_len] = '\0'; + if (strstr(data_copy, "\"streaming\":true") != NULL || strstr(data_copy, "\"streaming\": true") != NULL) { + /* Streaming mode — keep reading event lines. Exit + * 0 on EOF (server closed cleanly), non-zero on + * read error. */ + char ev[65536]; + while (read_line(fd, ev, sizeof(ev)) == 0) { + fputs(ev, stdout); + fputc('\n', stdout); + fflush(stdout); + } + close(fd); + return 0; + } + } + close(fd); return 0; } + close(fd); const char *code_v = json_value(resp, "code", &code_len); const char *msg_v = json_value(resp, "message", &msg_len); int exit_code = code_v ? (int)strtol(code_v, NULL, 10) : EX_SOFTWARE; diff --git a/test/cli/clide_cli_e2e_test.dart b/test/cli/clide_cli_e2e_test.dart index c6f69c36..3c44c319 100644 --- a/test/cli/clide_cli_e2e_test.dart +++ b/test/cli/clide_cli_e2e_test.dart @@ -11,6 +11,8 @@ library; import 'dart:convert'; import 'dart:io'; +import 'package:clide/kernel/src/events/bus.dart'; +import 'package:clide/kernel/src/events/types.dart'; import 'package:clide/kernel/src/log.dart'; import 'package:clide/src/cli/argv_dispatch.dart'; import 'package:clide/src/daemon/dispatcher.dart'; @@ -25,6 +27,7 @@ void main() { late final Directory workspaceRoot; late final IpcServer server; late final DaemonDispatcher dispatcher; + late final DaemonBus streamingBus; setUpAll(() async { final repoRoot = Directory.current.path; @@ -52,10 +55,12 @@ void main() { Directory('${workspaceRoot.path}/.git').createSync(); dispatcher = DaemonDispatcher(); registerArgvUnwrap(dispatcher); + streamingBus = DaemonBus(); server = IpcServer( dispatcher: dispatcher, workspaceRoot: workspaceRoot.path, log: Logger(minLevel: LogLevel.error, sinks: const []), + events: streamingBus, ); await server.start(); }); @@ -65,6 +70,7 @@ void main() { try { await server.stop(); } catch (_) {} + await streamingBus.dispose(); if (workspaceRoot.existsSync()) { workspaceRoot.deleteSync(recursive: true); } @@ -134,5 +140,36 @@ void main() { expect(r.exitCode, isNot(0)); expect(r.stderr.toString(), isNotEmpty); }); + + test('tail --events streams bus events to stdout (T-129)', () async { + if (!hasCC) { + markTestSkipped('cc not available'); + return; + } + final proc = await Process.start(binaryPath, ['tail', '--events', '--filter', 'pane'], workingDirectory: workspaceRoot.path); + addTearDown(() => proc.kill()); + final lines = []; + final sub = proc.stdout.transform(utf8.decoder).transform(const LineSplitter()).listen(lines.add); + addTearDown(sub.cancel); + // Wait for the ack so the server has registered us. + var attempts = 0; + while (lines.isEmpty && attempts < 50) { + await Future.delayed(const Duration(milliseconds: 20)); + attempts++; + } + expect(lines, isNotEmpty, reason: 'no ack received'); + // Emit two events. + streamingBus.emit(DaemonEvent(subsystem: 'pane', kind: 'spawned', data: const {'id': 'p1'}, ts: DateTime.now().toUtc())); + streamingBus.emit(DaemonEvent(subsystem: 'pane', kind: 'closed', data: const {'id': 'p1'}, ts: DateTime.now().toUtc())); + attempts = 0; + while (lines.length < 3 && attempts < 100) { + await Future.delayed(const Duration(milliseconds: 20)); + attempts++; + } + expect(lines.length, greaterThanOrEqualTo(3), reason: 'expected ack + 2 events, got: $lines'); + final concatenated = lines.skip(1).join('\n'); + expect(concatenated, contains('"kind":"spawned"')); + expect(concatenated, contains('"kind":"closed"')); + }); }); } diff --git a/test/ipc/server_streaming_test.dart b/test/ipc/server_streaming_test.dart new file mode 100644 index 00000000..319ddbe7 --- /dev/null +++ b/test/ipc/server_streaming_test.dart @@ -0,0 +1,223 @@ +/// T-129 — event streaming over the IPC socket. Tests the +/// `tail --events` subscription branch on the server: subscriber +/// registration, per-subsystem replay-buffer (D-6 / replayDepth=16), +/// filter matching, fanout on bus events, and broken-subscriber +/// cleanup. +library; + +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; + +import 'package:clide/kernel/src/events/bus.dart'; +import 'package:clide/kernel/src/events/types.dart'; +import 'package:clide/kernel/src/log.dart'; +import 'package:clide/src/daemon/dispatcher.dart'; +import 'package:clide/src/ipc/envelope.dart'; +import 'package:clide/src/ipc/server.dart'; +import 'package:test/test.dart'; + +Logger _silent() => Logger(minLevel: LogLevel.error, sinks: const []); + +Future _connect(IpcServer s) async => Socket.connect( + InternetAddress(s.socketPath, type: InternetAddressType.unix), + 0, + ); + +/// Wrap a Socket in a line iterator backed by a single broadcast +/// stream so the same connection can read multiple framed lines. +({Stream lines, Socket sock}) _lineReader(Socket s) { + final stream = s.cast>().transform(utf8.decoder).transform(const LineSplitter()).asBroadcastStream(); + return (lines: stream, sock: s); +} + +/// Thin wrapper over StreamIterator with a [next] convenience. +class _Lines { + _Lines(Stream s) : _it = StreamIterator(s); + final StreamIterator _it; + Future next({Duration timeout = const Duration(seconds: 2)}) async { + final ok = await _it.moveNext().timeout(timeout); + if (!ok) throw StateError('stream ended before next line'); + return _it.current; + } + + Future cancel() => _it.cancel(); +} + +Future _send(Socket s, IpcRequest req) async { + s.write('${req.encode()}\n'); + await s.flush(); +} + +IpcRequest _tailReq({String? filter, String id = 't'}) => IpcRequest( + id: id, + cmd: 'tail', + args: { + 'flags': { + 'events': true, + if (filter != null) 'filter': filter, + }, + }, + ); + +void main() { + late Directory ws; + late DaemonDispatcher dispatcher; + late DaemonBus bus; + late IpcServer server; + + setUp(() async { + ws = await Directory.systemTemp.createTemp('clide-stream-test-'); + dispatcher = DaemonDispatcher(); + bus = DaemonBus(); + server = + IpcServer(dispatcher: dispatcher, workspaceRoot: '${ws.path}/${DateTime.now().microsecondsSinceEpoch}', log: _silent(), events: bus, replayDepth: 4); + await server.start(); + }); + + tearDown(() async { + try { + await server.stop(); + } catch (_) {} + await bus.dispose(); + if (ws.existsSync()) ws.deleteSync(recursive: true); + }); + + test('subscribe → streaming ack with filter echoed', () async { + final s = await _connect(server); + addTearDown(s.close); + final r = _lineReader(s); + await _send(s, _tailReq(filter: 'pane')); + final line = await r.lines.first.timeout(const Duration(seconds: 2)); + final ack = IpcMessage.decode(line) as IpcResponse; + expect(ack.ok, isTrue); + expect(ack.data['streaming'], isTrue); + expect(ack.data['filter'], 'pane'); + }); + + test('subscribe with no filter → wildcard ack', () async { + final s = await _connect(server); + addTearDown(s.close); + final r = _lineReader(s); + await _send(s, _tailReq()); + final line = await r.lines.first.timeout(const Duration(seconds: 2)); + final ack = IpcMessage.decode(line) as IpcResponse; + expect(ack.data['filter'], '*'); + }); + + test('events emitted post-subscribe land on the subscriber', () async { + final s = await _connect(server); + addTearDown(s.close); + final r = _lineReader(s); + final lineQ = _Lines(r.lines); + await _send(s, _tailReq(filter: 'pane')); + await lineQ.next(); // ack + bus.emit(DaemonEvent(subsystem: 'pane', kind: 'spawned', data: const {'id': 'p1'}, ts: DateTime.now().toUtc())); + final evLine = await lineQ.next(); + final ev = IpcMessage.decode(evLine) as IpcEvent; + expect(ev.subsystem, 'pane'); + expect(ev.kind, 'spawned'); + expect(ev.data['id'], 'p1'); + await lineQ.cancel(); + }); + + test('filter excludes non-matching subsystems', () async { + final s = await _connect(server); + addTearDown(s.close); + final r = _lineReader(s); + final lineQ = _Lines(r.lines); + await _send(s, _tailReq(filter: 'pane')); + await lineQ.next(); // ack + // Emit a non-matching event first, then a matching one. The + // subscriber should only see the matching one. + bus.emit(DaemonEvent(subsystem: 'git', kind: 'changed', data: const {}, ts: DateTime.now().toUtc())); + bus.emit(DaemonEvent(subsystem: 'pane', kind: 'closed', data: const {'id': 'p1'}, ts: DateTime.now().toUtc())); + final ev = IpcMessage.decode(await lineQ.next()) as IpcEvent; + expect(ev.subsystem, 'pane'); + expect(ev.kind, 'closed'); + await lineQ.cancel(); + }); + + test('replay buffer surfaces pre-subscribe events on connect', () async { + // Push three events before any subscriber exists. + for (var i = 0; i < 3; i++) { + bus.emit(DaemonEvent(subsystem: 'pane', kind: 'spawned', data: {'i': i}, ts: DateTime.now().toUtc())); + } + await Future.delayed(const Duration(milliseconds: 20)); + + final s = await _connect(server); + addTearDown(s.close); + final r = _lineReader(s); + final lineQ = _Lines(r.lines); + await _send(s, _tailReq(filter: 'pane')); + await lineQ.next(); // ack + final replayed = []; + for (var i = 0; i < 3; i++) { + final ev = IpcMessage.decode(await lineQ.next()) as IpcEvent; + replayed.add(ev.data['i'] as int); + } + expect(replayed, [0, 1, 2]); + await lineQ.cancel(); + }); + + test('replay ring is bounded to replayDepth (4 for this test)', () async { + for (var i = 0; i < 10; i++) { + bus.emit(DaemonEvent(subsystem: 'pane', kind: 'spawned', data: {'i': i}, ts: DateTime.now().toUtc())); + } + await Future.delayed(const Duration(milliseconds: 20)); + + final s = await _connect(server); + addTearDown(s.close); + final r = _lineReader(s); + final lineQ = _Lines(r.lines); + await _send(s, _tailReq(filter: 'pane')); + await lineQ.next(); // ack + final replayed = []; + for (var i = 0; i < 4; i++) { + final ev = IpcMessage.decode(await lineQ.next()) as IpcEvent; + replayed.add(ev.data['i'] as int); + } + // Last 4 of 0..9 → 6,7,8,9. + expect(replayed, [6, 7, 8, 9]); + await lineQ.cancel(); + }); + + test('multiple subscribers each receive an event independently', () async { + final sA = await _connect(server); + addTearDown(sA.close); + final sB = await _connect(server); + addTearDown(sB.close); + final rA = _lineReader(sA); + final rB = _lineReader(sB); + final qA = _Lines(rA.lines); + final qB = _Lines(rB.lines); + await _send(sA, _tailReq(filter: 'pane', id: 'A')); + await _send(sB, _tailReq(filter: 'pane', id: 'B')); + await qA.next(); // ack + await qB.next(); // ack + bus.emit(DaemonEvent(subsystem: 'pane', kind: 'event', data: const {'tag': 'broadcast'}, ts: DateTime.now().toUtc())); + final evA = IpcMessage.decode(await qA.next()) as IpcEvent; + final evB = IpcMessage.decode(await qB.next()) as IpcEvent; + expect(evA.data['tag'], 'broadcast'); + expect(evB.data['tag'], 'broadcast'); + await qA.cancel(); + await qB.cancel(); + }); + + test('subscriber going away removes itself from fanout (no crash on emit)', () async { + final s = await _connect(server); + final r = _lineReader(s); + final q = _Lines(r.lines); + await _send(s, _tailReq(filter: 'pane')); + await q.next(); // ack + await q.cancel(); + await s.close(); + // Give the server's onDone a tick. + await Future.delayed(const Duration(milliseconds: 20)); + // Emitting should not throw or stall — covered by reaching the + // next assertion. + bus.emit(DaemonEvent(subsystem: 'pane', kind: 'orphan', data: const {}, ts: DateTime.now().toUtc())); + await Future.delayed(const Duration(milliseconds: 20)); + expect(server.isRunning, isTrue); + }); +}