Capture, decrypt, and restream toolkit with compiled-in app modules (RTE, TG4, BBC), on-device MITM proxy, streamd control plane, and www. BBC module.yaml is published (clear streams); other module values stay local.
245 lines
6.3 KiB
Go
245 lines
6.3 KiB
Go
// Package servecmd is the streamd control plane, exposed as a library so the
|
|
// single drm binary can host it as a subcommand.
|
|
package servecmd
|
|
|
|
import (
|
|
"flag"
|
|
"fmt"
|
|
"log"
|
|
"net/http"
|
|
"os"
|
|
"os/exec"
|
|
"os/signal"
|
|
"path/filepath"
|
|
"strings"
|
|
"syscall"
|
|
"time"
|
|
|
|
"drmdecryption/apps/streamd/internal/api"
|
|
"drmdecryption/apps/streamd/internal/db"
|
|
"drmdecryption/apps/streamd/internal/supervisor"
|
|
"drmdecryption/apps/streamd/internal/ui"
|
|
"drmdecryption/apps/streamd/internal/ws"
|
|
)
|
|
|
|
// Run dispatches a streamd subcommand. args[0] is the subcommand name.
|
|
func Run(args []string) error {
|
|
if len(args) < 1 {
|
|
Usage()
|
|
return fmt.Errorf("streamd: subcommand required")
|
|
}
|
|
sub, rest := args[0], args[1:]
|
|
switch sub {
|
|
case "serve":
|
|
serveCmd(rest)
|
|
case "worker":
|
|
return fmt.Errorf("worker: media is supervised in-process by `serve` now")
|
|
case "help", "-h", "--help":
|
|
Usage()
|
|
default:
|
|
Usage()
|
|
return fmt.Errorf("streamd: unknown subcommand %q", sub)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Usage prints the streamd subcommand help.
|
|
func Usage() {
|
|
fmt.Fprintf(os.Stderr, `streamd — Go control plane for DRM restreams
|
|
|
|
Usage:
|
|
streamd serve [flags]
|
|
|
|
serve flags:
|
|
--bind listen address (default 0.0.0.0:8083)
|
|
--data data dir for sqlite + www + work (default .cache/streamd)
|
|
--token shared bearer token for mutating API (or STREAMD_TOKEN)
|
|
--nre path to N_m3u8DL-RE (else NRE_PATH, bin/, PATH)
|
|
--ffmpeg path to ffmpeg (else FFMPEG_PATH, bin/, PATH)
|
|
--mp4decrypt path to mp4decrypt (else MP4DECRYPT_PATH, bin/, PATH)
|
|
`)
|
|
}
|
|
|
|
func serveCmd(args []string) {
|
|
fs := flag.NewFlagSet("serve", flag.ExitOnError)
|
|
bind := fs.String("bind", "0.0.0.0:8083", "listen address")
|
|
data := fs.String("data", "", "data directory (sqlite, www, work, logs)")
|
|
token := fs.String("token", "", "API bearer token (default: STREAMD_TOKEN env)")
|
|
nre := fs.String("nre", "", "N_m3u8DL-RE binary")
|
|
ffmpeg := fs.String("ffmpeg", "", "ffmpeg binary")
|
|
mp4decrypt := fs.String("mp4decrypt", "", "mp4decrypt binary")
|
|
_ = fs.Parse(args)
|
|
|
|
dataDir := *data
|
|
if dataDir == "" {
|
|
dataDir = filepath.Join(".", ".cache", "streamd")
|
|
}
|
|
tok := *token
|
|
if tok == "" {
|
|
tok = os.Getenv("STREAMD_TOKEN")
|
|
}
|
|
|
|
for _, sub := range []string{"www", "work", "logs"} {
|
|
if err := os.MkdirAll(filepath.Join(dataDir, sub), 0o755); err != nil {
|
|
log.Fatal(err)
|
|
}
|
|
}
|
|
|
|
store, err := db.Open(dataDir)
|
|
if err != nil {
|
|
log.Fatalf("db: %v", err)
|
|
}
|
|
defer store.Close()
|
|
|
|
nrePath := firstExisting(*nre,
|
|
os.Getenv("NRE_PATH"),
|
|
os.Getenv("N_M3U8DL_RE"),
|
|
filepath.Join("bin", "N_m3u8DL-RE.exe"),
|
|
"N_m3u8DL-RE",
|
|
)
|
|
ffPath := firstExisting(*ffmpeg,
|
|
os.Getenv("FFMPEG_PATH"),
|
|
os.Getenv("FFMPEG"),
|
|
filepath.Join("bin", "ffmpeg.exe"),
|
|
"ffmpeg",
|
|
)
|
|
decPath := firstExisting(*mp4decrypt,
|
|
os.Getenv("MP4DECRYPT_PATH"),
|
|
os.Getenv("MP4DECRYPT"),
|
|
filepath.Join("bin", "mp4decrypt.exe"),
|
|
"mp4decrypt",
|
|
)
|
|
|
|
sup := supervisor.NewMedia(supervisor.Options{
|
|
Store: store,
|
|
DataDir: dataDir,
|
|
NRE: nrePath,
|
|
FFmpeg: ffPath,
|
|
MP4Decrypt: decPath,
|
|
Log: log.Default(),
|
|
})
|
|
defer sup.Close()
|
|
|
|
apiH := &api.Handler{
|
|
Store: store,
|
|
Token: tok,
|
|
Supervisor: sup,
|
|
StartedAt: time.Now(),
|
|
}
|
|
hub := ws.NewHub(store, tok, log.Default())
|
|
|
|
mux := http.NewServeMux()
|
|
apiH.Mount(mux)
|
|
hub.Mount(mux)
|
|
mux.Handle("/", ui.Handler())
|
|
mux.Handle("/static/", http.StripPrefix("/static/", ui.Static()))
|
|
mux.Handle("/hls/", http.StripPrefix("/hls/", http.FileServer(http.Dir(filepath.Join(dataDir, "www")))))
|
|
|
|
srv := &http.Server{
|
|
Addr: *bind,
|
|
Handler: withCORS(mux),
|
|
ReadHeaderTimeout: 10 * time.Second,
|
|
}
|
|
|
|
log.Printf("streamd listening on http://%s data=%s", *bind, dataDir)
|
|
log.Printf("media: nre=%s", nrePath)
|
|
log.Printf("media: ffmpeg=%s", ffPath)
|
|
log.Printf("media: mp4decrypt=%s", decPath)
|
|
log.Printf("ws: agents=/ws/agent ui=/ws/ui")
|
|
if tok == "" {
|
|
log.Printf("warning: no --token / STREAMD_TOKEN set; mutating API is open")
|
|
}
|
|
|
|
stopBG := make(chan struct{})
|
|
go hub.RunExpireLoop(stopBG)
|
|
|
|
// Resume enabled streams that already have credentials (staggered).
|
|
go func() {
|
|
time.Sleep(500 * time.Millisecond)
|
|
list, err := store.ListStreams()
|
|
if err != nil {
|
|
return
|
|
}
|
|
for _, st := range list {
|
|
if !st.Enabled || strings.TrimSpace(st.MPD) == "" || strings.TrimSpace(st.Key) == "" {
|
|
continue
|
|
}
|
|
mpd := strings.ToLower(strings.TrimSpace(st.MPD))
|
|
if !strings.HasPrefix(mpd, "http://") && !strings.HasPrefix(mpd, "https://") {
|
|
log.Printf("auto-start skip %s: mpd is not an http(s) URL", st.Name)
|
|
continue
|
|
}
|
|
log.Printf("auto-start enabled stream %s", st.Name)
|
|
func(id int64, name string) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
log.Printf("auto-start %s panic: %v", name, r)
|
|
}
|
|
}()
|
|
if err := sup.Start(id); err != nil {
|
|
log.Printf("auto-start %s: %v", name, err)
|
|
}
|
|
}(st.ID, st.Name)
|
|
time.Sleep(1500 * time.Millisecond)
|
|
}
|
|
}()
|
|
|
|
go func() {
|
|
ch := make(chan os.Signal, 1)
|
|
signal.Notify(ch, os.Interrupt, syscall.SIGTERM)
|
|
<-ch
|
|
log.Printf("shutting down…")
|
|
close(stopBG)
|
|
sup.Close()
|
|
_ = srv.Close()
|
|
}()
|
|
|
|
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
|
log.Fatal(err)
|
|
}
|
|
}
|
|
|
|
func withCORS(next http.Handler) http.Handler {
|
|
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
w.Header().Set("Access-Control-Allow-Origin", "*")
|
|
w.Header().Set("Access-Control-Allow-Headers", "Authorization, Content-Type, X-Streamd-Token")
|
|
w.Header().Set("Access-Control-Allow-Methods", "GET, POST, PUT, PATCH, DELETE, OPTIONS")
|
|
if r.Method == http.MethodOptions {
|
|
w.WriteHeader(http.StatusNoContent)
|
|
return
|
|
}
|
|
next.ServeHTTP(w, r)
|
|
})
|
|
}
|
|
|
|
func firstExisting(explicit string, candidates ...string) string {
|
|
if explicit != "" {
|
|
if st, err := os.Stat(explicit); err == nil && !st.IsDir() {
|
|
return explicit
|
|
}
|
|
// still return explicit so exec fails clearly if user set a bad path
|
|
if filepath.IsAbs(explicit) || containsSep(explicit) {
|
|
return explicit
|
|
}
|
|
}
|
|
for _, c := range candidates {
|
|
if c == "" {
|
|
continue
|
|
}
|
|
if st, err := os.Stat(c); err == nil && !st.IsDir() {
|
|
return c
|
|
}
|
|
if p, err := execLookPath(c); err == nil {
|
|
return p
|
|
}
|
|
}
|
|
return explicit
|
|
}
|
|
|
|
func containsSep(s string) bool {
|
|
return filepath.Base(s) != s
|
|
}
|
|
|
|
func execLookPath(file string) (string, error) {
|
|
return exec.LookPath(file)
|
|
}
|