diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index 1b05a9b8b..b2f20ce60 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -140,6 +140,25 @@ build-docker-amd64: --destination "$CI_REGISTRY_IMAGE:$STREAMPLACE_BRANCH-amd64" --destination "$CI_REGISTRY_IMAGE:$STREAMPLACE_VERSION-amd64" +build-docker-mistserver: + stage: build + interruptible: true + needs: + - job: build + artifacts: true + image: + name: gcr.io/kaniko-project/executor:v1.14.0-debug + entrypoint: [""] + timeout: 2 hours + script: + - /kaniko/executor + --build-arg TARGETARCH=amd64 + --build-arg STREAMPLACE_URL=$STREAMPLACE_URL_LINUX_AMD64 + --context "${CI_PROJECT_DIR}" + --dockerfile "${CI_PROJECT_DIR}/docker/mistserver.Dockerfile" + --destination "$CI_REGISTRY_IMAGE:$STREAMPLACE_BRANCH-mistserver" + --destination "$CI_REGISTRY_IMAGE:$STREAMPLACE_VERSION-mistserver" + build-docker-arm64: stage: build interruptible: true diff --git a/Makefile b/Makefile index f7d8d95f7..eabbe0be8 100644 --- a/Makefile +++ b/Makefile @@ -471,14 +471,25 @@ DOCKER_OPTS?= in-container: docker-build-builder $(DOCKER_BIN) run $(DOCKER_OPTS) -v $$(pwd):$$(pwd) -w $$(pwd) --rm $(DOCKER_REF) bash -c "$(IN_CONTAINER_CMD)" +STREAMPLACE_URL?=https://git.stream.place/streamplace/streamplace/-/package_files/10122/download .PHONY: docker-release docker-release: cd docker \ && docker build -f release.Dockerfile \ --build-arg TARGETARCH=$(BUILDARCH) \ + --build-arg STREAMPLACE_URL=$(STREAMPLACE_URL) \ -t dist.stream.place/streamplace/streamplace \ . +.PHONY: docker-mistserver +docker-mistserver: + cd docker \ + && docker build -f mistserver.Dockerfile \ + --build-arg TARGETARCH=$(BUILDARCH) \ + --build-arg STREAMPLACE_URL=$(STREAMPLACE_URL) \ + -t dist.stream.place/streamplace/streamplace:mistserver \ + . + .PHONY: ci-upload ci-upload: ci-upload-node ci-upload-android diff --git a/docker/mistserver.Dockerfile b/docker/mistserver.Dockerfile new file mode 100644 index 000000000..739b55053 --- /dev/null +++ b/docker/mistserver.Dockerfile @@ -0,0 +1,13 @@ +ARG TARGETARCH +FROM --platform=linux/$TARGETARCH ubuntu:24.04 +RUN apt update && apt install -y curl +ARG STREAMPLACE_URL +ENV STREAMPLACE_URL $STREAMPLACE_URL +# strip the -cloudflare suffix from the url; we're on the git server we don't need to leave +RUN export LOCAL_URL="$(echo $STREAMPLACE_URL | sed 's/-cloudflare//')" && echo "downloading $LOCAL_URL" && cd /usr/local/bin && curl -L "$LOCAL_URL" | tar xzv + +RUN apt-get update && apt-get install -y curl +RUN curl -o - https://releases.mistserver.org/is/mistserver_64V3.6.1.tar.gz 2>/dev/null | sh +RUN mkdir -p /config +ADD ./docker/mistserver.json /config/mistserver.json +CMD ["MistController", "-c", "/config/mistserver.json"] diff --git a/docker/mistserver.json b/docker/mistserver.json new file mode 100644 index 000000000..dfb0f8302 --- /dev/null +++ b/docker/mistserver.json @@ -0,0 +1,149 @@ +{ + "account": { + "streamplace": { + "password": "9bc4bd49515c5ade1fa94f8301c24473" + } + }, + "auto_push": null, + "bandwidth": { + "exceptions": [ + "::1", + "127.0.0.0/8", + "10.0.0.0/8", + "192.168.0.0/16", + "172.16.0.0/12" + ] + }, + "config": { + "accesslog": "LOG", + "controller": { + "interface": "127.0.0.1", + "port": null, + "username": null + }, + "defaultStream": null, + "limits": null, + "prometheus": "", + "protocols": [ + { + "connector": "AAC" + }, + { + "connector": "CMAF" + }, + { + "connector": "EBML" + }, + { + "connector": "FLAC" + }, + { + "connector": "FLV" + }, + { + "connector": "H264" + }, + { + "connector": "HDS" + }, + { + "connector": "HLS" + }, + { + "connector": "HTTPTS" + }, + { + "connector": "JPG" + }, + { + "connector": "JSON" + }, + { + "connector": "MP3" + }, + { + "connector": "MP4" + }, + { + "connector": "OGG" + }, + { + "connector": "RTMP", + "interface": "127.0.0.1", + "port": 31935 + }, + { + "connector": "SDP" + }, + { + "connector": "SubRip" + }, + { + "connector": "WAV" + }, + { + "connector": "HTTP", + "interface": "127.0.0.1", + "port": 28080, + "pubaddr": [] + } + ], + "serverid": null, + "sessionInputMode": 15, + "sessionOutputMode": 15, + "sessionStreamInfoMode": 1, + "sessionUnspecifiedMode": 0, + "sessionViewerMode": 14, + "tknMode": 15, + "triggers": { + "PUSH_REWRITE": [ + { + "handler": "http://127.0.0.1:39090/mist-trigger", + "streams": [], + "sync": true + } + ] + }, + "trustedproxy": [] + }, + "extwriters": null, + "push_settings": { + "maxspeed": 0, + "wait": 3 + }, + "streams": { + "stream": { + "debug": 5, + "name": "stream", + "processes": [ + { + "debug": 5, + "exec": "streamplace live $wildcard", + "exit_unmask": false, + "inconsequential": false, + "process": "MKVExec", + "restart_type": "fixed" + } + ], + "source": "push://", + "stop_sessions": false, + "tags": [] + } + }, + "ui_settings": { + "HTTPUrl": "http://127.0.0.1:28080/", + "sort_autopushes": { + "by": "Stream", + "dir": 1 + }, + "sort_pushes": { + "by": "Statistics", + "dir": 1 + }, + "sortstreams": { + "by": "name", + "dir": 1 + } + }, + "variables": null +} diff --git a/js/app/components/live-dashboard/stream-key.tsx b/js/app/components/live-dashboard/stream-key.tsx index 6030113bc..4daedb556 100644 --- a/js/app/components/live-dashboard/stream-key.tsx +++ b/js/app/components/live-dashboard/stream-key.tsx @@ -10,7 +10,6 @@ import { useEffect, useState } from "react"; import { useAppDispatch, useAppSelector } from "store/hooks"; import { View, Paragraph, Button, Text } from "tamagui"; import { Redirect } from "components/aqlink"; -import Waiting from "./waiting"; const Row = ({ children }: { children: React.ReactNode }) => { return ( @@ -36,6 +35,7 @@ const Right = ({ children }: { children: React.ReactNode }) => { }; export default function StreamKeyScreen() { + const [protocol, setProtocol] = useState("whip"); const isReady = useAppSelector(selectIsReady); if (!isReady) { return ; @@ -49,6 +49,7 @@ export default function StreamKeyScreen() { if (!userProfile) { return ; } + return ( - - Service - - - WHIP - - - - - Server - - - {url} - - - - - Bearer Token - - - - + + + {protocol === "whip" && } + {protocol === "rtmp" && } Output Settings @@ -99,11 +98,74 @@ export default function StreamKeyScreen() { - ); } +export function WHIPDescription({ url }: { url: string }) { + return ( + <> + + + Service + + + WHIP + + + + + Server + + + {url} + + + + + Bearer Token + + + + + + + ); +} + +export function RTMPDescription({ url }: { url: string }) { + const u = new URL(url); + const rtmpUrl = `rtmps://${u.host}:1935/live`; + return ( + <> + + + Service + + + Custom... + + + + + Server + + + {rtmpUrl} + + + + + Stream Key + + + + + + + ); +} + export function StreamKey() { const dispatch = useAppDispatch(); const [generating, setGenerating] = useState(false); diff --git a/js/app/src/screens/live-dashboard.tsx b/js/app/src/screens/live-dashboard.tsx index cbbadf8f3..f7eb0d39f 100644 --- a/js/app/src/screens/live-dashboard.tsx +++ b/js/app/src/screens/live-dashboard.tsx @@ -12,6 +12,10 @@ import React, { useCallback, useState } from "react"; import { useLiveUser } from "hooks/useLiveUser"; import StreamKeyScreen from "components/live-dashboard/stream-key"; import { VideoElementProvider } from "contexts/VideoElementContext"; +import { Camera, FerrisWheel, X } from "@tamagui/lucide-icons"; +import { H6, Text } from "tamagui"; +import Waiting from "components/live-dashboard/waiting"; +import { selectTelemetry } from "features/streamplace/streamplaceSlice"; enum StreamSource { Start, @@ -100,10 +104,6 @@ export default function LiveDashboard() { ); } -import { Camera, FerrisWheel, X } from "@tamagui/lucide-icons"; -import { H6, Text } from "tamagui"; -import Waiting from "components/live-dashboard/waiting"; -import { selectTelemetry } from "features/streamplace/streamplaceSlice"; const elems = [ { title: "Stream your camera!", diff --git a/pkg/api/api.go b/pkg/api/api.go index e6d5a06fa..b4f9fb289 100644 --- a/pkg/api/api.go +++ b/pkg/api/api.go @@ -59,8 +59,10 @@ type StreamplaceAPI struct { connTracker *WebsocketTracker - limiters map[string]*rate.Limiter - limitersMu sync.Mutex + limiters map[string]*rate.Limiter + limitersMu sync.Mutex + SignerCache map[string]media.MediaSigner + SignerCacheMu sync.Mutex } type WebsocketTracker struct { @@ -87,6 +89,7 @@ func MakeStreamplaceAPI(cli *config.CLI, mod model.Model, signer *eip712.EIP712S Director: d, connTracker: NewWebsocketTracker(cli.RateLimitWebsocket), limiters: make(map[string]*rate.Limiter), + SignerCache: make(map[string]media.MediaSigner), } a.Mimes, err = updater.GetMimes() if err != nil { diff --git a/pkg/api/api_internal.go b/pkg/api/api_internal.go index 8d5e9194b..32e8b6bad 100644 --- a/pkg/api/api_internal.go +++ b/pkg/api/api_internal.go @@ -25,6 +25,7 @@ import ( "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/errors" "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/media" "stream.place/streamplace/pkg/mist/mistconfig" "stream.place/streamplace/pkg/mist/misttriggers" "stream.place/streamplace/pkg/model" @@ -55,14 +56,27 @@ func init() { func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, error) { router := httprouter.New() broker := misttriggers.NewTriggerBroker() - broker.OnPushOutStart(func(ctx context.Context, payload *misttriggers.PushOutStartPayload) (string, error) { - return payload.URL, nil - }) + broker.OnPushRewrite(func(ctx context.Context, payload *misttriggers.PushRewritePayload) (string, error) { log.Log(ctx, "got push out start", "streamName", payload.StreamName, "url", payload.URL.String()) + // Extract the last part of the URL path + urlPath := payload.URL.Path + parts := strings.Split(urlPath, "/") + lastPart := "" + if len(parts) > 0 { + lastPart = parts[len(parts)-1] + } + mediaSigner, err := a.MakeMediaSigner(ctx, lastPart) + if err != nil { + return "", err + } ms := time.Now().UnixMilli() - out := fmt.Sprintf("%s+%s_%d", mistconfig.STREAM_NAME, payload.StreamName, ms) + out := fmt.Sprintf("%s+%s_%d", mistconfig.STREAM_NAME, mediaSigner.Streamer(), ms) + a.SignerCacheMu.Lock() + a.SignerCache[mediaSigner.Streamer()] = mediaSigner + a.SignerCacheMu.Unlock() + log.Log(ctx, "added key to cache", "mist-stream", out, "streamer", mediaSigner.Streamer()) return out, nil }) @@ -239,8 +253,32 @@ func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, err }) handleIncomingStream := func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { + key := p.ByName("key") log.Log(ctx, "stream start") - err := a.MediaManager.IngestStream(ctx, r.Body, a.MediaSigner) + + var mediaSigner media.MediaSigner + var ok bool + var err error + parts := strings.Split(key, "_") + + if len(parts) == 2 { + a.SignerCacheMu.Lock() + mediaSigner, ok = a.SignerCache[parts[0]] + a.SignerCacheMu.Unlock() + if !ok { + log.Error(ctx, "couldn't find key in cache", "part", parts[0], "key", key) + errors.WriteHTTPUnauthorized(w, "invalid authorization key", nil) + return + } + } else { + mediaSigner, err = a.MakeMediaSigner(ctx, key) + if err != nil { + errors.WriteHTTPUnauthorized(w, "invalid authorization key", err) + return + } + } + + err = a.MediaManager.MKVIngest(ctx, r.Body, mediaSigner) if err != nil { log.Log(ctx, "stream error", "error", err) @@ -251,8 +289,8 @@ func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, err } // route to accept an incoming mkv stream from OBS, segment it, and push the segments back to this HTTP handler - router.POST("/stream/:key", handleIncomingStream) - router.PUT("/stream/:key", handleIncomingStream) + router.POST("/live/:key", handleIncomingStream) + router.PUT("/live/:key", handleIncomingStream) router.GET("/player-report/:id", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { id := p.ByName("id") diff --git a/pkg/api/playback.go b/pkg/api/playback.go index 02e78bc2a..03d3915ee 100644 --- a/pkg/api/playback.go +++ b/pkg/api/playback.go @@ -4,7 +4,6 @@ import ( "bufio" "bytes" "context" - "crypto" "fmt" "io" "net/http" @@ -12,19 +11,13 @@ import ( "strings" "time" - atcrypto "github.com/bluesky-social/indigo/atproto/crypto" - "github.com/decred/dcrd/dcrec/secp256k1" "github.com/julienschmidt/httprouter" - "github.com/mr-tron/base58" "github.com/pion/webrtc/v4" "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/aqtime" - "stream.place/streamplace/pkg/atproto" "stream.place/streamplace/pkg/constants" "stream.place/streamplace/pkg/errors" - apierrors "stream.place/streamplace/pkg/errors" "stream.place/streamplace/pkg/log" - "stream.place/streamplace/pkg/media" "stream.place/streamplace/pkg/spmetrics" ) @@ -208,89 +201,9 @@ func (a *StreamplaceAPI) HandleWebRTCIngest(ctx context.Context) httprouter.Hand encoded = strings.TrimSpace(encoded) } - if len(encoded) < 2 || encoded[0] != 'z' { - errors.WriteHTTPUnauthorized(w, "invalid authorization key (not a multibase base58btc string)", nil) - return - } - - var addrBytes []byte - var didBytes []byte - priv, err := atcrypto.ParsePrivateMultibase(encoded) - if err == nil { - addrBytes = priv.Bytes() - } else { - decoded, err := base58.Decode(encoded[1:]) - if err != nil { - errors.WriteHTTPUnauthorized(w, "invalid authorization key (not a base58btc string)", nil) - return - } - addrBytes = decoded[:32] - didBytes = decoded[32:] - priv, err = atcrypto.ParsePrivateBytesK256(addrBytes) - if err != nil { - errors.WriteHTTPUnauthorized(w, "invalid authorization key (not valid atcrypto)", err) - return - } - } - - key, _ := secp256k1.PrivKeyFromBytes(addrBytes) - if key == nil { - errors.WriteHTTPUnauthorized(w, "invalid authorization key (not valid secp256k1)", nil) - return - } - var signer crypto.Signer = key.ToECDSA() - pub, err := priv.PublicKey() - if err != nil { - apierrors.WriteHTTPUnauthorized(w, "invalid authorization key (could not parse as atcrypto)", err) - return - } - - did := string(didBytes) - - if did != "" { - repo, err := a.ATSync.SyncBlueskyRepo(ctx, did, a.Model) - if err != nil { - apierrors.WriteHTTPInternalServerError(w, "could not resolve streamplace key", err) - return - } - err = a.CLI.StreamIsAllowed(repo.DID) - if err != nil { - apierrors.WriteHTTPUnauthorized(w, "user is not allowed to stream", err) - return - } - signingKey, err := a.Model.GetSigningKey(ctx, pub.DIDKey(), repo.DID) - if err != nil { - apierrors.WriteHTTPUnauthorized(w, "signing key not found", err) - return - } - if signingKey == nil { - apierrors.WriteHTTPUnauthorized(w, "signing key not found", nil) - return - } - } else { - atkey, err := atproto.ParsePubKey(signer.Public()) - if err != nil { - apierrors.WriteHTTPUnauthorized(w, "invalid authorization key (not valid secp256k1)", err) - return - } - did = atkey.DIDKey() - err = a.CLI.StreamIsAllowed(did) - if err != nil { - apierrors.WriteHTTPUnauthorized(w, "user is not allowed to stream", err) - return - } - } - - ctx = log.WithLogValues(ctx, "did", did) - - var mediaSigner media.MediaSigner - if a.CLI.ExternalSigning { - mediaSigner, err = media.MakeMediaSignerExt(ctx, a.CLI, did, addrBytes) - } else { - mediaSigner, err = media.MakeMediaSigner(ctx, a.CLI, did, signer) - } + mediaSigner, err := a.MakeMediaSigner(ctx, encoded) if err != nil { - errors.WriteHTTPUnauthorized(w, "invalid authorization key (not valid secp256k1)", err) + errors.WriteHTTPUnauthorized(w, "invalid authorization key", err) return } diff --git a/pkg/api/stream_key.go b/pkg/api/stream_key.go new file mode 100644 index 000000000..3377fb10c --- /dev/null +++ b/pkg/api/stream_key.go @@ -0,0 +1,92 @@ +package api + +import ( + "context" + "crypto" + "fmt" + + atcrypto "github.com/bluesky-social/indigo/atproto/crypto" + "github.com/decred/dcrd/dcrec/secp256k1" + "github.com/mr-tron/base58" + "stream.place/streamplace/pkg/atproto" + "stream.place/streamplace/pkg/log" + "stream.place/streamplace/pkg/media" +) + +func (a *StreamplaceAPI) MakeMediaSigner(ctx context.Context, keyStr string) (media.MediaSigner, error) { + if len(keyStr) < 2 || keyStr[0] != 'z' { + return nil, fmt.Errorf("invalid authorization key (not a multibase base58btc string)") + } + + var addrBytes []byte + var didBytes []byte + priv, err := atcrypto.ParsePrivateMultibase(keyStr) + if err == nil { + addrBytes = priv.Bytes() + } else { + decoded, err := base58.Decode(keyStr[1:]) + if err != nil { + return nil, fmt.Errorf("invalid authorization key (not a base58btc string)") + } + addrBytes = decoded[:32] + didBytes = decoded[32:] + priv, err = atcrypto.ParsePrivateBytesK256(addrBytes) + if err != nil { + return nil, fmt.Errorf("invalid authorization key (not valid atproto): %w", err) + } + } + + key, _ := secp256k1.PrivKeyFromBytes(addrBytes) + if key == nil { + return nil, fmt.Errorf("invalid authorization key (not valid secp256k1)") + } + var signer crypto.Signer = key.ToECDSA() + pub, err := priv.PublicKey() + if err != nil { + return nil, fmt.Errorf("invalid authorization key (could not parse as atproto): %w", err) + } + + did := string(didBytes) + + if did != "" { + repo, err := a.ATSync.SyncBlueskyRepo(ctx, did, a.Model) + if err != nil { + return nil, fmt.Errorf("could not resolve streamplace key: %w", err) + } + err = a.CLI.StreamIsAllowed(repo.DID) + if err != nil { + return nil, fmt.Errorf("user is not allowed to stream: %w", err) + } + signingKey, err := a.Model.GetSigningKey(ctx, pub.DIDKey(), repo.DID) + if err != nil { + return nil, fmt.Errorf("signing key not found: %w", err) + } + if signingKey == nil { + return nil, fmt.Errorf("signing key not found") + } + } else { + atkey, err := atproto.ParsePubKey(signer.Public()) + if err != nil { + return nil, fmt.Errorf("invalid authorization key (not valid secp256k1): %w", err) + } + did = atkey.DIDKey() + err = a.CLI.StreamIsAllowed(did) + if err != nil { + return nil, fmt.Errorf("user is not allowed to stream: %w", err) + } + } + + ctx = log.WithLogValues(ctx, "did", did) + + var mediaSigner media.MediaSigner + if a.CLI.ExternalSigning { + mediaSigner, err = media.MakeMediaSignerExt(ctx, a.CLI, did, addrBytes) + } else { + mediaSigner, err = media.MakeMediaSigner(ctx, a.CLI, did, signer) + } + if err != nil { + return nil, fmt.Errorf("invalid authorization key (not valid secp256k1): %w", err) + } + + return mediaSigner, nil +} diff --git a/pkg/cmd/live.go b/pkg/cmd/live.go new file mode 100644 index 000000000..d66c01341 --- /dev/null +++ b/pkg/cmd/live.go @@ -0,0 +1,45 @@ +package cmd + +import ( + "fmt" + "io" + "net/http" + "os" +) + +func Live(streamKey string) error { + // Create the URL for the live stream endpoint + url := fmt.Sprintf("http://127.0.0.1:39090/live/%s", streamKey) + + // Create a new HTTP request with POST method + req, err := http.NewRequest("POST", url, os.Stdin) + if err != nil { + return fmt.Errorf("error creating request: %w", err) + } + + // Set appropriate headers if needed + req.Header.Set("Content-Type", "video/x-matroska") // Assuming MKV format, adjust if needed + + // Create HTTP client and send the request + client := &http.Client{} + resp, err := client.Do(req) + if err != nil { + return fmt.Errorf("error sending stream: %w", err) + } + defer resp.Body.Close() + + // Check response status + if resp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(resp.Body) + return fmt.Errorf("server returned non-OK status: %d %s - %s", + resp.StatusCode, resp.Status, string(body)) + } + + // Copy response to stdout (if any) + _, err = io.Copy(os.Stdout, resp.Body) + if err != nil { + return fmt.Errorf("error reading response: %w", err) + } + + return nil +} diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 961bb5738..4effc1b0f 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -28,6 +28,7 @@ import ( "stream.place/streamplace/pkg/notifications" "stream.place/streamplace/pkg/replication" "stream.place/streamplace/pkg/replication/boring" + "stream.place/streamplace/pkg/rtmps" v0 "stream.place/streamplace/pkg/schema/v0" "stream.place/streamplace/pkg/spmetrics" @@ -81,6 +82,14 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { return Stream(os.Args[2]) } + if len(os.Args) > 1 && os.Args[1] == "live" { + if len(os.Args) != 3 { + fmt.Println("usage: streamplace live [stream-key]") + os.Exit(1) + } + return Live(os.Args[2]) + } + if len(os.Args) > 1 && os.Args[1] == "sign" { return Sign(context.Background()) } @@ -153,6 +162,8 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { fs.IntVar(&cli.RateLimitPerSecond, "rate-limit-per-second", 10, "rate limit for requests per second per ip") fs.IntVar(&cli.RateLimitBurst, "rate-limit-burst", 10, "rate limit burst for requests per ip") fs.IntVar(&cli.RateLimitWebsocket, "rate-limit-websocket", 10, "number of concurrent websocket connections allowed per ip") + fs.StringVar(&cli.RTMPServerAddon, "rtmp-server-addon", "", "address of external RTMP server to forward streams to") + fs.StringVar(&cli.RtmpsAddr, "rtmps-addr", ":1935", "address to listen for RTMPS connections") version := fs.Bool("version", false, "print version and exit") if runtime.GOOS == "linux" { @@ -351,6 +362,11 @@ func start(build *config.BuildFlags, platformJobs []jobFunc) error { group.Go(func() error { return a.ServeHTTPRedirect(ctx) }) + if cli.RTMPServerAddon != "" { + group.Go(func() error { + return rtmps.ServeRTMPS(ctx, &cli) + }) + } } else { group.Go(func() error { return a.ServeHTTP(ctx) diff --git a/pkg/config/config.go b/pkg/config/config.go index ab0dee106..512489d15 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -56,6 +56,7 @@ type CLI struct { HttpAddr string HttpInternalAddr string HttpsAddr string + RtmpsAddr string Secure bool NoMist bool MistAdminPort int @@ -90,6 +91,7 @@ type CLI struct { Thumbnail bool SmearAudio bool ExternalSigning bool + RTMPServerAddon string TracingEndpoint string PublicHost string RateLimitPerSecond int diff --git a/pkg/media/gstreamer.go b/pkg/media/gstreamer.go index 9ccf34d68..c2354b685 100644 --- a/pkg/media/gstreamer.go +++ b/pkg/media/gstreamer.go @@ -157,70 +157,6 @@ func SelfTest(ctx context.Context) error { return nil } -func (mm *MediaManager) IngestStream(ctx context.Context, input io.Reader, ms MediaSigner) error { - ctx, cancel := context.WithCancel(ctx) - defer cancel() - pipelineSlice := []string{ - "appsrc name=streamsrc ! matroskademux name=demux", - "demux. ! queue ! h264parse name=parse", - "demux. ! queue ! aacparse name=audioparse", - } - pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) - if err != nil { - return fmt.Errorf("error creating IngestStream pipeline: %w", err) - } - defer runtime.KeepAlive(pipeline) - srcele, err := pipeline.GetElementByName("streamsrc") - if err != nil { - return err - } - // defer runtime.KeepAlive(srcele) - src := app.SrcFromElement(srcele) - src.SetCallbacks(&app.SourceCallbacks{ - NeedDataFunc: ReaderNeedData(ctx, input), - }) - parseEle, err := pipeline.GetElementByName("parse") - if err != nil { - return err - } - - signer, err := mm.SegmentAndSignElem(ctx, ms) - if err != nil { - return err - } - - err = pipeline.Add(signer) - if err != nil { - return err - } - err = parseEle.Link(signer) - if err != nil { - return err - } - audioparse, err := pipeline.GetElementByName("audioparse") - if err != nil { - return err - } - err = audioparse.Link(signer) - if err != nil { - return err - } - - go func() { - HandleBusMessages(ctx, pipeline) - cancel() - }() - - err = pipeline.SetState(gst.StatePlaying) - if err != nil { - return err - } - - <-ctx.Done() - - return nil -} - const TESTSRC_WIDTH = 1280 const TESTSRC_HEIGHT = 720 const QR_SIZE = 256 diff --git a/pkg/media/mkv_ingest.go b/pkg/media/mkv_ingest.go new file mode 100644 index 000000000..00b3a6294 --- /dev/null +++ b/pkg/media/mkv_ingest.go @@ -0,0 +1,86 @@ +package media + +import ( + "context" + "fmt" + "io" + "strings" + + "github.com/go-gst/go-gst/gst" + "github.com/go-gst/go-gst/gst/app" + "stream.place/streamplace/pkg/log" +) + +// ingest a H264+AAC MKV stream (prolly from an RTMP server) +func (mm *MediaManager) MKVIngest(ctx context.Context, input io.Reader, ms MediaSigner) error { + ctx, cancel := context.WithCancel(ctx) + defer cancel() + pipelineSlice := []string{ + "appsrc name=streamsrc ! matroskademux name=demux", + "demux. ! queue ! h264parse name=parse", + "demux. ! queue ! fdkaacdec ! audioresample ! opusenc name=audioenc", + } + pipeline, err := gst.NewPipelineFromString(strings.Join(pipelineSlice, "\n")) + if err != nil { + return fmt.Errorf("error creating MKVIngest pipeline: %w", err) + } + + srcele, err := pipeline.GetElementByName("streamsrc") + if err != nil { + return err + } + // defer runtime.KeepAlive(srcele) + src := app.SrcFromElement(srcele) + src.SetCallbacks(&app.SourceCallbacks{ + NeedDataFunc: ReaderNeedDataIncremental(ctx, input), + }) + parseEle, err := pipeline.GetElementByName("parse") + if err != nil { + return err + } + + signer, err := mm.SegmentAndSignElem(ctx, ms) + if err != nil { + return err + } + + err = pipeline.Add(signer) + if err != nil { + return err + } + err = parseEle.Link(signer) + if err != nil { + return err + } + audioenc, err := pipeline.GetElementByName("audioenc") + if err != nil { + return err + } + err = audioenc.Link(signer) + if err != nil { + return err + } + + busErr := make(chan error) + go func() { + err := HandleBusMessages(ctx, pipeline) + cancel() + busErr <- err + }() + + err = pipeline.SetState(gst.StatePlaying) + if err != nil { + return err + } + + defer func() { + err := pipeline.SetState(gst.StateNull) + if err != nil { + log.Error(ctx, "error setting pipeline to null state", "error", err) + } + }() + + <-busErr + + return nil +} diff --git a/pkg/rtmps/rtmps.go b/pkg/rtmps/rtmps.go new file mode 100644 index 000000000..a5bc324eb --- /dev/null +++ b/pkg/rtmps/rtmps.go @@ -0,0 +1,99 @@ +package rtmps + +import ( + "context" + "crypto/tls" + "errors" + "fmt" + "io" + "net" + "sync" + + "stream.place/streamplace/pkg/config" + "stream.place/streamplace/pkg/log" +) + +// passthrough RTMPS TLS terminator to external RTMP server +func ServeRTMPS(ctx context.Context, cli *config.CLI) error { + if cli.RTMPServerAddon == "" { + return fmt.Errorf("RTMP server address not configured") + } + + cert, err := tls.LoadX509KeyPair(cli.TLSCertPath, cli.TLSKeyPath) + if err != nil { + return fmt.Errorf("failed to load TLS certificate: %w", err) + } + + tlsConfig := &tls.Config{ + Certificates: []tls.Certificate{cert}, + MinVersion: tls.VersionTLS12, + } + + listener, err := tls.Listen("tcp", cli.RtmpsAddr, tlsConfig) + if err != nil { + return fmt.Errorf("failed to create RTMPS listener: %w", err) + } + + log.Log(ctx, "rtmps server starting", + "addr", cli.RtmpsAddr, + "forwarding_to", cli.RTMPServerAddon) + + go func() { + <-ctx.Done() + listener.Close() + }() + + for { + conn, err := listener.Accept() + if err != nil { + // Check if the context was canceled, which means we're shutting down + select { + case <-ctx.Done(): + return nil + default: + log.Error(ctx, "error accepting RTMPS connection", "error", err) + continue + } + } + + go func(clientConn net.Conn) { + defer clientConn.Close() + + rtmpConn, err := net.Dial("tcp", cli.RTMPServerAddon) + if err != nil { + log.Error(ctx, "failed to connect to RTMP server", "error", err) + return + } + defer rtmpConn.Close() + + // Create a wait group to wait for both copy operations to complete + var wg sync.WaitGroup + wg.Add(2) + + // Copy from client to RTMP server + go func() { + defer wg.Done() + _, err := io.Copy(rtmpConn, clientConn) + if err != nil && !errors.Is(err, io.EOF) { + log.Error(ctx, "error copying from client to RTMP server", "error", err) + } + // Signal the other goroutine to stop by closing the connection + rtmpConn.Close() + }() + + // Copy from RTMP server to client + go func() { + defer wg.Done() + _, err := io.Copy(clientConn, rtmpConn) + if err != nil && !errors.Is(err, io.EOF) { + log.Error(ctx, "error copying from RTMP server to client", "error", err) + } + // Signal the other goroutine to stop by closing the connection + clientConn.Close() + }() + + // Wait for both copy operations to complete + wg.Wait() + }(conn) + } +}