package api import ( "encoding/json" "errors" "io" "net/http" "strconv" "strings" "time" "drmdecryption/apps/streamd/internal/db" ) type Supervisor interface { Start(id int64) error Stop(id int64) error Restart(id int64) error } type Handler struct { Store *db.Store Token string Supervisor Supervisor StartedAt time.Time } func (h *Handler) Mount(mux *http.ServeMux) { mux.HandleFunc("/api/health", h.health) mux.HandleFunc("/api/streams", h.streams) mux.HandleFunc("/api/streams/", h.streamAction) mux.HandleFunc("/api/apps", h.apps) } func (h *Handler) auth(w http.ResponseWriter, r *http.Request) bool { if h.Token == "" { return true } if r.Method == http.MethodGet { return true } auth := r.Header.Get("Authorization") tok := strings.TrimPrefix(auth, "Bearer ") if tok == "" { tok = r.Header.Get("X-Streamd-Token") } if tok == "" { tok = r.URL.Query().Get("token") } if tok != h.Token { writeErr(w, http.StatusUnauthorized, "unauthorized") return false } return true } func (h *Handler) health(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { writeErr(w, http.StatusMethodNotAllowed, "method not allowed") return } writeJSON(w, http.StatusOK, map[string]any{ "ok": true, "service": "streamd", "uptime_s": time.Since(h.StartedAt).Seconds(), "started_at": h.StartedAt.UTC().Format(time.RFC3339), }) } func (h *Handler) streams(w http.ResponseWriter, r *http.Request) { switch r.Method { case http.MethodGet: list, err := h.Store.ListStreams() if err != nil { writeErr(w, http.StatusInternalServerError, err.Error()) return } if list == nil { list = []db.Stream{} } writeJSON(w, http.StatusOK, map[string]any{"streams": list}) case http.MethodPost: if !h.auth(w, r) { return } var in db.CreateStream if err := readJSON(r, &in); err != nil { writeErr(w, http.StatusBadRequest, err.Error()) return } st, err := h.Store.CreateStream(in) if err != nil { writeErr(w, http.StatusBadRequest, err.Error()) return } writeJSON(w, http.StatusCreated, st) default: writeErr(w, http.StatusMethodNotAllowed, "method not allowed") } } func (h *Handler) streamAction(w http.ResponseWriter, r *http.Request) { path := strings.TrimPrefix(r.URL.Path, "/api/streams/") path = strings.Trim(path, "/") if path == "" { h.streams(w, r) return } parts := strings.Split(path, "/") id, err := strconv.ParseInt(parts[0], 10, 64) if err != nil { writeErr(w, http.StatusBadRequest, "invalid id") return } action := "" if len(parts) > 1 { action = parts[1] } sub := "" if len(parts) > 2 { sub = parts[2] } switch { case action == "" && r.Method == http.MethodGet: st, err := h.Store.GetStream(id) if err != nil { writeErr(w, http.StatusNotFound, "not found") return } writeJSON(w, http.StatusOK, st) case action == "" && (r.Method == http.MethodPatch || r.Method == http.MethodPut): if !h.auth(w, r) { return } var in db.PatchStream if err := readJSON(r, &in); err != nil { writeErr(w, http.StatusBadRequest, err.Error()) return } before, _ := h.Store.GetStream(id) st, err := h.Store.PatchStream(id, in) if err != nil { writeErr(w, http.StatusBadRequest, err.Error()) return } // Editing mpd/key/headers must tear down the live media worker and // bring up a fresh instance with the new credentials. if h.Supervisor != nil && credentialsChanged(before, st, in) { if st.Enabled { _ = h.Supervisor.Restart(id) } else { _ = h.Supervisor.Stop(id) } st, _ = h.Store.GetStream(id) } writeJSON(w, http.StatusOK, st) case action == "" && r.Method == http.MethodDelete: if !h.auth(w, r) { return } if h.Supervisor != nil { _ = h.Supervisor.Stop(id) } if err := h.Store.DeleteStream(id); err != nil { writeErr(w, http.StatusInternalServerError, err.Error()) return } writeJSON(w, http.StatusOK, map[string]any{"ok": true}) case action == "start" && r.Method == http.MethodPost: if !h.auth(w, r) { return } st, err := h.Store.SetEnabled(id, true) if err != nil { writeErr(w, http.StatusNotFound, err.Error()) return } if h.Supervisor != nil { if err := h.Supervisor.Start(id); err != nil { writeErr(w, http.StatusInternalServerError, err.Error()) return } st, _ = h.Store.GetStream(id) } writeJSON(w, http.StatusOK, st) case action == "stop" && r.Method == http.MethodPost: if !h.auth(w, r) { return } if h.Supervisor != nil { _ = h.Supervisor.Stop(id) } st, err := h.Store.SetEnabled(id, false) if err != nil { writeErr(w, http.StatusNotFound, err.Error()) return } writeJSON(w, http.StatusOK, st) case action == "restart" && r.Method == http.MethodPost: if !h.auth(w, r) { return } st, err := h.Store.SetEnabled(id, true) if err != nil { writeErr(w, http.StatusNotFound, err.Error()) return } if h.Supervisor != nil { if err := h.Supervisor.Restart(id); err != nil { writeErr(w, http.StatusInternalServerError, err.Error()) return } st, _ = h.Store.GetStream(id) } writeJSON(w, http.StatusOK, st) case action == "credentials" && r.Method == http.MethodPost: if !h.auth(w, r) { return } var cred db.Credentials if err := readJSON(r, &cred); err != nil { writeErr(w, http.StatusBadRequest, err.Error()) return } st, err := h.Store.SetCredentials(id, cred) if err != nil { writeErr(w, http.StatusBadRequest, err.Error()) return } _ = h.Store.ReleaseClaim(id, "") // Always kill the old media worker; start a fresh one when enabled. if h.Supervisor != nil { if st.Enabled { _ = h.Supervisor.Restart(id) } else { _ = h.Supervisor.Stop(id) } st, _ = h.Store.GetStream(id) } writeJSON(w, http.StatusOK, st) case action == "claim" && sub == "" && r.Method == http.MethodPost: if !h.auth(w, r) { return } var body struct { AgentID string `json:"agent_id"` TTLSec int `json:"ttl_sec"` } _ = readJSON(r, &body) ttl := time.Duration(body.TTLSec) * time.Second st, err := h.Store.Claim(id, body.AgentID, ttl) if err != nil { if err.Error() == "claimed" { writeErr(w, http.StatusConflict, "already claimed") return } writeErr(w, http.StatusBadRequest, err.Error()) return } writeJSON(w, http.StatusOK, st) case action == "claim" && sub == "release" && r.Method == http.MethodPost: if !h.auth(w, r) { return } var body struct { AgentID string `json:"agent_id"` } _ = readJSON(r, &body) if err := h.Store.ReleaseClaim(id, body.AgentID); err != nil { writeErr(w, http.StatusInternalServerError, err.Error()) return } writeJSON(w, http.StatusOK, map[string]any{"ok": true}) default: writeErr(w, http.StatusNotFound, "not found") } } func readJSON(r *http.Request, dst any) error { defer r.Body.Close() b, err := io.ReadAll(io.LimitReader(r.Body, 1<<20)) if err != nil { return err } if len(b) == 0 { return errors.New("empty body") } return json.Unmarshal(b, dst) } func writeJSON(w http.ResponseWriter, code int, v any) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(code) enc := json.NewEncoder(w) enc.SetIndent("", " ") _ = enc.Encode(v) } func writeErr(w http.ResponseWriter, code int, msg string) { writeJSON(w, code, map[string]any{"error": msg}) } // credentialsChanged is true when the patch touched mpd, key, or headers that // the live NRE/ffmpeg worker is already using. func credentialsChanged(before, after db.Stream, in db.PatchStream) bool { if in.MPD != nil && strings.TrimSpace(before.MPD) != strings.TrimSpace(after.MPD) { return true } if in.Key != nil && strings.TrimSpace(before.Key) != strings.TrimSpace(after.Key) { return true } if in.HeadersJSON != nil && strings.TrimSpace(before.HeadersJSON) != strings.TrimSpace(after.HeadersJSON) { return true } return false }