diff --git a/.github/workflows/android.yml b/.github/workflows/android.yml index 881d127686..9a82759c2f 100644 --- a/.github/workflows/android.yml +++ b/.github/workflows/android.yml @@ -31,7 +31,6 @@ on: - "rs/moq-loc/**" - "rs/moq-msf/**" - "rs/moq-json/**" - - "rs/moq-binary/**" - "rs/moq-flate/**" - "rs/moq-pattern/**" - "rs/hang/**" diff --git a/Cargo.lock b/Cargo.lock index 49dd90fa96..c5b2f64b6b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4230,18 +4230,6 @@ dependencies = [ "usage-rs", ] -[[package]] -name = "moq-binary" -version = "0.1.7" -dependencies = [ - "bytes", - "kio 0.6.1", - "moq-flate", - "moq-net", - "thiserror 2.0.21", - "tracing", -] - [[package]] name = "moq-boy" version = "0.5.8" @@ -4374,7 +4362,10 @@ version = "0.2.0" dependencies = [ "bytes", "flate2", + "kio 0.6.1", + "moq-net", "thiserror 2.0.21", + "tracing", ] [[package]] @@ -4467,7 +4458,7 @@ dependencies = [ "hang", "kio 0.6.1", "memchr", - "moq-binary", + "moq-flate", "moq-json", "moq-loc", "moq-msf", diff --git a/Cargo.toml b/Cargo.toml index 0d602ce5fc..e58b93783a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -7,7 +7,6 @@ members = [ "rs/moq-audio", "rs/moq-auth", "rs/moq-bench", - "rs/moq-binary", "rs/moq-boy", "rs/moq-c", "rs/moq-cli", @@ -50,7 +49,6 @@ default-members = [ "rs/moq-boy", "rs/moq-audio", "rs/moq-bench", - "rs/moq-binary", "rs/moq-cli", "rs/moq-e2ee", # "rs/moq-ffi", # requires Python/maturin @@ -140,7 +138,6 @@ loom = { version = "0.7.2", features = ["futures"] } mdns-sd = { version = "0.21", features = ["async"] } moq-audio = { version = "0.1.7", path = "rs/moq-audio", default-features = false } moq-auth = { version = "0.1.5", path = "rs/moq-auth" } -moq-binary = { version = "0.1.7", path = "rs/moq-binary" } moq-flate = { version = "0.2.0", path = "rs/moq-flate" } moq-hls = { version = "0.5.8", path = "rs/moq-hls", default-features = false } moq-json = { version = "0.5.5", path = "rs/moq-json" } diff --git a/bun.lock b/bun.lock index 661ebc2605..be29b2b19b 100644 --- a/bun.lock +++ b/bun.lock @@ -93,20 +93,6 @@ "typescript": "7.0.2", }, }, - "js/binary": { - "name": "@moq/binary", - "version": "0.2.1", - "dependencies": { - "@moq/flate": "workspace:^", - "@moq/net": "workspace:^", - "@moq/signals": "workspace:^", - }, - "devDependencies": { - "@types/bun": "^1.4.2", - "rimraf": "^6.1.3", - "typescript": "7.0.2", - }, - }, "js/clock": { "name": "@moq/clock", "version": "0.1.0", @@ -127,6 +113,8 @@ "name": "@moq/flate", "version": "0.1.3", "dependencies": { + "@moq/net": "workspace:^", + "@moq/signals": "workspace:^", "pako": "^3.0.2", }, "devDependencies": { @@ -642,8 +630,6 @@ "@moq/auth": ["@moq/auth@workspace:js/auth"], - "@moq/binary": ["@moq/binary@workspace:js/binary"], - "@moq/boy": ["@moq/boy@workspace:js/moq-boy"], "@moq/clock": ["@moq/clock@workspace:js/clock"], diff --git a/dart/moq_ffi/lib/src/moq.dart b/dart/moq_ffi/lib/src/moq.dart index 1142719a15..785297f956 100644 --- a/dart/moq_ffi/lib/src/moq.dart +++ b/dart/moq_ffi/lib/src/moq.dart @@ -8,69 +8,12 @@ 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"; -class MoqBinaryConfig { - final bool compression; - final String? mime; - MoqBinaryConfig({this.compression = false, this.mime = null}); -} - -class FfiConverterMoqBinaryConfig { - static MoqBinaryConfig lift(RustBuffer buf) { - return FfiConverterMoqBinaryConfig.read(buf.asUint8List()).value; - } - - static LiftRetVal read(Uint8List buf) { - int new_offset = buf.offsetInBytes; - final compression_lifted = FfiConverterBool.read( - Uint8List.view(buf.buffer, new_offset), - ); - final compression = compression_lifted.value; - new_offset += compression_lifted.bytesRead; - final mime_lifted = FfiConverterOptionalString.read( - Uint8List.view(buf.buffer, new_offset), - ); - final mime = mime_lifted.value; - new_offset += mime_lifted.bytesRead; - return LiftRetVal( - MoqBinaryConfig(compression: compression, mime: mime), - new_offset - buf.offsetInBytes, - ); - } - - static RustBuffer lower(MoqBinaryConfig value) { - final total_length = - FfiConverterBool.allocationSize(value.compression) + - FfiConverterOptionalString.allocationSize(value.mime) + - 0; - final buf = Uint8List(total_length); - write(value, buf); - return toRustBuffer(buf); - } - - static int write(MoqBinaryConfig value, Uint8List buf) { - int new_offset = buf.offsetInBytes; - new_offset += FfiConverterBool.write( - value.compression, - Uint8List.view(buf.buffer, new_offset), - ); - new_offset += FfiConverterOptionalString.write( - value.mime, - Uint8List.view(buf.buffer, new_offset), - ); - return new_offset - buf.offsetInBytes; - } - - static int allocationSize(MoqBinaryConfig value) { - return FfiConverterBool.allocationSize(value.compression) + - FfiConverterOptionalString.allocationSize(value.mime) + - 0; - } -} - class MoqFetchGroupOptions { final int priority; MoqFetchGroupOptions({this.priority = 0}); @@ -301,6 +244,65 @@ class FfiConverterMoqProtocolError { } } +class MoqFlateConfig { + final bool compression; + final String? mime; + MoqFlateConfig({this.compression = false, this.mime = null}); +} + +class FfiConverterMoqFlateConfig { + static MoqFlateConfig lift(RustBuffer buf) { + return FfiConverterMoqFlateConfig.read(buf.asUint8List()).value; + } + + static LiftRetVal read(Uint8List buf) { + int new_offset = buf.offsetInBytes; + final compression_lifted = FfiConverterBool.read( + Uint8List.view(buf.buffer, new_offset), + ); + final compression = compression_lifted.value; + new_offset += compression_lifted.bytesRead; + final mime_lifted = FfiConverterOptionalString.read( + Uint8List.view(buf.buffer, new_offset), + ); + final mime = mime_lifted.value; + new_offset += mime_lifted.bytesRead; + return LiftRetVal( + MoqFlateConfig(compression: compression, mime: mime), + new_offset - buf.offsetInBytes, + ); + } + + static RustBuffer lower(MoqFlateConfig value) { + final total_length = + FfiConverterBool.allocationSize(value.compression) + + FfiConverterOptionalString.allocationSize(value.mime) + + 0; + final buf = Uint8List(total_length); + write(value, buf); + return toRustBuffer(buf); + } + + static int write(MoqFlateConfig value, Uint8List buf) { + int new_offset = buf.offsetInBytes; + new_offset += FfiConverterBool.write( + value.compression, + Uint8List.view(buf.buffer, new_offset), + ); + new_offset += FfiConverterOptionalString.write( + value.mime, + Uint8List.view(buf.buffer, new_offset), + ); + return new_offset - buf.offsetInBytes; + } + + static int allocationSize(MoqFlateConfig value) { + return FfiConverterBool.allocationSize(value.compression) + + FfiConverterOptionalString.allocationSize(value.mime) + + 0; + } +} + class MoqJsonSnapshotConfig { final int deltaRatio; final bool compression; @@ -4332,164 +4334,6 @@ class FfiConverterMoqReservation { } } -abstract class MoqBinarySnapshotProducerInterface { - void finish(); - void update({required Uint8List payload}); -} - -final _MoqBinarySnapshotProducerFinalizer = Finalizer>((ptr) { - rustCall( - (status) => uniffi_moq_ffi_fn_free_moqbinarysnapshotproducer(ptr, status), - ); -}); - -class MoqBinarySnapshotProducer implements MoqBinarySnapshotProducerInterface { - late final Pointer _ptr; - MoqBinarySnapshotProducer._(this._ptr) { - _MoqBinarySnapshotProducerFinalizer.attach(this, _ptr, detach: this); - } - factory MoqBinarySnapshotProducer.lift(Pointer ptr) { - return MoqBinarySnapshotProducer._(ptr); - } - Pointer uniffiClonePointer() { - return rustCall( - (status) => - uniffi_moq_ffi_fn_clone_moqbinarysnapshotproducer(_ptr, status), - ); - } - - void dispose() { - _MoqBinarySnapshotProducerFinalizer.detach(this); - rustCall( - (status) => - uniffi_moq_ffi_fn_free_moqbinarysnapshotproducer(_ptr, status), - ); - } - - void finish() { - return rustCall((status) { - uniffi_moq_ffi_fn_method_moqbinarysnapshotproducer_finish( - uniffiClonePointer(), - status, - ); - }, moqExceptionErrorHandler); - } - - void update({required Uint8List payload}) { - return rustCall((status) { - uniffi_moq_ffi_fn_method_moqbinarysnapshotproducer_update( - uniffiClonePointer(), - FfiConverterUint8List.lower(payload), - status, - ); - }, moqExceptionErrorHandler); - } -} - -class FfiConverterMoqBinarySnapshotProducer { - static MoqBinarySnapshotProducer lift(Pointer ptr) { - return MoqBinarySnapshotProducer.lift(ptr); - } - - static Pointer lower(MoqBinarySnapshotProducer value) { - return value.uniffiClonePointer(); - } - - static int allocationSize(MoqBinarySnapshotProducer value) { - return 8; - } - - static LiftRetVal read(Uint8List buf) { - final handle = buf.buffer.asByteData(buf.offsetInBytes).getInt64(0); - final pointer = Pointer.fromAddress(handle); - return LiftRetVal(MoqBinarySnapshotProducer.lift(pointer), 8); - } - - static int write(MoqBinarySnapshotProducer value, Uint8List buf) { - final handle = lower(value); - buf.buffer.asByteData(buf.offsetInBytes).setInt64(0, handle.address); - return 8; - } -} - -abstract class MoqBinaryStreamProducerInterface { - void append({required Uint8List payload}); - void finish(); -} - -final _MoqBinaryStreamProducerFinalizer = Finalizer>((ptr) { - rustCall( - (status) => uniffi_moq_ffi_fn_free_moqbinarystreamproducer(ptr, status), - ); -}); - -class MoqBinaryStreamProducer implements MoqBinaryStreamProducerInterface { - late final Pointer _ptr; - MoqBinaryStreamProducer._(this._ptr) { - _MoqBinaryStreamProducerFinalizer.attach(this, _ptr, detach: this); - } - factory MoqBinaryStreamProducer.lift(Pointer ptr) { - return MoqBinaryStreamProducer._(ptr); - } - Pointer uniffiClonePointer() { - return rustCall( - (status) => uniffi_moq_ffi_fn_clone_moqbinarystreamproducer(_ptr, status), - ); - } - - void dispose() { - _MoqBinaryStreamProducerFinalizer.detach(this); - rustCall( - (status) => uniffi_moq_ffi_fn_free_moqbinarystreamproducer(_ptr, status), - ); - } - - void append({required Uint8List payload}) { - return rustCall((status) { - uniffi_moq_ffi_fn_method_moqbinarystreamproducer_append( - uniffiClonePointer(), - FfiConverterUint8List.lower(payload), - status, - ); - }, moqExceptionErrorHandler); - } - - void finish() { - return rustCall((status) { - uniffi_moq_ffi_fn_method_moqbinarystreamproducer_finish( - uniffiClonePointer(), - status, - ); - }, moqExceptionErrorHandler); - } -} - -class FfiConverterMoqBinaryStreamProducer { - static MoqBinaryStreamProducer lift(Pointer ptr) { - return MoqBinaryStreamProducer.lift(ptr); - } - - static Pointer lower(MoqBinaryStreamProducer value) { - return value.uniffiClonePointer(); - } - - static int allocationSize(MoqBinaryStreamProducer value) { - return 8; - } - - static LiftRetVal read(Uint8List buf) { - final handle = buf.buffer.asByteData(buf.offsetInBytes).getInt64(0); - final pointer = Pointer.fromAddress(handle); - return LiftRetVal(MoqBinaryStreamProducer.lift(pointer), 8); - } - - static int write(MoqBinaryStreamProducer value, Uint8List buf) { - final handle = lower(value); - buf.buffer.asByteData(buf.offsetInBytes).setInt64(0, handle.address); - return 8; - } -} - abstract class MoqBroadcastConsumerInterface { Future fetchGroup({ required String name, @@ -5313,6 +5157,163 @@ class FfiConverterMoqTrackDemand { } } +abstract class MoqFlateSnapshotProducerInterface { + void finish(); + void update({required Uint8List payload}); +} + +final _MoqFlateSnapshotProducerFinalizer = Finalizer>((ptr) { + rustCall( + (status) => uniffi_moq_ffi_fn_free_moqflatesnapshotproducer(ptr, status), + ); +}); + +class MoqFlateSnapshotProducer implements MoqFlateSnapshotProducerInterface { + late final Pointer _ptr; + MoqFlateSnapshotProducer._(this._ptr) { + _MoqFlateSnapshotProducerFinalizer.attach(this, _ptr, detach: this); + } + factory MoqFlateSnapshotProducer.lift(Pointer ptr) { + return MoqFlateSnapshotProducer._(ptr); + } + Pointer uniffiClonePointer() { + return rustCall( + (status) => + uniffi_moq_ffi_fn_clone_moqflatesnapshotproducer(_ptr, status), + ); + } + + void dispose() { + _MoqFlateSnapshotProducerFinalizer.detach(this); + rustCall( + (status) => uniffi_moq_ffi_fn_free_moqflatesnapshotproducer(_ptr, status), + ); + } + + void finish() { + return rustCall((status) { + uniffi_moq_ffi_fn_method_moqflatesnapshotproducer_finish( + uniffiClonePointer(), + status, + ); + }, moqExceptionErrorHandler); + } + + void update({required Uint8List payload}) { + return rustCall((status) { + uniffi_moq_ffi_fn_method_moqflatesnapshotproducer_update( + uniffiClonePointer(), + FfiConverterUint8List.lower(payload), + status, + ); + }, moqExceptionErrorHandler); + } +} + +class FfiConverterMoqFlateSnapshotProducer { + static MoqFlateSnapshotProducer lift(Pointer ptr) { + return MoqFlateSnapshotProducer.lift(ptr); + } + + static Pointer lower(MoqFlateSnapshotProducer value) { + return value.uniffiClonePointer(); + } + + static int allocationSize(MoqFlateSnapshotProducer value) { + return 8; + } + + static LiftRetVal read(Uint8List buf) { + final handle = buf.buffer.asByteData(buf.offsetInBytes).getInt64(0); + final pointer = Pointer.fromAddress(handle); + return LiftRetVal(MoqFlateSnapshotProducer.lift(pointer), 8); + } + + static int write(MoqFlateSnapshotProducer value, Uint8List buf) { + final handle = lower(value); + buf.buffer.asByteData(buf.offsetInBytes).setInt64(0, handle.address); + return 8; + } +} + +abstract class MoqFlateStreamProducerInterface { + void append({required Uint8List payload}); + void finish(); +} + +final _MoqFlateStreamProducerFinalizer = Finalizer>((ptr) { + rustCall( + (status) => uniffi_moq_ffi_fn_free_moqflatestreamproducer(ptr, status), + ); +}); + +class MoqFlateStreamProducer implements MoqFlateStreamProducerInterface { + late final Pointer _ptr; + MoqFlateStreamProducer._(this._ptr) { + _MoqFlateStreamProducerFinalizer.attach(this, _ptr, detach: this); + } + factory MoqFlateStreamProducer.lift(Pointer ptr) { + return MoqFlateStreamProducer._(ptr); + } + Pointer uniffiClonePointer() { + return rustCall( + (status) => uniffi_moq_ffi_fn_clone_moqflatestreamproducer(_ptr, status), + ); + } + + void dispose() { + _MoqFlateStreamProducerFinalizer.detach(this); + rustCall( + (status) => uniffi_moq_ffi_fn_free_moqflatestreamproducer(_ptr, status), + ); + } + + void append({required Uint8List payload}) { + return rustCall((status) { + uniffi_moq_ffi_fn_method_moqflatestreamproducer_append( + uniffiClonePointer(), + FfiConverterUint8List.lower(payload), + status, + ); + }, moqExceptionErrorHandler); + } + + void finish() { + return rustCall((status) { + uniffi_moq_ffi_fn_method_moqflatestreamproducer_finish( + uniffiClonePointer(), + status, + ); + }, moqExceptionErrorHandler); + } +} + +class FfiConverterMoqFlateStreamProducer { + static MoqFlateStreamProducer lift(Pointer ptr) { + return MoqFlateStreamProducer.lift(ptr); + } + + static Pointer lower(MoqFlateStreamProducer value) { + return value.uniffiClonePointer(); + } + + static int allocationSize(MoqFlateStreamProducer value) { + return 8; + } + + static LiftRetVal read(Uint8List buf) { + final handle = buf.buffer.asByteData(buf.offsetInBytes).getInt64(0); + final pointer = Pointer.fromAddress(handle); + return LiftRetVal(MoqFlateStreamProducer.lift(pointer), 8); + } + + static int write(MoqFlateStreamProducer value, Uint8List buf) { + final handle = lower(value); + buf.buffer.asByteData(buf.offsetInBytes).setInt64(0, handle.address); + return 8; + } +} + abstract class MoqJsonSnapshotConsumerInterface { void cancel(); Future next(); @@ -6276,13 +6277,13 @@ class FfiConverterMoqBroadcastDynamic { } abstract class MoqBroadcastProducerInterface { - MoqBinarySnapshotProducer publishBinarySnapshot({ + MoqFlateSnapshotProducer publishFlateSnapshot({ required String name, - required MoqBinaryConfig config, + required MoqFlateConfig config, }); - MoqBinaryStreamProducer publishBinaryStream({ + MoqFlateStreamProducer publishFlateStream({ required String name, - required MoqBinaryConfig config, + required MoqFlateConfig config, }); MoqJsonSnapshotProducer publishJsonSnapshot({ required String name, @@ -6356,36 +6357,36 @@ class MoqBroadcastProducer implements MoqBroadcastProducerInterface { ); } - MoqBinarySnapshotProducer publishBinarySnapshot({ + MoqFlateSnapshotProducer publishFlateSnapshot({ required String name, - required MoqBinaryConfig config, + required MoqFlateConfig config, }) { return rustCallWithLifter( (status) => - uniffi_moq_ffi_fn_method_moqbroadcastproducer_publish_binary_snapshot( + uniffi_moq_ffi_fn_method_moqbroadcastproducer_publish_flate_snapshot( uniffiClonePointer(), FfiConverterString.lower(name), - FfiConverterMoqBinaryConfig.lower(config), + FfiConverterMoqFlateConfig.lower(config), status, ), - FfiConverterMoqBinarySnapshotProducer.lift, + FfiConverterMoqFlateSnapshotProducer.lift, moqExceptionErrorHandler, ); } - MoqBinaryStreamProducer publishBinaryStream({ + MoqFlateStreamProducer publishFlateStream({ required String name, - required MoqBinaryConfig config, + required MoqFlateConfig config, }) { return rustCallWithLifter( (status) => - uniffi_moq_ffi_fn_method_moqbroadcastproducer_publish_binary_stream( + uniffi_moq_ffi_fn_method_moqbroadcastproducer_publish_flate_stream( uniffiClonePointer(), FfiConverterString.lower(name), - FfiConverterMoqBinaryConfig.lower(config), + FfiConverterMoqFlateConfig.lower(config), status, ), - FfiConverterMoqBinaryStreamProducer.lift, + FfiConverterMoqFlateStreamProducer.lift, moqExceptionErrorHandler, ); } @@ -9755,72 +9756,6 @@ external void uniffi_moq_ffi_fn_method_moqreservation_update( Pointer uniffiStatus, ); -@Native Function(Pointer, Pointer)>( - assetId: _uniffiAssetId, -) -external Pointer uniffi_moq_ffi_fn_clone_moqbinarysnapshotproducer( - Pointer handle, - Pointer uniffiStatus, -); - -@Native, Pointer)>( - assetId: _uniffiAssetId, -) -external void uniffi_moq_ffi_fn_free_moqbinarysnapshotproducer( - Pointer handle, - Pointer uniffiStatus, -); - -@Native, Pointer)>( - assetId: _uniffiAssetId, -) -external void uniffi_moq_ffi_fn_method_moqbinarysnapshotproducer_finish( - Pointer ptr, - Pointer uniffiStatus, -); - -@Native, RustBuffer, Pointer)>( - assetId: _uniffiAssetId, -) -external void uniffi_moq_ffi_fn_method_moqbinarysnapshotproducer_update( - Pointer ptr, - RustBuffer payload, - Pointer uniffiStatus, -); - -@Native Function(Pointer, Pointer)>( - assetId: _uniffiAssetId, -) -external Pointer uniffi_moq_ffi_fn_clone_moqbinarystreamproducer( - Pointer handle, - Pointer uniffiStatus, -); - -@Native, Pointer)>( - assetId: _uniffiAssetId, -) -external void uniffi_moq_ffi_fn_free_moqbinarystreamproducer( - Pointer handle, - Pointer uniffiStatus, -); - -@Native, RustBuffer, Pointer)>( - assetId: _uniffiAssetId, -) -external void uniffi_moq_ffi_fn_method_moqbinarystreamproducer_append( - Pointer ptr, - RustBuffer payload, - Pointer uniffiStatus, -); - -@Native, Pointer)>( - assetId: _uniffiAssetId, -) -external void uniffi_moq_ffi_fn_method_moqbinarystreamproducer_finish( - Pointer ptr, - Pointer uniffiStatus, -); - @Native Function(Pointer, Pointer)>( assetId: _uniffiAssetId, ) @@ -10156,6 +10091,72 @@ external Pointer uniffi_moq_ffi_fn_method_moqtrackdemand_used( Pointer ptr, ); +@Native Function(Pointer, Pointer)>( + assetId: _uniffiAssetId, +) +external Pointer uniffi_moq_ffi_fn_clone_moqflatesnapshotproducer( + Pointer handle, + Pointer uniffiStatus, +); + +@Native, Pointer)>( + assetId: _uniffiAssetId, +) +external void uniffi_moq_ffi_fn_free_moqflatesnapshotproducer( + Pointer handle, + Pointer uniffiStatus, +); + +@Native, Pointer)>( + assetId: _uniffiAssetId, +) +external void uniffi_moq_ffi_fn_method_moqflatesnapshotproducer_finish( + Pointer ptr, + Pointer uniffiStatus, +); + +@Native, RustBuffer, Pointer)>( + assetId: _uniffiAssetId, +) +external void uniffi_moq_ffi_fn_method_moqflatesnapshotproducer_update( + Pointer ptr, + RustBuffer payload, + Pointer uniffiStatus, +); + +@Native Function(Pointer, Pointer)>( + assetId: _uniffiAssetId, +) +external Pointer uniffi_moq_ffi_fn_clone_moqflatestreamproducer( + Pointer handle, + Pointer uniffiStatus, +); + +@Native, Pointer)>( + assetId: _uniffiAssetId, +) +external void uniffi_moq_ffi_fn_free_moqflatestreamproducer( + Pointer handle, + Pointer uniffiStatus, +); + +@Native, RustBuffer, Pointer)>( + assetId: _uniffiAssetId, +) +external void uniffi_moq_ffi_fn_method_moqflatestreamproducer_append( + Pointer ptr, + RustBuffer payload, + Pointer uniffiStatus, +); + +@Native, Pointer)>( + assetId: _uniffiAssetId, +) +external void uniffi_moq_ffi_fn_method_moqflatestreamproducer_finish( + Pointer ptr, + Pointer uniffiStatus, +); + @Native Function(Pointer, Pointer)>( assetId: _uniffiAssetId, ) @@ -10596,7 +10597,7 @@ external Pointer uniffi_moq_ffi_fn_constructor_moqbroadcastproducer_new( ) >(assetId: _uniffiAssetId) external Pointer -uniffi_moq_ffi_fn_method_moqbroadcastproducer_publish_binary_snapshot( +uniffi_moq_ffi_fn_method_moqbroadcastproducer_publish_flate_snapshot( Pointer ptr, RustBuffer name, RustBuffer config, @@ -10612,7 +10613,7 @@ uniffi_moq_ffi_fn_method_moqbroadcastproducer_publish_binary_snapshot( ) >(assetId: _uniffiAssetId) external Pointer -uniffi_moq_ffi_fn_method_moqbroadcastproducer_publish_binary_stream( +uniffi_moq_ffi_fn_method_moqbroadcastproducer_publish_flate_stream( Pointer ptr, RustBuffer name, RustBuffer config, @@ -12115,18 +12116,6 @@ external int uniffi_moq_ffi_checksum_method_moqreservation_grant(); @Native(assetId: _uniffiAssetId) external int uniffi_moq_ffi_checksum_method_moqreservation_update(); -@Native(assetId: _uniffiAssetId) -external int uniffi_moq_ffi_checksum_method_moqbinarysnapshotproducer_finish(); - -@Native(assetId: _uniffiAssetId) -external int uniffi_moq_ffi_checksum_method_moqbinarysnapshotproducer_update(); - -@Native(assetId: _uniffiAssetId) -external int uniffi_moq_ffi_checksum_method_moqbinarystreamproducer_append(); - -@Native(assetId: _uniffiAssetId) -external int uniffi_moq_ffi_checksum_method_moqbinarystreamproducer_finish(); - @Native(assetId: _uniffiAssetId) external int uniffi_moq_ffi_checksum_method_moqbroadcastconsumer_fetch_group(); @@ -12220,6 +12209,18 @@ external int uniffi_moq_ffi_checksum_method_moqtrackdemand_unused(); @Native(assetId: _uniffiAssetId) external int uniffi_moq_ffi_checksum_method_moqtrackdemand_used(); +@Native(assetId: _uniffiAssetId) +external int uniffi_moq_ffi_checksum_method_moqflatesnapshotproducer_finish(); + +@Native(assetId: _uniffiAssetId) +external int uniffi_moq_ffi_checksum_method_moqflatesnapshotproducer_update(); + +@Native(assetId: _uniffiAssetId) +external int uniffi_moq_ffi_checksum_method_moqflatestreamproducer_append(); + +@Native(assetId: _uniffiAssetId) +external int uniffi_moq_ffi_checksum_method_moqflatestreamproducer_finish(); + @Native(assetId: _uniffiAssetId) external int uniffi_moq_ffi_checksum_method_moqjsonsnapshotconsumer_cancel(); @@ -12311,11 +12312,11 @@ uniffi_moq_ffi_checksum_method_moqbroadcastdynamic_requested_track(); @Native(assetId: _uniffiAssetId) external int -uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_binary_snapshot(); +uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_flate_snapshot(); @Native(assetId: _uniffiAssetId) external int -uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_binary_stream(); +uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_flate_stream(); @Native(assetId: _uniffiAssetId) external int @@ -12692,21 +12693,6 @@ void _checkApiChecksums() { if (uniffi_moq_ffi_checksum_method_moqreservation_update() != 9626) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } - if (uniffi_moq_ffi_checksum_method_moqbinarysnapshotproducer_finish() != - 10338) { - throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); - } - if (uniffi_moq_ffi_checksum_method_moqbinarysnapshotproducer_update() != - 56077) { - throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); - } - if (uniffi_moq_ffi_checksum_method_moqbinarystreamproducer_append() != 1645) { - throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); - } - if (uniffi_moq_ffi_checksum_method_moqbinarystreamproducer_finish() != - 60630) { - throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); - } if (uniffi_moq_ffi_checksum_method_moqbroadcastconsumer_fetch_group() != 18633) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); @@ -12803,6 +12789,20 @@ void _checkApiChecksums() { if (uniffi_moq_ffi_checksum_method_moqtrackdemand_used() != 18944) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } + if (uniffi_moq_ffi_checksum_method_moqflatesnapshotproducer_finish() != + 52594) { + throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); + } + if (uniffi_moq_ffi_checksum_method_moqflatesnapshotproducer_update() != + 28820) { + throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); + } + if (uniffi_moq_ffi_checksum_method_moqflatestreamproducer_append() != 38961) { + throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); + } + if (uniffi_moq_ffi_checksum_method_moqflatestreamproducer_finish() != 6982) { + throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); + } if (uniffi_moq_ffi_checksum_method_moqjsonsnapshotconsumer_cancel() != 45114) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); @@ -12897,12 +12897,12 @@ void _checkApiChecksums() { 24118) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } - if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_binary_snapshot() != - 6748) { + if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_flate_snapshot() != + 8575) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } - if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_binary_stream() != - 58418) { + if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_flate_stream() != + 56158) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } if (uniffi_moq_ffi_checksum_method_moqbroadcastproducer_publish_json_snapshot() != diff --git a/doc/.vitepress/config.ts b/doc/.vitepress/config.ts index 9c736e282c..bb6d10e60c 100644 --- a/doc/.vitepress/config.ts +++ b/doc/.vitepress/config.ts @@ -175,7 +175,7 @@ export default defineConfig({ { text: "moq-auth", link: "/lib/rs/moq-auth" }, { text: "moq-room", link: "/lib/rs/moq-room" }, { text: "moq-json", link: "/lib/rs/moq-json" }, - { text: "moq-binary", link: "/lib/rs/moq-binary" }, + { text: "moq-flate", link: "/lib/rs/moq-flate" }, ], }, { @@ -190,7 +190,7 @@ export default defineConfig({ { text: "@moq/auth", link: "/lib/js/auth" }, { text: "@moq/signals", link: "/lib/js/signals" }, { text: "@moq/json", link: "/lib/js/json" }, - { text: "@moq/binary", link: "/lib/js/binary" }, + { text: "@moq/flate", link: "/lib/js/flate" }, ], }, { text: "Swift", link: "/lib/swift/" }, diff --git a/doc/concept/hang.md b/doc/concept/hang.md index 39ec8b972f..c14d6685bf 100644 --- a/doc/concept/hang.md +++ b/doc/concept/hang.md @@ -117,7 +117,7 @@ retracts it when the producer drops. Read the config from `catalog.json.tracks` or `catalog.binary.tracks`, then pair its name and config with `moq_mux::catalog::Entry::new` to subscribe. In C, `moq_publish_json_*` and `moq_publish_binary_*` do the same, retracting on `_finish`. In the browser, read the same map, -subscribe by name, and hand the track to `@moq/json` or `@moq/binary`. +subscribe by name, and hand the track to `@moq/json` or `@moq/flate`. An application with its own per-track fields can list a data track in its own root section instead, flattening the JSON or binary entry beside those fields diff --git a/doc/lib/js/binary.md b/doc/lib/js/binary.md deleted file mode 100644 index d761026b52..0000000000 --- a/doc/lib/js/binary.md +++ /dev/null @@ -1,32 +0,0 @@ ---- -title: "@moq/binary" -description: Opaque binary payloads over MoQ tracks ---- - -# @moq/binary - -[![npm](https://img.shields.io/npm/v/@moq/binary)](https://www.npmjs.com/package/@moq/binary) - -Opaque payloads over [`@moq/net`](/lib/js/net) tracks, in two modes: - -- **Snapshot**: lossy latest-value, one payload per group. -- **Stream**: lossless append-log in a single group. - -The bytes are never inspected. Compression is [`@moq/flate`](https://www.npmjs.com/package/@moq/flate), -the same group-scoped DEFLATE `@moq/json` uses. - -`Config` is the codec options. `Producer.Config` / `Consumer.Config` add the -track. Compression is a shared `"none" | "deflate"` enum, not a boolean: both -sides set the same field. - -```ts -import { Snapshot } from "@moq/binary"; - -const producer = new Snapshot.Producer({ track, compression: "deflate" }); -producer.update(payload); -``` - -A payload is stamped when written, unless you pass its capture time: -`producer.update(payload, at)`. - -The Rust twin is [`moq-binary`](/lib/rs/moq-binary). diff --git a/doc/lib/js/flate.md b/doc/lib/js/flate.md new file mode 100644 index 0000000000..a128b85098 --- /dev/null +++ b/doc/lib/js/flate.md @@ -0,0 +1,33 @@ +--- +title: "@moq/flate" +description: Opaque MoQ tracks, optionally compressed with group-scoped DEFLATE +--- + +# @moq/flate + +[![npm](https://img.shields.io/npm/v/@moq/flate)](https://www.npmjs.com/package/@moq/flate) + +Opaque payloads over [`@moq/net`](/lib/js/net) tracks, in two modes: + +- **Snapshot**: lossy latest-value, one payload per group. +- **Stream**: lossless append-log in a single group. + +The bytes are never inspected. Compression is opt-in: `"none"` (the default) +or `"deflate"`, and both sides set the same field. With `"deflate"`, each +group is one raw DEFLATE stream sync-flushed at every frame boundary, the same +group-scoped DEFLATE `@moq/json` uses. + +```ts +import { Snapshot } from "@moq/flate"; + +const producer = new Snapshot.Producer({ track, compression: "deflate" }); +producer.update(payload); +``` + +A payload is stamped when written, unless you pass its capture time: +`producer.update(payload, at)`. + +The codec underneath is exported as `Encoder`/`Decoder`. Create one pair per +group and feed frames in order. + +The Rust twin is [`moq-flate`](/lib/rs/moq-flate). diff --git a/doc/lib/js/index.md b/doc/lib/js/index.md index 9de77dbcc8..1313a7c24b 100644 --- a/doc/lib/js/index.md +++ b/doc/lib/js/index.md @@ -21,8 +21,7 @@ and WebAudio. `@moq/net` also runs in Node, Bun, and Deno. | [@moq/auth](/lib/js/auth) | Validate a relay's request, build its grant, and mint and verify relay JWTs. | | [@moq/signals](/lib/js/signals) | The reactive primitives every package exposes its state through. | | [@moq/json](/lib/js/json) | JSON over tracks: snapshots with merge-patch deltas, or append logs. | -| [@moq/binary](/lib/js/binary) | Opaque payloads over tracks: snapshots or append logs. | -| [@moq/flate](https://www.npmjs.com/package/@moq/flate) | Group-scoped DEFLATE for any track. | +| [@moq/flate](/lib/js/flate) | Opaque payloads over tracks, optionally compressed with group-scoped DEFLATE: snapshots or append logs. | | [@moq/loc](https://www.npmjs.com/package/@moq/loc), [@moq/msf](https://www.npmjs.com/package/@moq/msf) | The IETF LOC container and MSF catalog. | | [@moq/boy](https://www.npmjs.com/package/@moq/boy) | The [MoQ Boy](/bin/demo) player element. | diff --git a/doc/lib/rs/index.md b/doc/lib/rs/index.md index 37feac3fab..bff79ce491 100644 --- a/doc/lib/rs/index.md +++ b/doc/lib/rs/index.md @@ -27,9 +27,8 @@ The reference implementation. Every crate is on | [moq-auth](/lib/rs/moq-auth) | The authorization contract: requests, grants, leases, the HTTP client, the reference server, JWT keys, signing, and verification, plus listing live sessions and pushing a re-check. | | [moq-room](/lib/rs/moq-room) | Headless rooms: announce-derived roster, token claims, and a chat track. | | [moq-json](/lib/rs/moq-json) | JSON over tracks: snapshots with merge-patch deltas, or append logs. | -| [moq-binary](/lib/rs/moq-binary) | Opaque payloads over tracks: snapshots or append logs. | +| [moq-flate](/lib/rs/moq-flate) | Opaque payloads over tracks, optionally compressed with group-scoped DEFLATE: snapshots or append logs. | | [moq-e2ee](https://docs.rs/moq-e2ee) | End-to-end encryption of groups, datagrams, and track names, scoped to a publisher epoch. | -| [moq-flate](https://docs.rs/moq-flate) | Group-scoped DEFLATE for any track. | | [moq-loc](https://docs.rs/moq-loc), [moq-msf](https://docs.rs/moq-msf) | The IETF LOC container and MSF catalog. | | [moq-stats](https://docs.rs/moq-stats) | Publish and consume relay traffic counters as tracks. | | [moq-hls](https://docs.rs/moq-hls), [moq-rtmp](https://docs.rs/moq-rtmp), [moq-srt](https://docs.rs/moq-srt), [moq-rtc](https://docs.rs/moq-rtc) | The [gateways](/bin/), as libraries you can embed with your own auth. | diff --git a/doc/lib/rs/moq-binary.md b/doc/lib/rs/moq-binary.md deleted file mode 100644 index f1b9a0d6e3..0000000000 --- a/doc/lib/rs/moq-binary.md +++ /dev/null @@ -1,36 +0,0 @@ ---- -title: moq-binary -description: Opaque binary payloads over MoQ tracks ---- - -# moq-binary - -[![crates.io](https://img.shields.io/crates/v/moq-binary)](https://crates.io/crates/moq-binary) -[![docs.rs](https://docs.rs/moq-binary/badge.svg)](https://docs.rs/moq-binary) - -Opaque payloads over [`moq-net`](/lib/rs/moq-net) tracks, in two modes: - -- **snapshot**: lossy latest-value, one payload per group. -- **stream**: lossless append-log in a single group. - -The bytes are never inspected. Compression is [`moq-flate`](https://docs.rs/moq-flate), -the same group-scoped DEFLATE `moq-json` uses. - -`Config` is the codec options. The track-owning pair is `producer::Config` / -`consumer::Config`. Compression is a shared `Compression` enum (`None` or -`Deflate`), not a bool: both sides set the same field. - -```rust -let mut config = moq_binary::snapshot::Config::default(); -config.compression = moq_binary::Compression::Deflate; - -let mut producer = moq_binary::snapshot::Producer::new(track, config); -producer.update(payload)?; -``` - -A payload is stamped when written, unless it carries its capture time: -`moq_net::Timed::from(bytes).at(captured)`. Writes return the encoded frame -size. - -The TypeScript twin is [`@moq/binary`](/lib/js/binary). API: -[docs.rs/moq-binary](https://docs.rs/moq-binary). diff --git a/doc/lib/rs/moq-flate.md b/doc/lib/rs/moq-flate.md new file mode 100644 index 0000000000..5996d62539 --- /dev/null +++ b/doc/lib/rs/moq-flate.md @@ -0,0 +1,38 @@ +--- +title: moq-flate +description: Opaque MoQ tracks, optionally compressed with group-scoped DEFLATE +--- + +# moq-flate + +[![crates.io](https://img.shields.io/crates/v/moq-flate)](https://crates.io/crates/moq-flate) +[![docs.rs](https://docs.rs/moq-flate/badge.svg)](https://docs.rs/moq-flate) + +Opaque payloads over [`moq-net`](/lib/rs/moq-net) tracks, in two modes: + +- **snapshot**: lossy latest-value, one payload per group. +- **stream**: lossless append-log in a single group. + +The bytes are never inspected. Compression is opt-in: `Compression` is `None` +(the default) or `Deflate`, and both sides set the same field. With `Deflate`, +each group is one raw DEFLATE stream sync-flushed at every frame boundary, so a +stream's payloads compress against the earlier ones in their group. + +```rust +let mut config = moq_flate::snapshot::Config::default(); +config.compression = moq_flate::Compression::Deflate; + +let mut producer = moq_flate::snapshot::Producer::new(track, config); +producer.update(payload)?; +``` + +A payload is stamped when written, unless it carries its capture time: +`moq_net::Timed::from(bytes).at(captured)`. Writes return the encoded frame +size. + +The group-scoped codec underneath is exported as `Encoder`/`Decoder`, which +[`moq-json`](/lib/rs/moq-json) reuses for its merge-patch deltas. Create one +pair per group and feed frames in order. + +The TypeScript twin is [`@moq/flate`](/lib/js/flate). API: +[docs.rs/moq-flate](https://docs.rs/moq-flate). diff --git a/doc/setup/upgrade.md b/doc/setup/upgrade.md index 12dfbcdb50..aea4de7e0a 100644 --- a/doc/setup/upgrade.md +++ b/doc/setup/upgrade.md @@ -15,6 +15,23 @@ A released flag, environment variable, or config key that was renamed is refused at startup with its replacement named, rather than ignored. Fix what the error lists and rerun. +## Unreleased + +These land with the next breaking release, not the 2026-09-23 train. + +- **moq-binary is moq-flate, and @moq/binary is @moq/flate.** The opaque + `snapshot` and `stream` tracks moved beside the codec; the wire and the + catalog's `binary` section are unchanged. In Rust, `moq_binary::X` is + `moq_flate::X`, and `moq_binary::Error::Flate(e)` is the matching + `moq_flate::Error` variant (`Decompress`, `TooLarge`). `moq_flate::Error` + now carries `moq_net::Error`, so it is no longer `PartialEq`. moq-mux's + `Error::Binary` is `Error::Flate`. In TypeScript, import `Snapshot` and + `Stream` from `@moq/flate`. +- **moq-ffi flate tracks.** `publish_binary_snapshot` / `publish_binary_stream` + are `publish_flate_snapshot` / `publish_flate_stream`, taking + `MoqFlateConfig` and returning `MoqFlateSnapshotProducer` / + `MoqFlateStreamProducer`. The C `moq_publish_binary_*` calls are unchanged. + ## Wire Older protocol versions still negotiate, so relays and clients can be upgraded @@ -181,7 +198,7 @@ PRs, so a minor rename may be missing. is `Watch.Player({ origin: connection.origin, ... })`. - **@moq/publish** takes `origin: connection.origin` instead of a `connection` signal. -- **@moq/json and @moq/binary** take one options object (#3640): +- **@moq/json and @moq/binary** (now `@moq/flate`) take one options object (#3640): `new Json.Snapshot.Consumer({ track })` instead of `(track, config)`. - **@moq/hang** reads the catalog `archive` entry instead of `timeline`. diff --git a/js/binary/README.md b/js/binary/README.md deleted file mode 100644 index c3448e284b..0000000000 --- a/js/binary/README.md +++ /dev/null @@ -1,15 +0,0 @@ -# @moq/binary - -Opaque binary payloads over MoQ tracks, in two modes: - -- `Snapshot`: **lossy**. One value updated over time; a consumer only gets the most recent one. -- `Stream`: **lossless**. An ordered append-log of self-contained payloads, nothing superseded. - -The bytes are opaque: this package frames them onto a track and optionally compresses them, and -never looks inside. For JSON documents reach for [`@moq/json`](../json) instead, which adds RFC 7396 -merge-patch deltas on top of the same two modes. - -`Config` is the codec options (`compression: "none" | "deflate"`). `Producer.Config` / -`Consumer.Config` add the track. Compression is [`@moq/flate`](../flate), the same group-scoped -DEFLATE `@moq/json` uses, so the two agree on the wire. Interoperable with the Rust `moq-binary` -crate. diff --git a/js/binary/package.json b/js/binary/package.json deleted file mode 100644 index fa1f4487fb..0000000000 --- a/js/binary/package.json +++ /dev/null @@ -1,32 +0,0 @@ -{ - "name": "@moq/binary", - "type": "module", - "version": "0.2.1", - "description": "Binary publishing over MoQ tracks: latest-value snapshots, or append-log streams.", - "license": "(MIT OR Apache-2.0)", - "repository": { - "type": "git", - "url": "git+https://github.com/moq-dev/moq.git", - "directory": "js/binary" - }, - "sideEffects": false, - "exports": { - ".": "./src/index.ts" - }, - "scripts": { - "build": "rimraf dist && tsc -b tsconfig.build.json && bun ../common/package.ts", - "check": "tsc --noEmit", - "test": "bun test --only-failures", - "release": "bun ../common/release.ts" - }, - "dependencies": { - "@moq/flate": "workspace:^", - "@moq/net": "workspace:^", - "@moq/signals": "workspace:^" - }, - "devDependencies": { - "@types/bun": "^1.4.2", - "rimraf": "^6.1.3", - "typescript": "7.0.2" - } -} diff --git a/js/binary/src/index.ts b/js/binary/src/index.ts deleted file mode 100644 index af19300fe1..0000000000 --- a/js/binary/src/index.ts +++ /dev/null @@ -1,29 +0,0 @@ -/** - * Opaque binary payloads over MoQ tracks, in two modes: - * - * - {@link Snapshot}: **lossy**. One value updated over time; a consumer only gets the most recent - * one. Older values are superseded and dropped. - * - {@link Stream}: **lossless**. An ordered append-log of self-contained payloads, delivered in - * order with nothing superseded. Bounded by the group cache: see {@link Stream} for what that - * costs a consumer that falls behind. - * - * Pick {@link Snapshot} when consumers care about "what is the value now" (a poster image, a - * serialized state blob) and {@link Stream} when they care about every payload (an event log, a - * sequence of samples). - * - * The bytes are opaque: this package frames them onto a track and optionally compresses them, and - * never looks inside. For JSON documents reach for `@moq/json` instead, which adds RFC 7396 - * merge-patch deltas on top of the same two modes. - * - * Compression is `@moq/flate`, the same group-scoped DEFLATE `@moq/json` uses, so the two agree on - * the wire: each group is one raw DEFLATE stream, sync-flushed at every frame boundary. A - * {@link Stream} therefore compresses each payload against the earlier ones in its group, while a - * {@link Snapshot} group holds a single self-contained value. Interoperable with the Rust - * `moq-binary` crate. - * - * @module - */ - -export type { Compression } from "./compression.ts"; -export * as Snapshot from "./snapshot/index.ts"; -export * as Stream from "./stream/index.ts"; diff --git a/js/binary/tsconfig.build.json b/js/binary/tsconfig.build.json deleted file mode 100644 index b0d79fdc2c..0000000000 --- a/js/binary/tsconfig.build.json +++ /dev/null @@ -1,7 +0,0 @@ -{ - // Build-only variant of tsconfig.json that keeps tests out of dist/, so they - // are neither published nor discovered by `bun test` at the package root. - // `tsc --noEmit` still uses tsconfig.json, so tests stay type-checked. - "extends": "./tsconfig.json", - "exclude": ["src/**/*.test.ts", "src/**/*.test.tsx"] -} diff --git a/js/binary/tsconfig.json b/js/binary/tsconfig.json deleted file mode 100644 index bb55d7c43f..0000000000 --- a/js/binary/tsconfig.json +++ /dev/null @@ -1,9 +0,0 @@ -{ - "extends": "../tsconfig.json", - "compilerOptions": { - "outDir": "dist", - "rootDir": "./src", - "types": ["bun"] - }, - "include": ["src"] -} diff --git a/js/flate/README.md b/js/flate/README.md index c013015d9e..de88ec1ac0 100644 --- a/js/flate/README.md +++ b/js/flate/README.md @@ -7,18 +7,34 @@ [![npm version](https://img.shields.io/npm/v/@moq/flate)](https://www.npmjs.com/package/@moq/flate) [![TypeScript](https://img.shields.io/badge/TypeScript-ready-blue.svg)](https://www.typescriptlang.org/) -Group-scoped DEFLATE: a stream of self-delimited frames sharing one compression window. +Opaque byte tracks over MoQ, optionally compressed with group-scoped DEFLATE, in two modes: -A sequence of frame payloads is compressed into a single raw DEFLATE ([RFC 1951](https://www.rfc-editor.org/rfc/rfc1951.html)) stream, sync-flushed at each frame boundary. Every frame is self-delimited (byte-aligned, the window retained) while later frames reuse the earlier ones as context, so a stream of similar payloads (a snapshot followed by deltas, repeated records, log lines) compresses far better than each payload alone. +- `Snapshot`: **lossy**. One value updated over time; a consumer only gets the most recent one. +- `Stream`: **lossless**. An ordered append-log of self-contained payloads, nothing superseded. -This is plain raw DEFLATE with a `Z_SYNC_FLUSH` after each frame, so it interoperates on the wire with any peer using the same primitive, including the Rust [`moq-flate`](https://crates.io/crates/moq-flate) crate. The fixed 4-byte sync-flush marker is stripped per frame ([RFC 7692](https://www.rfc-editor.org/rfc/rfc7692.html#section-7.2.1)'s permessage-deflate trick). There is no length prefix: the caller frames each slice ([`@moq/net`](../net) already does). +The bytes are opaque: the tracks frame them onto a [`@moq/net`](../net) track and optionally compress them, and never look inside. For JSON documents reach for [`@moq/json`](../json) instead, which adds RFC 7396 merge-patch deltas on top of the same two modes and codec. Interoperable with the Rust [`moq-flate`](https://crates.io/crates/moq-flate) crate. -## Quick Start +```ts +import { Snapshot, Stream } from "@moq/flate"; -```bash -npm add @moq/flate +const thumbnail = new Snapshot.Producer({ track, compression: "deflate" }); +thumbnail.update(jpeg); + +const log = new Stream.Consumer({ track: subscriber, compression: "deflate" }); +for (;;) { + const payload = await log.next(); + if (payload === undefined) break; +} ``` +Compression is opt-in per track (`compression: "none" | "deflate"`, default `"none"`); a consumer must set the same value as the producer. + +## Codec + +A sequence of frame payloads is compressed into a single raw DEFLATE ([RFC 1951](https://www.rfc-editor.org/rfc/rfc1951.html)) stream, sync-flushed at each frame boundary. Every frame is self-delimited (byte-aligned, the window retained) while later frames reuse the earlier ones as context, so a stream of similar payloads (a snapshot followed by deltas, repeated records, log lines) compresses far better than each payload alone. + +This is plain raw DEFLATE with a `Z_SYNC_FLUSH` after each frame, so it interoperates on the wire with any peer using the same primitive, including the Rust [`moq-flate`](https://crates.io/crates/moq-flate) crate. The fixed 4-byte sync-flush marker is stripped per frame ([RFC 7692](https://www.rfc-editor.org/rfc/rfc7692.html#section-7.2.1)'s permessage-deflate trick). There is no length prefix: the caller frames each slice ([`@moq/net`](../net) already does). + ```ts import { Encoder, Decoder } from "@moq/flate"; diff --git a/js/flate/package.json b/js/flate/package.json index 86f2c75783..7bd2397ea3 100644 --- a/js/flate/package.json +++ b/js/flate/package.json @@ -2,7 +2,7 @@ "name": "@moq/flate", "type": "module", "version": "0.1.3", - "description": "Group-scoped DEFLATE: a stream of self-delimited frames sharing one compression window.", + "description": "Opaque MoQ tracks, optionally compressed with group-scoped DEFLATE: latest-value snapshots or append-log streams.", "license": "(MIT OR Apache-2.0)", "repository": { "type": "git", @@ -20,6 +20,8 @@ "release": "bun ../common/release.ts" }, "dependencies": { + "@moq/net": "workspace:^", + "@moq/signals": "workspace:^", "pako": "^3.0.2" }, "devDependencies": { diff --git a/js/flate/src/index.test.ts b/js/flate/src/codec.test.ts similarity index 98% rename from js/flate/src/index.test.ts rename to js/flate/src/codec.test.ts index b4f5d11ea9..6145849397 100644 --- a/js/flate/src/index.test.ts +++ b/js/flate/src/codec.test.ts @@ -1,6 +1,6 @@ import { expect, test } from "bun:test"; import { Deflate, Inflate } from "fflate"; -import { DEFAULT_MAX_FRAME_SIZE, Decoder, Encoder } from "./index.ts"; +import { DEFAULT_MAX_FRAME_SIZE, Decoder, Encoder } from "./codec.ts"; const enc = new TextEncoder(); const dec = new TextDecoder(); diff --git a/js/flate/src/codec.ts b/js/flate/src/codec.ts new file mode 100644 index 0000000000..15c10ab6ee --- /dev/null +++ b/js/flate/src/codec.ts @@ -0,0 +1,145 @@ +// The group-scoped DEFLATE codec: see the module docs in index.ts. + +import * as pako from "pako"; + +/** The default DEFLATE level ({@link Encoder}): a good size/speed balance for small, repetitive payloads. */ +export const DEFAULT_LEVEL = 6; + +/** The default per-frame decompressed-size cap ({@link Decoder}): 64 MiB. */ +export const DEFAULT_MAX_FRAME_SIZE = 64 * 1024 * 1024; + +// The trailing bytes of a DEFLATE sync flush, stripped on the wire and re-appended to decode. +const SYNC_FLUSH_TAIL = new Uint8Array([0x00, 0x00, 0xff, 0xff]); + +// Concatenate chunks into one tight buffer (a single chunk passes through untouched). Safe only for +// output that is consumed before the next push, since a single chunk is pako's own reused buffer. +function concat(chunks: Uint8Array[], total: number): Uint8Array { + if (chunks.length === 1) return chunks[0]; + const out = new Uint8Array(total); + let offset = 0; + for (const chunk of chunks) { + out.set(chunk, offset); + offset += chunk.length; + } + return out; +} + +/** Options for an {@link Encoder}. */ +export interface EncoderOptions { + /** DEFLATE level, `0..=9` (higher is smaller and slower). Defaults to {@link DEFAULT_LEVEL}. */ + level?: number; +} + +/** + * Encodes a stream's frame payloads into one shared DEFLATE window, one self-delimited slice per + * frame. Hold one per stream; create a new one for each independent stream. + * + * @public + */ +export class Encoder { + #deflate: pako.Deflate; + #chunks: Uint8Array[] = []; + #total = 0; + + /** Start a fresh encoder with a cold window. */ + constructor(options: EncoderOptions = {}) { + // `raw`: no zlib header/trailer, matching the Rust side and the browser's `deflate-raw`. + // pako types `level` as a literal union; we accept a plain number and narrow here. + const level = (options.level ?? DEFAULT_LEVEL) as pako.DeflateOptions["level"]; + this.#deflate = new pako.Deflate({ raw: true, level }); + this.#deflate.onData = (chunk) => { + const bytes = chunk as Uint8Array; + this.#chunks.push(bytes); + this.#total += bytes.length; + }; + } + + /** + * Compress the next frame's `payload`, returning its slice of the stream: the DEFLATE bytes minus + * the fixed sync-flush marker. Empty in yields empty out. Slices must be produced in frame order. + */ + frame(payload: Uint8Array): Uint8Array { + if (payload.length === 0) return payload; + this.#chunks = []; + this.#total = 0; + this.#deflate.push(payload, pako.Z_SYNC_FLUSH); + + // Copy into one tight owned buffer, dropping the trailing sync-flush marker. We can't return + // pako's chunk views: a caller retains the reference and pako backs each chunk with a ~16 KB + // buffer, so a view would pin far more memory than the frame. + const out = new Uint8Array(this.#total - SYNC_FLUSH_TAIL.length); + let offset = 0; + for (const chunk of this.#chunks) { + if (offset >= out.length) break; + const take = Math.min(chunk.length, out.length - offset); + out.set(chunk.subarray(0, take), offset); + offset += take; + } + return out; + } +} + +/** Options for a {@link Decoder}. */ +export interface DecoderOptions { + /** + * Per-frame decompressed-size cap. A frame that inflates past this is rejected (zip-bomb guard). + * Defaults to {@link DEFAULT_MAX_FRAME_SIZE}. + */ + maxFrameSize?: number; +} + +/** + * Decodes a stream's frame slices back into the original payloads. Hold one per stream; feed slices + * in frame order (each frame builds on the earlier ones). + * + * @public + */ +export class Decoder { + #inflate = new pako.Inflate({ raw: true }); + #chunks: Uint8Array[] = []; + #total = 0; + #tooLarge = false; + #maxFrameSize: number; + + /** The per-frame decompressed-size cap in bytes, in effect for this decoder. */ + get maxFrameSize(): number { + return this.#maxFrameSize; + } + + /** Start a fresh decoder with a cold window. */ + constructor(options: DecoderOptions = {}) { + this.#maxFrameSize = options.maxFrameSize ?? DEFAULT_MAX_FRAME_SIZE; + this.#inflate.onData = (chunk) => { + const bytes = chunk as Uint8Array; + this.#total += bytes.length; + // Bound the inflated output as it is produced; a tiny slice can expand enormously. Stop + // retaining past the cap, then reject once the push returns. + if (this.#total > this.#maxFrameSize) { + this.#tooLarge = true; + return; + } + this.#chunks.push(bytes); + }; + } + + /** + * Decompress the next frame's `slice` back into its payload. Empty in yields empty out. Throws if + * the input is malformed or inflates past the per-frame size cap. + */ + frame(slice: Uint8Array): Uint8Array { + if (slice.length === 0) return slice; + + this.#chunks = []; + this.#total = 0; + this.#tooLarge = false; + + // Feed the slice then the re-appended sync-flush marker as two pushes, so no combined buffer is + // allocated. The marker delimits the frame and flushes its last bytes out of the inflate buffer. + this.#inflate.push(slice, false); + this.#inflate.push(SYNC_FLUSH_TAIL, pako.Z_SYNC_FLUSH); + if (this.#inflate.err) throw new Error(`decompression failed: ${this.#inflate.msg}`); + if (this.#tooLarge) throw new Error(`decompressed frame exceeded ${this.#maxFrameSize} bytes`); + + return concat(this.#chunks, this.#total); + } +} diff --git a/js/binary/src/compression.ts b/js/flate/src/compression.ts similarity index 75% rename from js/binary/src/compression.ts rename to js/flate/src/compression.ts index bf2b36d408..b76fe3ad4f 100644 --- a/js/binary/src/compression.ts +++ b/js/flate/src/compression.ts @@ -1,4 +1,4 @@ -/** How a binary track compresses its frames. `"none"` (the default) is uncompressed payloads. */ +/** How a `Snapshot` or `Stream` track compresses its frames. `"none"` (the default) is uncompressed payloads. */ export type Compression = "none" | "deflate"; /** Whether `compression` is group-scoped DEFLATE. */ diff --git a/js/flate/src/index.ts b/js/flate/src/index.ts index 9aa91ed9a6..aee9cb3298 100644 --- a/js/flate/src/index.ts +++ b/js/flate/src/index.ts @@ -1,13 +1,31 @@ /** - * Group-scoped DEFLATE: a stream of self-delimited frames sharing one compression window, using - * {@link https://github.com/nodeca/pako | pako}'s streaming deflate/inflate. + * Opaque byte tracks over MoQ, optionally compressed with group-scoped DEFLATE, in two modes: * - * A sequence of frame payloads is compressed into a single raw DEFLATE - * ([RFC 1951](https://www.rfc-editor.org/rfc/rfc1951.html)) stream, sync-flushed at each frame - * boundary, so every frame is self-delimited (byte-aligned, the window retained) while later frames - * reuse the earlier ones as context. A stream of similar payloads (a snapshot followed by deltas, - * repeated records, log lines) compresses far better than each payload alone. Create a fresh - * {@link Encoder}/{@link Decoder} pair per independent stream (in moq-net terms, per group). + * - {@link Snapshot}: **lossy**. One value updated over time; a consumer only gets the most recent + * one. Older values are superseded and dropped. + * - {@link Stream}: **lossless**. An ordered append-log of self-contained payloads, delivered in + * order with nothing superseded. Bounded by the group cache: see {@link Stream} for what that + * costs a consumer that falls behind. + * + * Pick {@link Snapshot} when consumers care about "what is the value now" (a poster image, a + * serialized state blob) and {@link Stream} when they care about every payload (an event log, a + * sequence of samples). The bytes are opaque: the tracks frame them and optionally compress them, + * and never look inside. For JSON documents reach for `@moq/json` instead, which adds RFC 7396 + * merge-patch deltas on top of the same two modes and codec. + * + * Compression is opt-in per track ({@link Compression}), so the package name is not a promise that + * every track is deflated: `"none"` writes the bytes through untouched. + * + * ## Codec + * + * Underneath, {@link Encoder}/{@link Decoder} compress a sequence of frame payloads into a single + * raw DEFLATE ([RFC 1951](https://www.rfc-editor.org/rfc/rfc1951.html)) stream, sync-flushed at each + * frame boundary, using {@link https://github.com/nodeca/pako | pako}. Every frame is self-delimited + * (byte-aligned, the window retained) while later frames reuse the earlier ones as context, so a + * stream of similar payloads compresses far better than each payload alone. Create a fresh pair per + * independent stream (in moq-net terms, per group). A {@link Stream} therefore compresses each + * payload against the earlier ones in its group, while a {@link Snapshot} group holds a single + * self-contained value. * * This is plain raw DEFLATE with a `Z_SYNC_FLUSH` after each frame, so any peer using the same * primitive (the Rust `moq-flate` crate, zlib's sync flush) interoperates on the wire. There is no @@ -17,153 +35,19 @@ * and {@link Decoder.frame} re-appends it, saving 4 bytes per frame, the same trick * [RFC 7692](https://www.rfc-editor.org/rfc/rfc7692.html#section-7.2.1) (permessage-deflate) uses. A * small slice can still inflate enormously, so {@link Decoder.frame} caps the inflated output as it - * is produced. - * - * pako is synchronous, so the whole codec is synchronous; it is a normal dependency. + * is produced. pako is synchronous, so the whole codec is synchronous. * * @module */ -import * as pako from "pako"; - -/** The default DEFLATE level ({@link Encoder}): a good size/speed balance for small, repetitive payloads. */ -export const DEFAULT_LEVEL = 6; - -/** The default per-frame decompressed-size cap ({@link Decoder}): 64 MiB. */ -export const DEFAULT_MAX_FRAME_SIZE = 64 * 1024 * 1024; - -// The trailing bytes of a DEFLATE sync flush, stripped on the wire and re-appended to decode. -const SYNC_FLUSH_TAIL = new Uint8Array([0x00, 0x00, 0xff, 0xff]); - -// Concatenate chunks into one tight buffer (a single chunk passes through untouched). Safe only for -// output that is consumed before the next push, since a single chunk is pako's own reused buffer. -function concat(chunks: Uint8Array[], total: number): Uint8Array { - if (chunks.length === 1) return chunks[0]; - const out = new Uint8Array(total); - let offset = 0; - for (const chunk of chunks) { - out.set(chunk, offset); - offset += chunk.length; - } - return out; -} - -/** Options for an {@link Encoder}. */ -export interface EncoderOptions { - /** DEFLATE level, `0..=9` (higher is smaller and slower). Defaults to {@link DEFAULT_LEVEL}. */ - level?: number; -} - -/** - * Encodes a stream's frame payloads into one shared DEFLATE window, one self-delimited slice per - * frame. Hold one per stream; create a new one for each independent stream. - * - * @public - */ -export class Encoder { - #deflate: pako.Deflate; - #chunks: Uint8Array[] = []; - #total = 0; - - /** Start a fresh encoder with a cold window. */ - constructor(options: EncoderOptions = {}) { - // `raw`: no zlib header/trailer, matching the Rust side and the browser's `deflate-raw`. - // pako types `level` as a literal union; we accept a plain number and narrow here. - const level = (options.level ?? DEFAULT_LEVEL) as pako.DeflateOptions["level"]; - this.#deflate = new pako.Deflate({ raw: true, level }); - this.#deflate.onData = (chunk) => { - const bytes = chunk as Uint8Array; - this.#chunks.push(bytes); - this.#total += bytes.length; - }; - } - - /** - * Compress the next frame's `payload`, returning its slice of the stream: the DEFLATE bytes minus - * the fixed sync-flush marker. Empty in yields empty out. Slices must be produced in frame order. - */ - frame(payload: Uint8Array): Uint8Array { - if (payload.length === 0) return payload; - this.#chunks = []; - this.#total = 0; - this.#deflate.push(payload, pako.Z_SYNC_FLUSH); - - // Copy into one tight owned buffer, dropping the trailing sync-flush marker. We can't return - // pako's chunk views: a caller retains the reference and pako backs each chunk with a ~16 KB - // buffer, so a view would pin far more memory than the frame. - const out = new Uint8Array(this.#total - SYNC_FLUSH_TAIL.length); - let offset = 0; - for (const chunk of this.#chunks) { - if (offset >= out.length) break; - const take = Math.min(chunk.length, out.length - offset); - out.set(chunk.subarray(0, take), offset); - offset += take; - } - return out; - } -} - -/** Options for a {@link Decoder}. */ -export interface DecoderOptions { - /** - * Per-frame decompressed-size cap. A frame that inflates past this is rejected (zip-bomb guard). - * Defaults to {@link DEFAULT_MAX_FRAME_SIZE}. - */ - maxFrameSize?: number; -} - -/** - * Decodes a stream's frame slices back into the original payloads. Hold one per stream; feed slices - * in frame order (each frame builds on the earlier ones). - * - * @public - */ -export class Decoder { - #inflate = new pako.Inflate({ raw: true }); - #chunks: Uint8Array[] = []; - #total = 0; - #tooLarge = false; - #maxFrameSize: number; - - /** The per-frame decompressed-size cap in bytes, in effect for this decoder. */ - get maxFrameSize(): number { - return this.#maxFrameSize; - } - - /** Start a fresh decoder with a cold window. */ - constructor(options: DecoderOptions = {}) { - this.#maxFrameSize = options.maxFrameSize ?? DEFAULT_MAX_FRAME_SIZE; - this.#inflate.onData = (chunk) => { - const bytes = chunk as Uint8Array; - this.#total += bytes.length; - // Bound the inflated output as it is produced; a tiny slice can expand enormously. Stop - // retaining past the cap, then reject once the push returns. - if (this.#total > this.#maxFrameSize) { - this.#tooLarge = true; - return; - } - this.#chunks.push(bytes); - }; - } - - /** - * Decompress the next frame's `slice` back into its payload. Empty in yields empty out. Throws if - * the input is malformed or inflates past the per-frame size cap. - */ - frame(slice: Uint8Array): Uint8Array { - if (slice.length === 0) return slice; - - this.#chunks = []; - this.#total = 0; - this.#tooLarge = false; - - // Feed the slice then the re-appended sync-flush marker as two pushes, so no combined buffer is - // allocated. The marker delimits the frame and flushes its last bytes out of the inflate buffer. - this.#inflate.push(slice, false); - this.#inflate.push(SYNC_FLUSH_TAIL, pako.Z_SYNC_FLUSH); - if (this.#inflate.err) throw new Error(`decompression failed: ${this.#inflate.msg}`); - if (this.#tooLarge) throw new Error(`decompressed frame exceeded ${this.#maxFrameSize} bytes`); - - return concat(this.#chunks, this.#total); - } -} +export { + DEFAULT_LEVEL, + DEFAULT_MAX_FRAME_SIZE, + Decoder, + type DecoderOptions, + Encoder, + type EncoderOptions, +} from "./codec.ts"; +export type { Compression } from "./compression.ts"; +export * as Snapshot from "./snapshot/index.ts"; +export * as Stream from "./stream/index.ts"; diff --git a/js/binary/src/snapshot/consumer.ts b/js/flate/src/snapshot/consumer.ts similarity index 94% rename from js/binary/src/snapshot/consumer.ts rename to js/flate/src/snapshot/consumer.ts index c6b0263613..9aa68a3429 100644 --- a/js/binary/src/snapshot/consumer.ts +++ b/js/flate/src/snapshot/consumer.ts @@ -1,15 +1,15 @@ -import { Decoder as Flate } from "@moq/flate"; import * as Moq from "@moq/net"; +import { Decoder as Flate } from "../codec.ts"; import { isDeflate } from "../compression.ts"; import type { Config as CodecConfig } from "./producer.ts"; /** - * Consumes a binary value from a track, yielding the newest one. + * Consumes an opaque value from a track, yielding the newest one. * * Jumps to the newest group and reads the value out of it, so a late joiner starts at the current * value rather than replaying superseded ones. Interoperable with the Rust - * `moq_binary::snapshot::Consumer`, which collapses the same backlog. + * `moq_flate::snapshot::Consumer`, which collapses the same backlog. */ export class Consumer { #track: Moq.Track.Ordered; diff --git a/js/binary/src/snapshot/index.ts b/js/flate/src/snapshot/index.ts similarity index 93% rename from js/binary/src/snapshot/index.ts rename to js/flate/src/snapshot/index.ts index 68fbf6165e..027e96f811 100644 --- a/js/binary/src/snapshot/index.ts +++ b/js/flate/src/snapshot/index.ts @@ -1,5 +1,5 @@ /** - * Lossy latest-value binary publishing over MoQ tracks. + * Lossy latest-value opaque publishing over MoQ tracks. * * One opaque value updated over time, for consumers that only care about the current state (a * poster image, a serialized state blob). This mode is **lossy** by design: a consumer yields only diff --git a/js/binary/src/snapshot/producer.ts b/js/flate/src/snapshot/producer.ts similarity index 93% rename from js/binary/src/snapshot/producer.ts rename to js/flate/src/snapshot/producer.ts index 75f4fe3881..60c8547aa9 100644 --- a/js/binary/src/snapshot/producer.ts +++ b/js/flate/src/snapshot/producer.ts @@ -1,6 +1,6 @@ -import { DEFAULT_MAX_FRAME_SIZE, Encoder as Flate } from "@moq/flate"; import * as Moq from "@moq/net"; import { Time } from "@moq/net"; +import { DEFAULT_MAX_FRAME_SIZE, Encoder as Flate } from "../codec.ts"; import { type Compression, isDeflate } from "../compression.ts"; @@ -17,7 +17,7 @@ export interface Config { } /** - * Publishes a binary value to a track, one value per group. + * Publishes an opaque value to a track, one value per group. * * Each {@link update} rolls a new group holding the whole value, so a consumer only ever needs the * newest group and older ones are dropped. For a log where every payload survives, use the `Stream` @@ -27,7 +27,7 @@ export class Producer { #track: Moq.Track.Producer; #compress: boolean; - /** Wrap a track to publish a binary value into it. */ + /** Wrap a track to publish an opaque value into it. */ constructor(config: Producer.Config) { this.#track = config.track; this.#compress = isDeflate(config.compression); diff --git a/js/binary/src/snapshot/snapshot.test.ts b/js/flate/src/snapshot/snapshot.test.ts similarity index 100% rename from js/binary/src/snapshot/snapshot.test.ts rename to js/flate/src/snapshot/snapshot.test.ts diff --git a/js/binary/src/stream/consumer.ts b/js/flate/src/stream/consumer.ts similarity index 97% rename from js/binary/src/stream/consumer.ts rename to js/flate/src/stream/consumer.ts index f4cc5cf441..0a313a3efa 100644 --- a/js/binary/src/stream/consumer.ts +++ b/js/flate/src/stream/consumer.ts @@ -1,6 +1,6 @@ -import { Decoder as Flate } from "@moq/flate"; import type * as Moq from "@moq/net"; import { race } from "@moq/signals"; +import { Decoder as Flate } from "../codec.ts"; import { isDeflate } from "../compression.ts"; import type { Config as CodecConfig } from "./producer.ts"; @@ -12,7 +12,7 @@ import type { Config as CodecConfig } from "./producer.ts"; * track rather than rolling. A second group therefore means whatever would have completed the * first is gone, so the read reports it instead of handing back the remainder as a continuous log. * - * Mirrors the Rust `moq_binary::Error::Rolled`. + * Mirrors the Rust `moq_flate::Error::Rolled`. */ export class Rolled extends Error { constructor() { @@ -22,7 +22,7 @@ export class Rolled extends Error { } /** - * Consumes an ordered log of binary payloads from a track, yielding every one in order. + * Consumes an ordered log of opaque payloads from a track, yielding every one in order. * * The log is a single group. That is what makes the mode lossless: rolling to a second group means * the payloads that would have completed the first are gone, so a publisher that cannot write ends diff --git a/js/binary/src/stream/index.ts b/js/flate/src/stream/index.ts similarity index 95% rename from js/binary/src/stream/index.ts rename to js/flate/src/stream/index.ts index 23272a845c..508439578c 100644 --- a/js/binary/src/stream/index.ts +++ b/js/flate/src/stream/index.ts @@ -1,5 +1,5 @@ /** - * Lossless append-log binary publishing over MoQ tracks. + * Lossless append-log opaque publishing over MoQ tracks. * * An ordered log of opaque payloads, for consumers that care about every one (an event log, a * sequence of samples). Nothing is ever superseded: a consumer yields each payload in the order it diff --git a/js/binary/src/stream/producer.ts b/js/flate/src/stream/producer.ts similarity index 96% rename from js/binary/src/stream/producer.ts rename to js/flate/src/stream/producer.ts index 8605954f6b..c63a14bd32 100644 --- a/js/binary/src/stream/producer.ts +++ b/js/flate/src/stream/producer.ts @@ -1,6 +1,6 @@ -import { DEFAULT_MAX_FRAME_SIZE, Encoder as Flate } from "@moq/flate"; import type * as Moq from "@moq/net"; import { Time } from "@moq/net"; +import { DEFAULT_MAX_FRAME_SIZE, Encoder as Flate } from "../codec.ts"; import { type Compression, isDeflate } from "../compression.ts"; @@ -15,7 +15,7 @@ export interface Config { } /** - * Publishes an ordered log of binary payloads to a track, one payload per frame in a single group. + * Publishes an ordered log of opaque payloads to a track, one payload per frame in a single group. */ export class Producer { #track: Moq.Track.Producer; diff --git a/js/binary/src/stream/stream.test.ts b/js/flate/src/stream/stream.test.ts similarity index 99% rename from js/binary/src/stream/stream.test.ts rename to js/flate/src/stream/stream.test.ts index 44a2430506..014a69a495 100644 --- a/js/binary/src/stream/stream.test.ts +++ b/js/flate/src/stream/stream.test.ts @@ -1,6 +1,6 @@ import { expect, spyOn, test } from "bun:test"; -import { DEFAULT_MAX_FRAME_SIZE } from "@moq/flate"; import { Time, Track } from "@moq/net"; +import { DEFAULT_MAX_FRAME_SIZE } from "../codec.ts"; import { Consumer, Producer, Rolled } from "./index.ts"; // Ask for a replay window, so the superseded first group is delivered rather than skipped by the diff --git a/js/binary/src/types.test.ts b/js/flate/src/types.test.ts similarity index 100% rename from js/binary/src/types.test.ts rename to js/flate/src/types.test.ts diff --git a/package.json b/package.json index 3a85e66642..49e3a47ede 100644 --- a/package.json +++ b/package.json @@ -33,7 +33,6 @@ "js/pattern", "js/auth", "js/flate", - "js/binary", "js/json", "js/hang", "js/loc", diff --git a/quest/m1/README.md b/quest/m1/README.md index 7f1d7e3615..fcc59dbd0a 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -22,7 +22,6 @@ transport, benchmark tooling); worktrees isolate commits, not semantics. - [Track tail interop](/quest/m1/track-tail-interop.md) - a Rust publisher ending a track with a group in flight is read to its end by the JS subscriber, and the reverse, in `just test interop` - [Worker socket count](/quest/m1/worker-socket-count.md) - the moq-tokio worker test counts only its own listener's sockets - [Binding audio delay](/quest/m1/binding-surface.md) - moq-ffi and every wrapper configure and observe audio playout delay -- [moq-binary folds into moq-flate](/quest/m1/flate-binary.md) - on dev, moq-flate and @moq/flate own the opaque snapshot and stream tracks and moq-binary is deleted - [FFI shape](/quest/m1/ffi-shape/README.md) - the bindings mirror Rust's layers: net at the root, then media, json, flate, audio, and video namespaces built from the handle below - [Track demand](/quest/m1/track-demand.md) - Rust and JS watch a track's subscribers through `demand()` alone - [Error messages](/quest/m1/error-display.md) - Python, Go, and Dart print `MoqError` with Rust's message, as Kotlin and Swift do diff --git a/quest/m1/data-capture-bindings.md b/quest/m1/data-capture-bindings.md index 95d320ac35..90f1f008ea 100644 --- a/quest/m1/data-capture-bindings.md +++ b/quest/m1/data-capture-bindings.md @@ -16,7 +16,7 @@ not moq-c. `Timed::from(value).at(t)`. The `moq-mux` data producers take `Timed<_, Instant>`, map it onto the broadcast clock, and refuse one ahead of now (`Error::InvalidCapture`). The moq-ffi producers in - `rs/moq-ffi/src/{json,binary}.rs` pass bare values. + `rs/moq-ffi/src/{json,flate}.rs` pass bare values. - Settled: the capture time is a media timestamp on the broadcast's timeline. moq-ffi exposes the broadcast clock's `now()` as a timestamp; callers stamp payloads with values taken from it, and moq refuses one ahead @@ -34,8 +34,8 @@ not moq-c. past capture time is accepted and a future one refused. - Go (`go/wrapper/moq`) and Python (`py/moq-rs`) wrap only the JSON producers; [#4137](https://github.com/moq-dev/moq/pull/4137) added - `publish_binary_snapshot` / `publish_binary_stream` to moq-ffi without - them. Add hand-written binary wrappers there, capture time included, so + `publish_binary_snapshot` / `publish_binary_stream` (now + `publish_flate_*`) to moq-ffi without them. Add hand-written flate wrappers there, capture time included, so the capture tests cover binary too. The maintainer asked for this. Update `doc/lib/{py,swift,kt,go,dart}`. diff --git a/quest/m1/ffi-shape/README.md b/quest/m1/ffi-shape/README.md index 476b57640c..5f10244392 100644 --- a/quest/m1/ffi-shape/README.md +++ b/quest/m1/ffi-shape/README.md @@ -20,8 +20,7 @@ Settled shape: binding never sees that split (catalog, import producers, container consumers); `json`, `flate`, `audio`, and `video` own their producers and consumers. `flate` holds the opaque snapshot and stream tracks moq-ffi - publishes as `publish_binary_*` today (#4137), named after the crate they - fold into in [moq-binary folds into moq-flate](/quest/m1/flate-binary.md). + publishes as `publish_flate_*` today, named after the `moq-flate` crate. - A layer's type is constructed from the handles its Rust constructor takes, not reached through an accessor on the broadcast: JSON wraps a track (`moq_json::snapshot::Producer::new(track, config)`), so it also works on a diff --git a/quest/m1/ffi-shape/json.md b/quest/m1/ffi-shape/json.md index 3cbf8969d6..1475f6217a 100644 --- a/quest/m1/ffi-shape/json.md +++ b/quest/m1/ffi-shape/json.md @@ -5,7 +5,7 @@ JSON tracks live under `json` and opaque tracks under `flate` in moq-ffi and every wrapper, constructed from a track producer or consumer as in `moq-json` and `moq-flate`, and `BroadcastProducer`/`BroadcastConsumer` lose -`publish_json_*`/`subscribe_json_*` and `publish_binary_*`. The per-language +`publish_json_*`/`subscribe_json_*` and `publish_flate_*`. The per-language namespace pattern this sets is what the other children copy. ## Plan @@ -26,10 +26,8 @@ keep doing so, and the rest may follow. Watch for Go import cycles: a subpackage takes the root's broadcast handle, so the root must not import it back. -`flate` is the same shape over opaque bytes: moq-ffi's `binary.rs` -(`publish_binary_snapshot`, `publish_binary_stream`, #4137) moves under it, -mirroring `moq_flate::{snapshot, stream}` once -[moq-binary folds into moq-flate](/quest/m1/flate-binary.md). If that fold has -not landed, name the namespace `flate` anyway rather than `binary`. +`flate` is the same shape over opaque bytes: moq-ffi's `flate.rs` +(`publish_flate_snapshot`, `publish_flate_stream`) moves under it, mirroring +`moq_flate::{snapshot, stream}`. Public API: breaking in every binding. Wire: none. diff --git a/quest/m1/flate-binary.md b/quest/m1/flate-binary.md deleted file mode 100644 index 6cab778cb7..0000000000 --- a/quest/m1/flate-binary.md +++ /dev/null @@ -1,50 +0,0 @@ -# [M] moq-binary folds into moq-flate - -## Goal - -Opaque binary tracks live in `moq-flate` and `@moq/flate`: the `snapshot` and -`stream` modes `moq-binary` and `@moq/binary` provide today move there beside -the group-scoped codec, and `moq-binary` and `@moq/binary` are deleted. One -package owns compressed and opaque tracks, so there is no second "compressed -track" wrapper to build. - -## Plan - -Decided in the 2026-09-28 quest audit: `moq-binary` already composes -`moq-flate` into per-group windows (each group one sync-flushed DEFLATE -stream), which is what the m2 flate line planned to add as a new track wrapper. -Folding the two removes the duplicate instead of building it. - -- Rust: move `rs/moq-binary/src/{snapshot,stream}` and `Compression` into - `moq-flate` as `moq_flate::{snapshot, stream}`, keeping the codec - (`Encoder`/`Decoder`) at the root. `moq-flate` gains the `moq-net` - dependency. Delete `rs/moq-binary` and its workspace member, and repoint - `moq-mux` (`src/binary.rs`, `src/error.rs`) and `rs/moq-c` if it still - exists. -- JS: move `js/binary/src/{snapshot,stream,compression.ts}` into `@moq/flate` - as `Snapshot` and `Stream`, adding the `@moq/net` and `@moq/signals` - dependencies, and delete `js/binary`. -- Wire and catalog: unchanged. The hang catalog's `binary` section and - `moq_mux::binary` keep their names; they describe the track's content, and - the catalog section is wire. -- moq-ffi: rename `binary.rs` and its `publish_binary_*`, `MoqBinaryConfig`, - and producer types after `flate`, so every binding names the crate it wraps. - If [FFI shape](/quest/m1/ffi-shape/README.md) has already given them a - `flate` namespace, follow it instead. -- Open: whether `Compression::None` survives the move. An uncompressed opaque - track still needs a home, so the recommendation is to keep it and document - that the crate name is not a promise every track is deflated. -- Docs: fold `doc/lib/rs/moq-binary.md` into a `moq-flate` page and - `doc/lib/js/binary.md` into a `@moq/flate` page, fix `doc/.vitepress/config.ts`, - `doc/lib/{rs,js}/index.md`, `doc/concept/hang.md`, the android workflow path - filter, and add an upgrade note in `doc/setup/upgrade.md`. Grep for - `moq-binary`, `moq_binary`, and `@moq/binary`. - -Public API: breaking. `moq-binary` and `@moq/binary` are published and -deleted, and moq-ffi's binary names change, so this lands on `dev`. -`moq-flate` and `@moq/flate` grow additively. Wire: none. - -## Related - -- [Compressed tracks](/quest/m2/flate/README.md) - the hand-written wrappers expose these tracks -- [FFI shape](/quest/m1/ffi-shape/README.md) - gives the flate tracks their binding namespace diff --git a/quest/m1/track-demand.md b/quest/m1/track-demand.md index 82fdcafe4c..499b5e6664 100644 --- a/quest/m1/track-demand.md +++ b/quest/m1/track-demand.md @@ -4,14 +4,14 @@ In Rust and JS, a track's subscribers are watched only through its `Demand`: `track::Producer` drops `is_used`/`used`/`unused` and every layer producer -built on a track (moq-json, moq-mux import, moq-audio, moq-binary) exposes +built on a track (moq-json, moq-mux import, moq-audio, moq-flate) exposes `demand()` instead of its own copies. Group and broadcast `used`/`unused` stay; `Demand` is a track concept, and group demand drives fetch coalescing. ## Plan About 150 call sites move, across moq-net, moq-mux, moq-json, moq-audio, -moq-binary, moq-relay, moq-transcode, moq-stats, and moq-c. Both waits +moq-flate, moq-relay, moq-transcode, moq-stats, and moq-c. Both waits already surface the track's abort reason, so callers keep their errors. Keep `abort_unused` if its race still needs an owner. JS mirrors the Rust shape in `js/net`. diff --git a/quest/m2/flate/README.md b/quest/m2/flate/README.md index 58936a70b3..d685a05710 100644 --- a/quest/m2/flate/README.md +++ b/quest/m2/flate/README.md @@ -8,9 +8,8 @@ are identical across all of them. ## Plan -`moq-flate` and `@moq/flate` absorb `moq-binary`'s snapshot and stream modes -in [moq-binary folds into moq-flate](/quest/m1/flate-binary.md), so the crate -already owns the per-group window a caller could otherwise desynchronize. The +`moq-flate` and `@moq/flate` own the opaque snapshot and stream modes, so the +crate already owns the per-group window a caller could otherwise desynchronize. The track wrapper this line once planned was dropped for that reason. What remains is reaching those tracks from the hand-written binding wrappers. diff --git a/quest/m2/flate/bindings.md b/quest/m2/flate/bindings.md index ef28b41d31..13c3c21db8 100644 --- a/quest/m2/flate/bindings.md +++ b/quest/m2/flate/bindings.md @@ -9,10 +9,8 @@ decodes in the browser with `@moq/flate` and vice versa. ## Plan -moq-ffi publishes opaque tracks today (`publish_binary_snapshot` and -`publish_binary_stream`, #4137), renamed after `flate` by -[moq-binary folds into moq-flate](/quest/m1/flate-binary.md). Only the -generated bindings reach them; no wrapper does. This quest binds the existing +moq-ffi publishes opaque tracks today (`publish_flate_snapshot` and +`publish_flate_stream`). Only the generated bindings reach them; no wrapper does. This quest binds the existing track modes, not the bare codec: a `frame()` call across the FFI boundary invites the window desync the track modes exist to prevent. @@ -30,7 +28,3 @@ invites the window desync the track modes exist to prevent. `just test interop --all`. Public API: additive on moq-ffi and every wrapper. Wire: none. - -## Required - -- [moq-binary folds into moq-flate](/quest/m1/flate-binary.md) - the flate snapshot and stream tracks and their moq-ffi names diff --git a/quest/m2/teleop/README.md b/quest/m2/teleop/README.md index fa7e5f5a4f..2f8062cec1 100644 --- a/quest/m2/teleop/README.md +++ b/quest/m2/teleop/README.md @@ -45,7 +45,7 @@ study on a long-RTT profile measured command staleness more than twice as bad for QUIC reliable streams as for DDS best-effort: correct and useless. The split is a framing decision, not a subscription flag, and `moq-json` and -`moq-binary` already implement both halves as their snapshot and stream +`moq-flate` already implement both halves as their snapshot and stream modes. What that means for the primitive is in [robot](/quest/m2/teleop/robot.md), and what it means for a protocol multiplexing many message rates onto one link is in diff --git a/quest/m2/teleop/mavlink.md b/quest/m2/teleop/mavlink.md index 5d545ee4d3..91a5098faf 100644 --- a/quest/m2/teleop/mavlink.md +++ b/quest/m2/teleop/mavlink.md @@ -42,8 +42,7 @@ them all in one long-lived group has the opposite failure: nothing can be skipped and head-of-line blocking is back. So the lossy class is a latest-value snapshot of opaque bytes: -`moq_binary::snapshot` (moving to `moq_flate::snapshot` in -[moq-binary folds into moq-flate](/quest/m1/flate-binary.md)), with the raw +`moq_flate::snapshot`, with the raw frame as the value. Every update is a self-contained group, so a newer value never waits behind an older one. One detail decides whether it actually delivers latest-value: diff --git a/quest/m2/teleop/robot.md b/quest/m2/teleop/robot.md index 5dd2735138..7417c46432 100644 --- a/quest/m2/teleop/robot.md +++ b/quest/m2/teleop/robot.md @@ -72,8 +72,7 @@ The framing is where the guarantee lives, not the subscription flags: - Announce-prefix fan-in, generalised from `rs/moq-boy/src/input.rs`. - The two delivery classes, as the snapshot and stream modes with the group structure and `Info::max_age` each one needs: `moq-json`'s for JSON, and the - opaque-bytes ones for binary frames (`moq-binary`, folding into `moq-flate` - per [moq-binary folds into moq-flate](/quest/m1/flate-binary.md)). + opaque-bytes ones for binary frames (`moq-flate`). - Per-stage timestamp instrumentation, generalised from moq-boy's `status` track. Check it against the publisher-reported stats broadcast ([client stats](/quest/m1/qos/stats/schema.md), moq#2734) before adding a diff --git a/rs/moq-binary/CHANGELOG.md b/rs/moq-binary/CHANGELOG.md deleted file mode 100644 index 23d61fa49e..0000000000 --- a/rs/moq-binary/CHANGELOG.md +++ /dev/null @@ -1,65 +0,0 @@ -# Changelog - -All notable changes to this project will be documented in this file. - -The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), -and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). - -## [Unreleased] - -## [0.1.7](https://github.com/moq-dev/moq/compare/moq-binary-v0.1.6...moq-binary-v0.1.7) - 2026-09-27 - -### Added - -- *(mux)* detect delay and jitter on JSON and binary tracks ([#4270](https://github.com/moq-dev/moq/pull/4270)) - -## [0.1.6](https://github.com/moq-dev/moq/compare/moq-binary-v0.1.5...moq-binary-v0.1.6) - 2026-09-26 - -### Other - -- updated the following local packages: kio, moq-net - -## [0.1.5](https://github.com/moq-dev/moq/compare/moq-binary-v0.1.4...moq-binary-v0.1.5) - 2026-09-26 - -### Other - -- updated the following local packages: moq-net - -## [0.1.4](https://github.com/moq-dev/moq/compare/moq-binary-v0.1.3...moq-binary-v0.1.4) - 2026-09-25 - -### Other - -- updated the following local packages: moq-net - -## [0.1.3](https://github.com/moq-dev/moq/compare/moq-binary-v0.1.2...moq-binary-v0.1.3) - 2026-09-25 - -### Other - -- updated the following local packages: moq-net - -## [0.1.2](https://github.com/moq-dev/moq/compare/moq-binary-v0.1.1...moq-binary-v0.1.2) - 2026-09-25 - -### Other - -- updated the following local packages: moq-flate, moq-net - -## [0.1.1](https://github.com/moq-dev/moq/compare/moq-binary-v0.1.0...moq-binary-v0.1.1) - 2026-09-24 - -### Other - -- fill release doc gaps ([#4027](https://github.com/moq-dev/moq/pull/4027)) - -## [0.1.0](https://github.com/moq-dev/moq/releases/tag/moq-binary-v0.1.0) - 2026-09-23 - -### Added - -- *(net)* [**breaking**] slim the moq-net public surface ([#3779](https://github.com/moq-dev/moq/pull/3779)) -- *(json)* [**breaking**] Config means the same thing in json and binary ([#3718](https://github.com/moq-dev/moq/pull/3718)) -- *(moq-net)* [**breaking**] abort an oversized group instead of shedding its head ([#3585](https://github.com/moq-dev/moq/pull/3585)) -- *(hang)* add json and binary data tracks to the catalog ([#3109](https://github.com/moq-dev/moq/pull/3109)) - -### Fixed - -- *(moq-json)* a lost snapshot group is not fatal ([#3907](https://github.com/moq-dev/moq/pull/3907)) -- *(net)* commit an unused-driven track teardown atomically ([#3455](https://github.com/moq-dev/moq/pull/3455)) -- *(stream)* detect a rolled log without waiting for the group to end ([#3282](https://github.com/moq-dev/moq/pull/3282)) diff --git a/rs/moq-binary/Cargo.toml b/rs/moq-binary/Cargo.toml deleted file mode 100644 index 439d2ea6eb..0000000000 --- a/rs/moq-binary/Cargo.toml +++ /dev/null @@ -1,21 +0,0 @@ -[package] -name = "moq-binary" -description = "Binary publishing over MoQ tracks: latest-value snapshots, or append-log streams." -authors = ["Luke Curley "] -repository = "https://github.com/moq-dev/moq" -license = "MIT OR Apache-2.0" - -version = "0.1.7" -edition = "2024" -rust-version.workspace = true - -keywords = ["quic", "http3", "webtransport", "binary", "live"] -categories = ["multimedia", "network-programming", "web-programming"] - -[dependencies] -bytes = { workspace = true } -kio = { workspace = true } -moq-flate = { workspace = true } -moq-net = { workspace = true } -thiserror = { workspace = true } -tracing = { workspace = true } diff --git a/rs/moq-binary/README.md b/rs/moq-binary/README.md deleted file mode 100644 index bcd4c2519f..0000000000 --- a/rs/moq-binary/README.md +++ /dev/null @@ -1,24 +0,0 @@ -[![Documentation](https://docs.rs/moq-binary/badge.svg)](https://docs.rs/moq-binary/) -[![Crates.io](https://img.shields.io/crates/v/moq-binary.svg)](https://crates.io/crates/moq-binary) -[![License: MIT](https://img.shields.io/badge/License-MIT-blue.svg)](https://github.com/moq-dev/moq/blob/main/LICENSE-MIT) - -# moq-binary - -Opaque binary payloads over [`moq-net`](https://docs.rs/moq-net) tracks, in two -modes: - -- **snapshot**: lossy latest value. A consumer gets only the most recent payload. -- **stream**: lossless append-log. Every payload is delivered in order. - -The bytes are never inspected. Compression is -[`moq-flate`](https://docs.rs/moq-flate), the same group-scoped DEFLATE -[`moq-json`](https://docs.rs/moq-json) uses, which adds merge-patch deltas for -JSON documents. The TypeScript twin is -[`@moq/binary`](https://www.npmjs.com/package/@moq/binary). - -```bash -cargo add moq-binary -``` - -See [doc.moq.dev](https://doc.moq.dev/lib/rs/moq-binary) and -[docs.rs/moq-binary](https://docs.rs/moq-binary). diff --git a/rs/moq-binary/src/lib.rs b/rs/moq-binary/src/lib.rs deleted file mode 100644 index a38495c204..0000000000 --- a/rs/moq-binary/src/lib.rs +++ /dev/null @@ -1,70 +0,0 @@ -//! Opaque binary payloads over [`moq-net`](moq_net) tracks, in two modes: -//! -//! - [`snapshot`]: **lossy**. One value updated over time; a consumer only gets the most recent -//! one. Older values are superseded and dropped. -//! - [`stream`]: **lossless**. An ordered append-log of self-contained payloads, delivered in order -//! with nothing superseded. Bounded by the group cache: see [`stream`] for what that costs a -//! consumer that falls behind. -//! -//! Pick [`snapshot`] when consumers care about "what is the value now" (a poster image, a -//! serialized state blob) and [`stream`] when they care about every payload (an event log, a -//! sequence of samples). -//! -//! The bytes are opaque: this crate frames them onto a track and optionally compresses them, and -//! never looks inside. For JSON documents reach for [`moq-json`](https://docs.rs/moq-json) instead, -//! which adds RFC 7396 merge-patch deltas on top of the same two modes. -//! -//! Compression is [`moq-flate`](moq_flate), the same group-scoped DEFLATE moq-json uses, so the two -//! agree on the wire: each group is one raw DEFLATE stream, sync-flushed at every frame boundary. A -//! [`stream`] therefore compresses each payload against the earlier ones in its group, while a -//! [`snapshot`] group holds a single self-contained value. - -// The browser transport is `!Send`, so on wasm the shared state behind these `Arc`s is -// too and clippy suggests `Rc`. The same code is genuinely cross-thread on native, so -// `Arc` stays and the lint is unactionable here. -#![cfg_attr(target_arch = "wasm32", allow(clippy::arc_with_non_send_sync))] - -pub mod snapshot; -pub mod stream; - -/// How a binary track compresses its frames. -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] -pub enum Compression { - /// Uncompressed payloads. - #[default] - None, - - /// Group-scoped raw DEFLATE, sync-flushed at each frame boundary. - Deflate, -} - -impl Compression { - pub(crate) const fn is_deflate(self) -> bool { - matches!(self, Self::Deflate) - } -} - -/// Errors produced while publishing or consuming binary payloads. -#[derive(thiserror::Error, Debug, Clone)] -#[non_exhaustive] -pub enum Error { - /// An error from the underlying track. - #[error(transparent)] - Net(#[from] moq_net::Error), - - /// A compressed frame could not be decoded (malformed, truncated, or oversized). - #[error(transparent)] - Flate(#[from] moq_flate::Error), - - /// A [`stream`] track carried a second group, which a lossless log cannot do. - /// - /// A stream is a single group by construction: a publisher that cannot write a payload ends - /// the track rather than rolling. A second group therefore means the records that would have - /// completed the first one are gone, so the read fails instead of presenting the remainder as - /// a continuous log. - #[error("stream rolled to a second group")] - Rolled, -} - -/// A [`Result`](std::result::Result) using this crate's [`Error`]. -pub type Result = std::result::Result; diff --git a/rs/moq-ffi/src/binary.rs b/rs/moq-ffi/src/flate.rs similarity index 72% rename from rs/moq-ffi/src/binary.rs rename to rs/moq-ffi/src/flate.rs index 3c231757f7..a447306ec1 100644 --- a/rs/moq-ffi/src/binary.rs +++ b/rs/moq-ffi/src/flate.rs @@ -1,8 +1,8 @@ -//! Binary data tracks over the FFI boundary, advertised in the catalog. +//! Opaque data tracks over the FFI boundary, advertised in the catalog. //! -//! The binary counterpart of [`crate::json`]: opaque payloads (for example a camera's latest JPEG -//! thumbnail) on a named track, in either mode — `snapshot` (each payload supersedes the last) or -//! `stream` (every payload preserved in order). The broadcast's catalog carries +//! The `moq-flate` counterpart of [`crate::json`]: opaque payloads (for example a camera's latest +//! JPEG thumbnail) on a named track, in either mode: `snapshot` (each payload supersedes the last) +//! or `stream` (every payload preserved in order). The broadcast's catalog carries //! `binary.tracks.` (mode, plus `mime` and `compression` when set) for as long as the track //! lives, so a consumer discovers it without knowing the application. @@ -13,9 +13,9 @@ use moq_mux::catalog::hang::Extra; use crate::error::MoqError; use crate::producer::MoqBroadcastProducer; -/// Options for a binary data track, in either mode (the mode is fixed by the constructor). +/// Options for an opaque data track, in either mode (the mode is fixed by the constructor). #[derive(Clone, uniffi::Record)] -pub struct MoqBinaryConfig { +pub struct MoqFlateConfig { /// DEFLATE-compress each payload, advertised in the catalog entry. #[uniffi(default = false)] pub compression: bool, @@ -25,8 +25,8 @@ pub struct MoqBinaryConfig { pub mime: Option, } -impl From for moq_mux::binary::Config { - fn from(config: MoqBinaryConfig) -> Self { +impl From for moq_mux::binary::Config { + fn from(config: MoqFlateConfig) -> Self { let mut out = moq_mux::binary::Config::default().with_compression(config.compression); if let Some(mime) = config.mime { out = out.with_mime(mime); @@ -37,41 +37,41 @@ impl From for moq_mux::binary::Config { #[uniffi::export] impl MoqBroadcastProducer { - /// Publish a binary snapshot track (lossy latest-value) by name, advertised in the catalog. + /// Publish an opaque snapshot track (lossy latest-value) by name, advertised in the catalog. /// /// Errors if the catalog already carries an entry under `name`. - pub fn publish_binary_snapshot( + pub fn publish_flate_snapshot( &self, name: String, - config: MoqBinaryConfig, - ) -> Result, MoqError> { + config: MoqFlateConfig, + ) -> Result, MoqError> { let _guard = crate::ffi::enter(); self.with_state(|state| { let track = state.broadcast.create_track(name, None)?; let producer = state .catalog .binary_snapshot(track, moq_mux::binary::Config::from(config))?; - Ok(Arc::new(MoqBinarySnapshotProducer { + Ok(Arc::new(MoqFlateSnapshotProducer { inner: std::sync::Mutex::new(Some(producer)), })) }) } - /// Publish a binary stream track (lossless append-log) by name, advertised in the catalog. + /// Publish an opaque stream track (lossless append-log) by name, advertised in the catalog. /// /// Errors if the catalog already carries an entry under `name`. - pub fn publish_binary_stream( + pub fn publish_flate_stream( &self, name: String, - config: MoqBinaryConfig, - ) -> Result, MoqError> { + config: MoqFlateConfig, + ) -> Result, MoqError> { let _guard = crate::ffi::enter(); self.with_state(|state| { let track = state.broadcast.create_track(name, None)?; let producer = state .catalog .binary_stream(track, moq_mux::binary::Config::from(config))?; - Ok(Arc::new(MoqBinaryStreamProducer { + Ok(Arc::new(MoqFlateStreamProducer { inner: std::sync::Mutex::new(Some(producer)), })) }) @@ -80,12 +80,12 @@ impl MoqBroadcastProducer { /// Publishes opaque payloads that consumers see as a single latest value. #[derive(uniffi::Object)] -pub struct MoqBinarySnapshotProducer { +pub struct MoqFlateSnapshotProducer { inner: std::sync::Mutex>>, } #[uniffi::export] -impl MoqBinarySnapshotProducer { +impl MoqFlateSnapshotProducer { /// Publish a new payload, superseding the last. pub fn update(&self, payload: Vec) -> Result<(), MoqError> { let _guard = crate::ffi::enter(); @@ -105,12 +105,12 @@ impl MoqBinarySnapshotProducer { /// Publishes an ordered log of opaque payloads, one per append. #[derive(uniffi::Object)] -pub struct MoqBinaryStreamProducer { +pub struct MoqFlateStreamProducer { inner: std::sync::Mutex>>, } #[uniffi::export] -impl MoqBinaryStreamProducer { +impl MoqFlateStreamProducer { /// Append one payload to the log. pub fn append(&self, payload: Vec) -> Result<(), MoqError> { let _guard = crate::ffi::enter(); diff --git a/rs/moq-ffi/src/lib.rs b/rs/moq-ffi/src/lib.rs index 8b2743fdef..d1a44519b8 100644 --- a/rs/moq-ffi/src/lib.rs +++ b/rs/moq-ffi/src/lib.rs @@ -17,11 +17,11 @@ mod android; #[cfg(all(feature = "audio", not(target_arch = "wasm32")))] pub mod audio; pub mod bandwidth; -pub mod binary; pub mod consumer; pub mod demand; pub mod error; mod ffi; +pub mod flate; pub mod json; #[cfg(not(target_arch = "wasm32"))] mod log; diff --git a/rs/moq-ffi/src/test.rs b/rs/moq-ffi/src/test.rs index f1430eafe1..446fdf4640 100644 --- a/rs/moq-ffi/src/test.rs +++ b/rs/moq-ffi/src/test.rs @@ -2,12 +2,12 @@ use super::origin::*; use super::producer::*; use super::server::MoqServer; use super::session::{MoqClient, MoqSession}; -use crate::binary::MoqBinaryConfig; use crate::consumer::MoqBroadcastConsumer; use crate::consumer::MoqFetchGroupOptions; use crate::consumer::MoqSubscription; use crate::consumer::MoqTrackConsumer; use crate::error::MoqError; +use crate::flate::MoqFlateConfig; use crate::json::{MoqJsonSnapshotConfig, MoqJsonStreamConfig}; use crate::media::{MoqAudioFormat, MoqAudioInit, MoqFrame, MoqVideoFormat, MoqVideoInit}; use crate::session::{MoqBackoff, MoqConnectionStatus}; @@ -4546,23 +4546,23 @@ async fn json_tracks_are_advertised_in_the_catalog() { assert!(published_catalog(&broadcast).json.tracks.is_empty()); } -/// Binary tracks carry their mode and (optional) media type in the catalog. +/// Flate tracks carry their mode and (optional) media type in the catalog. #[tokio::test] -async fn binary_tracks_are_advertised_in_the_catalog() { +async fn flate_tracks_are_advertised_in_the_catalog() { let broadcast = MoqBroadcastProducer::new().unwrap(); let thumb = broadcast - .publish_binary_snapshot( + .publish_flate_snapshot( "thumbnail".into(), - MoqBinaryConfig { + MoqFlateConfig { compression: false, mime: Some("image/jpeg".into()), }, ) .unwrap(); let log = broadcast - .publish_binary_stream( + .publish_flate_stream( "log".into(), - MoqBinaryConfig { + MoqFlateConfig { compression: false, mime: None, }, @@ -4604,9 +4604,9 @@ async fn data_track_names_cannot_collide() { .unwrap(); assert!( broadcast - .publish_binary_stream( + .publish_flate_stream( "state".into(), - MoqBinaryConfig { + MoqFlateConfig { compression: false, mime: None, }, diff --git a/rs/moq-flate/Cargo.toml b/rs/moq-flate/Cargo.toml index 4c9a396201..4c441f665a 100644 --- a/rs/moq-flate/Cargo.toml +++ b/rs/moq-flate/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "moq-flate" -description = "Group-scoped DEFLATE: a stream of self-delimited frames sharing one compression window." +description = "Opaque MoQ tracks, optionally compressed with group-scoped DEFLATE: latest-value snapshots or append-log streams." authors = ["Luke Curley "] repository = "https://github.com/moq-dev/moq" license = "MIT OR Apache-2.0" @@ -9,12 +9,15 @@ version = "0.2.0" edition = "2024" rust-version.workspace = true -keywords = ["deflate", "compression", "streaming", "live"] -categories = ["compression"] +keywords = ["deflate", "compression", "quic", "moq", "live"] +categories = ["compression", "network-programming"] [dependencies] bytes = { workspace = true } # zlib-rs, not the default miniz_oxide: miniz_oxide 0.9.1 can return from a sync flush with output # space left but the flush unfinished, silently truncating the frame. flate2 = { workspace = true, features = ["zlib-rs", "runtime_detection"] } +kio = { workspace = true } +moq-net = { workspace = true } thiserror = { workspace = true } +tracing = { workspace = true } diff --git a/rs/moq-flate/README.md b/rs/moq-flate/README.md new file mode 100644 index 0000000000..2e05b9e918 --- /dev/null +++ b/rs/moq-flate/README.md @@ -0,0 +1,25 @@ +[![Documentation](https://docs.rs/moq-flate/badge.svg)](https://docs.rs/moq-flate/) +[![Crates.io](https://img.shields.io/crates/v/moq-flate.svg)](https://crates.io/crates/moq-flate) +[![License: MIT](https://img.shields.io/badge/License-MIT-blue.svg)](https://github.com/moq-dev/moq/blob/main/LICENSE-MIT) + +# moq-flate + +Opaque byte tracks over [`moq-net`](https://docs.rs/moq-net), optionally +compressed with group-scoped DEFLATE, in two modes: + +- **snapshot**: lossy latest value. A consumer gets only the most recent payload. +- **stream**: lossless append-log. Every payload is delivered in order. + +The bytes are never inspected. Compression is opt-in per track: each group is +one raw DEFLATE stream, sync-flushed at every frame boundary, so later frames +reuse earlier ones as context. The bare `Encoder`/`Decoder` codec is exported +too, and [`moq-json`](https://docs.rs/moq-json) builds on it to add merge-patch +deltas for JSON documents. The TypeScript twin is +[`@moq/flate`](https://www.npmjs.com/package/@moq/flate). + +```bash +cargo add moq-flate +``` + +See [doc.moq.dev](https://doc.moq.dev/lib/rs/moq-flate) and +[docs.rs/moq-flate](https://docs.rs/moq-flate). diff --git a/rs/moq-flate/src/lib.rs b/rs/moq-flate/src/lib.rs index d6009e7766..717ae878c8 100644 --- a/rs/moq-flate/src/lib.rs +++ b/rs/moq-flate/src/lib.rs @@ -1,11 +1,32 @@ -//! Group-scoped DEFLATE: a stream of self-delimited frames sharing one compression window. +//! Opaque byte tracks over [`moq-net`](moq_net), optionally compressed with group-scoped DEFLATE. //! -//! A sequence of frame payloads is compressed into a single raw DEFLATE ([RFC 1951]) stream, -//! sync-flushed at each frame boundary. Every frame is therefore self-delimited (byte-aligned, the -//! window retained) while later frames reuse the earlier ones as context, so a stream of similar -//! payloads (a snapshot followed by deltas, repeated records, log lines) compresses far better than -//! each payload alone. The [`Encoder`]/[`Decoder`] hold that shared window; create a fresh pair per -//! independent stream (in moq-net terms, per group). +//! Two track modes carry the payloads: +//! +//! - [`snapshot`]: **lossy**. One value updated over time; a consumer only gets the most recent +//! one. Older values are superseded and dropped. +//! - [`stream`]: **lossless**. An ordered append-log of self-contained payloads, delivered in order +//! with nothing superseded. Bounded by the group cache: see [`stream`] for what that costs a +//! consumer that falls behind. +//! +//! Pick [`snapshot`] when consumers care about "what is the value now" (a poster image, a +//! serialized state blob) and [`stream`] when they care about every payload (an event log, a +//! sequence of samples). The bytes are opaque: the tracks frame them and optionally compress them, +//! and never look inside. For JSON documents reach for [`moq-json`](https://docs.rs/moq-json) +//! instead, which adds RFC 7396 merge-patch deltas on top of the same two modes and codec. +//! +//! Compression is opt-in per track ([`Compression`]), so the crate name is not a promise that every +//! track is deflated: [`Compression::None`] writes the bytes through untouched. +//! +//! # Codec +//! +//! Underneath, [`Encoder`]/[`Decoder`] compress a sequence of frame payloads into a single raw +//! DEFLATE ([RFC 1951]) stream, sync-flushed at each frame boundary. Every frame is therefore +//! self-delimited (byte-aligned, the window retained) while later frames reuse the earlier ones as +//! context, so a stream of similar payloads (a snapshot followed by deltas, repeated records, log +//! lines) compresses far better than each payload alone. The pair holds that shared window; create a +//! fresh pair per independent stream (in moq-net terms, per group). A [`stream`] track therefore +//! compresses each payload against the earlier ones in its group, while a [`snapshot`] group holds a +//! single self-contained value. //! //! This is plain raw DEFLATE with a `Z_SYNC_FLUSH` after each frame, so any peer using the same //! primitive (zlib's sync flush, the browser's `deflate-raw`) interoperates on the wire. There is no @@ -30,6 +51,14 @@ //! [RFC 1951]: https://www.rfc-editor.org/rfc/rfc1951.html //! [RFC 7692]: https://www.rfc-editor.org/rfc/rfc7692.html#section-7.2.1 +// The browser transport is `!Send`, so on wasm the shared state behind the track modes' `Arc`s is +// too and clippy suggests `Rc`. The same code is genuinely cross-thread on native, so `Arc` stays +// and the lint is unactionable here. +#![cfg_attr(target_arch = "wasm32", allow(clippy::arc_with_non_send_sync))] + +pub mod snapshot; +pub mod stream; + use bytes::Bytes; use flate2::{Compress, Decompress, FlushCompress, FlushDecompress, Status}; @@ -46,8 +75,25 @@ const SYNC_FLUSH_TAIL: [u8; 4] = [0x00, 0x00, 0xff, 0xff]; /// Scratch buffer size for the streaming (de)compress loops. const CHUNK: usize = 8 * 1024; -/// Errors produced while decoding a frame. -#[derive(thiserror::Error, Debug, Clone, PartialEq, Eq)] +/// How a [`snapshot`] or [`stream`] track compresses its frames. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub enum Compression { + /// Uncompressed payloads. + #[default] + None, + + /// Group-scoped raw DEFLATE, sync-flushed at each frame boundary. + Deflate, +} + +impl Compression { + pub(crate) const fn is_deflate(self) -> bool { + matches!(self, Self::Deflate) + } +} + +/// Errors produced while decoding a frame or publishing and consuming a track. +#[derive(thiserror::Error, Debug, Clone)] #[non_exhaustive] pub enum Error { /// A frame could not be decoded (malformed or truncated stream, or fed out of order). @@ -57,6 +103,19 @@ pub enum Error { /// A frame's decompressed size exceeded the configured limit (zip-bomb guard). #[error("decompressed frame exceeded {0} bytes")] TooLarge(u64), + + /// An error from the underlying track. + #[error(transparent)] + Net(#[from] moq_net::Error), + + /// A [`stream`] track carried a second group, which a lossless log cannot do. + /// + /// A stream is a single group by construction: a publisher that cannot write a payload ends + /// the track rather than rolling. A second group therefore means the records that would have + /// completed the first one are gone, so the read fails instead of presenting the remainder as + /// a continuous log. + #[error("stream rolled to a second group")] + Rolled, } /// A [`Result`](std::result::Result) using this crate's [`Error`]. @@ -325,7 +384,10 @@ mod test { #[test] fn decompress_rejects_garbage() { let mut dec = Decoder::new(); - assert_eq!(dec.frame(b"not a deflate stream at all"), Err(Error::Decompress)); + assert!(matches!( + dec.frame(b"not a deflate stream at all"), + Err(Error::Decompress) + )); } #[test] @@ -335,7 +397,7 @@ mod test { let slice = Encoder::new().frame(&payload); let mut dec = Decoder::with_max_frame_size(512); - assert_eq!(dec.frame(&slice), Err(Error::TooLarge(512))); + assert!(matches!(dec.frame(&slice), Err(Error::TooLarge(512)))); } #[test] diff --git a/rs/moq-binary/src/snapshot/consumer.rs b/rs/moq-flate/src/snapshot/consumer.rs similarity index 95% rename from rs/moq-binary/src/snapshot/consumer.rs rename to rs/moq-flate/src/snapshot/consumer.rs index 51d6b29d8a..9ce149d111 100644 --- a/rs/moq-binary/src/snapshot/consumer.rs +++ b/rs/moq-flate/src/snapshot/consumer.rs @@ -1,4 +1,4 @@ -//! Consuming a binary value from a track. +//! Consuming an opaque value from a track. use std::task::Poll; @@ -8,7 +8,7 @@ use crate::Result; pub use super::Config; -/// Consumes a binary value from a track, yielding the newest one. +/// Consumes an opaque value from a track, yielding the newest one. /// /// Jumps to the newest group and reads the value out of it, so a late joiner starts at the current /// value rather than replaying superseded ones. @@ -17,7 +17,7 @@ pub struct Consumer { group: Option, /// The DEFLATE decoder for the current group, `Some` while decompressing. A snapshot group is /// normally one frame, but the window is per group either way. - flate: Option, + flate: Option, compression: bool, } @@ -55,7 +55,7 @@ impl Consumer { match self.track.poll_next_group(waiter)? { Poll::Ready(Some(group)) => { self.group = Some(group); - self.flate = self.compression.then(moq_flate::Decoder::new); + self.flate = self.compression.then(crate::Decoder::new); } Poll::Ready(None) => break true, Poll::Pending => break false, diff --git a/rs/moq-binary/src/snapshot/mod.rs b/rs/moq-flate/src/snapshot/mod.rs similarity index 99% rename from rs/moq-binary/src/snapshot/mod.rs rename to rs/moq-flate/src/snapshot/mod.rs index f24ec05165..d1fda575f1 100644 --- a/rs/moq-binary/src/snapshot/mod.rs +++ b/rs/moq-flate/src/snapshot/mod.rs @@ -1,4 +1,4 @@ -//! Lossy latest-value binary publishing over [`moq-net`](moq_net) tracks. +//! Lossy latest-value opaque publishing over [`moq-net`](moq_net) tracks. //! //! One opaque value updated over time, for consumers that only care about the current state (a //! poster image, a serialized state blob). This mode is **lossy** by design: a consumer yields only diff --git a/rs/moq-binary/src/snapshot/producer.rs b/rs/moq-flate/src/snapshot/producer.rs similarity index 89% rename from rs/moq-binary/src/snapshot/producer.rs rename to rs/moq-flate/src/snapshot/producer.rs index 33d757ea51..cd817154a8 100644 --- a/rs/moq-binary/src/snapshot/producer.rs +++ b/rs/moq-flate/src/snapshot/producer.rs @@ -1,4 +1,4 @@ -//! Publishing a binary value over a track. +//! Publishing an opaque value over a track. use std::sync::{Arc, Mutex}; @@ -9,7 +9,7 @@ use crate::Result; pub use super::Config; -/// Publishes a binary value over a track, one value per group. +/// Publishes an opaque value over a track, one value per group. /// /// Each [`update`](Self::update) rolls a new group holding the whole value, so a consumer only ever /// needs the newest group and older ones are dropped. For a log where every payload survives, use @@ -75,12 +75,12 @@ impl Inner { let payload = match self.compression { true => { // Compression can take a large value under the group's frame limit, but every consumer - // decodes with moq-flate's default output cap, so publishing past it would advertise a + // decodes with the default output cap, so publishing past it would advertise a // value that always fails to read. Reject it here instead. - if payload.len() as u64 > moq_flate::DEFAULT_MAX_FRAME_SIZE { - return Err(moq_flate::Error::TooLarge(moq_flate::DEFAULT_MAX_FRAME_SIZE).into()); + if payload.len() as u64 > crate::DEFAULT_MAX_FRAME_SIZE { + return Err(crate::Error::TooLarge(crate::DEFAULT_MAX_FRAME_SIZE)); } - moq_flate::Encoder::new().frame(&payload) + crate::Encoder::new().frame(&payload) } false => payload, }; diff --git a/rs/moq-binary/src/stream/consumer.rs b/rs/moq-flate/src/stream/consumer.rs similarity index 94% rename from rs/moq-binary/src/stream/consumer.rs rename to rs/moq-flate/src/stream/consumer.rs index 049123b907..9ead657026 100644 --- a/rs/moq-binary/src/stream/consumer.rs +++ b/rs/moq-flate/src/stream/consumer.rs @@ -1,4 +1,4 @@ -//! Consuming an ordered log of binary payloads from a track. +//! Consuming an ordered log of opaque payloads from a track. use std::task::Poll; @@ -8,7 +8,7 @@ use crate::Result; pub use super::Config; -/// Consumes an ordered log of binary payloads from a track, yielding every one. +/// Consumes an ordered log of opaque payloads from a track, yielding every one. /// /// The log is a single group. That is what makes the mode lossless: rolling to a second group /// means the records that would have completed the first are gone, so a publisher that cannot @@ -28,7 +28,7 @@ pub struct Consumer { /// fails too rather than reporting the rest of the log as a whole one. rolled: bool, /// The DEFLATE decoder for the group, `Some` while decompressing. - flate: Option, + flate: Option, compression: bool, } @@ -68,7 +68,7 @@ impl Consumer { Poll::Ready(Some(_)) if self.taken => self.rolled = true, Poll::Ready(Some(group)) => { self.taken = true; - self.flate = self.compression.then(moq_flate::Decoder::new); + self.flate = self.compression.then(crate::Decoder::new); self.group = Some(group); } Poll::Ready(None) => return Poll::Ready(Ok(None)), diff --git a/rs/moq-binary/src/stream/mod.rs b/rs/moq-flate/src/stream/mod.rs similarity index 97% rename from rs/moq-binary/src/stream/mod.rs rename to rs/moq-flate/src/stream/mod.rs index aa1699fbde..5ed95d9bbf 100644 --- a/rs/moq-binary/src/stream/mod.rs +++ b/rs/moq-flate/src/stream/mod.rs @@ -1,4 +1,4 @@ -//! Lossless append-log binary publishing over [`moq-net`](moq_net) tracks. +//! Lossless append-log opaque publishing over [`moq-net`](moq_net) tracks. //! //! An ordered log of opaque payloads, for consumers that care about every one (an event log, a //! sequence of samples). Nothing is ever superseded: a consumer yields each payload in the order it @@ -231,11 +231,8 @@ mod test { let mut producer = Producer::new(track, cfg(true)); assert!(producer.is_used()); - let oversized = Bytes::from(vec![0u8; moq_flate::DEFAULT_MAX_FRAME_SIZE as usize + 1]); - assert!(matches!( - producer.append(oversized), - Err(crate::Error::Flate(moq_flate::Error::TooLarge(_))) - )); + let oversized = Bytes::from(vec![0u8; crate::DEFAULT_MAX_FRAME_SIZE as usize + 1]); + assert!(matches!(producer.append(oversized), Err(crate::Error::TooLarge(_)))); assert!(!producer.is_used()); @@ -312,7 +309,7 @@ mod test { for pair in payloads(4).chunks(2) { // Each group is its own DEFLATE stream, which is what a recovery roll would produce. - let mut flate = moq_flate::Encoder::new(); + let mut flate = crate::Encoder::new(); let mut group = track.append_group().unwrap(); for payload in pair { group @@ -387,7 +384,7 @@ mod test { // Publish sequence 1 before sequence 0, the way reordering delivers them. for sequence in [1u64, 0] { - let mut flate = moq_flate::Encoder::new(); + let mut flate = crate::Encoder::new(); let mut group = track.create_group(moq_net::group::Info { sequence }).unwrap(); group .write_frame(moq_net::Timestamp::now(), flate.frame(&[sequence as u8; 8])) diff --git a/rs/moq-binary/src/stream/producer.rs b/rs/moq-flate/src/stream/producer.rs similarity index 93% rename from rs/moq-binary/src/stream/producer.rs rename to rs/moq-flate/src/stream/producer.rs index d24e04ba32..9ccfa3b24d 100644 --- a/rs/moq-binary/src/stream/producer.rs +++ b/rs/moq-flate/src/stream/producer.rs @@ -1,4 +1,4 @@ -//! Publishing an ordered log of binary payloads over a track. +//! Publishing an ordered log of opaque payloads over a track. use std::sync::{Arc, Mutex}; @@ -9,7 +9,7 @@ use crate::Result; pub use super::Config; -/// Publishes an ordered log of binary payloads over a track, one payload per frame in a single +/// Publishes an ordered log of opaque payloads over a track, one payload per frame in a single /// group. /// /// Cheaply clonable: clones share one underlying track and publishing state, so multiple owners @@ -26,7 +26,7 @@ impl Producer { inner: Arc::new(Mutex::new(Inner { track, group: None, - flate: config.compression.is_deflate().then(moq_flate::Encoder::new), + flate: config.compression.is_deflate().then(crate::Encoder::new), })), } } @@ -73,7 +73,7 @@ struct Inner { group: Option, /// The DEFLATE encoder, one window for the whole group, `Some` while compressing. - flate: Option, + flate: Option, } impl Inner { @@ -85,9 +85,9 @@ impl Inner { // missing a record either way, and carrying on would present that gap as a complete log. // Checked before the group is opened, so nothing is published, and routed through the same // abort so a reader sees the failure rather than a clean end. - if self.flate.is_some() && payload.len() as u64 > moq_flate::DEFAULT_MAX_FRAME_SIZE { + if self.flate.is_some() && payload.len() as u64 > crate::DEFAULT_MAX_FRAME_SIZE { self.abort(moq_net::Error::FrameTooLarge); - return Err(moq_flate::Error::TooLarge(moq_flate::DEFAULT_MAX_FRAME_SIZE).into()); + return Err(crate::Error::TooLarge(crate::DEFAULT_MAX_FRAME_SIZE)); } // Open the group before compressing: a failure here must not leave the window ahead of a diff --git a/rs/moq-mux/Cargo.toml b/rs/moq-mux/Cargo.toml index 315fad5b99..3199311ce0 100644 --- a/rs/moq-mux/Cargo.toml +++ b/rs/moq-mux/Cargo.toml @@ -20,7 +20,7 @@ h264-parser = { version = "0.4.2" } hang = { workspace = true } kio = { workspace = true } memchr = "2" -moq-binary = { workspace = true } +moq-flate = { workspace = true } moq-json = { workspace = true } moq-loc = { workspace = true } moq-msf = { workspace = true } diff --git a/rs/moq-mux/src/binary.rs b/rs/moq-mux/src/binary.rs index 75b979df88..600c9b6977 100644 --- a/rs/moq-mux/src/binary.rs +++ b/rs/moq-mux/src/binary.rs @@ -136,7 +136,7 @@ fn prepare(config: &mut impl AsMut, mode: Mode) -> crate::Result { - inner: moq_binary::snapshot::Producer, + inner: moq_flate::snapshot::Producer, listing: Listing, /// Maps a payload's capture instant onto the broadcast timeline. clock: crate::Clock, @@ -150,11 +150,11 @@ impl Snapshot { rendition: crate::catalog::Rendition, mut config: C, ) -> crate::Result { - let mut binary = moq_binary::snapshot::Config::default(); + let mut binary = moq_flate::snapshot::Config::default(); if prepare(&mut config, Mode::Snapshot)? { - binary.compression = moq_binary::Compression::Deflate; + binary.compression = moq_flate::Compression::Deflate; } - let inner = moq_binary::snapshot::Producer::new(track, binary); + let inner = moq_flate::snapshot::Producer::new(track, binary); let clock = rendition.clock(); let listing = Listing::new(rendition, config)?; Ok(Self { @@ -200,7 +200,7 @@ impl Snapshot { /// Every [`append`](Self::append) is preserved and delivered in order. For a latest-value payload, /// use [`Snapshot`]. pub struct Stream { - inner: moq_binary::stream::Producer, + inner: moq_flate::stream::Producer, name: String, /// Cleared when a terminal failure ends the track, which retires the catalog entry with it. An @@ -219,11 +219,11 @@ impl Stream { rendition: crate::catalog::Rendition, mut config: C, ) -> crate::Result { - let mut binary = moq_binary::stream::Config::default(); + let mut binary = moq_flate::stream::Config::default(); if prepare(&mut config, Mode::Stream)? { - binary.compression = moq_binary::Compression::Deflate; + binary.compression = moq_flate::Compression::Deflate; } - let inner = moq_binary::stream::Producer::new(track, binary); + let inner = moq_flate::stream::Producer::new(track, binary); let clock = rendition.clock(); let listing = Listing::new(rendition, config)?; Ok(Self { @@ -251,7 +251,7 @@ impl Stream { /// Append one payload to the log. /// /// A payload that cannot be written ends the track (see - /// [`moq_binary::stream::Producer::append`]) and retires the catalog entry with it. A catalog + /// [`moq_flate::stream::Producer::append`]) and retires the catalog entry with it. A catalog /// error publishing the measured bitrate is returned after the payload was written, so the track /// stays open and a retry would duplicate it. pub fn append(&mut self, payload: impl Into>) -> crate::Result<()> { @@ -295,10 +295,10 @@ pub struct Consumer { mode: Mode, } -/// Which moq-binary consumer is doing the reading. Private: the caller sees one `Consumer`. +/// Which moq-flate consumer is doing the reading. Private: the caller sees one `Consumer`. enum Inner { - Snapshot(moq_binary::snapshot::Consumer), - Stream(moq_binary::stream::Consumer), + Snapshot(moq_flate::snapshot::Consumer), + Stream(moq_flate::stream::Consumer), } impl Consumer { @@ -315,18 +315,18 @@ impl Consumer { let inner = match &config.mode { Mode::Snapshot => { - let mut binary = moq_binary::snapshot::Config::default(); + let mut binary = moq_flate::snapshot::Config::default(); if compression { - binary.compression = moq_binary::Compression::Deflate; + binary.compression = moq_flate::Compression::Deflate; } - Inner::Snapshot(moq_binary::snapshot::Consumer::new(track, binary)) + Inner::Snapshot(moq_flate::snapshot::Consumer::new(track, binary)) } Mode::Stream => { - let mut binary = moq_binary::stream::Config::default(); + let mut binary = moq_flate::stream::Config::default(); if compression { - binary.compression = moq_binary::Compression::Deflate; + binary.compression = moq_flate::Compression::Deflate; } - Inner::Stream(moq_binary::stream::Consumer::new(track, binary)) + Inner::Stream(moq_flate::stream::Consumer::new(track, binary)) } other => return Err(crate::Error::UnsupportedMode(other.to_string())), }; diff --git a/rs/moq-mux/src/error.rs b/rs/moq-mux/src/error.rs index fd676051a9..046307d8a7 100644 --- a/rs/moq-mux/src/error.rs +++ b/rs/moq-mux/src/error.rs @@ -36,9 +36,9 @@ pub enum Error { #[error("json: {0}")] Json(#[from] moq_json::Error), - /// Error publishing or consuming binary payloads over a track. - #[error("binary: {0}")] - Binary(#[from] moq_binary::Error), + /// Error publishing or consuming opaque payloads over a track. + #[error("flate: {0}")] + Flate(#[from] moq_flate::Error), /// A catalog entry declares a track mode this build does not implement. #[error("unsupported track mode: {0}")]