From 46041a92141a91c565075ea903aaa15fd99a04e8 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sat, 28 Mar 2026 13:55:58 -0500 Subject: [PATCH] Revert "Revert "integrate MUXL segementation"" --- .ci/dockerfile-hash.yaml | 2 +- .github/workflows/build.yaml | 2 + .gitignore | 1 + Cargo.lock | 96 ++++++++++++ Cargo.toml | 2 +- Makefile | 18 ++- docker/build.Dockerfile | 2 + go.mod | 15 ++ go.sum | 34 ++++- pkg/config/config.go | 51 +++++++ pkg/director/s3_upload.go | 45 ++++++ pkg/director/stream_session.go | 7 + pkg/media/concat_demux.go | 3 +- pkg/media/io_helpers.go | 14 ++ pkg/media/muxl_segment.go | 120 +++++++++++++++ pkg/media/packetize.go | 1 + pkg/media/packetize_test.go | 19 ++- pkg/media/segmenter.go | 2 +- pkg/media/thumbnail.go | 112 +++++++++++--- pkg/media/thumbnail_test.go | 118 +++++++++------ pkg/model/segment.go | 1 - pkg/model/segment_test.go | 1 - pkg/muxl/muxl.go | 216 +++++++++++++++++++++++++++ pkg/s3/s3.go | 255 ++++++++++++++++++++++++++++++++ pkg/spxrpc/place_stream_live.go | 3 + rust/muxl-wasm/Cargo.toml | 11 ++ rust/muxl-wasm/src/main.rs | 3 + 27 files changed, 1074 insertions(+), 80 deletions(-) create mode 100644 pkg/director/s3_upload.go create mode 100644 pkg/media/muxl_segment.go delete mode 100644 pkg/model/segment.go delete mode 100644 pkg/model/segment_test.go create mode 100644 pkg/muxl/muxl.go create mode 100644 pkg/s3/s3.go create mode 100644 rust/muxl-wasm/Cargo.toml create mode 100644 rust/muxl-wasm/src/main.rs diff --git a/.ci/dockerfile-hash.yaml b/.ci/dockerfile-hash.yaml index 18b503a0..e55dc8fa 100644 --- a/.ci/dockerfile-hash.yaml +++ b/.ci/dockerfile-hash.yaml @@ -1,2 +1,2 @@ variables: - DOCKERFILE_HASH: 8ba2805e78d8d14ab757ef54f89a9778c925d6ad \ No newline at end of file + DOCKERFILE_HASH: 4dc3e0c009219b4ec68921e16fd9db640fcab34e \ No newline at end of file diff --git a/.github/workflows/build.yaml b/.github/workflows/build.yaml index 55dba125..55b62469 100644 --- a/.github/workflows/build.yaml +++ b/.github/workflows/build.yaml @@ -137,6 +137,8 @@ jobs: && corepack enable \ && export GOTOOLCHAIN=go1.25.1 \ && brew install go \ + && rustup target add wasm32-wasip1 \ + && rustup target add wasm32-unknown-unknown \ && python -m pip install virtualenv \ && python -m virtualenv ~/venv \ && source ~/venv/bin/activate \ diff --git a/.gitignore b/.gitignore index 4981419c..2ef79adb 100644 --- a/.gitignore +++ b/.gitignore @@ -25,3 +25,4 @@ my-release-key.keystore test.xml js/app/src/build-info.json *.env +muxl.wasm diff --git a/Cargo.lock b/Cargo.lock index 7bba713d..6ac7ccf4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -740,6 +740,12 @@ dependencies = [ "thiserror 1.0.69", ] +[[package]] +name = "cbor4ii" +version = "1.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "faed1a83001dc2c9201451030cc317e35bef36c84d3781d7c5bb9f343c397da8" + [[package]] name = "cc" version = "1.2.32" @@ -1184,12 +1190,49 @@ dependencies = [ "syn 2.0.105", ] +[[package]] +name = "dasl" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b59666035a4386b0fd272bd78da4cbc3ccb558941e97579ab00f0eb4639f2a49" +dependencies = [ + "blake3", + "cbor4ii", + "data-encoding", + "data-encoding-macro", + "scopeguard", + "serde", + "serde_bytes", + "sha2 0.10.9", + "thiserror 2.0.16", +] + [[package]] name = "data-encoding" version = "2.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2a2330da5de22e8a3cb63252ce2abb30116bf5265e89c0e01bc17015ce30a476" +[[package]] +name = "data-encoding-macro" +version = "0.1.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "47ce6c96ea0102f01122a185683611bd5ac8d99e62bc59dd12e6bda344ee673d" +dependencies = [ + "data-encoding", + "data-encoding-macro-internal", +] + +[[package]] +name = "data-encoding-macro-internal" +version = "0.1.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8d162beedaa69905488a8da94f5ac3edb4dd4788b732fadb7bd120b2625c1976" +dependencies = [ + "data-encoding", + "syn 2.0.105", +] + [[package]] name = "delegate" version = "0.8.0" @@ -3051,6 +3094,36 @@ dependencies = [ "thiserror 1.0.69", ] +[[package]] +name = "mp4-atom" +version = "0.10.1" +source = "git+https://github.com/streamplace/mp4-atom.git?branch=streamplace#126c49b8e3e3f8810089ed5cb157da94ce4f04b7" +dependencies = [ + "derive_more 2.0.1", + "num", + "paste", + "thiserror 1.0.69", + "tracing", +] + +[[package]] +name = "muxl" +version = "0.1.0" +source = "git+https://github.com/streamplace/s2pa-muxl.git?rev=5456732ab7844c7ffcba148583a08f4bf0b50dcd#5456732ab7844c7ffcba148583a08f4bf0b50dcd" +dependencies = [ + "dasl", + "mp4-atom", + "serde", + "serde_bytes", +] + +[[package]] +name = "muxl-wasm" +version = "0.1.0" +dependencies = [ + "muxl", +] + [[package]] name = "n0-future" version = "0.1.3" @@ -3338,6 +3411,20 @@ dependencies = [ "winapi", ] +[[package]] +name = "num" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35bd024e8b2ff75562e5f34e7f4905839deb4b22955ef5e73d2fea1b9813cb23" +dependencies = [ + "num-bigint", + "num-complex", + "num-integer", + "num-iter", + "num-rational", + "num-traits", +] + [[package]] name = "num-bigint" version = "0.4.6" @@ -3366,6 +3453,15 @@ dependencies = [ "zeroize", ] +[[package]] +name = "num-complex" +version = "0.4.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "73f88a1307638156682bada9d7604135552957b7818057dcef22705b4d509495" +dependencies = [ + "num-traits", +] + [[package]] name = "num-conv" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index 8c73d338..58ed954b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,4 +1,4 @@ [workspace] resolver = "3" -members = ["rust/export-c2pa-schema", "rust/iroh-streamplace"] +members = ["rust/export-c2pa-schema", "rust/iroh-streamplace", "rust/muxl-wasm"] exclude = ["subprojects/c2pa_go"] diff --git a/Makefile b/Makefile index 3a3f490f..5083a2cf 100644 --- a/Makefile +++ b/Makefile @@ -283,7 +283,7 @@ static-test: .PHONY: dev-setup dev-setup: - $(MAKE) -j16 app-cached dev-setup-meson + $(MAKE) -j16 app-cached dev-setup-meson muxl-wasm .PHONY: dev dev: app-cached $(LEXICON_STAMP) @@ -304,6 +304,11 @@ dev-setup-meson-configure: meson setup --default-library=shared $(BUILDDIR) $(SHARED_OPTS) meson configure --default-library=shared $(BUILDDIR) $(SHARED_OPTS) +.PHONY: muxl-wasm +muxl-wasm: + cargo build -p muxl-wasm --target wasm32-wasip1 --release + cp target/wasm32-wasip1/release/muxl-wasm.wasm pkg/muxl/muxl.wasm + .PHONY: dev-rust dev-rust: .build/bin/uniffi-bindgen-go-forked cargo build @@ -599,11 +604,11 @@ desktop-windows-amd64: && mv "js/desktop/out/make/squirrel.windows/x64/Streamplace-$(VERSION_ELECTRON) Setup.exe" ./bin/streamplace-desktop-$(VERSION)-windows-amd64.exe .PHONY: streamplace -streamplace: app-cached meson-setup-static +streamplace: app-cached meson-setup-static muxl-wasm meson compile -C $(BUILDDIR) streamplace | grep -v drectve .PHONY: archive -archive: app-cached meson-setup-static godeps +archive: app-cached meson-setup-static godeps muxl-wasm meson compile -C $(BUILDDIR) archive | grep -v drectve .PHONY: linux-amd64 @@ -711,6 +716,13 @@ link-ffmpeg: rm -rf subprojects/FFmpeg ln -s $$(realpath ../ffmpeg) ./subprojects/FFmpeg +.PHONY: build-muxl +build-muxl: + cd ../s2pa-muxl \ + && cargo build --target wasm32-wasip1 --release \ + && cd - \ + && cp ../s2pa-muxl/target/wasm32-wasip1/release/muxl.wasm pkg/muxl/muxl.wasm + # _____ ____ _____ _ ________ _____ # | __ \ / __ \ / ____| |/ / ____| __ \ # | | | | | | | | | ' /| |__ | |__) | diff --git a/docker/build.Dockerfile b/docker/build.Dockerfile index 8ba2805e..4dc3e0c0 100644 --- a/docker/build.Dockerfile +++ b/docker/build.Dockerfile @@ -74,6 +74,8 @@ RUN curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs > rustup.sh \ && rustup target add x86_64-pc-windows-gnu \ && rustup target add x86_64-apple-darwin \ && rustup target add aarch64-apple-darwin \ + && rustup target add wasm32-wasip1 \ + && rustup target add wasm32-unknown-unknown \ && rm rustup.sh RUN go env -w GOTOOLCHAIN=go$GO_VERSION diff --git a/go.mod b/go.mod index 81ecd07a..f2eca73f 100644 --- a/go.mod +++ b/go.mod @@ -19,6 +19,9 @@ require ( github.com/NYTimes/gziphandler v1.1.1 github.com/ThalesGroup/crypto11 v0.0.0-00010101000000-000000000000 github.com/acarl005/stripansi v0.0.0-20180116102854-5a71ef0e047d + github.com/aws/aws-sdk-go-v2 v1.41.4 + github.com/aws/aws-sdk-go-v2/credentials v1.19.12 + github.com/aws/aws-sdk-go-v2/service/s3 v1.97.1 github.com/bluenviron/gortmplib v0.1.2 github.com/bluenviron/gortsplib/v5 v5.2.1 github.com/bluesky-social/indigo v0.0.0-20251206005924-d49b45419635 @@ -36,6 +39,7 @@ require ( github.com/golangci/golangci-lint/v2 v2.1.6 github.com/google/uuid v1.6.0 github.com/gorilla/websocket v1.5.3 + github.com/hyphacoop/go-dasl v0.8.0 github.com/ipfs/go-cid v0.5.0 github.com/ipfs/go-ipld-cbor v0.2.0 github.com/ipld/go-car v0.6.1-0.20230509095817-92d28eb23ba4 @@ -65,6 +69,7 @@ require ( github.com/streamplace/oatproxy v0.0.0-20260318210219-4567e0e7926a github.com/stretchr/testify v1.11.1 github.com/tdewolff/canvas v0.0.0-20250728095813-50d4cb1eee71 + github.com/tetratelabs/wazero v1.11.0 github.com/urfave/cli/v3 v3.6.2 github.com/whyrusleeping/cbor-gen v0.3.1 github.com/whyrusleeping/go-did v0.0.0-20230824162731-404d1707d5d6 @@ -157,6 +162,15 @@ require ( github.com/asticode/go-astikit v0.30.0 // indirect github.com/asticode/go-astits v1.14.0 // indirect github.com/aws/aws-sdk-go v1.44.273 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.7 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.20 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.20 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.21 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.7 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.12 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.20 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.20 // indirect + github.com/aws/smithy-go v1.24.2 // indirect github.com/aymanbagabas/go-osc52/v2 v2.0.1 // indirect github.com/benoitkugler/textlayout v0.3.1 // indirect github.com/benoitkugler/textprocessing v0.0.3 // indirect @@ -285,6 +299,7 @@ require ( github.com/hashicorp/hcl v1.0.0 // indirect github.com/hexops/gotextdiff v1.0.3 // indirect github.com/holiman/uint256 v1.3.2 // indirect + github.com/hyphacoop/cbor/v2 v2.0.0-20251007204234-2a4fa83e606e // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect github.com/invopop/yaml v0.3.1 // indirect github.com/ipfs/bbloom v0.0.4 // indirect diff --git a/go.sum b/go.sum index 8311836c..698887c2 100644 --- a/go.sum +++ b/go.sum @@ -188,6 +188,30 @@ github.com/asticode/go-astits v1.14.0 h1:zkgnZzipx2XX5mWycqsSBeEyDH58+i4HtyF4j2R github.com/asticode/go-astits v1.14.0/go.mod h1:QSHmknZ51pf6KJdHKZHJTLlMegIrhega3LPWz3ND/iI= github.com/aws/aws-sdk-go v1.44.273 h1:CX8O0gK+cGrgUyv7bgJ6QQP9mQg7u5mweHdNzULH47c= github.com/aws/aws-sdk-go v1.44.273/go.mod h1:aVsgQcEevwlmQ7qHE9I3h+dtQgpqhFB+i8Phjh7fkwI= +github.com/aws/aws-sdk-go-v2 v1.41.4 h1:10f50G7WyU02T56ox1wWXq+zTX9I1zxG46HYuG1hH/k= +github.com/aws/aws-sdk-go-v2 v1.41.4/go.mod h1:mwsPRE8ceUUpiTgF7QmQIJ7lgsKUPQOUl3o72QBrE1o= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.7 h1:3kGOqnh1pPeddVa/E37XNTaWJ8W6vrbYV9lJEkCnhuY= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.7/go.mod h1:lyw7GFp3qENLh7kwzf7iMzAxDn+NzjXEAGjKS2UOKqI= +github.com/aws/aws-sdk-go-v2/credentials v1.19.12 h1:oqtA6v+y5fZg//tcTWahyN9PEn5eDU/Wpvc2+kJ4aY8= +github.com/aws/aws-sdk-go-v2/credentials v1.19.12/go.mod h1:U3R1RtSHx6NB0DvEQFGyf/0sbrpJrluENHdPy1j/3TE= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.20 h1:CNXO7mvgThFGqOFgbNAP2nol2qAWBOGfqR/7tQlvLmc= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.20/go.mod h1:oydPDJKcfMhgfcgBUZaG+toBbwy8yPWubJXBVERtI4o= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.20 h1:tN6W/hg+pkM+tf9XDkWUbDEjGLb+raoBMFsTodcoYKw= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.20/go.mod h1:YJ898MhD067hSHA6xYCx5ts/jEd8BSOLtQDL3iZsvbc= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.21 h1:SwGMTMLIlvDNyhMteQ6r8IJSBPlRdXX5d4idhIGbkXA= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.21/go.mod h1:UUxgWxofmOdAMuqEsSppbDtGKLfR04HGsD0HXzvhI1k= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.7 h1:5EniKhLZe4xzL7a+fU3C2tfUN4nWIqlLesfrjkuPFTY= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.7/go.mod h1:x0nZssQ3qZSnIcePWLvcoFisRXJzcTVvYpAAdYX8+GI= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.12 h1:qtJZ70afD3ISKWnoX3xB0J2otEqu3LqicRcDBqsj0hQ= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.12/go.mod h1:v2pNpJbRNl4vEUWEh5ytQok0zACAKfdmKS51Hotc3pQ= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.20 h1:2HvVAIq+YqgGotK6EkMf+KIEqTISmTYh5zLpYyeTo1Y= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.20/go.mod h1:V4X406Y666khGa8ghKmphma/7C0DAtEQYhkq9z4vpbk= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.20 h1:siU1A6xjUZ2N8zjTHSXFhB9L/2OY8Dqs0xXiLjF30jA= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.20/go.mod h1:4TLZCmVJDM3FOu5P5TJP0zOlu9zWgDWU7aUxWbr+rcw= +github.com/aws/aws-sdk-go-v2/service/s3 v1.97.1 h1:csi9NLpFZXb9fxY7rS1xVzgPRGMt7MSNWeQ6eo247kE= +github.com/aws/aws-sdk-go-v2/service/s3 v1.97.1/go.mod h1:qXVal5H0ChqXP63t6jze5LmFalc7+ZE7wOdLtZ0LCP0= +github.com/aws/smithy-go v1.24.2 h1:FzA3bu/nt/vDvmnkg+R8Xl46gmzEDam6mZ1hzmwXFng= +github.com/aws/smithy-go v1.24.2/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= github.com/aymanbagabas/go-osc52/v2 v2.0.1 h1:HwpRHbFMcZLEVr42D4p7XBqjyuxQH5SMiErDT4WkJ2k= github.com/aymanbagabas/go-osc52/v2 v2.0.1/go.mod h1:uYgXzlJ7ZpABp8OJ+exZzJJhRNQ2ASbcXHWsFqH8hp8= github.com/benbjohnson/clock v1.3.5 h1:VvXlSJBzZpA/zum6Sj74hxwYI2DIxRWuNIoXAzHZz5o= @@ -698,6 +722,10 @@ github.com/holiman/uint256 v1.3.2 h1:a9EgMPSC1AAaj1SZL5zIQD3WbwTuHrMGOerLjGmM/TA github.com/holiman/uint256 v1.3.2/go.mod h1:EOMSn4q6Nyt9P6efbI3bueV4e1b3dGlUCXeiRV4ng7E= github.com/huin/goupnp v1.3.0 h1:UvLUlWDNpoUdYzb2TCn+MuTWtcjXKSza2n6CBdQ0xXc= github.com/huin/goupnp v1.3.0/go.mod h1:gnGPsThkYa7bFi/KWmEysQRf48l2dvR5bxr2OFckNX8= +github.com/hyphacoop/cbor/v2 v2.0.0-20251007204234-2a4fa83e606e h1:4HqTNG0M8I/MYfra5Vn+4XAnuaivNqMuKpZu0B5hZjQ= +github.com/hyphacoop/cbor/v2 v2.0.0-20251007204234-2a4fa83e606e/go.mod h1:1ny0WdocVllO4iUxV5eXuJ9vMzmZ2nITeZg7uBxVEeU= +github.com/hyphacoop/go-dasl v0.8.0 h1:0K2nvGLUb8aBqG8GGUDbMgqCG9IVPT38hnpCVpzMCIU= +github.com/hyphacoop/go-dasl v0.8.0/go.mod h1:6Z8cAEAsv495mgxatpgIdEcJxDBc7Nbg/ZJX58iGJfU= github.com/ianlancetaylor/demangle v0.0.0-20181102032728-5e5cf60278f6/go.mod h1:aSSvb/t6k1mPoxDqO4vJh6VOCGPwU4O0C2/Eqndh1Sc= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= @@ -1364,6 +1392,8 @@ github.com/tenntenn/text/transform v0.0.0-20200319021203-7eef512accb3 h1:f+jULpR github.com/tenntenn/text/transform v0.0.0-20200319021203-7eef512accb3/go.mod h1:ON8b8w4BN/kE1EOhwT0o+d62W65a6aPw1nouo9LMgyY= github.com/tetafro/godot v1.5.1 h1:PZnjCol4+FqaEzvZg5+O8IY2P3hfY9JzRBNPv1pEDS4= github.com/tetafro/godot v1.5.1/go.mod h1:cCdPtEndkmqqrhiCfkmxDodMQJ/f3L1BCNskCUZdTwk= +github.com/tetratelabs/wazero v1.11.0 h1:+gKemEuKCTevU4d7ZTzlsvgd1uaToIDtlQlmNbwqYhA= +github.com/tetratelabs/wazero v1.11.0/go.mod h1:eV28rsN8Q+xwjogd7f4/Pp4xFxO7uOGbLcD/LzB1wiU= github.com/thales-e-security/pool v0.0.2 h1:RAPs4q2EbWsTit6tpzuvTFlgFRJ3S8Evf5gtvVDbmPg= github.com/thales-e-security/pool v0.0.2/go.mod h1:qtpMm2+thHtqhLzTwgDBj/OuNnMpupY8mv0Phz0gjhU= github.com/timakin/bodyclose v0.0.0-20241222091800-1db5c5ca4d67 h1:9LPGD+jzxMlnk5r6+hJnar67cgpDIz/iyD+rfl5r2Vk= @@ -2014,8 +2044,8 @@ mvdan.cc/gofumpt v0.8.0 h1:nZUCeC2ViFaerTcYKstMmfysj6uhQrA2vJe+2vwGU6k= mvdan.cc/gofumpt v0.8.0/go.mod h1:vEYnSzyGPmjvFkqJWtXkh79UwPWP9/HMxQdGEXZHjpg= mvdan.cc/unparam v0.0.0-20250301125049-0df0534333a4 h1:WjUu4yQoT5BHT1w8Zu56SP8367OuBV5jvo+4Ulppyf8= mvdan.cc/unparam v0.0.0-20250301125049-0df0534333a4/go.mod h1:rthT7OuvRbaGcd5ginj6dA2oLE7YNlta9qhBNNdCaLE= -pgregory.net/rapid v1.1.0 h1:CMa0sjHSru3puNx+J0MIAuiiEV4N0qj8/cMWGBBCsjw= -pgregory.net/rapid v1.1.0/go.mod h1:PY5XlDGj0+V1FCq0o192FdRhpKHGTRIWBgqjDBTrq04= +pgregory.net/rapid v1.2.0 h1:keKAYRcjm+e1F0oAuU5F5+YPAWcyxNNRK2wud503Gnk= +pgregory.net/rapid v1.2.0/go.mod h1:PY5XlDGj0+V1FCq0o192FdRhpKHGTRIWBgqjDBTrq04= rsc.io/binaryregexp v0.2.0/go.mod h1:qTv7/COck+e2FymRvadv62gMdZztPaShugOCi3I+8D8= rsc.io/pdf v0.1.1 h1:k1MczvYDUvJBe93bYd7wrZLLUEcLZAuF824/I4e5Xr4= rsc.io/pdf v0.1.1/go.mod h1:n8OzWcQ6Sp37PL01nO98y4iUCRdTGarVfzxY20ICaU4= diff --git a/pkg/config/config.go b/pkg/config/config.go index 65e30b4a..4e4c2e05 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -149,6 +149,12 @@ type CLI struct { PlayerTelemetry bool PlaybackWorkerURL string Ingests *placestream.IngestGetIngestUrls_Output + S3Endpoint string + S3Bucket string + S3AccessKeyID string + S3SecretAccessKey string + S3Region string + DisableSyndication bool } // ContentFilters represents the content filtering configuration @@ -808,6 +814,13 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { }, Sources: urfavecli.EnvVars("SP_SYNDICATE"), }, + &urfavecli.BoolFlag{ + Name: "disable-syndication", + Usage: `entirely disable syndication in both directions. useful for local development.`, + Value: false, + Destination: &cli.DisableSyndication, + Sources: urfavecli.EnvVars("SP_DISABLE_SYNDICATION"), + }, &urfavecli.BoolFlag{ Name: "player-telemetry", Usage: "enable player telemetry", @@ -833,6 +846,37 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { }, Sources: urfavecli.EnvVars("SP_INGESTS"), }, + &urfavecli.StringFlag{ + Name: "s3-endpoint", + Usage: "S3-compatible endpoint URL for segment archival uploads", + Destination: &cli.S3Endpoint, + Sources: urfavecli.EnvVars("SP_S3_ENDPOINT"), + }, + &urfavecli.StringFlag{ + Name: "s3-bucket", + Usage: "S3 bucket name for segment archival uploads", + Destination: &cli.S3Bucket, + Sources: urfavecli.EnvVars("SP_S3_BUCKET"), + }, + &urfavecli.StringFlag{ + Name: "s3-access-key-id", + Usage: "S3 access key ID for segment archival uploads", + Destination: &cli.S3AccessKeyID, + Sources: urfavecli.EnvVars("SP_S3_ACCESS_KEY_ID"), + }, + &urfavecli.StringFlag{ + Name: "s3-secret-access-key", + Usage: "S3 secret access key for segment archival uploads", + Destination: &cli.S3SecretAccessKey, + Sources: urfavecli.EnvVars("SP_S3_SECRET_ACCESS_KEY"), + }, + &urfavecli.StringFlag{ + Name: "s3-region", + Usage: "S3 region (default: us-east-1)", + Value: "us-east-1", + Destination: &cli.S3Region, + Sources: urfavecli.EnvVars("SP_S3_REGION"), + }, &urfavecli.BoolFlag{ Name: "external-signing", Usage: "DEPRECATED, does nothing.", @@ -1228,7 +1272,14 @@ func (cli *CLI) DumpDebugSegment(ctx context.Context, name string, r io.Reader) }() } +func (cli *CLI) S3Configured() bool { + return cli.S3Endpoint != "" && cli.S3Bucket != "" && cli.S3AccessKeyID != "" && cli.S3SecretAccessKey != "" +} + func (cli *CLI) ShouldSyndicate(did string) bool { + if cli.DisableSyndication { + return false + } for _, d := range cli.Syndicate { if d == "*" { return true diff --git a/pkg/director/s3_upload.go b/pkg/director/s3_upload.go new file mode 100644 index 00000000..d1ca7c23 --- /dev/null +++ b/pkg/director/s3_upload.go @@ -0,0 +1,45 @@ +package director + +import ( + "context" + "time" + + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/media" + "stream.place/streamplace/pkg/s3" +) + +func (ss *StreamSession) maybeStartS3Upload(ctx context.Context, repoDID string) { + if !ss.cli.S3Configured() { + return + } + cfg := s3.Config{ + Endpoint: ss.cli.S3Endpoint, + Bucket: ss.cli.S3Bucket, + AccessKeyID: ss.cli.S3AccessKeyID, + SecretAccessKey: ss.cli.S3SecretAccessKey, + Region: ss.cli.S3Region, + } + keyPrefix := repoDID + "/" + ss.s3Uploader = s3.NewS3Uploader(ctx, cfg, keyPrefix, time.Minute) + log.Log(ctx, "S3 upload enabled", "bucket", ss.cli.S3Bucket, "endpoint", ss.cli.S3Endpoint) +} + +func (ss *StreamSession) s3Upload(ctx context.Context, notif *media.NewSegmentNotification) { + if ss.s3Uploader == nil { + return + } + ss.Go(ctx, func() error { + return ss.s3Uploader.AddSegment(ctx, notif.Data) + }) +} + +func (ss *StreamSession) s3Close(ctx context.Context) { + if ss.s3Uploader == nil { + return + } + err := ss.s3Uploader.Close(ctx) + if err != nil { + log.Error(ctx, "error closing S3 upload", "error", err) + } +} diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index 31c15761..d8d49f29 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -29,6 +29,7 @@ import ( "stream.place/streamplace/pkg/model" "stream.place/streamplace/pkg/renditions" "stream.place/streamplace/pkg/replication" + "stream.place/streamplace/pkg/s3" "stream.place/streamplace/pkg/spmetrics" "stream.place/streamplace/pkg/statedb" "stream.place/streamplace/pkg/streamplace" @@ -66,6 +67,7 @@ type StreamSession struct { lastLivestreamTime time.Time lastViewCountTime time.Time + s3Uploader *s3.S3Uploader } func (ss *StreamSession) Start(ctx context.Context, notif *media.NewSegmentNotification) error { @@ -107,6 +109,8 @@ func (ss *StreamSession) Start(ctx context.Context, notif *media.NewSegmentNotif allRenditions = append(allRenditions, renditions.AudioRendition) ss.hls = media.NewM3U8(allRenditions) + ss.maybeStartS3Upload(ctx, notif.Segment.RepoDID) + close(ss.started) // Start background workers for status, origin, and livestream updates @@ -135,6 +139,7 @@ func (ss *StreamSession) Start(ctx context.Context, notif *media.NewSegmentNotif // reset timer case <-ctx.Done(): // Signal all background workers to stop + ss.s3Close(ctx) return ss.g.Wait() // case <-time.After(time.Minute * 1): case <-time.After(ss.cli.StreamSessionTimeout): @@ -189,6 +194,8 @@ func (ss *StreamSession) NewSegment(ctx context.Context, notif *media.NewSegment return fmt.Errorf("could not convert segment to streamplace segment: %w", err) } + ss.s3Upload(ctx, notif) + ss.bus.Publish(spseg.Creator, spseg) ss.Go(ctx, func() error { return ss.AddPlaybackSegment(ctx, spseg, "source", &bus.Seg{ diff --git a/pkg/media/concat_demux.go b/pkg/media/concat_demux.go index 2d5e6bc1..0317f2b7 100644 --- a/pkg/media/concat_demux.go +++ b/pkg/media/concat_demux.go @@ -173,7 +173,6 @@ func ConcatDemuxBin(ctx context.Context, seg *bus.Seg, doH264Parse bool) (*gst.B var padAdded func(self *gst.Element, pad *gst.Pad) // the defer funcs are needed to avoid leaking pads for some reason padAdded = func(self *gst.Element, pad *gst.Pad) { - log.Debug(ctx, "demux pad-added", "name", pad.GetName(), "direction", pad.GetDirection()) var downstreamPad *gst.Pad if strings.HasPrefix(pad.GetName(), "video_") { downstreamPad = mqVideoSink @@ -224,7 +223,7 @@ func ConcatDemuxBin(ctx context.Context, seg *bus.Seg, doH264Parse bool) (*gst.B src := app.SrcFromElement(appSrc) src.SetCallbacks(&app.SourceCallbacks{ - NeedDataFunc: ReaderNeedData(ctx, bytes.NewReader(seg.Data)), + NeedDataFunc: ReaderNeedDataIncremental(ctx, bytes.NewReader(seg.Data)), }) return bin, nil diff --git a/pkg/media/io_helpers.go b/pkg/media/io_helpers.go index 9cad9f6e..886f58f9 100644 --- a/pkg/media/io_helpers.go +++ b/pkg/media/io_helpers.go @@ -16,11 +16,21 @@ func ReaderNeedData(ctx context.Context, input io.Reader) func(self *app.Source, if err != nil { log.Error(ctx, "error reading from input", "error", err) } + wrote := false return func(self *app.Source, length uint) { if ctx.Err() != nil { self.EndStream() return } + if wrote { + log.Error(ctx, "ReaderNeedData: called after already wrote") + ret := self.EndStream() + if ret != gst.FlowOK { + log.Error(ctx, "failed to end stream after already wrote", "error", ret.String()) + } + return + } + wrote = true buffer := gst.NewBufferWithSize(int64(len(bsCopy))) buffer.Map(gst.MapWrite).WriteData(bsCopy) defer buffer.Unmap() @@ -30,6 +40,10 @@ func ReaderNeedData(ctx context.Context, input io.Reader) func(self *app.Source, } else { log.Debug(ctx, "pushed buffer", "length", len(bsCopy)) } + ret = self.EndStream() + if ret != gst.FlowOK { + log.Error(ctx, "failed to end stream", "error", ret.String()) + } } } diff --git a/pkg/media/muxl_segment.go b/pkg/media/muxl_segment.go new file mode 100644 index 00000000..e4957316 --- /dev/null +++ b/pkg/media/muxl_segment.go @@ -0,0 +1,120 @@ +package media + +import ( + "bytes" + "time" + + "github.com/go-gst/go-gst/gst" + "github.com/go-gst/go-gst/gst/app" + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/muxl" + + "context" + _ "embed" + "fmt" + "io" +) + +func MuxlSegmentElem(ctx context.Context, cli *config.CLI, streamer string, doH264Parse bool, cb func(ctx context.Context, buf []byte, now int64) error) (*gst.Element, error) { + bin := gst.NewBin("muxl-segment-bin") + elem, err := gst.NewElementWithProperties("mp4mux", map[string]any{ + "name": "fmp4mux", + "fragment-mode": 0, + "fragment-duration": 1, + }) + if err != nil { + return nil, err + } + + err = bin.Add(elem) + if err != nil { + return nil, fmt.Errorf("failed to add mp4mux to bin: %w", err) + } + + videoPad := elem.GetRequestPad("video_%u") + if videoPad == nil { + return nil, fmt.Errorf("failed to get video pad") + } + videoGhost := gst.NewGhostPad("video_0", videoPad) + if videoGhost == nil { + return nil, fmt.Errorf("failed to create video ghost pad") + } + audioPad := elem.GetRequestPad("audio_%u") + if audioPad == nil { + return nil, fmt.Errorf("failed to get audio pad") + } + audioGhost := gst.NewGhostPad("audio_0", audioPad) + if audioGhost == nil { + return nil, fmt.Errorf("failed to create audio ghost pad") + } + + ok := bin.AddPad(videoGhost.Pad) + if !ok { + return nil, fmt.Errorf("failed to add video ghost pad to bin") + } + + ok = bin.AddPad(audioGhost.Pad) + if !ok { + return nil, fmt.Errorf("failed to add audio ghost pad to bin") + } + + appsink, err := gst.NewElementWithProperties("appsink", map[string]any{ + "name": "muxl-appsink", + }) + if err != nil { + return nil, fmt.Errorf("failed to create appsink element: %w", err) + } + err = bin.Add(appsink) + if err != nil { + return nil, fmt.Errorf("failed to add appsink to bin: %w", err) + } + + err = elem.Link(appsink) + if err != nil { + return nil, fmt.Errorf("failed to link mp4mux to appsink: %w", err) + } + + initCh := make(chan []byte) + segCh := make(chan []byte) + r, w := io.Pipe() + go func() { + err := muxl.RunMuxlSegmenter(ctx, r, initCh, segCh) + if err != nil { + log.Error(ctx, "error running muxl segmenter", "error", err) + } + }() + + go func() { + var initSeg []byte + select { + case <-ctx.Done(): + return + case initSeg = <-initCh: + log.Debug(ctx, "got init segment", "size", len(initSeg)) + } + for { + select { + case <-ctx.Done(): + return + case seg := <-segCh: + log.Debug(ctx, "got segment", "size", len(seg)) + fullSeg := []byte{} + fullSeg = append(fullSeg, initSeg...) + fullSeg = append(fullSeg, seg...) + cli.DumpDebugSegment(ctx, "muxl_segment_input.fmp4", bytes.NewReader(fullSeg)) + err := cb(ctx, fullSeg, time.Now().UnixMilli()) + if err != nil { + log.Error(ctx, "error calling callback", "error", err) + } + } + } + }() + + sink := app.SinkFromElement(appsink) + sink.SetCallbacks(&app.SinkCallbacks{ + NewSampleFunc: WriterNewSample(ctx, w), + }) + + return bin.Element, nil +} diff --git a/pkg/media/packetize.go b/pkg/media/packetize.go index ef755606..6ebaf940 100644 --- a/pkg/media/packetize.go +++ b/pkg/media/packetize.go @@ -137,6 +137,7 @@ func Packetize(ctx context.Context, seg *bus.Seg) (*bus.PacketizedSegment, error } samples := buffer.Bytes() + // log.Warn(ctx, "audioappsink NewSampleFunc", "sample", len(samples)) audioOutput = append(audioOutput, samples) diff --git a/pkg/media/packetize_test.go b/pkg/media/packetize_test.go index 414c7d8e..8b0a06e3 100644 --- a/pkg/media/packetize_test.go +++ b/pkg/media/packetize_test.go @@ -11,6 +11,7 @@ import ( "github.com/stretchr/testify/require" "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/bus" + "stream.place/streamplace/test/remote" ) func TestPacketize(t *testing.T) { @@ -18,7 +19,7 @@ func TestPacketize(t *testing.T) { g, _ := errgroup.WithContext(context.Background()) for range streamplaceTestCount { g.Go(func() error { - innerTestPacketize(t) + innerTestPacketize(t, getFixture("sample-segment.mp4"), 49, 40, time.Duration(800*time.Millisecond)) return nil }) } @@ -27,8 +28,14 @@ func TestPacketize(t *testing.T) { }) } -func innerTestPacketize(t *testing.T) { - filename := getFixture("sample-segment.mp4") +func TestPacketizeMuxl(t *testing.T) { + withNoGSTLeaks(t, func() { + filename := remote.RemoteFixture("c6b57a53fc5a2234dbdd388922f0e293d8063d2b30620321e974b7c85640f228/2026-03-17T19-02-08-607Z-muxl_segment_input.fmp4") + innerTestPacketize(t, filename, 60, 50, time.Duration(1000*time.Millisecond)) + }) +} + +func innerTestPacketize(t *testing.T, filename string, expectedVideo int, expectedAudio int, expectedDuration time.Duration) { inputFile, err := os.Open(filename) require.NoError(t, err) defer inputFile.Close() @@ -44,9 +51,9 @@ func innerTestPacketize(t *testing.T) { packet, err := Packetize(context.Background(), testSeg) require.NoError(t, err) require.NotNil(t, packet) - require.Equal(t, 49, len(packet.Video)) - require.Equal(t, 40, len(packet.Audio)) - require.Equal(t, time.Duration(800*time.Millisecond), packet.Duration) + require.Equal(t, expectedVideo, len(packet.Video)) + require.Equal(t, expectedAudio, len(packet.Audio)) + require.Equal(t, expectedDuration, packet.Duration) } func TestPacketizeInvalid(t *testing.T) { diff --git a/pkg/media/segmenter.go b/pkg/media/segmenter.go index 8cda85f3..b58edd42 100644 --- a/pkg/media/segmenter.go +++ b/pkg/media/segmenter.go @@ -179,7 +179,7 @@ func SegmentElem(ctx context.Context, cli *config.CLI, streamer string, doH264Pa } func (mm *MediaManager) SegmentAndSignElem(ctx context.Context, ms MediaSigner) (*gst.Element, error) { - return SegmentElem(ctx, mm.cli, ms.Streamer(), false, func(ctx context.Context, bs []byte, now int64) error { + return MuxlSegmentElem(ctx, mm.cli, ms.Streamer(), false, func(ctx context.Context, bs []byte, now int64) error { if mm.cli.SmearAudio { smearedBuf := &bytes.Buffer{} err := RewriteAudioTimestamps(ctx, mm.cli, bytes.NewReader(bs), smearedBuf, true) diff --git a/pkg/media/thumbnail.go b/pkg/media/thumbnail.go index 1e3ef9ca..0503950b 100644 --- a/pkg/media/thumbnail.go +++ b/pkg/media/thumbnail.go @@ -8,6 +8,7 @@ import ( "github.com/go-gst/go-gst/gst" "github.com/go-gst/go-gst/gst/app" + "stream.place/streamplace/pkg/bus" "stream.place/streamplace/pkg/log" ) @@ -27,10 +28,21 @@ func Thumbnail(ctx context.Context, r io.Reader, w io.Writer, format string) err encoder = "pngenc snapshot=true" } + // Read all data from the reader to create a Seg for ConcatDemuxBin + data, err := io.ReadAll(r) + if err != nil { + return fmt.Errorf("error reading input data: %w", err) + } + + seg := &bus.Seg{ + Data: data, + } + pipelineSlice := []string{ - "appsrc name=appsrc ! qtdemux name=demux ! decodebin ! videoconvert ! videoscale ! videorate ! capsfilter name=capsfilter caps=video/x-raw,width=[1,1280],height=[1,720],pixel-aspect-ratio=1/1,framerate=1/999999 ! ", + "decodebin name=decode ! videoconvert ! videoscale ! videorate ! capsfilter name=capsfilter caps=video/x-raw,width=[1,1280],height=[1,720],pixel-aspect-ratio=1/1,framerate=1/999999 ! ", encoder, " ! appsink name=appsink", + "fakesink name=audiofakesink sync=false", } pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) @@ -38,19 +50,56 @@ func Thumbnail(ctx context.Context, r io.Reader, w io.Writer, format string) err return fmt.Errorf("error creating Thumbnail pipeline: %w", err) } - defer func() { - cancel() - err = pipeline.BlockSetState(gst.StateNull) - }() - appsrc, err := pipeline.GetElementByName("appsrc") + demuxBin, err := ConcatDemuxBin(ctx, seg, false) if err != nil { - return err + return fmt.Errorf("failed to create demux bin: %w", err) } - src := app.SrcFromElement(appsrc) - src.SetCallbacks(&app.SourceCallbacks{ - NeedDataFunc: ReaderNeedData(ctx, r), - }) + err = pipeline.Add(demuxBin.Element) + if err != nil { + return fmt.Errorf("failed to add demux bin to pipeline: %w", err) + } + + demuxBinPadVideoSrc := demuxBin.GetStaticPad("video_0") + if demuxBinPadVideoSrc == nil { + return fmt.Errorf("failed to get demux bin video src pad") + } + + demuxBinPadAudioSrc := demuxBin.GetStaticPad("audio_0") + if demuxBinPadAudioSrc == nil { + return fmt.Errorf("failed to get demux bin audio src pad") + } + + decode, err := pipeline.GetElementByName("decode") + if err != nil { + return fmt.Errorf("failed to get decodebin element: %w", err) + } + + audioFakeSink, err := pipeline.GetElementByName("audiofakesink") + if err != nil { + return fmt.Errorf("failed to get audio fakesink element: %w", err) + } + + linked := demuxBinPadVideoSrc.Link(decode.GetStaticPad("sink")) + if linked != gst.PadLinkOK { + return fmt.Errorf("failed to link demux bin video src to decodebin: %v", linked) + } + + linked = demuxBinPadAudioSrc.Link(audioFakeSink.GetStaticPad("sink")) + if linked != gst.PadLinkOK { + return fmt.Errorf("failed to link demux bin audio src to fakesink: %v", linked) + } + + defer func() { + err := pipeline.SetState(gst.StateNull) + if err != nil { + log.Error(ctx, "failed to set pipeline state to null", "error", err) + } + err = pipeline.Remove(demuxBin.Element) + if err != nil { + log.Error(ctx, "failed to remove demux bin from pipeline", "error", err) + } + }() appsink, err := pipeline.GetElementByName("appsink") if err != nil { @@ -60,25 +109,50 @@ func Thumbnail(ctx context.Context, r io.Reader, w io.Writer, format string) err errCh := make(chan error) go func() { err := HandleBusMessages(ctx, pipeline) - cancel() errCh <- err close(errCh) }() + thumbCh := make(chan struct{}) + sink := app.SinkFromElement(appsink) sink.SetCallbacks(&app.SinkCallbacks{ - NewSampleFunc: WriterNewSample(ctx, w), + NewSampleFunc: func(sink *app.Sink) gst.FlowReturn { + sample := sink.PullSample() + if sample == nil { + return gst.FlowOK + } + + // Retrieve the buffer from the sample. + buffer := sample.GetBuffer() + bs := buffer.Map(gst.MapRead).Bytes() + defer buffer.Unmap() + + _, err := w.Write(bs) + log.Debug(ctx, "wrote buffer", "length", len(bs)) + + if err != nil { + log.Error(ctx, "error writing to output", "error", err) + return gst.FlowError + } + + close(thumbCh) + + return gst.FlowOK + }, }) - if err := pipeline.BlockSetState(gst.StatePlaying); err != nil { + if err := pipeline.SetState(gst.StatePlaying); err != nil { return fmt.Errorf("error setting pipeline state: %w", err) } - <-ctx.Done() + <-thumbCh + // signals the pipeline to clean up cleanly + pipeline.Error(ErrConcatDone.Error(), ErrConcatDone) - if err := pipeline.BlockSetState(gst.StateNull); err != nil { - return fmt.Errorf("error setting pipeline state: %w", err) - } + busErr := <-errCh + + log.Debug(ctx, "thumbnail done") - return <-errCh + return busErr } diff --git a/pkg/media/thumbnail_test.go b/pkg/media/thumbnail_test.go index cfe1da9e..6e37622e 100644 --- a/pkg/media/thumbnail_test.go +++ b/pkg/media/thumbnail_test.go @@ -7,60 +7,92 @@ import ( "io" "os" "testing" + "time" "github.com/stretchr/testify/require" "golang.org/x/sync/errgroup" + "stream.place/streamplace/pkg/log" "stream.place/streamplace/test/remote" ) +var thumbnailTestCases = []struct { + name string + fixtureFn func() string +}{ + { + name: "SampleSegment", + fixtureFn: func() string { + return getFixture("sample-segment.mp4") + }, + }, + { + name: "MuxlSegment", + fixtureFn: func() string { + return remote.RemoteFixture("c6b57a53fc5a2234dbdd388922f0e293d8063d2b30620321e974b7c85640f228/2026-03-17T19-02-08-607Z-muxl_segment_input.fmp4") + }, + }, +} + func TestThumbnail(t *testing.T) { - withNoGSTLeaks(t, func() { - // Open input file - inputFile, err := os.Open(getFixture("sample-segment.mp4")) - require.NoError(t, err) - defer inputFile.Close() - bs, err := io.ReadAll(inputFile) - require.NoError(t, err) + for _, tc := range thumbnailTestCases { + t.Run(tc.name, func(t *testing.T) { + withNoGSTLeaks(t, func() { + inputFile, err := os.Open(tc.fixtureFn()) + require.NoError(t, err) + defer inputFile.Close() + bs, err := io.ReadAll(inputFile) + require.NoError(t, err) - ctx := context.Background() - g, ctx := errgroup.WithContext(ctx) + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + ctx = log.WithDebugValue(ctx, map[string]map[string]int{"function": {"Thumbnail": 9}}) + g, ctx := errgroup.WithContext(ctx) - for i := 0; i < streamplaceTestCount; i++ { - g.Go(func() error { - thumbnail := bytes.Buffer{} - // thumbnailCtx = log.WithDebugValue(ctx, map[string]map[string]int{"function": {"Thumbnail": 9}}) - err := Thumbnail(ctx, bytes.NewReader(bs), &thumbnail, "png") - if err != nil { - return err - } - if thumbnail.Len() == 0 { - return fmt.Errorf("thumbnail buffer is empty") + for i := 0; i < streamplaceTestCount; i++ { + // g.Go(func() error { + // thumbnail := bytes.Buffer{} + // err := Thumbnail(ctx, bytes.NewReader(bs), &thumbnail, "png") + // if err != nil { + // return err + // } + // if thumbnail.Len() == 0 { + // return fmt.Errorf("thumbnail buffer is empty") + // } + // // No strict length checks for muxl variant, but keep sample-segment's as before. + // if tc.name == "sample-segment" { + // require.Equal(t, 1418910, thumbnail.Len()) + // } else { + // require.Greater(t, thumbnail.Len(), 50000) + // } + // return nil + // }) + g.Go(func() error { + thumbnail := bytes.Buffer{} + err := Thumbnail(ctx, bytes.NewReader(bs), &thumbnail, "jpeg") + if err != nil { + return err + } + if thumbnail.Len() == 0 { + return fmt.Errorf("thumbnail buffer is empty") + } + // For jpeg, apply broad range checking for muxl, strict for sample-segment + if tc.name == "sample-segment" { + require.Greater(t, thumbnail.Len(), 140000) + require.Less(t, thumbnail.Len(), 150000) + require.Equal(t, 140969, thumbnail.Len()) + } else { + require.Greater(t, thumbnail.Len(), 10000) + require.Less(t, thumbnail.Len(), 150000) + } + return nil + }) } - require.Equal(t, 1418910, thumbnail.Len()) - return nil - }) - g.Go(func() error { - thumbnail := bytes.Buffer{} - // thumbnailCtx = log.WithDebugValue(ctx, map[string]map[string]int{"function": {"Thumbnail": 9}}) - err := Thumbnail(ctx, bytes.NewReader(bs), &thumbnail, "jpeg") - if err != nil { - return err - } - if thumbnail.Len() == 0 { - return fmt.Errorf("thumbnail buffer is empty") - } - // jpeg thumbnails aren't deterministic, so let's give a range instead - // testing gave 140969 bytes, but it can vary a bit - require.Greater(t, thumbnail.Len(), 140000) - require.Less(t, thumbnail.Len(), 150000) - require.Equal(t, 140969, thumbnail.Len()) - return nil - }) - } - err = g.Wait() - require.NoError(t, err) - }) + err = g.Wait() + require.NoError(t, err) + }) + }) + } } // This segment once caused a segfault in gst-libav. diff --git a/pkg/model/segment.go b/pkg/model/segment.go deleted file mode 100644 index 8b537907..00000000 --- a/pkg/model/segment.go +++ /dev/null @@ -1 +0,0 @@ -package model diff --git a/pkg/model/segment_test.go b/pkg/model/segment_test.go deleted file mode 100644 index 8b537907..00000000 --- a/pkg/model/segment_test.go +++ /dev/null @@ -1 +0,0 @@ -package model diff --git a/pkg/muxl/muxl.go b/pkg/muxl/muxl.go new file mode 100644 index 00000000..36471e55 --- /dev/null +++ b/pkg/muxl/muxl.go @@ -0,0 +1,216 @@ +package muxl + +import ( + "context" + "errors" + "fmt" + "io" + "sync/atomic" + + _ "embed" + + "github.com/hyphacoop/go-dasl/drisl" + "github.com/tetratelabs/wazero" + "github.com/tetratelabs/wazero/imports/wasi_snapshot_preview1" + "stream.place/streamplace/pkg/log" +) + +var moduleCounter atomic.Uint64 + +// MuxlEvent represents an event from the muxl segmenter. +type MuxlEvent struct { + Type string // "INIT" or "SEGM" + Number uint32 // segment number (only for SEGM) + Tracks map[string][]byte + Data []byte +} + +//go:embed muxl.wasm +var wasmBytes []byte + +var wasmRuntime wazero.Runtime + +func init() { + wasmRuntime = wazero.NewRuntime(context.Background()) + wasi_snapshot_preview1.MustInstantiate(context.Background(), wasmRuntime) +} + +func getCompiledModule(ctx context.Context) (wazero.CompiledModule, error) { + compiledModule, err := wasmRuntime.CompileModule(ctx, wasmBytes) + if err != nil { + return nil, fmt.Errorf("error compiling module: %w", err) + } + return compiledModule, nil +} + +// Segment arbitrary fMP4 input into MUXL-compatible init and segment chunks. +func RunMuxlSegmenter(ctx context.Context, input io.Reader, initCh chan []byte, segCh chan []byte) error { + return runMuxl(ctx, []string{"muxl", "segment", "-", "--stdout"}, input, initCh, segCh) +} + +// Given a bunch of MUXL-compatible fMP4 archives containing init and segment chunks, concatenate them into a single fMP4 archive. +// If the init segment changes, you'll get a new init segment in the output. +func RunMuxlConcatenator(ctx context.Context, input io.Reader, initCh chan []byte, segCh chan []byte) error { + return runMuxl(ctx, []string{"muxl", "concat"}, input, initCh, segCh) +} + +// Concatenator accepts full fMP4 archives (init+segments) and produces +// deduplicated output: init segments are emitted only when they change, +// and segment data is emitted without the init header, suitable for +// concatenation into a single fMP4 stream. +// +// Usage: +// +// cat := muxl.NewConcatenator(ctx) +// go func() { cat.Write(fullFmp4Archive); cat.Close() }() +// initSeg := <-cat.InitCh +// for seg := range cat.SegCh { /* append to output */ } +type Concatenator struct { + stdinWriter *io.PipeWriter + InitCh chan []byte + SegCh chan []byte + done chan error +} + +// NewConcatenator starts the WASM concat process in the background. +// Write full fMP4 archives via Write(), receive processed output on InitCh and SegCh. +// InitCh receives a new init segment only when the track configuration changes. +// SegCh receives raw segment data (moof+mdat) that can be concatenated after an init. +// Both channels are closed when the concatenator finishes (after Close + WASM exit). +func NewConcatenator(ctx context.Context) *Concatenator { + initCh := make(chan []byte, 1) + segCh := make(chan []byte, 16) + stdinReader, stdinWriter := io.Pipe() + done := make(chan error, 1) + + c := &Concatenator{ + stdinWriter: stdinWriter, + InitCh: initCh, + SegCh: segCh, + done: done, + } + + go func() { + err := RunMuxlConcatenator(ctx, stdinReader, initCh, segCh) + close(initCh) + close(segCh) + done <- err + }() + + return c +} + +// Write feeds a full fMP4 archive (init+segments) to the concatenator. +func (c *Concatenator) Write(data []byte) error { + _, err := c.stdinWriter.Write(data) + return err +} + +// Close signals that no more data will be written. The WASM process will +// finish processing and the output channels will be closed. +func (c *Concatenator) Close() error { + c.stdinWriter.Close() + return <-c.done +} + +// logWriter adapts WASM stderr output to log.Debug calls, emitting one +// log message per line. +type logWriter struct { + ctx context.Context + instanceID uint64 + buf []byte +} + +func (w *logWriter) Write(p []byte) (int, error) { + w.buf = append(w.buf, p...) + for { + i := 0 + for i < len(w.buf) && w.buf[i] != '\n' { + i++ + } + if i >= len(w.buf) { + break + } + line := string(w.buf[:i]) + w.buf = w.buf[i+1:] + log.Debug(w.ctx, "muxl wasm", "instance", w.instanceID, "msg", line) + } + return len(p), nil +} + +func runMuxl(ctx context.Context, args []string, input io.Reader, initCh chan []byte, segCh chan []byte) error { + compiledModule, err := getCompiledModule(ctx) + if err != nil { + return fmt.Errorf("error getting compiled module: %w", err) + } + + // Set up stdin/stdout pipes + stdinReader, stdinWriter := io.Pipe() + stdoutReader, stdoutWriter := io.Pipe() + + instanceID := moduleCounter.Add(1) + config := wazero.NewModuleConfig(). + WithName(fmt.Sprintf("muxl-%d", instanceID)). + WithStdin(stdinReader). + WithStdout(stdoutWriter). + WithStderr(&logWriter{ctx: ctx, instanceID: instanceID}). + WithArgs(args...) + + // Run the module in a goroutine + errCh := make(chan error, 1) + go func() { + _, err := wasmRuntime.InstantiateModule(ctx, compiledModule, config) + if err != nil { + log.Error(ctx, "error instantiating module", "error", err) + } + stdoutWriter.Close() + errCh <- err + }() + + // Feed input to stdin in a goroutine + go func() { + _, err := io.Copy(stdinWriter, input) + if err != nil { + log.Error(ctx, "error copying input to stdin", "error", err) + } + stdinWriter.Close() + }() + + // Parse framed events from stdout + err = ParseMuxlEvents(stdoutReader, initCh, segCh) + if err != nil { + return fmt.Errorf("parsing events: %w", err) + } + + // Wait for WASM module to finish + if wasmErr := <-errCh; wasmErr != nil { + return fmt.Errorf("wasm execution: %w", wasmErr) + } + + return nil +} + +func ParseMuxlEvents(r io.Reader, initCh chan []byte, segCh chan []byte) error { + decoder := drisl.NewDecoder(r) + + for { + var ev MuxlEvent + err := decoder.Decode(&ev) + if errors.Is(err, io.EOF) { + break + } + if ev.Type == "init" { + initCh <- ev.Data + } else if ev.Type == "segment" { + combined := []byte{} + for _, data := range ev.Tracks { + combined = append(combined, data...) + } + segCh <- combined + } else { + return fmt.Errorf("unknown event type: %s", ev.Type) + } + } + + return nil +} diff --git a/pkg/s3/s3.go b/pkg/s3/s3.go new file mode 100644 index 00000000..7725bc16 --- /dev/null +++ b/pkg/s3/s3.go @@ -0,0 +1,255 @@ +package s3 + +import ( + "bytes" + "context" + "fmt" + "time" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/credentials" + "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/aws/aws-sdk-go-v2/service/s3/types" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/muxl" +) + +// Config holds the configuration for an S3-compatible upload target. +type Config struct { + Endpoint string + Bucket string + AccessKeyID string + SecretAccessKey string + Region string +} + +// S3Uploader manages streaming multipart uploads to an S3-compatible endpoint. +// Full fMP4 archives are fed via AddSegment. They are run through a muxl +// Concatenator to strip duplicate init segments, then uploaded as a +// multipart upload. Every cutoverEvery, the current upload is completed +// and a new one begins. +type S3Uploader struct { + client *s3.Client + bucket string + cutoverEvery time.Duration + keyPrefix string // e.g. "did:plc:abc123/" + concat *muxl.Concatenator + done chan error +} + +// S3 requires each part except the last to be at least 5MB. +const minPartSize = 5 * 1024 * 1024 + +type activeUpload struct { + key string + uploadID string + parts []types.CompletedPart + partNum int32 + started time.Time + buf []byte // accumulates segments until we hit minPartSize +} + +// NewS3Uploader creates a new S3Uploader. keyPrefix is prepended to every +// object key (typically the streamer DID + "/"). Starts the muxl Concatenator +// and a background goroutine that reads processed segments and uploads them. +func NewS3Uploader(ctx context.Context, cfg Config, keyPrefix string, cutoverEvery time.Duration) *S3Uploader { + client := s3.New(s3.Options{ + Region: cfg.Region, + Credentials: credentials.NewStaticCredentialsProvider( + cfg.AccessKeyID, + cfg.SecretAccessKey, + "", + ), + BaseEndpoint: aws.String(cfg.Endpoint), + UsePathStyle: true, + }) + if cutoverEvery == 0 { + cutoverEvery = time.Minute + } + concat := muxl.NewConcatenator(ctx) + u := &S3Uploader{ + client: client, + bucket: cfg.Bucket, + cutoverEvery: cutoverEvery, + keyPrefix: keyPrefix, + concat: concat, + done: make(chan error, 1), + } + go u.uploadLoop(ctx) + return u +} + +// AddSegment feeds a full fMP4 archive (init+segments) to the concatenator +// for processing and upload. +func (u *S3Uploader) AddSegment(ctx context.Context, data []byte) error { + return u.concat.Write(data) +} + +// Close signals that no more segments will be added, waits for all +// in-flight uploads to complete, and returns any error. +func (u *S3Uploader) Close(ctx context.Context) error { + closeErr := u.concat.Close() + uploadErr := <-u.done + if uploadErr != nil { + return uploadErr + } + return closeErr +} + +// uploadLoop reads init and segment events from the concatenator and manages +// multipart uploads. Runs until the concatenator's channels are closed. +func (u *S3Uploader) uploadLoop(ctx context.Context) { + ctx = log.WithLogValues(ctx, "func", "s3.uploadLoop") + var initSeg []byte + var current *activeUpload + + // Helper: prepend init to buffer when starting a new upload + startUpload := func() error { + now := time.Now() + key := fmt.Sprintf("%s%s.mp4", u.keyPrefix, now.UTC().Format("2006-01-02T15-04-05")) + + resp, err := u.client.CreateMultipartUpload(ctx, &s3.CreateMultipartUploadInput{ + Bucket: aws.String(u.bucket), + Key: aws.String(key), + ContentType: aws.String("video/mp4"), + }) + if err != nil { + return fmt.Errorf("creating multipart upload for %s: %w", key, err) + } + + current = &activeUpload{ + key: key, + uploadID: *resp.UploadId, + started: now, + } + // Prepend init segment to the buffer so the file starts valid + if initSeg != nil { + current.buf = append(current.buf, initSeg...) + } + log.Log(ctx, "started S3 multipart upload", "key", key) + return nil + } + + flushBuffer := func() error { + if current == nil || len(current.buf) == 0 { + return nil + } + current.partNum++ + partNum := current.partNum + + resp, err := u.client.UploadPart(ctx, &s3.UploadPartInput{ + Bucket: aws.String(u.bucket), + Key: aws.String(current.key), + UploadId: aws.String(current.uploadID), + PartNumber: aws.Int32(partNum), + Body: bytes.NewReader(current.buf), + }) + if err != nil { + return fmt.Errorf("uploading part %d: %w", partNum, err) + } + log.Debug(ctx, "uploaded S3 part", "key", current.key, "part", partNum, "size", len(current.buf)) + current.parts = append(current.parts, types.CompletedPart{ + ETag: resp.ETag, + PartNumber: aws.Int32(partNum), + }) + current.buf = current.buf[:0] + return nil + } + + completeUpload := func() error { + if current == nil { + return nil + } + if err := flushBuffer(); err != nil { + return err + } + if len(current.parts) == 0 { + _, err := u.client.AbortMultipartUpload(ctx, &s3.AbortMultipartUploadInput{ + Bucket: aws.String(u.bucket), + Key: aws.String(current.key), + UploadId: aws.String(current.uploadID), + }) + if err != nil { + log.Error(ctx, "aborting empty multipart upload", "key", current.key, "error", err) + } + current = nil + return nil + } + _, err := u.client.CompleteMultipartUpload(ctx, &s3.CompleteMultipartUploadInput{ + Bucket: aws.String(u.bucket), + Key: aws.String(current.key), + UploadId: aws.String(current.uploadID), + MultipartUpload: &types.CompletedMultipartUpload{ + Parts: current.parts, + }, + }) + if err != nil { + return fmt.Errorf("completing multipart upload %s: %w", current.key, err) + } + log.Log(ctx, "completed S3 multipart upload", "key", current.key, "parts", len(current.parts)) + current = nil + return nil + } + + handleSegment := func(seg []byte) error { + now := time.Now() + + // Cut over if needed + if current != nil && now.Sub(current.started) >= u.cutoverEvery { + if err := completeUpload(); err != nil { + return err + } + } + + // Start a new upload if needed + if current == nil { + if err := startUpload(); err != nil { + return err + } + } + + // Append segment data to buffer + current.buf = append(current.buf, seg...) + + // Flush if buffer is large enough for a part + if len(current.buf) >= minPartSize { + if err := flushBuffer(); err != nil { + return err + } + } + + return nil + } + + var err error + for err == nil { + select { + case init, ok := <-u.concat.InitCh: + if !ok { + u.concat.InitCh = nil + continue + } + initSeg = init + log.Log(ctx, "received init segment for S3 upload", "size", len(init)) + + case seg, ok := <-u.concat.SegCh: + log.Log(ctx, "received segment for S3 upload", "size", len(seg)) + if !ok { + // Concatenator is done, complete any in-progress upload + err = completeUpload() + u.done <- err + return + } + if err = handleSegment(seg); err != nil { + log.Error(ctx, "error handling segment", "error", err) + } + + case <-ctx.Done(): + _ = completeUpload() + u.done <- ctx.Err() + return + } + } + + u.done <- err +} diff --git a/pkg/spxrpc/place_stream_live.go b/pkg/spxrpc/place_stream_live.go index 0caa6cce..c00ee9e6 100644 --- a/pkg/spxrpc/place_stream_live.go +++ b/pkg/spxrpc/place_stream_live.go @@ -209,6 +209,9 @@ func (s *Server) handlePlaceStreamLiveGetLiveUsers(ctx context.Context, before s } func (s *Server) handlePlaceStreamLiveSubscribeSegments(c echo.Context) error { + if s.cli.DisableSyndication { + return echo.NewHTTPError(http.StatusNotImplemented, "Syndication is disabled") + } user := c.QueryParam("streamer") if user == "" { return echo.NewHTTPError(http.StatusBadRequest, "User DID is required") diff --git a/rust/muxl-wasm/Cargo.toml b/rust/muxl-wasm/Cargo.toml new file mode 100644 index 00000000..2a5fb7e4 --- /dev/null +++ b/rust/muxl-wasm/Cargo.toml @@ -0,0 +1,11 @@ +[package] +name = "muxl-wasm" +version = "0.1.0" +edition = "2024" +license = "Apache-2.0" + +[dependencies] +# For local development, use the path dependency: +# muxl = { path = "../../../s2pa-muxl" } +# To use the published git version instead, comment the line above and uncomment: +muxl = { git = "https://github.com/streamplace/s2pa-muxl.git", rev = "5456732ab7844c7ffcba148583a08f4bf0b50dcd" } diff --git a/rust/muxl-wasm/src/main.rs b/rust/muxl-wasm/src/main.rs new file mode 100644 index 00000000..2197c04d --- /dev/null +++ b/rust/muxl-wasm/src/main.rs @@ -0,0 +1,3 @@ +fn main() { + muxl::cli_main(); +} -- 2.51.2