serve static method

  1. @visibleForTesting
Future<void> serve(
  1. Stream<String> input,
  2. void output(
    1. String message
    ),
  3. FutureOr<void> body(
    1. World w
    )
)

run, over input and output rather than the owner's socket — answers once the world has closed.

Implementation

@visibleForTesting
static Future<void> serve(
  Stream<String> input,
  void Function(String line) output,
  FutureOr<void> Function(World w) body,
) async {
  void send(String type, [Map<String, Object?> fields = const {}]) =>
      output(encodeWorldMessage(type, fields));

  World? world;
  // Once, however many ways it is asked for: a `close` and the owner's
  // hanging up usually arrive together.
  Future<void>? closing;
  var asked = Completer<void>();
  Future<void> close() {
    if (!asked.isCompleted) asked.complete();
    return closing ??= () async {
      await world?._close();
      send(WorldMessage.closed);
    }();
  }

  var subscription = input.listen(
    (line) {
      var message = decodeWorldMessage(line);
      switch (message['type']) {
        case WorldMessage.open when world == null:
          var knobs = (message['knobs'] as Map? ?? const {})
              .cast<String, Object?>();
          var opened = world = World._(_newId(), knobs, send);
          send(WorldMessage.hello, {
            'protocol': worldProtocolVersion,
            'id': opened.id,
            'pid': pid,
          });
          unawaited(opened._setUp(body));
        case WorldMessage.ready:
          if (!(world?._ready.isCompleted ?? true)) world!._ready.complete();
        case WorldMessage.invoke:
          world?._invoke(
            message['action']! as String,
            message['run']! as int,
            step: message['step'] as String?,
          );
        case WorldMessage.cancel:
          world?._runs[message['run']]?._cancel();
        case WorldMessage.close:
          unawaited(close());
      }
    },
    // The owner went away without closing: close anyway, so whatever the
    // script started is stopped rather than orphaned.
    onDone: () => unawaited(close()),
  );
  await asked.future;
  await closing;
  await subscription.cancel();
}