Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions dart/moq/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,10 @@ Idiomatic Dart and Flutter bindings for Media over QUIC.
import 'package:moq/moq.dart';

final connection = await Moq.connect('https://relay.example.com');
await for (final event in connection.announcements(
final announced = connection.announced(
options: const AnnounceOptions(prefix: 'live/', filter: '*/camera'),
)) {
);
await for (final event in announced.updates()) {
if (event is AnnounceEventStart) {
// Prefix stays origin-relative; captures reports wildcard matches.
print(event.announce.prefix);
Expand Down
142 changes: 73 additions & 69 deletions dart/moq/lib/src/client.dart
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import 'package:moq_ffi/moq_ffi.dart';

import 'aliases.dart';

/// Scope for discovering announcements.
Expand All @@ -18,6 +20,9 @@ final class AnnounceOptions {
}

/// Everything [Moq.connect] can be told beyond the URL.
///
/// Every field is optional; a null keeps the native default. A value the native
/// side cannot use fails [Moq.connect] with `MoqException.Config`.
final class ConnectOptions {
/// Set false to skip certificate verification (local dev only).
final bool tlsVerify;
Expand All @@ -40,6 +45,10 @@ final class ConnectOptions {
/// Local socket address to bind, e.g. `0.0.0.0:0`.
final String? bind;

/// Protocol versions to offer, most preferred first, e.g. `moq-lite-03`.
/// Null offers every supported version.
final List<String>? versions;

/// Cap on the concurrent QUIC streams the peer may open toward this
/// connection. MoQ opens one stream per group, so a subscriber to many
/// tracks may want this raised.
Expand All @@ -55,7 +64,7 @@ final class ConnectOptions {

/// Set false for a one-shot dial. By default the session redials with
/// backoff whenever the transport drops.
final bool? reconnect;
final bool reconnect;

/// Retry pacing for the automatic reconnect.
final Backoff? backoff;
Expand All @@ -64,7 +73,7 @@ final class ConnectOptions {
final OriginProducer? publish;

/// Origin to discover broadcasts through; auto-created when null.
final OriginProducer? subscribe;
final OriginProducer? consume;

const ConnectOptions({
this.tlsVerify = true,
Expand All @@ -74,14 +83,47 @@ final class ConnectOptions {
this.tlsCert,
this.tlsKey,
this.bind,
this.versions,
this.maxStreams,
this.websocketEnabled,
this.websocketDelay,
this.reconnect,
this.reconnect = true,
this.backoff,
this.publish,
this.subscribe,
this.consume,
});

MoqClientConfig get _ffi {
final delay = websocketDelay;
if (delay != null && delay.isNegative) {
throw ArgumentError.value(
delay,
'websocketDelay',
'must not be negative',
);
}
return MoqClientConfig(
bind: bind,
versions: versions ?? const [],
tls: MoqClientTls(
insecure: !tlsVerify,
roots: tlsRoots ?? const [],
systemRoots: tlsSystemRoots,
fingerprints: tlsFingerprints ?? const [],
cert: tlsCert,
key: tlsKey,
),
quic: MoqQuicConfig(maxStreams: maxStreams),
websocket: MoqWebSocketConfig(
enabled: websocketEnabled,
delayUs: delay?.inMicroseconds,
),
once: !reconnect,
backoff: backoff ?? MoqBackoff(),
publish: publish,
consume: consume,
);
}
}

/// A connected MoQ session with publishing and subscription conveniences.
Expand All @@ -95,58 +137,16 @@ final class Moq {

/// Connect to a relay at [url].
///
/// With neither [ConnectOptions.publish] nor [ConnectOptions.subscribe]
/// With neither [ConnectOptions.publish] nor [ConnectOptions.consume]
/// given, both sides of the session share one origin, so a broadcast
/// announced here is discoverable through [announcements]. Wiring either
/// announced here is discoverable through [announced]. Wiring either
/// side opts out and isolates the two directions.
static Future<Moq> connect(
String url, {
ConnectOptions options = const ConnectOptions(),
}) async {
final websocketDelay = options.websocketDelay;
if (websocketDelay != null && websocketDelay.isNegative) {
throw ArgumentError.value(
websocketDelay,
'websocketDelay',
'must not be negative',
);
}

final client = Client();
final client = Client(config: options._ffi);
try {
if (!options.tlsVerify) client.setTlsVerify(verify: false);
if (options.tlsRoots != null) {
client.setTlsRoots(paths: options.tlsRoots!);
}
if (options.tlsSystemRoots != null) {
client.setTlsSystemRoots(systemRoots: options.tlsSystemRoots!);
}
if (options.tlsFingerprints != null) {
client.setTlsFingerprints(fingerprints: options.tlsFingerprints!);
}
if (options.tlsCert != null) client.setTlsCert(path: options.tlsCert);
if (options.tlsKey != null) client.setTlsKey(path: options.tlsKey);
if (options.bind != null) client.setBind(addr: options.bind!);
if (options.maxStreams != null) {
client.setQuicMaxStreams(maxStreams: options.maxStreams!);
}
if (options.websocketEnabled != null) {
client.setWebsocketEnabled(enabled: options.websocketEnabled!);
}
if (websocketDelay != null) {
client.setWebsocketDelay(delayUs: websocketDelay.inMicroseconds);
}
if (options.reconnect != null) {
client.setReconnect(enabled: options.reconnect!);
}
if (options.backoff != null) {
client.setBackoff(backoff: options.backoff!);
}
if (options.publish != null) client.setPublish(origin: options.publish);
if (options.subscribe != null) {
client.setConsume(origin: options.subscribe);
}

final session = await client.connect(url: url);
return Moq._(session, client);
} catch (_) {
Expand All @@ -162,27 +162,10 @@ final class Moq {
BroadcastProducer createBroadcast(String path) =>
session.publish().createBroadcast(path: path);

/// Stream announce events matching [options]; prefixes stay relative to the origin.
/// Discover routes matching [options]; prefixes stay relative to the origin.
///
/// A [AnnounceEventLive] follows the routes live at subscribe time, so a
/// listener can collect what is live and stop there.
Stream<AnnounceEvent> announcements({
AnnounceOptions options = const AnnounceOptions(),
}) async* {
final announced = session.consume().announced(config: options._ffi);
try {
while (true) {
final event = await announced.next();
if (event == null) return;
yield event;
}
} finally {
announced.cancel();
announced.dispose();
}
}

/// Return the raw cursor for [options].
/// Listen to `announced(...).updates()` for a [Stream] of [AnnounceEvent]
/// that releases the cursor when the subscription ends.
AnnounceConsumer announced({
AnnounceOptions options = const AnnounceOptions(),
}) => session.consume().announced(config: options._ffi);
Expand Down Expand Up @@ -212,3 +195,24 @@ final class Moq {
_client.cancel();
}
}

/// A [Stream] view over an announcement cursor.
extension AnnounceConsumerUpdates on AnnounceConsumer {
/// Stream announce events until the cursor ends.
///
/// A [AnnounceEventLive] follows the routes live at subscribe time, so a
/// listener can collect what is live and stop there. Listen once: the cursor
/// is cancelled and released when the subscription ends.
Stream<AnnounceEvent> updates() async* {
try {
while (true) {
final event = await next();
if (event == null) return;
yield event;
}
} finally {
cancel();
dispose();
}
}
}
15 changes: 9 additions & 6 deletions dart/moq/lib/src/durations.dart
Original file line number Diff line number Diff line change
Expand Up @@ -13,14 +13,17 @@ extension ConnectionStatsDuration on MoqConnectionStats {

/// Duration views over the reconnect pacing.
extension BackoffDuration on MoqBackoff {
/// Delay before the first reconnect attempt.
Duration get initial => Duration(microseconds: initialUs);
/// Delay before the first reconnect attempt, or null for the default.
Duration? get initial =>
initialUs == null ? null : Duration(microseconds: initialUs!);

/// Maximum delay between reconnect attempts.
Duration get max => Duration(microseconds: maxUs);
/// Maximum delay between reconnect attempts, or null for the default.
Duration? get max => maxUs == null ? null : Duration(microseconds: maxUs!);

/// Time spent retrying before giving up. [Duration.zero] retries forever.
Duration get timeout => Duration(microseconds: timeoutUs);
/// Time spent retrying before giving up, or null for the default.
/// [Duration.zero] retries forever.
Duration? get timeout =>
timeoutUs == null ? null : Duration(microseconds: timeoutUs!);
}

/// Duration views over the subscription knobs.
Expand Down
47 changes: 31 additions & 16 deletions dart/moq/lib/src/server.dart
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,9 @@ import 'package:moq_ffi/moq_ffi.dart';
import 'aliases.dart';

/// Everything [Server.listen] can be told.
///
/// A value the native side cannot use fails [Server.listen] with
/// `MoqException.Config`.
final class ListenOptions {
/// Local socket address to listen on, e.g. `127.0.0.1:4443` or `[::]:443`.
final String bind;
Expand All @@ -16,19 +19,28 @@ final class ListenOptions {
/// Hostnames to generate a self-signed certificate for.
final List<String>? tlsGenerate;

/// Protocol versions to accept, e.g. `moq-lite-03`. Null accepts every
/// supported version.
final List<String>? versions;

/// Cap on the concurrent QUIC streams each peer may open toward this server.
final int? maxStreams;

/// Origin whose broadcasts are served to incoming sessions; auto-created when null.
final OriginProducer? publish;

/// Origin that receives broadcasts published by incoming sessions; auto-created when null.
final OriginProducer? subscribe;
final OriginProducer? consume;

const ListenOptions({
this.bind = '[::]:443',
this.tlsCert,
this.tlsKey,
this.tlsGenerate,
this.versions,
this.maxStreams,
this.publish,
this.subscribe,
this.consume,
});
}

Expand All @@ -52,30 +64,33 @@ final class Server {

/// Bind a server at [ListenOptions.bind] and start accepting.
///
/// With neither [ListenOptions.publish] nor [ListenOptions.subscribe] given,
/// With neither [ListenOptions.publish] nor [ListenOptions.consume] given,
/// both sides share one origin, so a broadcast created here is also visible
/// to sessions publishing into this server. Wiring either side opts out and
/// isolates the two directions.
static Future<Server> listen({
ListenOptions options = const ListenOptions(),
}) async {
final shared = options.publish == null && options.subscribe == null
final shared = options.publish == null && options.consume == null
? OriginProducer(config: OriginConfig())
: null;
final publishOrigin = options.publish ?? shared;
final subscribeOrigin = options.subscribe ?? shared;

final server = MoqServer();
final server = MoqServer(
config: MoqServerConfig(
bind: options.bind,
versions: options.versions ?? const [],
tls: MoqServerTls(
cert: options.tlsCert ?? const [],
key: options.tlsKey ?? const [],
generate: options.tlsGenerate ?? const [],
),
quic: MoqQuicConfig(maxStreams: options.maxStreams),
publish: publishOrigin,
consume: options.consume ?? shared,
),
);
try {
server.setBind(addr: options.bind);
if (options.tlsCert != null) server.setTlsCert(paths: options.tlsCert!);
if (options.tlsKey != null) server.setTlsKey(paths: options.tlsKey!);
if (options.tlsGenerate != null) {
server.setTlsGenerate(hostnames: options.tlsGenerate!);
}
if (publishOrigin != null) server.setPublish(origin: publishOrigin);
if (subscribeOrigin != null) server.setConsume(origin: subscribeOrigin);

final localAddr = await server.listen();
return Server._(server, localAddr, publishOrigin);
} catch (_) {
Expand All @@ -88,7 +103,7 @@ final class Server {
/// Create a broadcast at [path], served to incoming sessions.
///
/// Advertise it with `announce` after populating tracks. Throws when [listen]
/// was given a [ListenOptions.subscribe] origin but no
/// was given a [ListenOptions.consume] origin but no
/// [ListenOptions.publish] one, since there is then nothing to serve from.
BroadcastProducer createBroadcast(String path) {
final origin = _publishOrigin;
Expand Down
Loading
Loading