From 93e8d15fc014aabe0fe6f17e3699bda4ab214b20 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Mon, 25 May 2026 17:03:15 -0700 Subject: [PATCH] =?UTF-8?q?spxrpc:=20XRPC=20live=20HLS=20=E2=80=94=20getLi?= =?UTF-8?q?vePlaylist=20+=20getLiveSegment?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Serve live HLS over XRPC (documented + auth-capable), reading the in-memory per-stream window via MediaManager.GetLiveWindow: - getLivePlaylist: master playlist (track omitted) or per-track media playlist. - getLiveSegment: a track's init (seg=init) or a windowed canonical .m4s (seg=), HTTP Range via http.ServeContent. Segment bytes are the verbatim signed .m4s, so provenance travels with playback. New lexicons under lexicons/place/stream/playback/ (make lexicons regenerates the stubs + OpenAPI docs + client types). Custom echo routes override the lexgen stubs for Range + the mpegurl MIME, mirroring getVideoBlob/Playlist; segment URLs end in a cosmetic ".m4s" so ffmpeg-based players fetch them. Open playback for now (account-ban gated); the XRPC framing lets auth middleware layer on without a wire change. View-count logging and mid-stream init changes (resolution switches) are follow-ups. Co-Authored-By: Claude Opus 4.7 (1M context) --- .../content/docs/lex-reference/openapi.json | 167 +++++++++++++++++ .../place-stream-playback-getliveplaylist.md | 89 +++++++++ .../place-stream-playback-getlivesegment.md | 89 +++++++++ .../stream/playback/getLivePlaylist.json | 46 +++++ .../place/stream/playback/getLiveSegment.json | 46 +++++ pkg/spxrpc/place_stream_playback_getlive.go | 173 ++++++++++++++++++ .../place_stream_playback_getlive_test.go | 38 ++++ pkg/spxrpc/spxrpc.go | 4 + pkg/spxrpc/stubs.go | 35 ++++ pkg/streamplace/playbackgetLivePlaylist.go | 35 ++++ pkg/streamplace/playbackgetLiveSegment.go | 35 ++++ 11 files changed, 757 insertions(+) create mode 100644 js/docs/src/content/docs/lex-reference/playback/place-stream-playback-getliveplaylist.md create mode 100644 js/docs/src/content/docs/lex-reference/playback/place-stream-playback-getlivesegment.md create mode 100644 lexicons/place/stream/playback/getLivePlaylist.json create mode 100644 lexicons/place/stream/playback/getLiveSegment.json create mode 100644 pkg/spxrpc/place_stream_playback_getlive.go create mode 100644 pkg/spxrpc/place_stream_playback_getlive_test.go create mode 100644 pkg/streamplace/playbackgetLivePlaylist.go create mode 100644 pkg/streamplace/playbackgetLiveSegment.go diff --git a/js/docs/src/content/docs/lex-reference/openapi.json b/js/docs/src/content/docs/lex-reference/openapi.json index 8d760e24..927e08b8 100644 --- a/js/docs/src/content/docs/lex-reference/openapi.json +++ b/js/docs/src/content/docs/lex-reference/openapi.json @@ -641,6 +641,173 @@ } } }, + "/xrpc/place.stream.playback.getLivePlaylist": { + "get": { + "summary": "Get an HLS CMAF playlist for a live stream. Returns a master playlist when `track` is omitted, or a single-track media playlist when `track` is supplied. The playlist references each segment + per-track init segment via getLiveSegment. Segments come from an in-memory sliding window fed as the stream is ingested (or replicated to this node), so a playlist is only available while the stream is live here.", + "operationId": "place.stream.playback.getLivePlaylist", + "tags": ["place.stream.playback"], + "responses": { + "200": { + "description": "Success", + "content": { + "*/*": { + "schema": {} + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "StreamNotLive" + }, + { + "const": "TrackNotFound" + }, + { + "const": "StreamUnavailable" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "parameters": [ + { + "name": "did", + "in": "query", + "required": true, + "description": "DID of the streamer whose live stream to play back.", + "schema": { + "type": "string", + "description": "DID of the streamer whose live stream to play back.", + "format": "did" + } + }, + { + "name": "track", + "in": "query", + "required": false, + "description": "Track ID (stringified u32 matching the MUXL container) for a single-track media playlist. Omit for the master playlist.", + "schema": { + "type": "string", + "description": "Track ID (stringified u32 matching the MUXL container) for a single-track media playlist. Omit for the master playlist." + } + }, + { + "name": "sid", + "in": "query", + "required": false, + "description": "Opaque playback session identifier. Omit on the master playlist request; the server generates one and threads it through the sub-playlist + segment URLs it returns, for view-count correlation.", + "schema": { + "type": "string", + "description": "Opaque playback session identifier. Omit on the master playlist request; the server generates one and threads it through the sub-playlist + segment URLs it returns, for view-count correlation." + } + } + ] + } + }, + "/xrpc/place.stream.playback.getLiveSegment": { + "get": { + "summary": "Fetch a single live HLS segment, or a track's init segment, from the in-memory live window. `seg` is `init` for the EXT-X-MAP init segment, otherwise the segment's media-sequence number; a cosmetic `.m4s` suffix is accepted (and ignored) so ffmpeg-based HLS players will fetch it. HTTP Range is honored. Segments are the verbatim signed canonical .m4s, so provenance travels with playback.", + "operationId": "place.stream.playback.getLiveSegment", + "tags": ["place.stream.playback"], + "responses": { + "200": { + "description": "Success", + "content": { + "video/mp4": { + "schema": {} + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "type": "object", + "required": ["error", "message"], + "properties": { + "error": { + "type": "string", + "oneOf": [ + { + "const": "StreamNotLive" + }, + { + "const": "SegmentNotFound" + } + ] + }, + "message": { + "type": "string" + } + } + } + } + } + } + }, + "parameters": [ + { + "name": "did", + "in": "query", + "required": true, + "description": "DID of the streamer.", + "schema": { + "type": "string", + "description": "DID of the streamer.", + "format": "did" + } + }, + { + "name": "track", + "in": "query", + "required": true, + "description": "Track ID (stringified u32 matching the MUXL container).", + "schema": { + "type": "string", + "description": "Track ID (stringified u32 matching the MUXL container)." + } + }, + { + "name": "seg", + "in": "query", + "required": true, + "description": "`init` for the track's init segment, or the segment's media-sequence number. A trailing `.m4s` is accepted and ignored.", + "schema": { + "type": "string", + "description": "`init` for the track's init segment, or the segment's media-sequence number. A trailing `.m4s` is accepted and ignored." + } + }, + { + "name": "sid", + "in": "query", + "required": false, + "description": "Opaque playback session identifier, propagated from the media playlist that referenced this segment. Logged for view-count correlation; not used for access control.", + "schema": { + "type": "string", + "description": "Opaque playback session identifier, propagated from the media playlist that referenced this segment. Logged for view-count correlation; not used for access control." + } + } + ] + } + }, "/xrpc/place.stream.playback.getPlaybackServer": { "get": { "summary": "Get available playback servers for a livestream.", diff --git a/js/docs/src/content/docs/lex-reference/playback/place-stream-playback-getliveplaylist.md b/js/docs/src/content/docs/lex-reference/playback/place-stream-playback-getliveplaylist.md new file mode 100644 index 00000000..9b59d8eb --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/playback/place-stream-playback-getliveplaylist.md @@ -0,0 +1,89 @@ +--- +title: place.stream.playback.getLivePlaylist +description: Reference for the place.stream.playback.getLivePlaylist lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `query` + +Get an HLS CMAF playlist for a live stream. Returns a master playlist when `track` is omitted, or a single-track media playlist when `track` is supplied. The playlist references each segment + per-track init segment via getLiveSegment. Segments come from an in-memory sliding window fed as the stream is ingested (or replicated to this node), so a playlist is only available while the stream is live here. + +**Parameters:** + +| Name | Type | Req'd | Description | Constraints | +| ------- | -------- | ----- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | ------------- | +| `did` | `string` | ✅ | DID of the streamer whose live stream to play back. | Format: `did` | +| `track` | `string` | ❌ | Track ID (stringified u32 matching the MUXL container) for a single-track media playlist. Omit for the master playlist. | | +| `sid` | `string` | ❌ | Opaque playback session identifier. Omit on the master playlist request; the server generates one and threads it through the sub-playlist + segment URLs it returns, for view-count correlation. | | + +**Output:** + +- **Encoding:** `*/*` +- **Schema:** + +_Schema not defined._ +**Possible Errors:** + +- `StreamNotLive`: No live segments are currently windowed for this streamer on this node. +- `TrackNotFound`: The requested track ID is not present in the live stream. +- `StreamUnavailable`: The streamer's account is unavailable (e.g. banned). + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.playback.getLivePlaylist", + "defs": { + "main": { + "type": "query", + "description": "Get an HLS CMAF playlist for a live stream. Returns a master playlist when `track` is omitted, or a single-track media playlist when `track` is supplied. The playlist references each segment + per-track init segment via getLiveSegment. Segments come from an in-memory sliding window fed as the stream is ingested (or replicated to this node), so a playlist is only available while the stream is live here.", + "parameters": { + "type": "params", + "required": ["did"], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "DID of the streamer whose live stream to play back." + }, + "track": { + "type": "string", + "description": "Track ID (stringified u32 matching the MUXL container) for a single-track media playlist. Omit for the master playlist." + }, + "sid": { + "type": "string", + "description": "Opaque playback session identifier. Omit on the master playlist request; the server generates one and threads it through the sub-playlist + segment URLs it returns, for view-count correlation." + } + } + }, + "output": { + "encoding": "*/*" + }, + "errors": [ + { + "name": "StreamNotLive", + "description": "No live segments are currently windowed for this streamer on this node." + }, + { + "name": "TrackNotFound", + "description": "The requested track ID is not present in the live stream." + }, + { + "name": "StreamUnavailable", + "description": "The streamer's account is unavailable (e.g. banned)." + } + ] + } + } +} +``` diff --git a/js/docs/src/content/docs/lex-reference/playback/place-stream-playback-getlivesegment.md b/js/docs/src/content/docs/lex-reference/playback/place-stream-playback-getlivesegment.md new file mode 100644 index 00000000..19b12a97 --- /dev/null +++ b/js/docs/src/content/docs/lex-reference/playback/place-stream-playback-getlivesegment.md @@ -0,0 +1,89 @@ +--- +title: place.stream.playback.getLiveSegment +description: Reference for the place.stream.playback.getLiveSegment lexicon +--- + +**Lexicon Version:** 1 + +## Definitions + + + +### `main` + +**Type:** `query` + +Fetch a single live HLS segment, or a track's init segment, from the in-memory live window. `seg` is `init` for the EXT-X-MAP init segment, otherwise the segment's media-sequence number; a cosmetic `.m4s` suffix is accepted (and ignored) so ffmpeg-based HLS players will fetch it. HTTP Range is honored. Segments are the verbatim signed canonical .m4s, so provenance travels with playback. + +**Parameters:** + +| Name | Type | Req'd | Description | Constraints | +| ------- | -------- | ----- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------- | +| `did` | `string` | ✅ | DID of the streamer. | Format: `did` | +| `track` | `string` | ✅ | Track ID (stringified u32 matching the MUXL container). | | +| `seg` | `string` | ✅ | `init` for the track's init segment, or the segment's media-sequence number. A trailing `.m4s` is accepted and ignored. | | +| `sid` | `string` | ❌ | Opaque playback session identifier, propagated from the media playlist that referenced this segment. Logged for view-count correlation; not used for access control. | | + +**Output:** + +- **Encoding:** `video/mp4` +- **Schema:** + +_Schema not defined._ +**Possible Errors:** + +- `StreamNotLive`: No live window for this streamer on this node. +- `SegmentNotFound`: The requested segment has slid out of the window, or never existed. + +--- + +## Lexicon Source + +```json +{ + "lexicon": 1, + "id": "place.stream.playback.getLiveSegment", + "defs": { + "main": { + "type": "query", + "description": "Fetch a single live HLS segment, or a track's init segment, from the in-memory live window. `seg` is `init` for the EXT-X-MAP init segment, otherwise the segment's media-sequence number; a cosmetic `.m4s` suffix is accepted (and ignored) so ffmpeg-based HLS players will fetch it. HTTP Range is honored. Segments are the verbatim signed canonical .m4s, so provenance travels with playback.", + "parameters": { + "type": "params", + "required": ["did", "track", "seg"], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "DID of the streamer." + }, + "track": { + "type": "string", + "description": "Track ID (stringified u32 matching the MUXL container)." + }, + "seg": { + "type": "string", + "description": "`init` for the track's init segment, or the segment's media-sequence number. A trailing `.m4s` is accepted and ignored." + }, + "sid": { + "type": "string", + "description": "Opaque playback session identifier, propagated from the media playlist that referenced this segment. Logged for view-count correlation; not used for access control." + } + } + }, + "output": { + "encoding": "video/mp4" + }, + "errors": [ + { + "name": "StreamNotLive", + "description": "No live window for this streamer on this node." + }, + { + "name": "SegmentNotFound", + "description": "The requested segment has slid out of the window, or never existed." + } + ] + } + } +} +``` diff --git a/lexicons/place/stream/playback/getLivePlaylist.json b/lexicons/place/stream/playback/getLivePlaylist.json new file mode 100644 index 00000000..7894072e --- /dev/null +++ b/lexicons/place/stream/playback/getLivePlaylist.json @@ -0,0 +1,46 @@ +{ + "lexicon": 1, + "id": "place.stream.playback.getLivePlaylist", + "defs": { + "main": { + "type": "query", + "description": "Get an HLS CMAF playlist for a live stream. Returns a master playlist when `track` is omitted, or a single-track media playlist when `track` is supplied. The playlist references each segment + per-track init segment via getLiveSegment. Segments come from an in-memory sliding window fed as the stream is ingested (or replicated to this node), so a playlist is only available while the stream is live here.", + "parameters": { + "type": "params", + "required": ["did"], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "DID of the streamer whose live stream to play back." + }, + "track": { + "type": "string", + "description": "Track ID (stringified u32 matching the MUXL container) for a single-track media playlist. Omit for the master playlist." + }, + "sid": { + "type": "string", + "description": "Opaque playback session identifier. Omit on the master playlist request; the server generates one and threads it through the sub-playlist + segment URLs it returns, for view-count correlation." + } + } + }, + "output": { + "encoding": "*/*" + }, + "errors": [ + { + "name": "StreamNotLive", + "description": "No live segments are currently windowed for this streamer on this node." + }, + { + "name": "TrackNotFound", + "description": "The requested track ID is not present in the live stream." + }, + { + "name": "StreamUnavailable", + "description": "The streamer's account is unavailable (e.g. banned)." + } + ] + } + } +} diff --git a/lexicons/place/stream/playback/getLiveSegment.json b/lexicons/place/stream/playback/getLiveSegment.json new file mode 100644 index 00000000..86303b1f --- /dev/null +++ b/lexicons/place/stream/playback/getLiveSegment.json @@ -0,0 +1,46 @@ +{ + "lexicon": 1, + "id": "place.stream.playback.getLiveSegment", + "defs": { + "main": { + "type": "query", + "description": "Fetch a single live HLS segment, or a track's init segment, from the in-memory live window. `seg` is `init` for the EXT-X-MAP init segment, otherwise the segment's media-sequence number; a cosmetic `.m4s` suffix is accepted (and ignored) so ffmpeg-based HLS players will fetch it. HTTP Range is honored. Segments are the verbatim signed canonical .m4s, so provenance travels with playback.", + "parameters": { + "type": "params", + "required": ["did", "track", "seg"], + "properties": { + "did": { + "type": "string", + "format": "did", + "description": "DID of the streamer." + }, + "track": { + "type": "string", + "description": "Track ID (stringified u32 matching the MUXL container)." + }, + "seg": { + "type": "string", + "description": "`init` for the track's init segment, or the segment's media-sequence number. A trailing `.m4s` is accepted and ignored." + }, + "sid": { + "type": "string", + "description": "Opaque playback session identifier, propagated from the media playlist that referenced this segment. Logged for view-count correlation; not used for access control." + } + } + }, + "output": { + "encoding": "video/mp4" + }, + "errors": [ + { + "name": "StreamNotLive", + "description": "No live window for this streamer on this node." + }, + { + "name": "SegmentNotFound", + "description": "The requested segment has slid out of the window, or never existed." + } + ] + } + } +} diff --git a/pkg/spxrpc/place_stream_playback_getlive.go b/pkg/spxrpc/place_stream_playback_getlive.go new file mode 100644 index 00000000..fcaaf322 --- /dev/null +++ b/pkg/spxrpc/place_stream_playback_getlive.go @@ -0,0 +1,173 @@ +package spxrpc + +import ( + "bytes" + "context" + "io" + "net/http" + "net/url" + "strconv" + "strings" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/labstack/echo/v4" +) + +// --- stubs for the auto-generated wrappers in stubs.go ------------------ +// +// Like getVideoBlob/getVideoPlaylist, the lexgen stubs hard-code status 200 + +// a fixed Content-Type, neither of which suits these endpoints (Range/206 on +// getLiveSegment, vnd.apple.mpegurl on getLivePlaylist). NewServer registers +// custom echo routes that override them; these exist only to satisfy the build. + +func (s *Server) handlePlaceStreamPlaybackGetLivePlaylist(ctx context.Context, did string, sid string, track string) (io.Reader, error) { + return nil, stubMisrouted("getLivePlaylist") +} + +func (s *Server) handlePlaceStreamPlaybackGetLiveSegment(ctx context.Context, did string, seg string, sid string, track string) (io.Reader, error) { + return nil, stubMisrouted("getLiveSegment") +} + +// --- getLivePlaylist ---------------------------------------------------- + +// HandleGetLivePlaylist serves a live HLS master playlist (track omitted) or a +// single-track media playlist out of the streamer's in-memory live window. The +// window is fed by ValidateMP4 for every segment that flows through this node, +// so a playlist exists only while the stream is live here. Open playback, +// gated only on an account ban (auth middleware can layer on later). +func (s *Server) HandleGetLivePlaylist(c echo.Context) error { + parsedDID, err := syntax.ParseDID(c.QueryParam("did")) + if err != nil { + return echo.NewHTTPError(http.StatusBadRequest, "did is required and must be a valid DID") + } + did := parsedDID.String() + + if banned, err := s.accountBanned(did); err != nil { + return echo.NewHTTPError(http.StatusInternalServerError, err.Error()) + } else if banned { + return echo.NewHTTPError(http.StatusForbidden, "StreamUnavailable") + } + + w := s.mm.GetLiveWindow(did) + if w == nil { + return echo.NewHTTPError(http.StatusNotFound, "StreamNotLive") + } + + // Reuse the caller's sid when present so master/media/segment requests of + // one playback session share an identifier; mint one otherwise. + sid, err := sessionIDOrNew(c.QueryParam("sid")) + if err != nil { + return echo.NewHTTPError(http.StatusBadRequest, err.Error()) + } + + track := c.QueryParam("track") + var body string + if track == "" { + body = w.MasterPlaylist(func(tid string) string { + return liveTrackPlaylistURL(did, tid, sid) + }) + } else { + initURL := liveSegmentURL(did, track, "init", sid) + segURI := func(seq uint64) string { + return liveSegmentURL(did, track, strconv.FormatUint(seq, 10), sid) + } + body = w.MediaPlaylist(track, initURL, segURI) + if body == "" { + return echo.NewHTTPError(http.StatusNotFound, "TrackNotFound") + } + } + + h := c.Response().Header() + h.Set("Content-Type", "application/vnd.apple.mpegurl") + // The live media playlist changes as the window slides; never cache it. + h.Set("Cache-Control", "no-cache") + c.Response().WriteHeader(http.StatusOK) + _, err = c.Response().Writer.Write([]byte(body)) + return err +} + +// --- getLiveSegment ----------------------------------------------------- + +// HandleGetLiveSegment serves a track's init segment (seg=init) or a windowed +// canonical .m4s (seg=) from the live window, with HTTP Range. +// The bytes are the verbatim signed segment, so provenance travels with +// playback. +func (s *Server) HandleGetLiveSegment(c echo.Context) error { + parsedDID, err := syntax.ParseDID(c.QueryParam("did")) + if err != nil { + return echo.NewHTTPError(http.StatusBadRequest, "did is required and must be a valid DID") + } + did := parsedDID.String() + track := c.QueryParam("track") + if track == "" { + return echo.NewHTTPError(http.StatusBadRequest, "track is required") + } + // The cosmetic ".m4s" lets ffmpeg-based players fetch the URL; ignore it. + seg := strings.TrimSuffix(c.QueryParam("seg"), ".m4s") + if seg == "" { + return echo.NewHTTPError(http.StatusBadRequest, "seg is required") + } + + w := s.mm.GetLiveWindow(did) + if w == nil { + return echo.NewHTTPError(http.StatusNotFound, "StreamNotLive") + } + + var data []byte + isInit := seg == "init" + if isInit { + data = w.InitSegment(track) + } else { + seq, err := strconv.ParseUint(seg, 10, 64) + if err != nil { + return echo.NewHTTPError(http.StatusBadRequest, "seg must be 'init' or a media-sequence number") + } + data = w.SegmentData(track, seq) + } + if data == nil { + return echo.NewHTTPError(http.StatusNotFound, "SegmentNotFound") + } + + h := c.Response().Header() + h.Set("Content-Type", "video/mp4") + if isInit { + // The init can change mid-stream (e.g. a resolution change mints a new + // catalog), so don't let it be cached long. + h.Set("Cache-Control", "no-cache") + } else { + // A numbered segment's bytes are fixed once minted (signed), so it's + // safely immutable for as long as it stays in the window. + h.Set("Cache-Control", "public, max-age=31536000, immutable") + } + // http.ServeContent honors Range / sets Accept-Ranges + Content-Length and + // keeps the Content-Type we set above. + http.ServeContent(c.Response().Writer, c.Request(), "segment.m4s", time.Time{}, bytes.NewReader(data)) + return nil +} + +// --- url builders ------------------------------------------------------- + +// liveTrackPlaylistURL is the URL to a single-track live media playlist served +// by this same handler. sid is propagated so a player's playlist + segment +// requests share an identifier. +func liveTrackPlaylistURL(did, track, sid string) string { + q := url.Values{"did": {did}, "track": {track}} + if sid != "" { + q.Set("sid", sid) + } + return "/xrpc/place.stream.playback.getLivePlaylist?" + q.Encode() +} + +// liveSegmentURL builds a getLiveSegment URL. seg is "init" or a media +// sequence number; the trailing ".m4s" is appended last (so seg is the final +// query token) to satisfy ffmpeg's segment-extension allowlist — the handler +// strips it back off. +func liveSegmentURL(did, track, seg, sid string) string { + q := url.Values{"did": {did}, "track": {track}} + if sid != "" { + q.Set("sid", sid) + } + return "/xrpc/place.stream.playback.getLiveSegment?" + q.Encode() + + "&seg=" + url.QueryEscape(seg) + ".m4s" +} diff --git a/pkg/spxrpc/place_stream_playback_getlive_test.go b/pkg/spxrpc/place_stream_playback_getlive_test.go new file mode 100644 index 00000000..33e413b0 --- /dev/null +++ b/pkg/spxrpc/place_stream_playback_getlive_test.go @@ -0,0 +1,38 @@ +package spxrpc + +import ( + "strings" + "testing" + + "github.com/stretchr/testify/require" +) + +// TestLiveSegmentURL pins the ffmpeg-compatibility contract: the segment URL +// must END in ".m4s" (ffmpeg's HLS demuxer checks the URL extension against +// its allowlist before fetching), and the handler must be able to recover the +// did/track/seg from it. +func TestLiveSegmentURL(t *testing.T) { + u := liveSegmentURL("did:plc:abc123", "1", "42", "3kabc") + require.True(t, strings.HasSuffix(u, ".m4s"), "segment URL must end in .m4s, got %q", u) + require.True(t, strings.HasPrefix(u, "/xrpc/place.stream.playback.getLiveSegment?")) + require.Contains(t, u, "track=1") + require.Contains(t, u, "sid=3kabc") + require.Contains(t, u, "seg=42.m4s") + // DID is percent-encoded in the query. + require.Contains(t, u, "did=did%3Aplc%3Aabc123") + + // init segment URL (no sid). + ui := liveSegmentURL("did:plc:abc123", "2", "init", "") + require.True(t, strings.HasSuffix(ui, "seg=init.m4s"), "init URL must end in seg=init.m4s, got %q", ui) + require.NotContains(t, ui, "sid=") +} + +// TestLiveTrackPlaylistURL confirms the master playlist's per-track sub-URLs +// carry did + track + sid back to getLivePlaylist. +func TestLiveTrackPlaylistURL(t *testing.T) { + u := liveTrackPlaylistURL("did:plc:xyz", "1", "3ksid") + require.True(t, strings.HasPrefix(u, "/xrpc/place.stream.playback.getLivePlaylist?")) + require.Contains(t, u, "track=1") + require.Contains(t, u, "sid=3ksid") + require.Contains(t, u, "did=did%3Aplc%3Axyz") +} diff --git a/pkg/spxrpc/spxrpc.go b/pkg/spxrpc/spxrpc.go index 12d66d0b..216b8849 100644 --- a/pkg/spxrpc/spxrpc.go +++ b/pkg/spxrpc/spxrpc.go @@ -104,6 +104,10 @@ func NewServer(ctx context.Context, cli *config.CLI, model model.Model, stateful // echo's last-write-wins for exact-match routes. e.GET("/xrpc/place.stream.playback.getVideoBlob", s.HandleGetVideoBlob) e.GET("/xrpc/place.stream.playback.getVideoPlaylist", s.HandleGetVideoPlaylist) + // Live HLS, same override rationale (Range + mpegurl), served from the + // in-memory live window instead of a stored metafile. + e.GET("/xrpc/place.stream.playback.getLivePlaylist", s.HandleGetLivePlaylist) + e.GET("/xrpc/place.stream.playback.getLiveSegment", s.HandleGetLiveSegment) e.GET("/xrpc/*", s.HandleWildcard) e.POST("/xrpc/*", s.HandleWildcard) return s, nil diff --git a/pkg/spxrpc/stubs.go b/pkg/spxrpc/stubs.go index 508dfc92..b6ab8c69 100644 --- a/pkg/spxrpc/stubs.go +++ b/pkg/spxrpc/stubs.go @@ -318,6 +318,8 @@ func (s *Server) RegisterHandlersPlaceStream(e *echo.Echo) error { e.POST("/xrpc/place.stream.multistream.deleteTarget", s.HandlePlaceStreamMultistreamDeleteTarget) e.GET("/xrpc/place.stream.multistream.listTargets", s.HandlePlaceStreamMultistreamListTargets) e.POST("/xrpc/place.stream.multistream.putTarget", s.HandlePlaceStreamMultistreamPutTarget) + e.GET("/xrpc/place.stream.playback.getLivePlaylist", s.HandlePlaceStreamPlaybackGetLivePlaylist) + e.GET("/xrpc/place.stream.playback.getLiveSegment", s.HandlePlaceStreamPlaybackGetLiveSegment) e.GET("/xrpc/place.stream.playback.getPlaybackServer", s.HandlePlaceStreamPlaybackGetPlaybackServer) e.GET("/xrpc/place.stream.playback.getVideoBlob", s.HandlePlaceStreamPlaybackGetVideoBlob) e.GET("/xrpc/place.stream.playback.getVideoPlaylist", s.HandlePlaceStreamPlaybackGetVideoPlaylist) @@ -974,6 +976,39 @@ func (s *Server) HandlePlaceStreamMultistreamPutTarget(c echo.Context) error { return c.JSON(200, out) } +func (s *Server) HandlePlaceStreamPlaybackGetLivePlaylist(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamPlaybackGetLivePlaylist") + defer span.End() + did := c.QueryParam("did") + sid := c.QueryParam("sid") + track := c.QueryParam("track") + var out io.Reader + var handleErr error + // func (s *Server) handlePlaceStreamPlaybackGetLivePlaylist(ctx context.Context,did string,sid string,track string) (io.Reader, error) + out, handleErr = s.handlePlaceStreamPlaybackGetLivePlaylist(ctx, did, sid, track) + if handleErr != nil { + return handleErr + } + return c.Stream(200, "application/octet-stream", out) +} + +func (s *Server) HandlePlaceStreamPlaybackGetLiveSegment(c echo.Context) error { + ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamPlaybackGetLiveSegment") + defer span.End() + did := c.QueryParam("did") + seg := c.QueryParam("seg") + sid := c.QueryParam("sid") + track := c.QueryParam("track") + var out io.Reader + var handleErr error + // func (s *Server) handlePlaceStreamPlaybackGetLiveSegment(ctx context.Context,did string,seg string,sid string,track string) (io.Reader, error) + out, handleErr = s.handlePlaceStreamPlaybackGetLiveSegment(ctx, did, seg, sid, track) + if handleErr != nil { + return handleErr + } + return c.Stream(200, "video/mp4", out) +} + func (s *Server) HandlePlaceStreamPlaybackGetPlaybackServer(c echo.Context) error { ctx, span := otel.Tracer("server").Start(c.Request().Context(), "HandlePlaceStreamPlaybackGetPlaybackServer") defer span.End() diff --git a/pkg/streamplace/playbackgetLivePlaylist.go b/pkg/streamplace/playbackgetLivePlaylist.go new file mode 100644 index 00000000..4ce2e143 --- /dev/null +++ b/pkg/streamplace/playbackgetLivePlaylist.go @@ -0,0 +1,35 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +// Lexicon schema: place.stream.playback.getLivePlaylist + +package streamplace + +import ( + "bytes" + "context" + + lexutil "github.com/bluesky-social/indigo/lex/util" +) + +// PlaybackGetLivePlaylist calls the XRPC method "place.stream.playback.getLivePlaylist". +// +// did: DID of the streamer whose live stream to play back. +// sid: Opaque playback session identifier. Omit on the master playlist request; the server generates one and threads it through the sub-playlist + segment URLs it returns, for view-count correlation. +// track: Track ID (stringified u32 matching the MUXL container) for a single-track media playlist. Omit for the master playlist. +func PlaybackGetLivePlaylist(ctx context.Context, c lexutil.LexClient, did string, sid string, track string) ([]byte, error) { + buf := new(bytes.Buffer) + + params := map[string]interface{}{} + params["did"] = did + if sid != "" { + params["sid"] = sid + } + if track != "" { + params["track"] = track + } + if err := c.LexDo(ctx, lexutil.Query, "", "place.stream.playback.getLivePlaylist", params, nil, buf); err != nil { + return nil, err + } + + return buf.Bytes(), nil +} diff --git a/pkg/streamplace/playbackgetLiveSegment.go b/pkg/streamplace/playbackgetLiveSegment.go new file mode 100644 index 00000000..fa914256 --- /dev/null +++ b/pkg/streamplace/playbackgetLiveSegment.go @@ -0,0 +1,35 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +// Lexicon schema: place.stream.playback.getLiveSegment + +package streamplace + +import ( + "bytes" + "context" + + lexutil "github.com/bluesky-social/indigo/lex/util" +) + +// PlaybackGetLiveSegment calls the XRPC method "place.stream.playback.getLiveSegment". +// +// did: DID of the streamer. +// seg: `init` for the track's init segment, or the segment's media-sequence number. A trailing `.m4s` is accepted and ignored. +// sid: Opaque playback session identifier, propagated from the media playlist that referenced this segment. Logged for view-count correlation; not used for access control. +// track: Track ID (stringified u32 matching the MUXL container). +func PlaybackGetLiveSegment(ctx context.Context, c lexutil.LexClient, did string, seg string, sid string, track string) ([]byte, error) { + buf := new(bytes.Buffer) + + params := map[string]interface{}{} + params["did"] = did + params["seg"] = seg + if sid != "" { + params["sid"] = sid + } + params["track"] = track + if err := c.LexDo(ctx, lexutil.Query, "", "place.stream.playback.getLiveSegment", params, nil, buf); err != nil { + return nil, err + } + + return buf.Bytes(), nil +} -- 2.51.2