diff --git a/cpp/obs/src/moq-output.cpp b/cpp/obs/src/moq-output.cpp index bbe39e6879..7a5e071951 100644 --- a/cpp/obs/src/moq-output.cpp +++ b/cpp/obs/src/moq-output.cpp @@ -301,10 +301,10 @@ void MoQOutput::Reset() } audio_tracks.clear(); - // Finish the broadcast so the origin unpublishes it immediately; Start() + // Close the broadcast so the origin retracts it immediately; Start() // creates a fresh one on restart. if (broadcast > 0) { - moq_publish_finish(broadcast); + moq_publish_close(broadcast); broadcast = 0; } } diff --git a/cpp/obs/test/moq-output-test.cpp b/cpp/obs/test/moq-output-test.cpp index e28039ea60..c20c9d30c8 100644 --- a/cpp/obs/test/moq-output-test.cpp +++ b/cpp/obs/test/moq-output-test.cpp @@ -233,7 +233,7 @@ int32_t moq_publish_announce(uint32_t, const moq_route *) return 0; } -int32_t moq_publish_finish(uint32_t) +int32_t moq_publish_close(uint32_t) { return 0; } diff --git a/dart/moq/test/moq_test.dart b/dart/moq/test/moq_test.dart index 8fcc824ab2..278456ee19 100644 --- a/dart/moq/test/moq_test.dart +++ b/dart/moq/test/moq_test.dart @@ -102,7 +102,8 @@ void main() { client.close(); serverSession.cancel(code: 0); track.finish(); - broadcast.finish(); + broadcast.close(); + broadcast.close(); // a second close is a no-op server.close(); }); diff --git a/dart/moq_ffi/lib/src/moq.dart b/dart/moq_ffi/lib/src/moq.dart index d4635f31c0..9dc5359082 100644 --- a/dart/moq_ffi/lib/src/moq.dart +++ b/dart/moq_ffi/lib/src/moq.dart @@ -8,9 +8,7 @@ import "dart:ffi"; import "dart:io" show Platform, File, Directory; import "dart:isolate"; import "dart:typed_data"; - import "package:ffi/ffi.dart"; - import "uniffi_runtime.dart"; export "uniffi_runtime.dart"; @@ -5900,6 +5898,7 @@ abstract class MoqBroadcastProducerInterface { required MoqJsonStreamConfig config, }); void announce({required MoqRoute route}); + void close(); MoqBroadcastConsumer consume(); MoqBroadcastDynamic dynamic_(); void finish(); @@ -6007,6 +6006,15 @@ class MoqBroadcastProducer implements MoqBroadcastProducerInterface { }, moqExceptionErrorHandler); } + void close() { + return rustCall((status) { + uniffi_moq_ffi_fn_method_moqbroadcastproducer_close( + uniffiClonePointer(), + status, + ); + }, moqExceptionErrorHandler); + } + MoqBroadcastConsumer consume() { return rustCallWithLifter( (status) => uniffi_moq_ffi_fn_method_moqbroadcastproducer_consume( @@ -10173,6 +10181,14 @@ external void uniffi_moq_ffi_fn_method_moqbroadcastproducer_announce( Pointer uniffiStatus, ); +@Native, Pointer)>( + assetId: _uniffiAssetId, +) +external void uniffi_moq_ffi_fn_method_moqbroadcastproducer_close( + Pointer ptr, + Pointer uniffiStatus, +); + @Native Function(Pointer, Pointer)>( assetId: _uniffiAssetId, ) @@ -11825,6 +11841,9 @@ uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_json_stream(); @Native(assetId: _uniffiAssetId) external int uniffi_moq_ffi_checksum_method_moqbroadcastproducer_announce(); +@Native(assetId: _uniffiAssetId) +external int uniffi_moq_ffi_checksum_method_moqbroadcastproducer_close(); + @Native(assetId: _uniffiAssetId) external int uniffi_moq_ffi_checksum_method_moqbroadcastproducer_consume(); @@ -12335,7 +12354,7 @@ void _checkApiChecksums() { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } if (uniffi_moq_ffi_checksum_method_moqannouncedbroadcast_available() != - 42497) { + 37458) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } if (uniffi_moq_ffi_checksum_method_moqannouncedbroadcast_cancel() != 63175) { @@ -12399,13 +12418,16 @@ void _checkApiChecksums() { if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_announce() != 13700) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } + if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_close() != 19191) { + throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); + } if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_consume() != 27634) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_dynamic() != 55635) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } - if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_finish() != 7183) { + if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_finish() != 29562) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_audio() != diff --git a/doc/lib/c/index.md b/doc/lib/c/index.md index b77ae8a4bd..5bb5aaa84b 100644 --- a/doc/lib/c/index.md +++ b/doc/lib/c/index.md @@ -42,7 +42,7 @@ and `target/include/moq.h`. - **Server.** `moq_server_listen` binds before it returns (a bad address or certificate fails there) and hands each incoming session to `on_request` as a request handle. Read `moq_session_request_path` and `_query` to route and authenticate, then `moq_session_request_accept` (a session handle, with origins like `moq_session_connect`) or `moq_session_request_reject` with an HTTP-style code (401 and 403 become the protocol's unauthorized close). An accepted session reports `1` once SETUP completes and never reconnects. `moq_server_addr` reports an ephemeral port and `moq_server_fingerprints` the hashes a client pins for a `tls_generate` certificate. `moq_server_close` stops listening; its terminal callback fires once the sockets are released. - **Demand.** A watcher on a published track (`moq_publish_track_demand`, `moq_publish_media_demand`, `moq_encode_video_demand`, `moq_encode_audio_demand`) calls `on_demand` with `MOQ_DEMAND_USED` or `MOQ_DEMAND_UNUSED` right away and again on every change, so an encoder on a battery-powered device runs only while someone is watching. The first call is the current state, so a track that went unused before the watcher existed still reports it. `moq_publish_demand_cancel` stops it; the terminal callback still fires. A container has no single demand and is refused. Demand follows the last real subscriber: an origin that served the track drops its source copy on the unused edge and keeps only the finished groups it already cached warm for 30 seconds, so the cache linger does not delay the unused edge. - **Requests.** `moq_publish_dynamic` serves subscriptions to tracks the broadcast never declared: each arrives as a request handle, read its name with `moq_track_request_name`, then `moq_track_request_accept` (a raw track handle), `moq_track_request_video` / `_audio` (the media handle `moq_publish_video` / `_audio` return), or `moq_track_request_abort` with an application code the subscriber sees. Without a live handler an unknown name is refused. `moq_publish_track_dynamic` does the same for fetches of groups a track no longer has cached, delivered as `moq_group_request_*` (`sequence`, `priority`, `frame_start`); `moq_group_request_accept` starts the producer at `frame_start` so written frames keep their group indices. Register it with `moq_track_request_dynamic` before accepting a track that was itself requested by a fetch, so that pending group survives the transition. Both handlers stop with `moq_publish_dynamic_cancel`. -- **Everything the bindings can do** ([list](/lib/#what-every-binding-can-do)): media publish and consume with the catalog managed for you, raw pixels and PCM with the codec inside (`moq_encode_video`, `moq_encode_audio`, and the `moq_decode_*` mirrors), raw tracks with timestamps and datagrams, JSON and binary data tracks (snapshot or stream, each advertised in the catalog for as long as it lives), group fetch, catalog sections, shared video properties, and stalled hints. The three advertising operations are `moq_origin_create_broadcast` (unannounced producer, invisible to everyone), `moq_publish_announce` / `moq_publish_unannounce` (exact-path advertisement), and `moq_origin_dynamic` (a claim over a path prefix and everything beneath it; `""` for everything). A route is a capability, not an inventory. `moq_origin_announced` takes a literal prefix and an optional relative pattern filter; `moq_announce_update.prefix` stays relative to the origin, while `captures` reports what each wildcard matched when `has_captures` is true. Paths with a `.`-prefixed segment below the prefix are [hidden](/concept/moq-lite#hidden-broadcasts); name the dot segment in `prefix` to list them. +- **Everything the bindings can do** ([list](/lib/#what-every-binding-can-do)): media publish and consume with the catalog managed for you, raw pixels and PCM with the codec inside (`moq_encode_video`, `moq_encode_audio`, and the `moq_decode_*` mirrors), raw tracks with timestamps and datagrams, JSON and binary data tracks (snapshot or stream, each advertised in the catalog for as long as it lives), group fetch, catalog sections, shared video properties, and stalled hints. The three advertising operations are `moq_origin_create_broadcast` (unannounced producer, invisible to everyone), `moq_publish_announce` / `moq_publish_unannounce` (exact-path advertisement), `moq_publish_close` (ends the broadcast for good and releases its handle; `moq_publish_finish` is its deprecated alias), and `moq_origin_dynamic` (a claim over a path prefix and everything beneath it; `""` for everything). A route is a capability, not an inventory. `moq_origin_announced` takes a literal prefix and an optional relative pattern filter; `moq_announce_update.prefix` stays relative to the origin, while `captures` reports what each wildcard matched when `has_captures` is true. Paths with a `.`-prefixed segment below the prefix are [hidden](/concept/moq-lite#hidden-broadcasts); name the dot segment in `prefix` to list them. ```c moq_client_config config; diff --git a/doc/lib/dart/index.md b/doc/lib/dart/index.md index 564d4e06ed..d7d31ee23c 100644 --- a/doc/lib/dart/index.md +++ b/doc/lib/dart/index.md @@ -65,7 +65,8 @@ await for (final request in server.requests()) { The three advertising operations: `moq.createBroadcast(path)` (or `origin.createBroadcast`) returns an unannounced producer, invisible to everyone; `broadcast.announce(route:)` / `broadcast.unannounce()` own that exact-path -advertisement; `origin.dynamic_(prefix:, route:)` claims `prefix` and +advertisement, and `broadcast.close()` ends the broadcast for good (a second +call is a no-op; `finish()` is its deprecated alias); `origin.dynamic_(prefix:, route:)` claims `prefix` and every path beneath it (`''` for everything; Dart spells the origin method `dynamic_` because `dynamic` is reserved). Hold the returned handle while the claim should stay advertised, and reject the requests you will not serve. A diff --git a/doc/lib/go/index.md b/doc/lib/go/index.md index 001f5a50e4..5b3ffe9e47 100644 --- a/doc/lib/go/index.md +++ b/doc/lib/go/index.md @@ -66,7 +66,7 @@ video, _ := broadcast.EncodeVideo( ) _ = video.Write(moq.VideoFrame{TimestampUs: pts, Data: rgba}) _ = broadcast.Announce(moq.Route{}) -broadcast.Finish() // keep the producer reachable while publishing, then finish explicitly +broadcast.Close() // keep the producer reachable while publishing, then close explicitly ``` For locally encoded media, call `MediaProducer.Flush(timestampUs)` after `WriteFrame` with the same broadcast-clock PTS. It measures catalog jitter at the transport handoff. File, pipe, and network imports should omit `Flush`; built-in encoders observe their own output. @@ -74,7 +74,8 @@ For locally encoded media, call `MediaProducer.Flush(timestampUs)` after `WriteF The three advertising operations: `client.CreateBroadcast(path)` (or `origin.CreateBroadcast`) returns an unannounced producer, invisible to everyone; `broadcast.Announce(route)` / `broadcast.Unannounce()` own that exact-path -advertisement; `origin.Dynamic(prefix, route)` claims `prefix` and every +advertisement, and `broadcast.Close()` ends the broadcast for good (a second +call is a no-op; `Finish` is its deprecated alias); `origin.Dynamic(prefix, route)` claims `prefix` and every path beneath it (`""` for everything). Hold the returned `OriginDynamic` while the claim should stay advertised, and reject the requests you will not serve. A route is a capability, not an inventory. `Announced(options)` combines diff --git a/doc/lib/kt/index.md b/doc/lib/kt/index.md index a241ad6b0c..2b8e9d5795 100644 --- a/doc/lib/kt/index.md +++ b/doc/lib/kt/index.md @@ -56,7 +56,8 @@ Moq.connect("https://relay.example.com").use { moq -> The three advertising operations: `moq.createBroadcast(path)` (or `origin.createBroadcast`) returns an unannounced producer, invisible to everyone; `broadcast.announce(route)` / `broadcast.unannounce()` own that exact-path -advertisement; `origin.dynamic(prefix, route)` claims `prefix` and every +advertisement, and `broadcast.close()` (or `use { }`) releases the producer, +ending the broadcast once no `dynamic()` handle remains; `origin.dynamic(prefix, route)` claims `prefix` and every path beneath it (`""` for everything). Hold the returned `OriginDynamic` while the claim should stay advertised, and reject the requests you will not serve. A route is a capability, not an inventory. `announcements(config)` takes diff --git a/doc/lib/py/index.md b/doc/lib/py/index.md index a41cb7255b..10b2edd914 100644 --- a/doc/lib/py/index.md +++ b/doc/lib/py/index.md @@ -74,7 +74,9 @@ For already-encoded live output, call `audio.flush(timestamp_us)` after each `au The three advertising operations, as the other bindings spell them: `client.create_broadcast(path)` (or `OriginProducer.create_broadcast`) returns an unannounced producer, invisible to everyone; `broadcast.announce(route)` / -`broadcast.unannounce()` own that exact-path advertisement; +`broadcast.unannounce()` own that exact-path advertisement, and +`broadcast.close()` ends the broadcast for good (a second call is a no-op; +`finish()` is its deprecated alias); `origin.dynamic(prefix, route)` claims `prefix` and every path beneath it (`""` for everything). Hold the returned handle while the claim should stay advertised, and reject the requests you will not serve. A route is a diff --git a/doc/lib/swift/index.md b/doc/lib/swift/index.md index a9aa8b3c44..413b9d7554 100644 --- a/doc/lib/swift/index.md +++ b/doc/lib/swift/index.md @@ -59,7 +59,9 @@ For already-encoded live output, call `audio.flush(timestampUs:)` after `writeFr The three advertising operations: `session.publish.createBroadcast(path:)` returns an unannounced producer, invisible to everyone; `broadcast.announce(route:)` / -`broadcast.unannounce()` own that exact-path advertisement; +`broadcast.unannounce()` own that exact-path advertisement, and +`broadcast.close()` ends the broadcast for good (a second call is a no-op; +`finish()` is its deprecated alias); `session.publish.dynamic(prefix:route:)` claims `prefix` and every path beneath it (`""` for everything). Hold the returned `OriginDynamic` while the claim should stay advertised, and reject the requests you will not serve. A diff --git a/go/wrapper/example_test.go b/go/wrapper/example_test.go index 84eef29901..60c3c40d29 100644 --- a/go/wrapper/example_test.go +++ b/go/wrapper/example_test.go @@ -53,8 +53,8 @@ func ExampleClient_CreateBroadcast() { if err != nil { log.Fatal(err) } - // Finishing unpublishes the broadcast immediately. - defer broadcast.Finish() + // Closing ends the broadcast for good. + defer broadcast.Close() media, err := broadcast.PublishAudio(moq.AudioFormatOpus, opusHead()) if err != nil { @@ -91,7 +91,7 @@ func ExampleBroadcastProducer_PublishVideo_videoHint() { if err != nil { log.Fatal(err) } - defer broadcast.Finish() + defer broadcast.Close() media, err := broadcast.PublishVideo(moq.VideoFormatAvc3, nil, moq.WithVideoHint(moq.VideoHint{})) if err != nil { diff --git a/go/wrapper/moq_test.go b/go/wrapper/moq_test.go index 0a91ad7927..936fd7e5df 100644 --- a/go/wrapper/moq_test.go +++ b/go/wrapper/moq_test.go @@ -115,7 +115,7 @@ func TestDynamicBroadcastRequest(t *testing.T) { if err := track.Finish(); err != nil { t.Fatal(err) } - if err := served.Finish(); err != nil { + if err := served.Close(); err != nil { t.Fatal(err) } } @@ -135,11 +135,27 @@ func TestPublishAudioLifecycle(t *testing.T) { if err := media.Finish(); err != nil { t.Fatal(err) } - if err := broadcast.Finish(); err != nil { + if err := broadcast.Close(); err != nil { t.Fatal(err) } } +func TestBroadcastCloseTwiceIsNoop(t *testing.T) { + broadcast, err := moq.NewBroadcastProducer() + if err != nil { + t.Fatal(err) + } + if err := broadcast.Close(); err != nil { + t.Fatal(err) + } + if err := broadcast.Close(); err != nil { + t.Fatalf("second Close: %v", err) + } + if _, err := broadcast.PublishTrack("events", nil); err == nil { + t.Fatal("PublishTrack after Close succeeded") + } +} + func TestEncodeAudioWithOpusObject(t *testing.T) { // The producer retains the codec, so releasing the codec and the // config in either order must still encode. @@ -179,7 +195,7 @@ func TestEncodeAudioWithOpusObject(t *testing.T) { if err := producer.Finish(); err != nil { t.Fatal(err) } - if err := broadcast.Finish(); err != nil { + if err := broadcast.Close(); err != nil { t.Fatal(err) } } @@ -221,7 +237,7 @@ func TestEncodeAudioFrameDurations(t *testing.T) { t.Fatalf("err = %v, want ErrAudio: 2 ms is not an opus frame duration", err) } - if err := broadcast.Finish(); err != nil { + if err := broadcast.Close(); err != nil { t.Fatal(err) } } @@ -235,7 +251,7 @@ func TestVideoPropertiesUseDefaultedFields(t *testing.T) { if err := broadcast.SetVideoProperties(moq.VideoProperties{Rotation: &rotation}); err != nil { t.Fatal(err) } - if err := broadcast.Finish(); err != nil { + if err := broadcast.Close(); err != nil { t.Fatal(err) } } @@ -350,7 +366,7 @@ func TestDecodeVideoFormat(t *testing.T) { if err := video.Finish(); err != nil { t.Fatal(err) } - if err := broadcast.Finish(); err != nil { + if err := broadcast.Close(); err != nil { t.Fatal(err) } } @@ -806,7 +822,7 @@ func TestDynamicTrackRequest(t *testing.T) { if err != nil { t.Fatal(err) } - defer broadcast.Finish() + defer broadcast.Close() dynamic, err := broadcast.Dynamic() if err != nil { @@ -882,7 +898,7 @@ func TestDynamicTrackRequestCanPublishAudio(t *testing.T) { if err != nil { t.Fatal(err) } - defer broadcast.Finish() + defer broadcast.Close() dynamic, err := broadcast.Dynamic() if err != nil { @@ -974,7 +990,7 @@ func TestRecvGroupCancelRace(t *testing.T) { if err != nil { t.Fatal(err) } - defer broadcast.Finish() + defer broadcast.Close() var wg sync.WaitGroup for i := 0; i < 16; i++ { @@ -1007,7 +1023,7 @@ func TestConsumerCancelConcurrent(t *testing.T) { if err != nil { t.Fatal(err) } - defer broadcast.Finish() + defer broadcast.Close() track, err := broadcast.PublishTrack("x", nil) if err != nil { @@ -1075,7 +1091,7 @@ func TestRequestBroadcastCancelKeepsTheOrigin(t *testing.T) { if err != nil { t.Fatal(err) } - defer served.Finish() + defer served.Close() resolved := make(chan error, 1) go func() { @@ -1108,7 +1124,7 @@ func TestSubscribeTrackCancelKeepsTheBroadcast(t *testing.T) { if err != nil { t.Fatal(err) } - defer broadcast.Finish() + defer broadcast.Close() dynamic, err := broadcast.Dynamic() if err != nil { @@ -1177,7 +1193,7 @@ func TestUsedCancelKeepsTheTrack(t *testing.T) { if err != nil { t.Fatal(err) } - defer broadcast.Finish() + defer broadcast.Close() track, err := broadcast.PublishTrack("status", nil) if err != nil { @@ -1224,7 +1240,7 @@ func TestCancelDoesNotLeakGoroutines(t *testing.T) { if err != nil { t.Fatal(err) } - defer broadcast.Finish() + defer broadcast.Close() dynamic, err := broadcast.Dynamic() if err != nil { @@ -1338,7 +1354,7 @@ func TestBroadcastIsReachableOnlyWhileAnnounced(t *testing.T) { if _, err := consumer.RequestBroadcast(ctx, "live"); err != nil { t.Fatal(err) } - if err := broadcast.Finish(); err != nil { + if err := broadcast.Close(); err != nil { t.Fatal(err) } } diff --git a/go/wrapper/origin.go b/go/wrapper/origin.go index 67eb77e040..78642147c1 100644 --- a/go/wrapper/origin.go +++ b/go/wrapper/origin.go @@ -49,9 +49,8 @@ func (o *OriginProducer) Dynamic(prefix string, route Route) (*OriginDynamic, er // // The broadcast is invisible and unroutable, for this origin's consumers and // peers alike, until [BroadcastProducer.Announce]. Announce it after -// populating tracks. Finish unpublishes immediately, while dropping the -// producer without finishing also unpublishes but reads to subscribers as a -// failure rather than a deliberate end. +// populating tracks. [BroadcastProducer.Close] ends it for good; dropping the +// last handle does the same. func (o *OriginProducer) CreateBroadcast(path string) (*BroadcastProducer, error) { inner, err := o.inner.CreateBroadcast(path) if err != nil { diff --git a/go/wrapper/publish.go b/go/wrapper/publish.go index 6eccc6d8f6..27c7ff2087 100644 --- a/go/wrapper/publish.go +++ b/go/wrapper/publish.go @@ -262,9 +262,17 @@ func (b *BroadcastProducer) RemoveCatalogSection(name string) error { return b.inner.RemoveCatalogSection(name) } -// Finish closes the broadcast. +// Close ends the broadcast for good: it retracts and serves no new tracks. +// Tracks already subscribed carry on to their own end. Closing again is a no-op. +func (b *BroadcastProducer) Close() error { + return b.inner.Close() +} + +// Finish ends the broadcast. +// +// Deprecated: use [BroadcastProducer.Close]; a broadcast end carries no cause. func (b *BroadcastProducer) Finish() error { - return b.inner.Finish() + return b.inner.Close() } // BroadcastDynamic is a stream of subscriber-requested tracks. diff --git a/go/wrapper/reconnect_test.go b/go/wrapper/reconnect_test.go index 72e71536d2..762a46c38d 100644 --- a/go/wrapper/reconnect_test.go +++ b/go/wrapper/reconnect_test.go @@ -166,7 +166,7 @@ func TestReconnectAcrossRelayRestart(t *testing.T) { if err != nil { t.Fatal(err) } - defer broadcast.Finish() + defer broadcast.Close() track, err := broadcast.PublishTrack("data", nil) if err != nil { diff --git a/kt/moq/src/jvmAndAndroidMain/kotlin/dev/moq/Moq.kt b/kt/moq/src/jvmAndAndroidMain/kotlin/dev/moq/Moq.kt index 728c381c8e..dec946bc2f 100644 --- a/kt/moq/src/jvmAndAndroidMain/kotlin/dev/moq/Moq.kt +++ b/kt/moq/src/jvmAndAndroidMain/kotlin/dev/moq/Moq.kt @@ -22,7 +22,8 @@ class Moq internal constructor( /** * Create an unannounced broadcast at [path], invisible to everyone until announced. * - * Advertise it with `announce` after populating tracks. `finish()` unpublishes immediately. + * Advertise it with `announce` after populating tracks. `close()` (or `use`) ends it once no + * `dynamic()` handle remains. */ fun createBroadcast(path: String): BroadcastProducer = session.publish().createBroadcast(path) diff --git a/kt/moq/src/jvmAndAndroidMain/kotlin/dev/moq/Server.kt b/kt/moq/src/jvmAndAndroidMain/kotlin/dev/moq/Server.kt index e085291a69..9257ba51dd 100644 --- a/kt/moq/src/jvmAndAndroidMain/kotlin/dev/moq/Server.kt +++ b/kt/moq/src/jvmAndAndroidMain/kotlin/dev/moq/Server.kt @@ -33,8 +33,8 @@ class Server internal constructor( * Create a live broadcast at [path], served to incoming sessions. * * The origin announces the path so subscribers can discover it, becoming visible - * Advertise it with `announce` after populating tracks. `finish()` - * unpublishes immediately. + * Advertise it with `announce` after populating tracks. `close()` (or `use`) + * ends it once no `dynamic()` handle remains. */ fun createBroadcast(path: String): BroadcastProducer { val origin = publishOrigin ?: throw IllegalStateException("no publish origin configured") diff --git a/kt/moq/src/jvmAndAndroidTest/kotlin/dev/moq/SmokeTest.kt b/kt/moq/src/jvmAndAndroidTest/kotlin/dev/moq/SmokeTest.kt index 048ee430c2..9cd8e32517 100644 --- a/kt/moq/src/jvmAndAndroidTest/kotlin/dev/moq/SmokeTest.kt +++ b/kt/moq/src/jvmAndAndroidTest/kotlin/dev/moq/SmokeTest.kt @@ -132,6 +132,16 @@ class SmokeTest { assertEquals(ConnectionStatus.CONNECTED, status) } + @Test + fun `closing a broadcast ends it, and closing again is a no-op`() = runTest { + val broadcast = BroadcastProducer() + val consumer = broadcast.consume() + // AutoCloseable.close() releases the last producer handle, which ends the broadcast. + broadcast.close() + broadcast.close() + assertFailsWith { consumer.subscribeTrack("events", null) } + } + @Test fun `broadcast updates shared video properties`() { BroadcastProducer().use { broadcast -> @@ -302,7 +312,6 @@ class SmokeTest { server.createBroadcast("live").use { broadcast -> broadcast.announce(Route()) broadcast.unannounce() - broadcast.finish() } } } diff --git a/py/moq-rs/README.md b/py/moq-rs/README.md index 6add3e7adb..0902031a5a 100644 --- a/py/moq-rs/README.md +++ b/py/moq-rs/README.md @@ -68,7 +68,7 @@ async def main(): # Clean up audio.finish() - broadcast.finish() + broadcast.close() asyncio.run(main()) @@ -132,7 +132,7 @@ client = moq.Client( - **`Server(bind="[::]:443", *, tls_cert=(), tls_key=(), tls_generate=(), publish=None, subscribe=None)`**. Async context manager + async iterator of incoming `Request`s. - `.local_addr`. The bound address (useful when binding to port `0`). - `.cert_fingerprints()`. SHA-256 fingerprints of the configured TLS certificates, for `serverCertificateHashes` browser cert pinning. - - `.create_broadcast(path) → BroadcastProducer`. Create an unannounced broadcast, invisible to everyone; `announce()` makes it discoverable and reachable; `finish()` unpublishes it. + - `.create_broadcast(path) → BroadcastProducer`. Create an unannounced broadcast, invisible to everyone; `announce()` makes it discoverable and reachable; `close()` ends it. - **`Request`**. An incoming session, yielded by `async for request in server`. - `.url`, `.path`, `.query`, `.transport`. The query-free path is uniform across transports; the root or missing path is `""`. The encoded query may contain credentials. - `.set_publish(origin)`, `.set_consume(origin)`. Per-request overrides, captured at `accept()`. Raise if the request is already answered, cancelled, or currently accepting. @@ -155,7 +155,7 @@ client = moq.Client( - `.publish_video(format, init=b"", *, label=None, hint=None, track=None) → MediaProducer`. `init` may be empty for a format that resolves in band; a `VideoHint` pins catalog fields the stream can't reveal (bitrate) or publishes the catalog before the first keyframe. `track` names the track as in `publish_audio`. - `.encode_video(input, output, *, bandwidth=None) → VideoProducer`. Encode raw `VideoFrame`s inside the binding; `.write(frame)` each one. - `.encode_audio(name, input, output, *, bandwidth=None) → AudioProducer`. Encode raw PCM `AudioFrame`s; the codec is `output.codec`, e.g. `AudioCodec.opus()`, with `output.frame_duration_us` setting the Opus frame length. - - `.finish()` + - `.close()` ends the broadcast for good; a second call is a no-op. `.finish()` is its deprecated alias. - **`BroadcastDynamic`**. Async source of tracks requested by subscribers. - `await .requested_track() → TrackRequest`. Call `.accept()` on it for a `TrackProducer`, or `.abort(code)` to reject. - Async iterator yielding `TrackRequest` diff --git a/py/moq-rs/moq/origin.py b/py/moq-rs/moq/origin.py index f4a7c1886b..517ebe1210 100644 --- a/py/moq-rs/moq/origin.py +++ b/py/moq-rs/moq/origin.py @@ -248,8 +248,7 @@ def create_broadcast(self, path: str) -> BroadcastProducer: and peers alike, until :meth:`BroadcastProducer.announce`. Announce it after populating tracks. Create, :meth:`dynamic` if tracks are served on demand, populate, then announce. - ``finish()`` unpublishes immediately, while dropping the producer without - finishing also unpublishes but reads to subscribers as a failure rather - than a deliberate end. + :meth:`BroadcastProducer.close` ends it for good; dropping the last + handle does the same. """ return BroadcastProducer._from_inner(self._inner.create_broadcast(path)) diff --git a/py/moq-rs/moq/publish.py b/py/moq-rs/moq/publish.py index 2a5929954a..c6e0b4b633 100644 --- a/py/moq-rs/moq/publish.py +++ b/py/moq-rs/moq/publish.py @@ -3,6 +3,7 @@ from __future__ import annotations import json +import warnings from typing import TYPE_CHECKING, Any from moq_ffi import ( @@ -837,6 +838,14 @@ def consume(self) -> BroadcastConsumer: return BroadcastConsumer(self._inner.consume()) + def close(self) -> None: + """End the broadcast for good: retract it and serve no new tracks. + + Tracks already subscribed carry on to their own end. Closing again is a no-op. + """ + self._inner.close() + def finish(self) -> None: - """Finish the broadcast, closing its tracks and unpublishing it.""" - self._inner.finish() + """Deprecated: use :meth:`close`. A broadcast end carries no cause.""" + warnings.warn("use close(); a broadcast end carries no cause", DeprecationWarning, stacklevel=2) + self._inner.close() diff --git a/py/moq-rs/tests/test_local.py b/py/moq-rs/tests/test_local.py index ad004a566f..d8b23fb367 100644 --- a/py/moq-rs/tests/test_local.py +++ b/py/moq-rs/tests/test_local.py @@ -127,7 +127,7 @@ def test_publish_media_lifecycle(): media = broadcast.publish_audio(moq.AudioFormat.OPUS, opus_head()) media.write_frame(b"opus frame", 1000) media.finish() - broadcast.finish() + broadcast.close() def test_publish_media_cut_and_seek(): @@ -148,7 +148,7 @@ def test_publish_media_cut_and_seek(): media.seek(42) media.finish() - broadcast.finish() + broadcast.close() def test_video_properties_use_defaulted_fields(): @@ -157,7 +157,7 @@ def test_video_properties_use_defaulted_fields(): assert properties.display is None assert properties.flip is None broadcast.set_video_properties(properties) - broadcast.finish() + broadcast.close() def test_audio_rejects_bad_init_bytes(): @@ -303,12 +303,19 @@ async def test_catalog_update_on_new_track(): break -def test_finish_closes_producer(): +def test_close_twice_is_a_noop(): broadcast = moq.BroadcastProducer() _media = broadcast.publish_audio(moq.AudioFormat.OPUS, opus_head()) - broadcast.finish() + broadcast.close() + broadcast.close() with pytest.raises(Exception): + broadcast.publish_audio(moq.AudioFormat.OPUS, opus_head()) + + +def test_finish_is_deprecated(): + broadcast = moq.BroadcastProducer() + with pytest.deprecated_call(): broadcast.finish() @@ -330,7 +337,7 @@ def test_publish_lifecycle(): track = broadcast.publish_track("status") track.write_frame(b'{"cmd": "ready"}', 0) track.finish() - broadcast.finish() + broadcast.close() async def test_publish_track_info_and_subscription(): @@ -520,7 +527,7 @@ async def test_dynamic_broadcast_request(): assert frame.payload == payload assert frame.timestamp_us == 20_000 track.finish() - served.finish() + served.close() async def test_dynamic_broadcast_request_can_reject(): @@ -1075,7 +1082,7 @@ def test_encode_audio_with_opus_object(): producer.write(moq.AudioFrame(timestamp_us=0, data=bytes(960 * 4))) assert producer.name == "mic" producer.finish() - broadcast.finish() + broadcast.close() async def test_decode_video_format(): @@ -1129,7 +1136,7 @@ async def test_decode_video_format(): i420.cancel() packed.cancel() video.finish() - broadcast.finish() + broadcast.close() async def test_broadcast_is_reachable_only_while_announced(): @@ -1159,7 +1166,7 @@ async def test_broadcast_is_reachable_only_while_announced(): await asyncio.wait_for(consumer.request_broadcast("live"), timeout=5.0) announced.cancel() track.finish() - broadcast.finish() + broadcast.close() async def test_announced_pattern_captures(): @@ -1180,8 +1187,8 @@ async def test_announced_pattern_captures(): announced.cancel() dynamic.cancel() - audio.finish() - chat.finish() + audio.close() + chat.close() async def test_dynamic_serves_a_request_under_a_prefix(): @@ -1196,7 +1203,7 @@ async def test_dynamic_serves_a_request_under_a_prefix(): request.accept(served) await asyncio.wait_for(pending, timeout=5.0) dynamic.cancel() - served.finish() + served.close() async def test_dynamic_and_json_handles_are_async_context_managers(): @@ -1235,4 +1242,4 @@ async def assert_cancelled(awaitable) -> None: snapshot.finish() stream.finish() track.finish() - broadcast.finish() + broadcast.close() diff --git a/py/moq-rs/tests/test_server.py b/py/moq-rs/tests/test_server.py index 170129306b..3dd07bc8a9 100644 --- a/py/moq-rs/tests/test_server.py +++ b/py/moq-rs/tests/test_server.py @@ -72,7 +72,7 @@ async def accept_loop() -> None: except asyncio.CancelledError: pass media.finish() - broadcast.finish() + broadcast.close() async def test_client_reconnects_and_resumes_announcements(): @@ -127,7 +127,7 @@ async def accept_loop() -> None: async for announcement in client.announced(): assert announcement.prefix == "after-reconnect" break - broadcast.finish() + broadcast.close() finally: accept_task.cancel() try: @@ -287,7 +287,7 @@ async def test_serve_helper_accepts_clients(): await serve_task except asyncio.CancelledError: pass - broadcast.finish() + broadcast.close() async def test_broadcast_route_over_wire(): @@ -317,7 +317,7 @@ async def test_broadcast_route_over_wire(): await serve_task except asyncio.CancelledError: pass - broadcast.finish() + broadcast.close() async def test_route_update_observes_restart(): diff --git a/quest/m1/broadcast-close/README.md b/quest/m1/broadcast-close/README.md index 3ad63bad52..99f38a846d 100644 --- a/quest/m1/broadcast-close/README.md +++ b/quest/m1/broadcast-close/README.md @@ -34,11 +34,10 @@ object ended", not "the path went offline": announcements say whether a path is live, and a path can be announced again by a new object. Caches and moq-hls use it to tell a live object from a replaced one at the same path. -Stage it as the children below: the bindings, then the `dev` removal. +Stage the `dev` removal as the child below. ## Quests -- [Binding close](/quest/m1/broadcast-close/bindings.md) - moq-ffi, libmoq, and every wrapper expose `close()` and deprecate `finish` - [Remove finish](/quest/m1/broadcast-close/remove.md) - on dev, the deprecated broadcast end APIs are gone and `closed()` carries no cause ## Related diff --git a/quest/m1/broadcast-close/bindings.md b/quest/m1/broadcast-close/bindings.md deleted file mode 100644 index 5f1444c897..0000000000 --- a/quest/m1/broadcast-close/bindings.md +++ /dev/null @@ -1,25 +0,0 @@ -# [M] Binding close - -## Goal - -Every binding ends a broadcast with `close()`, mirroring Rust, and its -`finish` is deprecated. - -## Plan - -- moq-ffi: add `MoqBroadcastProducer::close()` next to `finish()` in - `rs/moq-ffi/src/producer.rs`. It closes the broadcast, then the catalog, as - `finish` does, so a catalog error cannot stop the broadcast from ending. - Deprecate `finish`. -- libmoq: export `moq_publish_close` and deprecate `moq_publish_finish`. Move - `cpp/obs/src/moq-output.cpp` and the C tests over. -- Wrappers: `py/moq-rs` `BroadcastProducer.close()` (its `finish` docstring - wrongly says it closes the tracks), swift `Broadcast.close()`, and the Go - `Publish.Close()`. Kotlin and Dart only alias the generated type, so they - pick `close` up from moq-ffi. -- Move the binding tests over, keeping one test per binding that a second - `close` is a no-op, as in Rust. -- Update `doc/lib/{py,swift,kt,go,dart,c}`, including `doc/lib/go/index.md`'s - `broadcast.Finish()` sample. -- Fix the moq-ffi `origin.rs` doc comment that points users at a - `broadcast.closed()` the bindings don't have. diff --git a/quest/m1/broadcast-close/remove.md b/quest/m1/broadcast-close/remove.md index 805896a6ee..6da1df879c 100644 --- a/quest/m1/broadcast-close/remove.md +++ b/quest/m1/broadcast-close/remove.md @@ -15,7 +15,7 @@ This is a published API break, so it targets `dev`. `()` rather than an `Error` cause. - Remove the deprecated binding `finish` methods, `moq_publish_finish`, and JS's `close(abort)` parameter. - -## Required - -- [Binding close](/quest/m1/broadcast-close/bindings.md) - every binding already has `close()` to move to +- Kotlin has no generated `close()`: it would collide with `AutoCloseable.close()`, + so `rs/moq-ffi/uniffi.toml` excludes it. Kotlin's `close()` releases the handle, + which ends the broadcast only once no `dynamic()` handle remains. Removing + `finish` leaves Kotlin without a forced end; decide whether it needs one. diff --git a/rs/libmoq/README.md b/rs/libmoq/README.md index 635024d478..2ea1d55bb9 100644 --- a/rs/libmoq/README.md +++ b/rs/libmoq/README.md @@ -63,7 +63,7 @@ int32_t moq_origin_announced_cancel(uint32_t announced); // Publishing int32_t moq_publish_announce(uint32_t broadcast, const moq_route *route); int32_t moq_publish_unannounce(uint32_t broadcast); -int32_t moq_publish_finish(uint32_t broadcast); +int32_t moq_publish_close(uint32_t broadcast); int32_t moq_publish_audio(uint32_t broadcast, const moq_audio_init *config); int32_t moq_publish_video(uint32_t broadcast, const moq_video_init *config); int32_t moq_publish_container(uint32_t broadcast, const moq_container_init *config); diff --git a/rs/libmoq/c-tests/decode-output.c b/rs/libmoq/c-tests/decode-output.c index 585f107bf2..057f33e62d 100644 --- a/rs/libmoq/c-tests/decode-output.c +++ b/rs/libmoq/c-tests/decode-output.c @@ -284,8 +284,8 @@ int main(void) { fail("error: moq_consume_close failed (%s)\n", moq_error()); if (moq_encode_video_finish((uint32_t)producer) < 0) fail("error: moq_encode_video_finish failed (%s)\n", moq_error()); - if (moq_publish_finish((uint32_t)broadcast) < 0) - fail("error: moq_publish_finish failed (%s)\n", moq_error()); + if (moq_publish_close((uint32_t)broadcast) < 0) + fail("error: moq_publish_close failed (%s)\n", moq_error()); if (moq_origin_close((uint32_t)origin) < 0) fail("error: moq_origin_close failed (%s)\n", moq_error()); diff --git a/rs/libmoq/src/api.rs b/rs/libmoq/src/api.rs index 42bfcecbe7..181d63dd2d 100644 --- a/rs/libmoq/src/api.rs +++ b/rs/libmoq/src/api.rs @@ -1540,7 +1540,7 @@ pub extern "C" fn moq_origin_create() -> i32 { /// /// The broadcast is invisible and unroutable, on this origin and its peers /// alike, until [moq_publish_announce]. Fill it with the `moq_publish_*` -/// functions, then announce it. [moq_publish_finish] unpublishes immediately. +/// functions, then announce it. [moq_publish_close] ends it for good. /// /// Returns a non-zero broadcast handle on success, or a negative code on failure. /// @@ -1920,20 +1920,29 @@ pub extern "C" fn moq_publish_unannounce(broadcast: u32) -> i32 { }) } -/// Finish a broadcast and release it, ending its catalog cleanly. +/// End a broadcast for good and release its handle. /// -/// Subscribers see a normal end of stream rather than an error, and the origin unpublishes -/// the path immediately. +/// The origin retracts the path immediately and serves no new tracks; tracks already +/// subscribed carry on to their own end. The handle is invalid afterwards, so closing +/// it again fails like any unknown handle. /// /// Returns a zero on success, or a negative code on failure. #[unsafe(no_mangle)] -pub extern "C" fn moq_publish_finish(broadcast: u32) -> i32 { +pub extern "C" fn moq_publish_close(broadcast: u32) -> i32 { ffi::enter(move || { let broadcast = ffi::parse_id(broadcast)?; - State::lock().publish.finish(broadcast) + State::lock().publish.close(broadcast) }) } +/// Deprecated: use [moq_publish_close]. A broadcast end carries no cause. +/// +/// Returns a zero on success, or a negative code on failure. +#[unsafe(no_mangle)] +pub extern "C" fn moq_publish_finish(broadcast: u32) -> i32 { + moq_publish_close(broadcast) +} + /// Publish one audio codec as a new media track. /// /// The track is named after the format (`0.opus`), so a subscriber finds it diff --git a/rs/libmoq/src/publish.rs b/rs/libmoq/src/publish.rs index e83a4e842a..7b93c083e0 100644 --- a/rs/libmoq/src/publish.rs +++ b/rs/libmoq/src/publish.rs @@ -162,9 +162,8 @@ impl Publish { Ok((&mut broadcast.producer, &mut broadcast.catalog)) } - /// Cleanly finish the broadcast and finalize the catalog stream, so subscribers - /// see a normal end rather than [`moq_net::Error::Dropped`]. - pub fn finish(&mut self, broadcast: Id) -> Result<(), Error> { + /// End the broadcast for good and release it, finalizing the catalog stream. + pub fn close(&mut self, broadcast: Id) -> Result<(), Error> { let Broadcast { producer, mut catalog, diff --git a/rs/libmoq/src/test.rs b/rs/libmoq/src/test.rs index 535d3574d8..d48de75362 100644 --- a/rs/libmoq/src/test.rs +++ b/rs/libmoq/src/test.rs @@ -367,7 +367,7 @@ fn publish_media_lifecycle() { let origin = id(moq_origin_create()); let broadcast = publish_broadcast(origin, b"publish-media-lifecycle"); let _guard = Guard(Some(|| { - moq_publish_finish(broadcast); + moq_publish_close(broadcast); })); let init = opus_head(); @@ -379,7 +379,7 @@ fn publish_media_lifecycle() { assert_eq!(ret, 0, "moq_publish_media_frame should succeed"); assert_eq!(moq_publish_media_finish(media), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); } #[test] @@ -389,7 +389,7 @@ fn publish_media_rejects_a_null_config() { assert!(unsafe { moq_publish_audio(broadcast, std::ptr::null()) } < 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -424,7 +424,7 @@ fn container_and_media_handles_are_not_interchangeable() { assert_eq!(moq_publish_container_finish(container), 0); assert_eq!(moq_publish_media_finish(media), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -511,7 +511,7 @@ fn publish_media_labels_config_without_naming_track() { assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_publish_media_finish(media1), 0); assert_eq!(moq_publish_media_finish(media2), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -621,7 +621,7 @@ fn publish_media_owns_its_rendition_before_the_first_keyframe() { "finishing the media track releases its rendition name" ); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -667,7 +667,7 @@ fn publish_video_config_replaces_its_own_rendition() { "the name is free once the caller removes its rendition" ); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -690,7 +690,7 @@ fn publish_catalog_config_null_pointer() { -6, "null config should return InvalidPointer (-6)" ); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); } #[test] @@ -891,7 +891,7 @@ fn publish_catalog_roundtrip() { assert_eq!(moq_consume_catalog_cancel(catalog_task), 0); assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0"); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -963,13 +963,13 @@ fn a_half_specified_coded_size_round_trips() { assert_eq!(moq_consume_catalog_cancel(forwarded_task), 0); assert_eq!(forwarded_cb.recv_catalog_terminal(), 0); assert_eq!(moq_consume_close(forwarded), 0); - assert_eq!(moq_publish_finish(forward), 0); + assert_eq!(moq_publish_close(forward), 0); assert_eq!(moq_consume_catalog_free(catalog), 0); assert_eq!(moq_consume_catalog_cancel(catalog_task), 0); assert_eq!(catalog_cb.recv_catalog_terminal(), 0); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1065,7 +1065,7 @@ fn raw_loc_video_uses_the_declared_catalog_container() { assert_eq!(catalog_cb.recv_terminal(), 0); assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_publish_track_finish(track), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1129,7 +1129,7 @@ fn cmaf_catalog_container_carries_its_init_segment() { assert_eq!(moq_consume_catalog_cancel(catalog_task), 0); assert_eq!(catalog_cb.recv_terminal(), 0); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1189,7 +1189,7 @@ fn unpublishable_catalog_containers_are_rejected() { "cmaf without an init segment should return InvalidPointer (-6)" ); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1344,7 +1344,7 @@ fn catalog_section_roundtrip() { assert_eq!(moq_consume_catalog_cancel(catalog_task), 0); assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0"); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1391,7 +1391,7 @@ fn publish_track_with_info_rejects_invalid_timescale() { }; assert!(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), &info) } < 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); } #[test] @@ -1491,7 +1491,7 @@ fn raw_track_publish_consume() { assert_eq!(moq_publish_track_finish(track), 0); assert!(moq_publish_track_finish(track) < 0, "double-close should fail"); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1553,7 +1553,7 @@ fn raw_track_datagram_publish_consume() { assert!(moq_consume_datagrams_cancel(consumer) < 0, "double-close should fail"); assert_eq!(moq_publish_track_finish(track), 0); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1572,7 +1572,7 @@ fn raw_track_sparse_groups_and_known_end() { assert_eq!(moq_publish_group_finish(group), 0); assert!(moq_publish_track_group_at(track, 5) < 0); assert_eq!(moq_publish_track_finish(track), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); } #[test] @@ -1587,7 +1587,7 @@ fn raw_track_and_group_abort_consume_their_handles() { assert!(moq_publish_group_finish(group) < 0); assert_eq!(moq_publish_track_abort(track, 410), 0); assert!(moq_publish_track_finish(track) < 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); } #[test] @@ -1673,7 +1673,7 @@ fn raw_track_subscription_options_and_update() { assert_eq!(frame_cb.recv_terminal(), 0); assert_eq!(moq_publish_track_finish(track), 0); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1739,7 +1739,7 @@ fn json_snapshot_publish_consume() { "double-close should fail" ); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1799,7 +1799,7 @@ fn json_stream_publish_consume() { assert_eq!(moq_publish_json_stream_finish(producer), 0); assert!(moq_publish_json_stream_finish(producer) < 0, "double-close should fail"); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1807,13 +1807,26 @@ fn json_stream_publish_consume() { fn close_invalid_or_zero_ids() { assert!(moq_origin_close(9999) < 0); assert!(moq_session_close(9999) < 0); - assert!(moq_publish_finish(9999) < 0); + assert!(moq_publish_close(9999) < 0); assert!(moq_consume_close(9999) < 0); assert!(moq_consume_frame_free(9999) < 0); assert!(moq_origin_close(0) < 0); assert!(moq_session_close(0) < 0); - assert!(moq_publish_finish(0) < 0); + assert!(moq_publish_close(0) < 0); +} + +#[test] +fn publish_close_releases_the_handle() { + let origin = id(moq_origin_create()); + let broadcast = publish_broadcast(origin, b"close/twice"); + assert_eq!(moq_publish_close(broadcast), 0); + assert!(moq_publish_close(broadcast) < 0, "a closed handle is released"); + + // The deprecated alias still ends a broadcast. + let broadcast = publish_broadcast(origin, b"close/finish"); + assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_origin_close(origin), 0); } #[test] @@ -1865,7 +1878,7 @@ fn announced_free_lifecycle() { ann_cb.recv_terminal(); assert_eq!(moq_origin_close(origin), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); } #[test] @@ -1882,7 +1895,7 @@ fn double_close_all_resource_types() { assert_eq!(moq_publish_media_finish(media), 0); assert!(moq_publish_media_finish(media) < 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); let origin = id(moq_origin_create()); let path = b"double-close-test"; @@ -1922,7 +1935,7 @@ fn double_close_all_resource_types() { assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_publish_media_finish(media), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1961,7 +1974,7 @@ fn media_cut_bounds_audio_groups() { assert!(moq_publish_media_seek(9999, 0) < 0); assert_eq!(moq_publish_media_finish(media), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -1970,7 +1983,7 @@ fn unknown_format() { let origin = id(moq_origin_create()); let broadcast = publish_broadcast(origin, b"unknown-format"); let _guard = Guard(Some(|| { - moq_publish_finish(broadcast); + moq_publish_close(broadcast); })); // A format is an enum now, so the only bad value C can still supply is an out-of-range @@ -2025,7 +2038,7 @@ fn local_announce() { assert_eq!(moq_origin_announced_cancel(announced_task), 0); assert_eq!(cb.recv_terminal(), 0, "announced close delivers terminal 0"); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2072,8 +2085,8 @@ fn announced_filters_patterns_and_reports_captures() { assert_eq!(moq_origin_announced_free(announced_id), 0); assert_eq!(moq_origin_announced_cancel(announced_task), 0); assert_eq!(cb.recv_terminal(), 0); - assert_eq!(moq_publish_finish(audio), 0); - assert_eq!(moq_publish_finish(chat), 0); + assert_eq!(moq_publish_close(audio), 0); + assert_eq!(moq_publish_close(chat), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2125,7 +2138,7 @@ fn announced_deactivation() { assert_eq!(moq_origin_announced_cancel(announced_task), 0); assert_eq!(cb.recv_terminal(), 0, "announced close delivers terminal 0"); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2166,7 +2179,7 @@ fn create_broadcast_is_unroutable_until_announced() { assert_eq!(moq_origin_announced_cancel(announced_task), 0); assert_eq!(cb.recv_terminal(), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2184,7 +2197,7 @@ fn announce_accepts_an_anonymous_hop() { has_cold: false, }; assert_eq!(unsafe { moq_publish_announce(broadcast, &route) }, 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2233,7 +2246,7 @@ fn dynamic_serves_a_request_under_a_prefix() { assert_eq!(moq_origin_dynamic_cancel(dynamic), 0); assert!(moq_origin_dynamic_cancel(dynamic) < 0, "double-cancel should fail"); assert_eq!(cb.recv_terminal(), 0); - assert_eq!(moq_publish_finish(served), 0); + assert_eq!(moq_publish_close(served), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2363,7 +2376,7 @@ fn track_demand_follows_subscribers() { ); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2396,7 +2409,7 @@ fn track_demand_reports_current_state_before_close() { ); assert_eq!(moq_publish_track_finish(track), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2419,7 +2432,7 @@ fn track_demand_reports_an_abort() { assert_eq!(protocol.kind, moq_protocol_kind::MOQ_PROTOCOL_KIND_APP as u32); assert_eq!(protocol.code, 64 + 7); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2456,7 +2469,7 @@ fn media_demand_refuses_a_container() { assert_eq!(cb.recv_terminal(), 0); assert_eq!(moq_publish_container_finish(container), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2559,7 +2572,7 @@ fn dynamic_serves_track_requests() { assert_eq!(moq_publish_track_finish(track), 0); assert_eq!(demand_cb.recv_terminal(), 0); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2603,7 +2616,7 @@ fn dynamic_track_request_publishes_media() { assert_eq!(demand_cb.recv(), moq_demand::MOQ_DEMAND_UNUSED as i32); assert_eq!(moq_publish_media_finish(media), 0); assert_eq!(demand_cb.recv_terminal(), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(request_cb.recv_terminal(), 0); assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_origin_close(origin), 0); @@ -2688,7 +2701,7 @@ fn track_dynamic_serves_a_fetch_miss() { assert_eq!(moq_publish_dynamic_cancel(dynamic), 0); assert_eq!(group_cb.recv_terminal(), 0); assert_eq!(moq_publish_track_finish(track), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2723,7 +2736,7 @@ fn track_dynamic_serves_a_fetch_from_frame_start() { assert_eq!(moq_publish_dynamic_cancel(dynamic), 0); assert_eq!(group_cb.recv_terminal(), 0); assert_eq!(moq_publish_track_finish(track), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2757,7 +2770,7 @@ fn track_request_dynamic_survives_accept() { assert_eq!(moq_publish_dynamic_cancel(track_dynamic), 0); assert_eq!(group_cb.recv_terminal(), 0); assert_eq!(moq_publish_track_finish(track), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(request_cb.recv_terminal(), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2855,7 +2868,7 @@ fn local_publish_consume() { assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0"); assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_publish_media_finish(media), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -2912,7 +2925,7 @@ fn consume_announced_local() { assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0"); assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_publish_media_finish(media), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3014,8 +3027,8 @@ fn consume_audio_follows_a_sibling_broadcast_reference() { assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0"); assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_publish_media_finish(media), 0); - assert_eq!(moq_publish_finish(broadcast), 0); - assert_eq!(moq_publish_finish(source), 0); + assert_eq!(moq_publish_close(broadcast), 0); + assert_eq!(moq_publish_close(source), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3142,7 +3155,7 @@ fn video_publish_consume() { assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0"); assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_publish_media_finish(media), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3197,7 +3210,7 @@ fn audio_raw_publish() { "a finished producer should take no more frames" ); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3252,7 +3265,7 @@ fn audio_raw_publish_frame_durations() { assert_eq!(moq_encode_audio_finish(id(encode(b"default", 0))), 0, "0 = 20 ms"); assert!(encode(b"rounded", 2_000) < 0, "2 ms is not an opus frame duration"); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3355,7 +3368,7 @@ fn video_raw_publish_consume() { assert_eq!(catalog_cb.recv_catalog_terminal(), 0); assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_encode_video_finish(producer), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3493,7 +3506,7 @@ fn decode_first_frame(output: &moq_video_decoder_output) -> (u32, u32, usize) { assert_eq!(catalog_cb.recv_catalog_terminal(), 0); assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_encode_video_finish(producer), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); result } @@ -3578,7 +3591,7 @@ fn video_raw_publish_from_many_threads() { .join() .unwrap(); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3674,7 +3687,7 @@ fn a_stalled_encode_does_not_block_unrelated_calls() { assert_eq!(moq_origin_close(id(created)), 0); assert_eq!(moq_encode_video_finish(stalled), 0); assert_eq!(moq_encode_video_finish(other), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3712,7 +3725,7 @@ fn video_raw_publish_rejects_frame_size_mismatch() { assert!(unsafe { moq_encode_video_frame(producer, &frame) } < 0); assert_eq!(moq_encode_video_finish(producer), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3801,7 +3814,7 @@ fn video_raw_publish_rejects_invalid_config() { assert!(moq_encode_video_bitrate(0, 1_000_000) < 0); assert!(moq_encode_video_finish(0) < 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3901,7 +3914,7 @@ fn video_raw_decode() { } assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_publish_media_finish(media), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -3961,7 +3974,7 @@ fn multiple_frames_ordering() { ); assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_publish_media_finish(media), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -4010,7 +4023,7 @@ fn catalog_update_on_new_track() { assert_eq!(moq_consume_close(consume), 0); assert_eq!(moq_publish_media_finish(media1), 0); assert_eq!(moq_publish_media_finish(media2), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -4549,7 +4562,7 @@ fn bandwidth_reservations_split_the_estimate() { let _ = first_cb.recv_terminal(); let _ = second_cb.recv_terminal(); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -4592,7 +4605,7 @@ fn bandwidth_handles_share_the_registry() { let _ = first_cb.recv_terminal(); let _ = second_cb.recv_terminal(); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -4642,7 +4655,7 @@ fn encode_video_bitrate_caps_the_reservation() { assert_eq!(moq_encode_video_finish(producer), 0); assert_eq!(moq_bandwidth_close(bandwidth), 0); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -4761,7 +4774,7 @@ fn server_accepts_a_session() { "the server handle is gone after its terminal callback" ); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(served), 0); assert_eq!(moq_origin_close(received), 0); } @@ -4931,7 +4944,7 @@ fn json_tracks_are_advertised_in_the_catalog() { assert_eq!(moq_publish_json_stream_finish(stream), 0); assert!(published_catalog(broadcast).json.tracks.is_empty()); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -4979,7 +4992,7 @@ fn binary_snapshot_is_advertised_and_delivered() { assert!(published_catalog(broadcast).binary.tracks.is_empty()); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -5025,7 +5038,7 @@ fn binary_stream_is_advertised_and_delivered() { assert!(published_catalog(broadcast).binary.tracks.is_empty()); assert_eq!(moq_consume_close(consume), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } @@ -5087,6 +5100,6 @@ fn data_track_names_cannot_collide() { ); assert_eq!(moq_publish_json_snapshot_finish(first), 0); - assert_eq!(moq_publish_finish(broadcast), 0); + assert_eq!(moq_publish_close(broadcast), 0); assert_eq!(moq_origin_close(origin), 0); } diff --git a/rs/moq-ffi/src/origin.rs b/rs/moq-ffi/src/origin.rs index 1f325f1692..253d5a749e 100644 --- a/rs/moq-ffi/src/origin.rs +++ b/rs/moq-ffi/src/origin.rs @@ -504,7 +504,7 @@ impl MoqAnnounceUpdate { impl MoqAnnouncedBroadcast { /// Wait until the broadcast is announced. Returns `Closed` if cancelled or the origin is closed. /// - /// Use `broadcast.closed()` to learn when the broadcast ends. + /// Its end arrives as an inactive [`MoqAnnounceUpdate`] on the origin's announcements. pub async fn available(&self) -> Result, MoqError> { self.task.run(|mut state| async move { state.available().await }).await } diff --git a/rs/moq-ffi/src/producer.rs b/rs/moq-ffi/src/producer.rs index 6210300fe4..c1117c2eeb 100644 --- a/rs/moq-ffi/src/producer.rs +++ b/rs/moq-ffi/src/producer.rs @@ -170,7 +170,7 @@ impl MoqBroadcastProducer { } /// Run `f` against the open broadcast and catalog. Errors with - /// [`MoqError::Closed`] if `finish()` has already run. Used by + /// [`MoqError::Closed`] if `close()` has already run. Used by /// sibling modules (e.g. `audio`) that need joint access. pub(crate) fn with_state( &self, @@ -470,17 +470,27 @@ impl MoqBroadcastProducer { })) } - /// Finish this publisher, finalizing the catalog stream and cleanly closing the - /// broadcast so subscribers see a normal end rather than `Error::Dropped`. - pub fn finish(&self) -> Result<(), MoqError> { + /// End the broadcast for good: retract it, serve no new tracks, and finalize the catalog. + /// + /// Tracks already subscribed carry on to their own end. Every later call on this + /// producer fails with `Closed`; closing again is a no-op. + pub fn close(&self) -> Result<(), MoqError> { let _guard = crate::ffi::enter(); + // Hold the lock through shutdown so a concurrent close() returns only once it is done. let mut guard = self.state.lock().unwrap(); - let mut state = guard.take().ok_or(MoqError::Closed)?; + let Some(mut state) = guard.take() else { + return Ok(()); + }; // Close the broadcast first so it ends even if finalizing the catalog fails. state.broadcast.close(); state.catalog.finish()?; Ok(()) } + + /// Deprecated: use `close()`. A broadcast end carries no cause. + pub fn finish(&self) -> Result<(), MoqError> { + self.close() + } } // ---- Dynamic Broadcast Producer ---- diff --git a/rs/moq-ffi/src/test.rs b/rs/moq-ffi/src/test.rs index 4963dff98f..4ec19b991e 100644 --- a/rs/moq-ffi/src/test.rs +++ b/rs/moq-ffi/src/test.rs @@ -331,7 +331,7 @@ async fn announced_route_keeps_cold_cost_on_reannounce() { assert_eq!(back.cost, moq_net::origin::Cost { warm: 0, cold: 9 }); assert_eq!(MoqRoute::from(back), route); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } #[test] @@ -346,7 +346,7 @@ fn publish_media_lifecycle() { }) .unwrap(); media.finish().unwrap(); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } #[tokio::test] @@ -466,7 +466,7 @@ async fn raw_audio_activity() { assert_eq!(resumed.timestamp_us, RESUMED_TIMESTAMP_US); audio.finish().unwrap(); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } /// `frame_duration_us` is microseconds so Opus' 2.5 ms frame survives the trip, where @@ -508,7 +508,7 @@ async fn raw_audio_frame_durations() { "2 ms is not an opus frame duration" ); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } #[tokio::test] @@ -1481,7 +1481,7 @@ async fn create_broadcast_is_invisible_until_announced() { .expect("the cursor is still open"); assert_eq!(update.prefix(), "live"); assert!(update.active()); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } /// Waiting for an exact path must hand the broadcast back named by that path, the base a @@ -1508,7 +1508,7 @@ async fn announced_broadcast_keeps_the_requested_path() { .unwrap(); assert_eq!(requested.inner().info().path.as_str(), "a/pub"); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } /// A catalog rendition may name a sibling broadcast (`./source`), and the track then lives @@ -1549,8 +1549,8 @@ async fn decode_audio_follows_a_sibling_broadcast_reference() { "the catalog broadcast does not serve the track itself" ); - catalog.finish().unwrap(); - source.finish().unwrap(); + catalog.close().unwrap(); + source.close().unwrap(); } /// Announcement filters do not re-root the origin, so relative broadcast references @@ -1604,8 +1604,8 @@ async fn announced_broadcasts_resolve_siblings_under_the_prefix() { .expect("timed out subscribing on the resolved broadcast") .unwrap(); - catalog.finish().unwrap(); - source.finish().unwrap(); + catalog.close().unwrap(); + source.close().unwrap(); } #[tokio::test] @@ -1641,7 +1641,7 @@ async fn announced_filters_patterns_and_reports_captures() { assert_eq!(update.captures(), Some(vec!["alice".into()])); assert!(update.active()); - chat.finish().unwrap(); + chat.close().unwrap(); } /// A `.`-named broadcast is listed only when the config opts in or the prefix names it. @@ -1730,8 +1730,8 @@ async fn resolve_returns_a_broadcast_that_resolves_further_references() { .unwrap(); assert_eq!(back.inner().info().path.as_str(), "a/pub"); - catalog.finish().unwrap(); - source.finish().unwrap(); + catalog.close().unwrap(); + source.close().unwrap(); } #[tokio::test] @@ -1775,7 +1775,7 @@ async fn announce_and_unannounce_toggles_discovery() { .expect("timed out requesting the reannounced broadcast") .expect("a reannounced broadcast resolves"); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } #[tokio::test] @@ -1792,7 +1792,7 @@ async fn finish_unpublishes() { // A graceful finish detaches immediately; the path stops resolving. Removal is // asynchronous, so poll until it takes effect. - broadcast.finish().unwrap(); + broadcast.close().unwrap(); let removed = tokio::time::timeout(TIMEOUT, async { loop { if consumer.request_broadcast("live".into()).await.is_err() { @@ -1861,7 +1861,7 @@ async fn local_publish_consume_audio() { assert_eq!(frame.payload, payload); assert_eq!(frame.timestamp_us, 1_000_000); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } #[tokio::test] @@ -1923,7 +1923,7 @@ async fn video_publish_consume() { assert_eq!(frame.timestamp_us, 0); assert!(!frame.payload.is_empty(), "frame should have payload data"); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } /// The raw-video publish path: hand mid-gray RGBA to `encode_video` and check @@ -2036,7 +2036,7 @@ async fn video_raw_publish_consume() { assert!(!frame.payload.is_empty(), "frame should carry encoded video"); video.finish().unwrap(); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } /// The decode side picks its CPU pixel layout: an unset `format` delivers I420, @@ -2144,7 +2144,7 @@ async fn video_decode_format() { i420.cancel(); rgba_out.cancel(); video.finish().unwrap(); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } /// Regression: a `MoqVideoProducer` is shared, so its calls land on whichever @@ -2230,7 +2230,7 @@ async fn video_raw_publish_from_many_threads() { let closer = video.clone(); std::thread::spawn(move || closer.finish()).join().unwrap().unwrap(); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } /// A raw video producer rejects a buffer that isn't one picture at the @@ -2294,7 +2294,7 @@ async fn video_raw_publish_rejects_bad_frames() { Err(MoqError::Closed) )); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } #[tokio::test] @@ -2349,7 +2349,7 @@ async fn multiple_frames_ordering() { assert_eq!(frame.payload, expected.as_bytes(), "frame {i} has wrong payload"); } - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } #[tokio::test] @@ -2391,17 +2391,22 @@ async fn catalog_update_on_new_track() { assert_eq!(catalog2.audio["0.opus"].label.as_deref(), Some("English")); assert_eq!(catalog2.audio["1.opus"].label, None); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } #[test] -fn finish_closes_producer() { +fn close_twice_is_a_noop() { let broadcast = MoqBroadcastProducer::new().unwrap(); let init = opus_head(); - let _media = broadcast.publish_audio(audio_init(MoqAudioFormat::Opus, init)).unwrap(); - broadcast.finish().unwrap(); + let _media = broadcast + .publish_audio(audio_init(MoqAudioFormat::Opus, init.clone())) + .unwrap(); + broadcast.close().unwrap(); + broadcast.close().unwrap(); - let err = broadcast.finish().unwrap_err(); + let Err(err) = broadcast.publish_audio(audio_init(MoqAudioFormat::Opus, init)) else { + panic!("publishing after close succeeded"); + }; assert!( matches!(err, crate::error::MoqError::Closed), "expected Closed error, got {err}" @@ -2430,7 +2435,7 @@ async fn announced_broadcast() { .unwrap(); // Finish so consumers observe a deliberate end (the canonical end for a // publisher; dropping without finish reads as a failure). - _broadcast.finish().unwrap(); + _broadcast.close().unwrap(); } fn serve(origin: &MoqOriginProducer, prefix: &str) -> Arc { @@ -2483,7 +2488,7 @@ async fn dynamic_broadcast_request() { assert_eq!(frame.timestamp_us, 20_000); track.finish().unwrap(); - served.finish().unwrap(); + served.close().unwrap(); } /// A prefix serves requests beneath it; cancelling the handle @@ -2523,7 +2528,7 @@ async fn dynamic_serves_a_request_under_a_prefix() { crate::error::MoqProtocolKind::Unroutable, ); - served.finish().unwrap(); + served.close().unwrap(); } /// Tearing the origin down ends every handler with `Closed`. A parked request @@ -3265,7 +3270,7 @@ fn without_runtime() { announced.cancel(); client.cancel(); media.finish().unwrap(); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); drop(client); drop(consumer); drop(announcement); @@ -3366,7 +3371,7 @@ async fn server_client_roundtrip() { // Clean up. Exercise `shutdown()` on the client side and the underlying // `cancel(code)` on the server side, so both shutdown paths run. media.finish().unwrap(); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); cs.shutdown(); server_session.cancel(0); server.cancel(); @@ -3438,10 +3443,10 @@ async fn server_client_roundtrip_auto_origin() { .await .expect("timed out waiting for the loopback broadcast") .expect("an auto-origin session should discover its own announcement"); - local_broadcast.finish().unwrap(); + local_broadcast.close().unwrap(); media.finish().unwrap(); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); cs.shutdown(); server_session.cancel(0); server.cancel(); @@ -3647,7 +3652,7 @@ async fn request_per_session_publish_override() { .expect("expected an announcement"); assert_eq!(announcement.prefix(), "override-only"); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); cs.cancel(0); server_session.cancel(0); server.cancel(); @@ -3775,7 +3780,7 @@ async fn client_reconnects_and_resumes_announcements() { .expect("expected an announcement"); assert_eq!(announcement.prefix(), "after-reconnect"); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); cs.cancel(0); server_session.cancel(0); server.cancel(); @@ -4039,7 +4044,7 @@ async fn video_encoder_follows_a_shrinking_grant() { } video.finish().unwrap(); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } #[cfg(feature = "video")] @@ -4116,7 +4121,7 @@ async fn set_bitrate_caps_a_later_bandwidth_grant() { assert_eq!(video.applied_bitrate(), 1_000_000); video.finish().unwrap(); - broadcast.finish().unwrap(); + broadcast.close().unwrap(); } async fn one_shot_peers() -> (Arc, Arc, Arc) { diff --git a/rs/moq-ffi/uniffi.toml b/rs/moq-ffi/uniffi.toml new file mode 100644 index 0000000000..ddf0c572fe --- /dev/null +++ b/rs/moq-ffi/uniffi.toml @@ -0,0 +1,5 @@ +# Kotlin objects are AutoCloseable, and `close()` there releases the handle, +# which ends the broadcast like dropping the last Rust producer. A generated +# `close()` would collide with it. +[bindings.kotlin] +exclude = ["MoqBroadcastProducer.close"] diff --git a/swift/Sources/Moq/Broadcast.swift b/swift/Sources/Moq/Broadcast.swift index cf52ddba1b..15a29613c6 100644 --- a/swift/Sources/Moq/Broadcast.swift +++ b/swift/Sources/Moq/Broadcast.swift @@ -344,8 +344,15 @@ public final class BroadcastProducer: Sendable { try ffi.removeCatalogSection(name: name) } - /// Finish the broadcast, finalizing the catalog stream. + /// End the broadcast for good: retract it and serve no new tracks. + /// + /// Tracks already subscribed carry on to their own end. Closing again is a no-op. + public func close() throws { + try ffi.close() + } + + @available(*, deprecated, renamed: "close", message: "A broadcast end carries no cause.") public func finish() throws { - try ffi.finish() + try ffi.close() } } diff --git a/swift/Sources/Moq/Origin.swift b/swift/Sources/Moq/Origin.swift index 6a5f4fe2d8..8c8aec69b6 100644 --- a/swift/Sources/Moq/Origin.swift +++ b/swift/Sources/Moq/Origin.swift @@ -34,10 +34,8 @@ public final class OriginProducer: Sendable { /// /// The broadcast is invisible and unroutable, for this origin's consumers /// and peers alike, until `BroadcastProducer.announce(route:)`. Announce it - /// after populating tracks. `finish()` - /// unpublishes immediately, while releasing the producer without finishing - /// also unpublishes but reads to subscribers as a failure rather than a - /// deliberate end. + /// after populating tracks. `BroadcastProducer.close()` ends it for good; + /// releasing the last handle does the same. public func createBroadcast(path: String) throws -> BroadcastProducer { BroadcastProducer(try ffi.createBroadcast(path: path)) } diff --git a/swift/Tests/MoqTests/SmokeTests.swift b/swift/Tests/MoqTests/SmokeTests.swift index 84fa239c41..38485871e1 100644 --- a/swift/Tests/MoqTests/SmokeTests.swift +++ b/swift/Tests/MoqTests/SmokeTests.swift @@ -135,7 +135,14 @@ final class SmokeTests: XCTestCase { let track = try broadcast.publishTrack(name: "events") XCTAssertEqual(try track.name, "events") try track.finish() - try broadcast.finish() + try broadcast.close() + } + + func testBroadcastCloseTwiceIsNoop() throws { + let broadcast = try BroadcastProducer() + try broadcast.close() + try broadcast.close() + XCTAssertThrowsError(try broadcast.publishTrack(name: "events")) } func testVideoHintsReachMediaPublishApi() throws { @@ -148,13 +155,13 @@ final class SmokeTests: XCTestCase { ) let media = try broadcast.publishVideo(format: .avc3, hint: hint) try media.finish() - try broadcast.finish() + try broadcast.close() } func testVideoPropertiesUseDefaultedFields() throws { let broadcast = try BroadcastProducer() try broadcast.setVideoProperties(VideoProperties(rotation: 315)) - try broadcast.finish() + try broadcast.close() } func testBroadcastConsumerFetchesCachedGroup() async throws { @@ -199,7 +206,7 @@ final class SmokeTests: XCTestCase { consumer.cancel() try producer.finish() - try broadcast.finish() + try broadcast.close() } func testJsonStreamRoundTrip() async throws { @@ -220,7 +227,7 @@ final class SmokeTests: XCTestCase { consumer.cancel() try producer.finish() - try broadcast.finish() + try broadcast.close() } func testJsonProducersReportDemand() async throws { @@ -244,7 +251,7 @@ final class SmokeTests: XCTestCase { try await snapshotDemand.unused() try await streamDemand.unused() - try broadcast.finish() + try broadcast.close() } func testRawTrackTimestamps() async throws { @@ -270,7 +277,7 @@ final class SmokeTests: XCTestCase { XCTAssertEqual(groupFrame?.timestampUs, 23_456) try track.finish() - try broadcast.finish() + try broadcast.close() } func testReadFrameSkipsEmptyThenPopulatedGroups() async throws { @@ -287,7 +294,7 @@ final class SmokeTests: XCTestCase { XCTAssertEqual(frame?.timestampUs, 2_000) try track.finish() - try broadcast.finish() + try broadcast.close() } func testSparseGroupsAndKnownEnd() throws { @@ -301,7 +308,7 @@ final class SmokeTests: XCTestCase { try track.createGroup(sequence: 4).finish() XCTAssertThrowsError(try track.createGroup(sequence: 5)) try track.finish() - try broadcast.finish() + try broadcast.close() } /// `frameDurationUs` is microseconds so Opus' 2.5 ms frame is expressible at @@ -331,7 +338,7 @@ final class SmokeTests: XCTestCase { XCTFail("2 ms is not an opus frame duration: \(error)") } - try broadcast.finish() + try broadcast.close() } /// The decode side picks its CPU layout: an unset `format` is I420, and RGBA @@ -388,7 +395,7 @@ final class SmokeTests: XCTestCase { XCTAssertTrue(stride(from: 3, to: frame.data.count, by: 4).allSatisfy { frame.data[$0] == 0xFF }) try video.finish() - try broadcast.finish() + try broadcast.close() } func testEncodeAudioWithOpusObject() throws { @@ -408,7 +415,7 @@ final class SmokeTests: XCTestCase { try producer.write(silence) XCTAssertEqual(try producer.name, "mic") try producer.finish() - try broadcast.finish() + try broadcast.close() } // Release the config before finishing: the producer retains what it needs. @@ -422,7 +429,7 @@ final class SmokeTests: XCTestCase { } try producer.write(silence) try producer.finish() - try broadcast.finish() + try broadcast.close() } } }