pullServerEvents method

Stream<PullEvent> pullServerEvents(
  1. PullInfoDto pullInfo
)

A typed pull stream. Cancelling the subscription closes the underlying HTTP request, which in turn triggers daemon-side pull cancellation.

Implementation

Stream<PullEvent> pullServerEvents(PullInfoDto pullInfo) async* {
  final Stream<String> sseStream;
  try {
    sseStream = await serverSse.request(
      '/pull',
      requestBody: pullInfo.toJson(),
    );
  } catch (error, stackTrace) {
    final normalized = _normalizeError(error);
    if (identical(normalized, error)) rethrow;
    Error.throwWithStackTrace(normalized, stackTrace);
  }

  var terminalReceived = false;
  try {
    await for (final chunk in sseStream) {
      final events = <PullEvent>[];
      SseClient.parse(chunk, (event, data) {
        switch (event) {
          case EventType.START:
            events.add(PullStarted(PullStartDto.fromJson(data)));
          case EventType.DATA:
            events.add(PullProgress(PullDownloadDto.fromJson(data)));
          case EventType.DONE:
            terminalReceived = true;
            events.add(PullCompleted(PullDoneDto.fromJson(data)));
          case EventType.ERROR:
            terminalReceived = true;
            events.add(PullFailed(PullErrorDto.fromJson(data)));
        }
      });
      for (final event in events) {
        yield event;
      }
      if (terminalReceived) return;
    }
  } catch (error, stackTrace) {
    final normalized = _normalizeError(error);
    if (identical(normalized, error)) rethrow;
    Error.throwWithStackTrace(normalized, stackTrace);
  }
  if (!terminalReceived) {
    throw const SseProtocolException(
      'Pull stream closed before a terminal done or error event',
    );
  }
}