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
20 changes: 10 additions & 10 deletions cpp/obs/src/moq-settings.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -196,16 +196,16 @@ const DefaultValues &LibraryDefaults()
return d;
}

d.connect_timeout_ms = (long long)config.connect_timeout_ms;
d.failover_delay_ms = (long long)config.failover_delay_ms;
d.connect_timeout_ms = (long long)(config.connect_timeout_us / 1000);
d.failover_delay_ms = (long long)(config.failover_delay_us / 1000);
d.backoff_initial_ms = (long long)(config.backoff_initial_us / 1000);
d.backoff_max_ms = (long long)(config.backoff_max_us / 1000);
d.backoff_timeout_ms = (long long)(config.backoff_timeout_us / 1000);
d.quic_max_streams = (long long)config.quic_max_streams;
d.quic_idle_timeout_ms = (long long)config.quic_idle_timeout_ms;
d.quic_idle_timeout_ms = (long long)(config.quic_idle_timeout_us / 1000);
// Absent means "no keep-alive", which the UI shows as zero.
d.quic_keep_alive_ms = config.has_quic_keep_alive ? (long long)config.quic_keep_alive_ms : 0;
d.websocket_delay_ms = config.has_websocket_delay ? (long long)config.websocket_delay_ms : 0;
d.quic_keep_alive_ms = config.has_quic_keep_alive ? (long long)(config.quic_keep_alive_us / 1000) : 0;
d.websocket_delay_ms = config.has_websocket_delay ? (long long)(config.websocket_delay_us / 1000) : 0;
d.websocket_enabled = config.websocket_enabled;
d.loaded = true;
return d;
Expand Down Expand Up @@ -403,9 +403,9 @@ bool BuildConfig(obs_data_t *settings, Config *out)
borrow(OptionalString(settings, BACKEND), &out->backend, &config.backend, &config.backend_len);
borrow(OptionalString(settings, BIND), &out->bind, &config.bind, &config.bind_len);

config.connect_timeout_ms = (uint64_t)Amount(settings, CONNECT_TIMEOUT);
config.connect_timeout_us = (uint64_t)Amount(settings, CONNECT_TIMEOUT) * 1000;
config.has_connect_timeout = true;
config.failover_delay_ms = (uint64_t)Amount(settings, FAILOVER_DELAY);
config.failover_delay_us = (uint64_t)Amount(settings, FAILOVER_DELAY) * 1000;
config.has_failover_delay = true;

config.tls_disable_verify = obs_data_get_bool(settings, TLS_DISABLE_VERIFY);
Expand Down Expand Up @@ -436,9 +436,9 @@ bool BuildConfig(obs_data_t *settings, Config *out)

config.quic_max_streams = (uint64_t)Amount(settings, QUIC_MAX_STREAMS);
config.has_quic_max_streams = true;
config.quic_idle_timeout_ms = (uint64_t)Amount(settings, QUIC_IDLE_TIMEOUT);
config.quic_idle_timeout_us = (uint64_t)Amount(settings, QUIC_IDLE_TIMEOUT) * 1000;
config.has_quic_idle_timeout = true;
config.quic_keep_alive_ms = (uint64_t)Amount(settings, QUIC_KEEP_ALIVE);
config.quic_keep_alive_us = (uint64_t)Amount(settings, QUIC_KEEP_ALIVE) * 1000;
config.has_quic_keep_alive = true;

// The tri-states stay unset when the user left them on Automatic, which is what
Expand All @@ -465,7 +465,7 @@ bool BuildConfig(obs_data_t *settings, Config *out)

config.websocket_enabled = obs_data_get_bool(settings, WEBSOCKET_ENABLED);
config.has_websocket_enabled = true;
config.websocket_delay_ms = (uint64_t)Amount(settings, WEBSOCKET_DELAY);
config.websocket_delay_us = (uint64_t)Amount(settings, WEBSOCKET_DELAY) * 1000;
config.has_websocket_delay = true;

return true;
Expand Down
26 changes: 13 additions & 13 deletions cpp/obs/src/moq-source.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -225,7 +225,7 @@ struct subscription_ref {
subscription_ref &operator=(const subscription_ref &) = delete;
};

// user_data for a single moq_origin_consume_announced. The generation must travel
// user_data for a single moq_origin_announced_broadcast. The generation must travel
// with the request rather than live on ctx: a reconnect can issue a new request while
// an older one still has a delivery in flight, and a single slot on ctx would let
// that stale delivery read the new generation and pass the staleness check.
Expand Down Expand Up @@ -686,14 +686,14 @@ static void moq_source_subscribe_video(struct moq_source *ctx, int32_t catalog,
ctx->video_track = track;
pthread_mutex_unlock(&ctx->mutex);
if (old_track >= 0)
moq_consume_video_close(old_track);
moq_consume_video_cancel(old_track);
LOG_INFO("Subscribed to video track successfully");
} else {
// Stale or shutting down: close the track we just created; its terminal
// callback releases the reference added above.
pthread_mutex_unlock(&ctx->mutex);
if (!video_state->terminal.load())
moq_consume_video_close(track);
moq_consume_video_cancel(track);
}
}

Expand Down Expand Up @@ -794,12 +794,12 @@ static void moq_source_subscribe_audio(struct moq_source *ctx, int32_t catalog,
ctx->audio_track = track;
pthread_mutex_unlock(&ctx->mutex);
if (old_track >= 0)
moq_consume_audio_close(old_track);
moq_consume_audio_cancel(old_track);
LOG_INFO("Subscribed to audio track successfully (%u Hz, %u ch)", sample_rate, channels);
} else {
pthread_mutex_unlock(&ctx->mutex);
if (!audio_state->terminal.load())
moq_consume_audio_close(track);
moq_consume_audio_cancel(track);
}
}

Expand Down Expand Up @@ -1013,7 +1013,7 @@ static void moq_source_start_consume(struct moq_source *ctx, uint32_t expected_g
// so it need not outlive this call, and delivers the broadcast handle
// asynchronously to on_broadcast.
int32_t request =
moq_origin_consume_announced(origin, broadcast_copy, strlen(broadcast_copy), on_broadcast, req);
moq_origin_announced_broadcast(origin, broadcast_copy, strlen(broadcast_copy), on_broadcast, req);
if (request < 0) {
LOG_ERROR("Failed to request broadcast '%s': %d", broadcast_copy, request);
bfree(broadcast_copy);
Expand Down Expand Up @@ -1041,13 +1041,13 @@ static void moq_source_start_consume(struct moq_source *ctx, uint32_t expected_g
} else {
// Stale or shutting down: close it; its terminal releases the reference.
pthread_mutex_unlock(&ctx->mutex);
moq_origin_consume_announced_close(request);
moq_origin_announced_broadcast_cancel(request);
}
}

// Receives the announced broadcast: a positive handle once announced, then exactly
// once more with a terminal code (0 = finished, including after
// moq_origin_consume_announced_close; < 0 = error). The terminal is the last touch
// moq_origin_announced_broadcast_cancel; < 0 = error). The terminal is the last touch
// of user_data, so it both frees the request context and releases the request's
// lifetime reference via subscription_ref.
static void on_broadcast(void *user_data, int32_t broadcast)
Expand Down Expand Up @@ -1137,7 +1137,7 @@ static void on_broadcast(void *user_data, int32_t broadcast)
// Stale or shutting down: close it; its terminal releases the reference.
pthread_mutex_unlock(&ctx->mutex);
if (!state->terminal.load())
moq_consume_catalog_close(catalog_handle);
moq_consume_catalog_cancel(catalog_handle);
}
}

Expand All @@ -1152,7 +1152,7 @@ static void moq_source_clear_video_locked(struct moq_source *ctx)
{
ctx->video_attempt++;
if (ctx->video_track >= 0) {
moq_consume_video_close(ctx->video_track);
moq_consume_video_cancel(ctx->video_track);
ctx->video_track = -1;
}
moq_source_destroy_decoder_locked(ctx);
Expand All @@ -1170,15 +1170,15 @@ static void moq_source_disconnect_locked(struct moq_source *ctx)
moq_source_clear_audio_locked(ctx);

if (ctx->catalog_handle >= 0) {
moq_consume_catalog_close(ctx->catalog_handle);
moq_consume_catalog_cancel(ctx->catalog_handle);
ctx->catalog_handle = -1;
}

// An unresolved wait still owes a terminal on_broadcast; closing it makes that
// fire (with 0) instead of leaving it pending until the source dies. This is the
// path that ends a wait for a broadcast that is never announced.
if (ctx->request >= 0) {
moq_origin_consume_announced_close(ctx->request);
moq_origin_announced_broadcast_cancel(ctx->request);
ctx->request = -1;
}

Expand Down Expand Up @@ -1636,7 +1636,7 @@ static void moq_source_clear_audio_locked(struct moq_source *ctx)
{
ctx->audio_attempt++;
if (ctx->audio_track >= 0) {
moq_consume_audio_close(ctx->audio_track);
moq_consume_audio_cancel(ctx->audio_track);
ctx->audio_track = -1;
}
moq_source_destroy_audio_decoder_locked(ctx);
Expand Down
14 changes: 7 additions & 7 deletions cpp/obs/test/moq-source-test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -557,7 +557,7 @@ std::atomic<int> g_origin_closes{0};
std::atomic<int> g_session_connects{0};
std::atomic<int> g_announced_calls{0};
// moq_origin_request resolves only broadcasts that are already announced. The
// source must wait with moq_origin_consume_announced, so this stays zero.
// source must wait with moq_origin_announced_broadcast, so this stays zero.
std::atomic<int> g_request_calls{0};
std::atomic<int> g_catalog_calls{0};
std::atomic<int> g_video_calls{0};
Expand Down Expand Up @@ -735,8 +735,8 @@ int32_t moq_session_close(uint32_t session)
return closeSub(static_cast<int32_t>(session));
}

int32_t moq_origin_consume_announced(uint32_t, const char *, uintptr_t, void (*on_broadcast)(void *, int32_t),
void *user_data)
int32_t moq_origin_announced_broadcast(uint32_t, const char *, uintptr_t, void (*on_broadcast)(void *, int32_t),
void *user_data)
{
if (g_announced_result < 0)
return g_announced_result;
Expand All @@ -746,7 +746,7 @@ int32_t moq_origin_consume_announced(uint32_t, const char *, uintptr_t, void (*o
return handle;
}

int32_t moq_origin_consume_announced_close(uint32_t task)
int32_t moq_origin_announced_broadcast_cancel(uint32_t task)
{
return closeSub(static_cast<int32_t>(task));
}
Expand All @@ -769,7 +769,7 @@ int32_t moq_consume_catalog(uint32_t, void (*on_catalog)(void *, int32_t), void
return handle;
}

int32_t moq_consume_catalog_close(uint32_t catalog)
int32_t moq_consume_catalog_cancel(uint32_t catalog)
{
return closeSub(static_cast<int32_t>(catalog));
}
Expand Down Expand Up @@ -832,7 +832,7 @@ int32_t moq_consume_video(uint32_t catalog, uint32_t, uint64_t, void (*on_frame)
return handle;
}

int32_t moq_consume_video_close(uint32_t track)
int32_t moq_consume_video_cancel(uint32_t track)
{
return closeSub(static_cast<int32_t>(track));
}
Expand Down Expand Up @@ -882,7 +882,7 @@ int32_t moq_consume_audio(uint32_t catalog, uint32_t, uint64_t, void (*on_frame)
return handle;
}

int32_t moq_consume_audio_close(uint32_t track)
int32_t moq_consume_audio_cancel(uint32_t track)
{
g_audio_closes++;
return closeSub(static_cast<int32_t>(track));
Expand Down
3 changes: 2 additions & 1 deletion dart/moq/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,8 @@ import 'package:moq/moq.dart';

final connection = await Moq.connect('https://relay.example.com');
await for (final announcement in connection.announcements()) {
print(announcement.path());
// The covered prefix is relative to the requested announcements prefix.
print(announcement.prefix());
}
```

Expand Down
6 changes: 3 additions & 3 deletions dart/moq/lib/moq.dart
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ final class Moq {
}) async {
final client = MoqClient();
try {
if (!tlsVerify) client.setTlsDisableVerify(disable: true);
if (!tlsVerify) client.setTlsVerify(verify: false);
if (tlsRoots != null) client.setTlsRoots(paths: tlsRoots);
if (tlsSystemRoots != null) {
client.setTlsSystemRoots(systemRoots: tlsSystemRoots);
Expand Down Expand Up @@ -60,7 +60,7 @@ final class Moq {
MoqBroadcastProducer createBroadcast(String path) =>
session.publish().createBroadcast(path: path);

/// Stream announcements whose paths begin with [prefix].
/// Stream routes under requested [prefix]; updates return relative covered prefixes.
Stream<MoqAnnounceUpdate> announcements({String prefix = ''}) async* {
final announced = session.consume().announced(prefix: prefix);
try {
Expand All @@ -75,7 +75,7 @@ final class Moq {
}
}

/// Return the raw announcement cursor for [prefix].
/// Return the raw cursor for requested [prefix]; updates return relative covered prefixes.
MoqAnnounceConsumer announced({String prefix = ''}) =>
session.consume().announced(prefix: prefix);

Expand Down
8 changes: 4 additions & 4 deletions dart/moq/test/moq_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -34,10 +34,10 @@ void main() {
final track = broadcast.publishTrack(name: 'events', info: null);
broadcast.announce(route: MoqRoute());
final announced = await announcement.timeout(timeout);
expect(announced.path(), 'live');
expect(announced.prefix(), 'live');

final requested = await client
.requestBroadcast(announced.path())
.requestBroadcast(announced.prefix())
.timeout(timeout);
final consumer = await requested
.subscribeTrack(name: 'events', subscription: null)
Expand Down Expand Up @@ -71,12 +71,12 @@ void main() {

final announced = origin.consume().announced(prefix: '');
final first = await announced.next().timeout(timeout);
expect(first?.path(), 'live');
expect(first?.prefix(), 'live');
expect(first?.active(), isTrue);

broadcast.unannounce();
final retracted = await announced.next().timeout(timeout);
expect(retracted?.path(), 'live');
expect(retracted?.prefix(), 'live');
expect(retracted?.active(), isFalse);
await origin.consume().requestBroadcast(path: 'live').timeout(timeout);
announced.cancel();
Expand Down
Loading
Loading