serve static method
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();
}