Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
20 kB · 576 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577package api
import ( "context" "encoding/base64" "encoding/json" "fmt" "io" "log/slog" "net/http" "net/http/pprof" "os" "runtime" rtpprof "runtime/pprof" "strconv" "strings" "time"
"github.com/juju/ratelimit" "github.com/julienschmidt/httprouter" "github.com/pion/webrtc/v4" "github.com/prometheus/client_golang/prometheus/promhttp" sloghttp "github.com/samber/slog-http" "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" notificationpkg "stream.place/streamplace/pkg/notifications" "stream.place/streamplace/pkg/rtcrec" v0 "stream.place/streamplace/pkg/schema/v0")
func (a *StreamplaceAPI) ServeInternalHTTP(ctx context.Context) error { handler, err := a.InternalHandler(ctx) if err != nil { return err } return a.ServerWithShutdown(ctx, handler, func(s *http.Server) error { s.Addr = a.CLI.HTTPInternalAddr log.Log(ctx, "http server starting", "addr", s.Addr) return s.ListenAndServe() })}
func (a *StreamplaceAPI) InternalHandler(ctx context.Context) (http.Handler, error) { router := httprouter.New() broker := misttriggers.NewTriggerBroker()
// serverCtx outlives any single trigger request — the Mist pull ingests // spawned below run for the life of their stream, not the life of the // PUSH_REWRITE request that announced it. serverCtx := ctx broker.OnPushRewrite(func(ctx context.Context, payload *misttriggers.PushRewritePayload) (string, error) { log.Log(ctx, "got push rewrite", "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.StreamName, 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())
// The push is authed and named — ingest it by pulling Mist's live fMP4 // output for the stream we just named. This replaces the old Mist-side // MKVExec process (`streamplace live` POSTing MKV back to /live): fMP4 // carries real decode timestamps, so ingest no longer reconstructs DTS. // Mist accepts the push right after this trigger returns, so the pull // retries briefly while the stream boots — including waiting for a // header with tracks, since an encoder can connect before its media // registers (mistPullConnect). go func() { if perr := a.MediaManager.MistPullIngest(serverCtx, out, mediaSigner); perr != nil { log.Error(serverCtx, "mist pull ingest ended", "mist-stream", out, "streamer", mediaSigner.Streamer(), "error", perr) } else { log.Log(serverCtx, "mist pull ingest ended cleanly", "mist-stream", out, "streamer", mediaSigner.Streamer()) } }()
return out, nil }) triggerCollection := misttriggers.NewMistCallbackHandlersCollection(a.CLI, broker) router.POST("/mist-trigger", triggerCollection.Trigger()) router.HandlerFunc("GET", "/healthz", a.HandleHealthz(ctx))
// Add pprof handlers router.HandlerFunc("GET", "/debug/pprof/", pprof.Index) router.HandlerFunc("GET", "/debug/pprof/cmdline", pprof.Cmdline) router.HandlerFunc("GET", "/debug/pprof/profile", pprof.Profile) router.HandlerFunc("GET", "/debug/pprof/symbol", pprof.Symbol) router.HandlerFunc("GET", "/debug/pprof/trace", pprof.Trace) router.Handler("GET", "/debug/pprof/goroutine", pprof.Handler("goroutine")) router.Handler("GET", "/debug/pprof/heap", pprof.Handler("heap")) router.Handler("GET", "/debug/pprof/threadcreate", pprof.Handler("threadcreate")) router.Handler("GET", "/debug/pprof/block", pprof.Handler("block")) router.Handler("GET", "/debug/pprof/allocs", pprof.Handler("allocs")) router.Handler("GET", "/debug/pprof/mutex", pprof.Handler("mutex"))
router.POST("/gc", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { runtime.GC() w.WriteHeader(204) })
// Pull a VOD blob from another Streamplace node into our playback // store and attest to it via place.stream.media.origin. router.POST("/vod-transfer", a.HandleVODTransfer(ctx))
// Rebuild the local media.origin index from our own server repo, for // blobs we host but never indexed (a dropped firehose event, or a // --secure node whose self-subscription never connected). router.POST("/reindex-origins", a.HandleReindexOrigins(ctx))
router.Handler("GET", "/metrics", promhttp.Handler())
// Legacy disk-served HLS (ffconcat -> latest.mp4 -> segment/:file, all reading // .m4s off disk) has been removed — live playback is all MUXL HLS now, served // from the in-memory window in pkg/livehls / place_stream_playback_getlive.
router.HEAD("/playback/:user/:rendition/stream.mkv", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { user := p.ByName("user") if user == "" { errors.WriteHTTPBadRequest(w, "user required", nil) return } w.Header().Set("Content-Type", "video/x-matroska") w.Header().Set("Transfer-Encoding", "chunked") w.WriteHeader(200) })
router.POST("/http-pipe/:uuid", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { uu := p.ByName("uuid") if uu == "" { errors.WriteHTTPBadRequest(w, "uuid required", nil) return } pr := a.MediaManager.GetHTTPPipeWriter(uu) if pr == nil { errors.WriteHTTPNotFound(w, "http-pipe not found", nil) return } if _, err := io.Copy(pr, r.Body); err != nil { errors.WriteHTTPInternalServerError(w, "failed to copy response", nil) } })
// self-destruct code, useful for dumping goroutines on windows router.POST("/abort", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { if err := rtpprof.Lookup("goroutine").WriteTo(os.Stderr, 2); err != nil { log.Log(ctx, "error writing rtpprof", "error", err) } log.Log(ctx, "got POST /abort, self-destructing") os.Exit(1) })
handleIncomingStream := func(w http.ResponseWriter, httpReq *http.Request, p httprouter.Params) { reqCtx := httpReq.Context() var r io.Reader = httpReq.Body key := p.ByName("key") limitStr := httpReq.URL.Query().Get("ratelimit") if limitStr != "" { limit, err := strconv.Atoi(limitStr) if err != nil { errors.WriteHTTPBadRequest(w, "invalid ratelimit", err) return } bucket := ratelimit.NewBucketWithRate(float64(limit), int64(limit)) // 2 Mbps r = ratelimit.Reader(r, bucket) } log.Log(reqCtx, "stream start")
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(reqCtx, "couldn't find key in cache", "part", parts[0], "key", key) errors.WriteHTTPUnauthorized(w, "invalid authorization key", nil) return } } else { mediaSigner, err = a.MakeMediaSigner(reqCtx, key) if err != nil { errors.WriteHTTPUnauthorized(w, "invalid authorization key", err) return } }
reqCtx = log.WithLogValues(reqCtx, "streamer", mediaSigner.Streamer())
err = a.checkBanned(reqCtx, mediaSigner.Streamer()) if err != nil { errors.WriteHTTPUnauthorized(w, err.Error(), err) return }
if a.CLI.IsolatedIngest { // Zero-downtime path: hijack the authed push connection and hand it to a // DETACHED worker that owns the connection (so it survives a main // restart) and serves signed segments back over its socket. if hj, ok := w.(http.Hijacker); ok { conn, bufrw, herr := hj.Hijack() if herr != nil { log.Error(reqCtx, "ingest hijack failed", "error", herr) return } var prebuf []byte if n := bufrw.Reader.Buffered(); n > 0 { prebuf = make([]byte, n) _, _ = io.ReadFull(bufrw.Reader, prebuf) } chunked := len(httpReq.TransferEncoding) > 0 && httpReq.TransferEncoding[0] == "chunked" if derr := a.MediaManager.MP4IngestDetached(reqCtx, conn, prebuf, chunked, mediaSigner); derr != nil { log.Log(reqCtx, "isolated stream ended", "error", derr) } return // connection hijacked; the HTTP response is ours now } // The isolated path needs a hijackable HTTP/1.1 connection (which the // real /live clients — the `streamplace live` CLI and tests pushing // over localhost — always are; the Mist ingest itself now arrives via // MistPullIngest, not this route). We don't support a non-hijack // fallback: it couldn't receive mid-stream manifest updates and would // stay stuck pre-live, so refuse. Such a client can use WHIP instead. log.Error(reqCtx, "isolated ingest requires a hijackable HTTP/1.1 connection; refusing push") errors.WriteHTTPInternalServerError(w, "isolated ingest requires a hijackable HTTP/1.1 connection; use WHIP", fmt.Errorf("connection is not hijackable")) return } else { err = a.MediaManager.MP4Ingest(reqCtx, r, mediaSigner) }
if err != nil { log.Log(reqCtx, "stream error", "error", err) errors.WriteHTTPInternalServerError(w, "stream error", err) return } log.Log(reqCtx, "stream success", "url", httpReq.URL.String()) }
// route to accept an incoming fragmented-MP4 stream (the `streamplace live` // CLI piping from stdin), segment it, and validate the signed segments. // The co-located MistServer's streams are ingested by pulling its fMP4 // output instead (MistPullIngest, kicked off from PUSH_REWRITE above). 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") if id == "" { errors.WriteHTTPBadRequest(w, "id required", nil) return } events, err := a.Model.PlayerReport(id) if err != nil { errors.WriteHTTPBadRequest(w, err.Error(), err) return } bs, err := json.Marshal(events) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to marhsal json", err) return } if _, err := w.Write(bs); err != nil { log.Error(ctx, "error writing response", "error", err) } })
router.GET("/segment/:id", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { id := p.ByName("id") if id == "" { errors.WriteHTTPBadRequest(w, "id required", nil) return } segment, err := a.LocalDB.GetSegment(id) if err != nil { errors.WriteHTTPBadRequest(w, err.Error(), err) return } if segment == nil { errors.WriteHTTPNotFound(w, "segment not found", nil) return } spSeg, err := segment.ToStreamplaceSegment() if err != nil { errors.WriteHTTPInternalServerError(w, "unable to convert segment to streamplace segment", err) return } bs, err := json.Marshal(spSeg) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to marhsal json", err) return } if _, err := w.Write(bs); err != nil { log.Error(ctx, "error writing response", "error", err) } })
router.DELETE("/player-events", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { err := a.Model.ClearPlayerEvents() if err != nil { errors.WriteHTTPInternalServerError(w, "unable to delete player events", err) return } w.WriteHeader(204) })
router.GET("/followers/:user", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { user := p.ByName("user") if user == "" { errors.WriteHTTPBadRequest(w, "user required", nil) return }
followers, err := a.Model.GetUserFollowers(ctx, user) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to get followers", err) return } bs, err := json.Marshal(followers) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to marshal json", err) return } if _, err := w.Write(bs); err != nil { log.Error(ctx, "error writing response", "error", err) } })
router.GET("/following/:user", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { user := p.ByName("user") if user == "" { errors.WriteHTTPBadRequest(w, "user required", nil) return }
followers, err := a.Model.GetUserFollowing(ctx, user) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to get followers", err) return } bs, err := json.Marshal(followers) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to marshal json", err) return } if _, err := w.Write(bs); err != nil { log.Error(ctx, "error writing response", "error", err) } })
router.GET("/notifications", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { notifications, err := a.StatefulDB.ListNotifications() if err != nil { errors.WriteHTTPInternalServerError(w, "unable to get notifications", err) return } bs, err := json.Marshal(notifications) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to marshal json", err) return } if _, err := w.Write(bs); err != nil { log.Error(ctx, "error writing response", "error", err) } })
router.GET("/chat-posts", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { posts, err := a.Model.ListFeedPosts() if err != nil { errors.WriteHTTPInternalServerError(w, "unable to get chat posts", err) return } bs, err := json.Marshal(posts) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to marshal json", err) return } if _, err := w.Write(bs); err != nil { log.Error(ctx, "error writing response", "error", err) } })
router.GET("/chat/:uri", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { uri := p.ByName("uri") if uri == "" { errors.WriteHTTPBadRequest(w, "uri required", nil) return } msg, err := a.Model.GetChatMessage(uri) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to get chat posts", err) return } spmsg, err := msg.ToStreamplaceMessageView() if err != nil { errors.WriteHTTPInternalServerError(w, "unable to convert chat message to streamplace message view", err) return } if spmsg.Author.Handle == "" || spmsg.Author.Handle == "handle.invalid" { spmsg.Author.Handle = a.ATSync.ResolveAuthorHandle(ctx, spmsg.Author.Did) } bs, err := json.Marshal(spmsg) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to marshal json", err) return } if _, err := w.Write(bs); err != nil { log.Error(ctx, "error writing response", "error", err) } })
router.GET("/oauth-sessions", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { sessions, err := a.StatefulDB.ListOAuthSessions() if err != nil { errors.WriteHTTPInternalServerError(w, "unable to get oauth sessions", err) return } bs, err := json.Marshal(sessions) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to marshal oauth sessions", err) return } if _, err := w.Write(bs); err != nil { log.Error(ctx, "error writing response", "error", err) } })
router.POST("/notification-blast", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { var payload notificationpkg.NotificationBlast if err := json.NewDecoder(r.Body).Decode(&payload); err != nil { errors.WriteHTTPBadRequest(w, "invalid request body", err) return } notifications, err := a.StatefulDB.ListNotifications() if err != nil { errors.WriteHTTPInternalServerError(w, "unable to get notifications", err) return } if a.Notifier == nil { errors.WriteHTTPInternalServerError(w, "notifier not initialized", nil) return } targets := make([]notificationpkg.NotificationTarget, len(notifications)) for i, not := range notifications { targets[i] = notificationpkg.NotificationTarget{Token: not.Token, Type: not.Type} } err = a.Notifier.Blast(ctx, targets, &payload) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to blast notifications", err) return } w.WriteHeader(http.StatusNoContent) })
router.PUT("/settings/:id", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { w.Header().Set("Access-Control-Allow-Origin", "*") w.Header().Set("Access-Control-Allow-Methods", "PUT") w.Header().Set("Access-Control-Allow-Headers", "Content-Type")
id := p.ByName("id") if id == "" { errors.WriteHTTPBadRequest(w, "id required", nil) return }
var ident model.Identity if err := json.NewDecoder(r.Body).Decode(&ident); err != nil { errors.WriteHTTPBadRequest(w, "invalid request body", err) return } ident.ID = id
if err := a.Model.UpdateIdentity(&ident); err != nil { errors.WriteHTTPInternalServerError(w, "unable to update settings", err) return }
w.WriteHeader(http.StatusNoContent) })
router.POST("/replay/:streamKey", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { reqCtx := r.Context() key := p.ByName("streamKey") if key == "" { errors.WriteHTTPBadRequest(w, "streamKey required", nil) return } mediaSigner, err := a.MakeMediaSigner(reqCtx, key) if err != nil { errors.WriteHTTPUnauthorized(w, "invalid authorization key", err) return } pc, err := rtcrec.NewReplayPeerConnection(reqCtx, r.Body) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to create replay peer connection", err) return } answer, err := a.MediaManager.WebRTCIngest(reqCtx, &webrtc.SessionDescription{SDP: "placeholder"}, mediaSigner, pc, make(chan error, 1)) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to ingest web rtc", err) return } w.WriteHeader(200) if _, err := w.Write([]byte(answer.SDP)); err != nil { errors.WriteHTTPInternalServerError(w, "unable to write response", err) log.Error(reqCtx, "error writing response", "error", err) } })
router.GET("/clip/:did/clip.mp4", func(w http.ResponseWriter, r *http.Request, p httprouter.Params) { did := p.ByName("did") if did == "" { errors.WriteHTTPBadRequest(w, "did required", nil) return } user, err := a.NormalizeUser(ctx, did) if err != nil { errors.WriteHTTPBadRequest(w, "invalid user", err) return } secsStr := r.URL.Query().Get("secs") secs := 60 // Default to 60 seconds if secsStr != "" { parsedSecs, err := strconv.Atoi(secsStr) if err != nil { errors.WriteHTTPBadRequest(w, "invalid secs parameter", err) return } secs = parsedSecs } after := time.Now().Add(-time.Duration(secs) * time.Second) w.Header().Set("Content-Type", "video/mp4") err = a.MediaManager.ClipUser(ctx, user, w, nil, &after) if err != nil { errors.WriteHTTPInternalServerError(w, "unable to clip user", err) return } })
handler := sloghttp.Recovery(router) if log.Level(4) { handler = sloghttp.New(slog.Default())(handler) } return handler, nil}
func (a *StreamplaceAPI) keyToUser(ctx context.Context, key string) (string, error) { payload, err := base64.URLEncoding.DecodeString(key) if err != nil { return "", err } signed, err := a.Signer.Verify(payload) if err != nil { return "", err } _, ok := signed.Data().(*v0.StreamKey) if !ok { return "", fmt.Errorf("got signed data but it wasn't a stream key") } return strings.ToLower(signed.Signer()), nil}