start method

Future<void> start()

Binds the event source, starts queue polling, and begins serving. Resolves once the server is listening.

Implementation

Future<void> start() async {
  await state.bind(source);
  _server = await HttpServer.bind(host, port);
  if (broker != null) {
    _queuePoll = Timer.periodic(queuePollInterval, (_) => _pollQueues());
    unawaited(_pollQueues());
  }
  // Fan live events out to connected SSE clients.
  state.tap.listen((event) {
    if (_sseClients.isEmpty) return;
    final frame = 'data: ${event.encode()}\n\n';
    for (final client in _sseClients.toList()) {
      try {
        client.write(frame);
      } catch (_) {
        _sseClients.remove(client);
      }
    }
  });
  _serve(_server!);
}