diff --git a/CHANGELOG.md b/CHANGELOG.md index 645d5cda..5c171fbf 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -18,6 +18,10 @@ heading, and (b) bumping `pubspec.yaml` `version:` in the same commit. ### Added +- `clide events --since [--filter X]` reads events after a cursor and + returns them plus a next-cursor — the pull-based complement to the + `tail --events` stream, made for agent poll loops. Reports `gap: true` when + the cursor has aged out of the in-memory ring (D-85). (T-223) - "Install 'clide' command in PATH" command (`clide.installCli`) copies the bundled C client to `~/.local/bin`, VS Code style. On launch clide warns when `clide` is missing from PATH or points at the GUI bundle instead of the CLI diff --git a/lib/src/cli/argv_to_request.dart b/lib/src/cli/argv_to_request.dart index 7154f472..ff327b5f 100644 --- a/lib/src/cli/argv_to_request.dart +++ b/lib/src/cli/argv_to_request.dart @@ -24,8 +24,10 @@ import 'package:clide/src/ipc/envelope.dart'; import 'package:clide/src/ipc/schema_v1.dart'; /// Umbrella commands — single-token names with no subsystem.verb -/// split. Match the IDs the dispatcher exposes directly. -const Set _umbrellaCommands = {'status', 'tail', 'version', 'ping', 'capabilities'}; +/// split. Match the IDs the dispatcher exposes directly. `tail` and +/// `events` are handled by the IPC server itself (streaming / cursor-pull +/// event reads, T-129 / T-223) rather than the dispatcher. +const Set _umbrellaCommands = {'status', 'tail', 'events', 'version', 'ping', 'capabilities'}; /// Sealed result of translating argv. Caller (the IPC server, or the /// C client wrapper in T-126) handles either branch. diff --git a/lib/src/ipc/server.dart b/lib/src/ipc/server.dart index 31443d4a..8a368218 100644 --- a/lib/src/ipc/server.dart +++ b/lib/src/ipc/server.dart @@ -29,6 +29,7 @@ class IpcServer { required this.log, this.events, this.replayDepth = 16, + this.eventLogDepth = 1024, }); final DaemonDispatcher dispatcher; @@ -62,6 +63,21 @@ class IpcServer { /// Per-subsystem ring buffer of recent events for replay. final Map> _replay = {}; + /// Bound on the cursor log that serves `clide events --since` (T-223). + /// Larger than [replayDepth]: a polling agent reads at its own cadence, + /// so a deeper window means fewer gaps between polls. + final int eventLogDepth; + + /// Global, arrival-ordered log of events keyed by a monotonic cursor — + /// the pull-based read surface (T-223 / D-85). Drop-oldest; never blocks + /// the producer. [_lastCursor] is the high-water mark handed back as the + /// next cursor; [_droppedThrough] is the highest evicted cursor, so a pull + /// whose cursor predates it is told there's a gap rather than silently + /// missing the dropped events. + final Queue<_LoggedEvent> _eventLog = Queue<_LoggedEvent>(); + int _lastCursor = 0; + int _droppedThrough = 0; + String get socketPath => _socketPath ?? workspaceSocketPath(workspaceRoot); bool get isRunning => _socket != null; @@ -117,6 +133,9 @@ class IpcServer { _busSub = null; _subscribers.clear(); _replay.clear(); + _eventLog.clear(); + _lastCursor = 0; + _droppedThrough = 0; await _accepts?.cancel(); _accepts = null; for (final c in List.from(_clients)) { @@ -209,7 +228,7 @@ class IpcServer { await _enterStreamingMode(client, req); return; } - response = await dispatcher.dispatch(req); + response = _isEventsPull(req) ? _eventsSince(req) : await dispatcher.dispatch(req); } } on FormatException catch (e) { response = IpcResponse.err( @@ -312,6 +331,59 @@ class IpcServer { return _replay[filter] ?? const []; } + // -- event pull (T-223) --------------------------------------------------- + + /// `clide events [--since ] [--filter X]` — a one-shot, cursor-based + /// read of the event log, the request/response complement to the + /// never-returning `tail --events` stream. + bool _isEventsPull(IpcRequest req) => req.cmd == 'events'; + + /// Build the pull response: every logged event with cursor > `since` + /// (optionally filtered by subsystem), the high-water `cursor` to poll + /// from next, and `gap: true` when `since` predates the retained window + /// (events between `since` and the oldest retained entry were dropped). + IpcResponse _eventsSince(IpcRequest req) { + final flags = req.args['flags']; + final flagMap = flags is Map ? flags : const {}; + final filter = (flagMap['filter'] as String?) ?? '*'; + final since = _parseSince(flagMap['since']); + if (since == null) { + return IpcResponse.err( + id: req.id, + error: IpcError( + code: IpcExitCode.userError, + kind: IpcErrorKind.userError, + message: '--since must be a non-negative integer cursor', + ), + ); + } + final out = >[]; + for (final logged in _eventLog) { + if (logged.cursor <= since) continue; + if (filter != '*' && logged.event.subsystem != filter) continue; + out.add({...logged.event.toJson(), 'cursor': logged.cursor}); + } + // A gap only means something when the caller had a prior position + // (since > 0); a first read (since 0) just gets whatever's retained. + final gap = since > 0 && since < _droppedThrough; + return IpcResponse.ok(id: req.id, data: { + 'events': out, + 'cursor': _lastCursor, + 'gap': gap, + if (gap) 'oldestCursor': _eventLog.isEmpty ? _lastCursor : _eventLog.first.cursor, + }); + } + + /// Parse the `--since` flag (a string from argv or an int from a typed + /// request). Null on an invalid (non-integer / negative) value; absent + /// means 0 (read from the beginning of the retained window). + int? _parseSince(Object? raw) { + if (raw == null) return 0; + if (raw is int) return raw < 0 ? null : raw; + final n = int.tryParse('$raw'); + return (n == null || n < 0) ? null : n; + } + void _onBusEvent(DaemonEvent e) { final ev = IpcEvent( subsystem: e.subsystem, @@ -325,6 +397,13 @@ class IpcServer { while (ring.length > replayDepth) { ring.removeFirst(); } + // Push to the cursor log (T-223). Drop-oldest; record the highest + // evicted cursor so a later `--since` below it reports a gap. + final cursor = ++_lastCursor; + _eventLog.addLast(_LoggedEvent(cursor, ev)); + while (_eventLog.length > eventLogDepth) { + _droppedThrough = _eventLog.removeFirst().cursor; + } // Fan out to live subscribers whose filter matches. final stale = []; for (final entry in _subscribers.entries) { @@ -368,3 +447,11 @@ class IpcServer { } } } + +/// One entry in the cursor log (T-223): an [IpcEvent] tagged with the +/// monotonic [cursor] assigned when it was recorded. +class _LoggedEvent { + _LoggedEvent(this.cursor, this.event); + final int cursor; + final IpcEvent event; +} diff --git a/test/cli/argv_to_request_test.dart b/test/cli/argv_to_request_test.dart index 9e3ff2bb..13707519 100644 --- a/test/cli/argv_to_request_test.dart +++ b/test/cli/argv_to_request_test.dart @@ -94,6 +94,13 @@ void main() { expect(_expectOk(parseArgv(['ping'], requestId: 'p')).cmd, 'ping'); expect(_expectOk(parseArgv(['version'], requestId: 'v')).cmd, 'version'); }); + + test('"events --since 5 --filter pane" parses to the events command (T-223)', () { + final req = _expectOk(parseArgv(['events', '--since', '5', '--filter', 'pane'], requestId: 'e')); + expect(req.cmd, 'events'); + expect((req.args['flags'] as Map)['since'], '5'); + expect((req.args['flags'] as Map)['filter'], 'pane'); + }); }); group('parseArgv — errors', () { diff --git a/test/ipc/server_events_pull_test.dart b/test/ipc/server_events_pull_test.dart new file mode 100644 index 00000000..2a1976c9 --- /dev/null +++ b/test/ipc/server_events_pull_test.dart @@ -0,0 +1,188 @@ +/// T-223 — cursor-based pull events over the IPC socket +/// (`clide events --since `). Tests the one-shot read complement +/// to `tail --events`: events-after-cursor, the next-cursor high-water mark, +/// no-drop/no-duplicate across polls, subsystem filtering, and the gap marker +/// when a cursor predates the retained window (drop-oldest, D-85). +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/schema_v1.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, + ); + +void _emit(DaemonBus bus, String sub, String kind, [Map data = const {}]) => + bus.emit(DaemonEvent(subsystem: sub, kind: kind, data: data, ts: DateTime.now().toUtc())); + +/// Let the bus drain into the server's cursor log before a pull. +Future _drain() => Future.delayed(const Duration(milliseconds: 20)); + +void main() { + late Directory ws; + late DaemonBus bus; + late IpcServer server; + + setUp(() async { + ws = await Directory.systemTemp.createTemp('clide-events-pull-'); + bus = DaemonBus(); + server = IpcServer( + dispatcher: DaemonDispatcher(), + workspaceRoot: '${ws.path}/${DateTime.now().microsecondsSinceEpoch}', + log: _silent(), + events: bus, + eventLogDepth: 4, + ); + await server.start(); + }); + + tearDown(() async { + try { + await server.stop(); + } catch (_) {} + if (ws.existsSync()) ws.deleteSync(recursive: true); + }); + + /// Open a connection, send one `events` request, read the single response. + Future pull({Object? since, String? filter}) async { + final s = await _connect(server); + try { + final lines = s.cast>().transform(utf8.decoder).transform(const LineSplitter()); + final it = StreamIterator(lines); + s.write('${IpcRequest(id: 'e', cmd: 'events', args: { + 'flags': { + if (since != null) 'since': since, + if (filter != null) 'filter': filter, + }, + }).encode()}\n'); + await s.flush(); + if (!await it.moveNext().timeout(const Duration(seconds: 2))) { + throw StateError('no response'); + } + final resp = IpcMessage.decode(it.current) as IpcResponse; + await it.cancel(); + return resp; + } finally { + await s.close(); + } + } + + List cursors(IpcResponse r) => [for (final e in r.data['events'] as List) (e as Map)['cursor'] as int]; + + test('since 0 returns all retained events with a monotonic cursor + high-water', () async { + _emit(bus, 'pane', 'spawned', {'id': 'p1'}); + _emit(bus, 'git', 'changed'); + _emit(bus, 'pane', 'closed', {'id': 'p1'}); + await _drain(); + + final r = await pull(since: 0); + expect(r.ok, isTrue); + expect(cursors(r), [1, 2, 3]); + expect(r.data['cursor'], 3); // next poll uses --since 3 + expect(r.data['gap'], isFalse); + // Each event keeps its wire shape plus the cursor. + final first = (r.data['events'] as List).first as Map; + expect(first['subsystem'], 'pane'); + expect(first['kind'], 'spawned'); + }); + + test('since returns only events after it', () async { + _emit(bus, 'pane', 'a'); + _emit(bus, 'pane', 'b'); + _emit(bus, 'pane', 'c'); + await _drain(); + final r = await pull(since: 1); + expect(cursors(r), [2, 3]); + }); + + test('repeated polls neither drop nor duplicate', () async { + _emit(bus, 'pane', 'a'); + _emit(bus, 'pane', 'b'); + await _drain(); + final first = await pull(since: 0); + expect(cursors(first), [1, 2]); + final next = first.data['cursor'] as int; + + _emit(bus, 'pane', 'c'); + _emit(bus, 'pane', 'd'); + await _drain(); + final second = await pull(since: next); + expect(cursors(second), [3, 4]); // no overlap with the first batch + expect(second.data['cursor'], 4); + + // Polling again with the latest cursor yields nothing new. + final third = await pull(since: 4); + expect(third.data['events'], isEmpty); + expect(third.data['cursor'], 4); + }); + + test('filter restricts to a single subsystem', () async { + _emit(bus, 'pane', 'a'); + _emit(bus, 'git', 'changed'); + _emit(bus, 'pane', 'b'); + await _drain(); + final r = await pull(since: 0, filter: 'pane'); + expect(cursors(r), [1, 3]); // git event (cursor 2) excluded + }); + + test('a cursor that aged out of the ring is reported as a gap', () async { + // Depth is 4; emit 6 so cursors 1,2 are evicted (retained: 3,4,5,6). + for (var i = 0; i < 6; i++) { + _emit(bus, 'pane', 'e$i'); + } + await _drain(); + + final r = await pull(since: 1); // 1 predates the retained window + expect(r.data['gap'], isTrue); + expect(r.data['oldestCursor'], 3); + expect(cursors(r), [3, 4, 5, 6]); + + // A cursor at/after the dropped watermark is not a gap. + final r2 = await pull(since: 2); + expect(r2.data['gap'], isFalse); + }); + + test('a first read (since 0) is never a gap even after eviction', () async { + for (var i = 0; i < 6; i++) { + _emit(bus, 'pane', 'e$i'); + } + await _drain(); + final r = await pull(since: 0); + expect(r.data['gap'], isFalse); + }); + + test('bare events with no flags reads from the start', () async { + _emit(bus, 'pane', 'a'); + _emit(bus, 'pane', 'b'); + await _drain(); + final r = await pull(); + expect(cursors(r), [1, 2]); + }); + + test('a string --since (argv shape) parses', () async { + _emit(bus, 'pane', 'a'); + _emit(bus, 'pane', 'b'); + await _drain(); + final r = await pull(since: '1'); + expect(cursors(r), [2]); + }); + + test('an invalid --since is a user error', () async { + final r = await pull(since: 'abc'); + expect(r.ok, isFalse); + expect(r.error!.kind, IpcErrorKind.userError); + }); +}