stremio/node/media-node/cmd/media-node/main.go

485 lines
13 KiB
Go

package main
import (
"bytes"
"encoding/json"
"fmt"
"io"
"log"
"net"
"net/http"
"os"
"os/exec"
"os/signal"
"path/filepath"
"runtime"
"strings"
"syscall"
"time"
"github.com/sthmedia/media-node/internal/config"
"github.com/sthmedia/media-node/internal/control"
"github.com/sthmedia/media-node/internal/database"
"github.com/sthmedia/media-node/internal/scanner"
"github.com/sthmedia/media-node/internal/streaming"
)
var version = "1.3.12"
func main() {
if len(os.Args) < 2 {
printUsage()
os.Exit(1)
}
switch os.Args[1] {
case "install":
runInstall()
case "uninstall":
runUninstall()
case "register":
runRegister()
case "start":
runSystemctl("start")
case "stop":
runSystemctl("stop")
case "restart":
runSystemctl("restart")
case "status":
runSystemctl("status")
case "scan":
runScan()
case "config":
runConfig()
case "diagnostics":
runDiagnostics()
case "version":
fmt.Printf("media-node %s (%s/%s)\n", version, runtime.GOOS, runtime.GOARCH)
case "run":
runService()
default:
runService()
}
}
func printUsage() {
fmt.Println(`media-node - Distributed Stremio Media Node
Usage:
media-node install Interactive installer + systemd
media-node uninstall Remove systemd service
media-node register Register with Master using enrollment token
media-node start|stop|restart|status
media-node run Run the node service (foreground)
media-node scan Trigger a full media scan
media-node config Show current configuration
media-node diagnostics Run connectivity diagnostics
media-node version Show version`)
}
func runService() {
configPath := envOr("MEDIA_NODE_CONFIG", config.DefaultConfigPath)
dataDir := envOr("MEDIA_NODE_DATA", config.DefaultDataDir)
credsPath := envOr("MEDIA_NODE_CREDENTIALS", config.DefaultCredentialsPath)
ensureDirs(filepath.Dir(configPath), dataDir, filepath.Dir(credsPath))
cfg, err := config.Load(configPath)
if err != nil {
log.Fatalf("Failed to load config: %v", err)
}
creds, err := loadCredentials(credsPath)
if err != nil {
token := strings.TrimSpace(os.Getenv("MEDIA_NODE_ENROLLMENT_TOKEN"))
if token == "" {
log.Fatalf("Failed to load credentials: %v (set MEDIA_NODE_ENROLLMENT_TOKEN for first boot)", err)
}
log.Printf("No credentials at %s — enrolling with Master...", credsPath)
creds, err = enroll(cfg.Master.URL, token, cfg.Node.Name, cfg.Stream.PublicURL)
if err != nil {
log.Fatalf("Enrollment failed: %v", err)
}
data, _ := json.MarshalIndent(creds, "", " ")
if writeErr := os.WriteFile(credsPath, data, 0o600); writeErr != nil {
log.Fatalf("Failed to save credentials: %v", writeErr)
}
log.Printf("Enrolled as node %s — credentials saved", creds.NodeID)
cfg.Node.ID = creds.NodeID
_ = config.Save(configPath, cfg)
}
store, err := database.Open(dataDir)
if err != nil {
log.Fatalf("Failed to open database: %v", err)
}
defer store.Close()
streamer := streaming.New(store, cfg.Stream.Listen, cfg.Security.MaxConcurrentPerIP)
var ctrlClient *control.Client
sc := scanner.New(cfg, store, func(events []scanner.LibraryEvent) {
if ctrlClient != nil {
ctrlClient.SendLibraryEvents(events)
}
})
ctrlClient = control.NewClient(cfg, creds, store, sc, streamer, version, configPath)
sc.SetOnScanComplete(func() {
ctrlClient.RequestFullLibrarySync()
})
streamer.SetSessionEndHandler(func(sessionID, reason string) {
ctrlClient.NotifySessionEnded(sessionID, reason)
})
streamer.SetSessionProgressHandler(func(sessionID string, bytesSent, lastByteOffset, fileSizeBytes int64) {
ctrlClient.NotifySessionProgress(sessionID, bytesSent, lastByteOffset, fileSizeBytes)
})
go sc.Start()
go ctrlClient.Run()
go func() {
if err := streamer.ListenAndServe(); err != nil {
log.Fatalf("Stream server error: %v", err)
}
}()
log.Printf("media-node %s started as %s", version, cfg.Node.Name)
sig := make(chan os.Signal, 1)
signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM)
<-sig
log.Println("Shutting down...")
ctrlClient.Stop()
}
func runInstall() {
fmt.Println("=== Media Node Installer ===")
fmt.Println()
masterURL := prompt("Master URL (https://master.media.example.com): ")
nodeName := prompt("Node name: ")
publicURL := prompt("Public stream URL (https://node01.media.example.com): ")
moviesDir := prompt("Movies directory: ")
seriesDir := prompt("Series directory (optional): ")
token := prompt("Enrollment token: ")
cfg := &config.Config{
Node: config.NodeConfig{Name: nodeName},
Master: config.MasterConfig{URL: strings.TrimRight(masterURL, "/")},
Stream: config.StreamConfig{
Listen: "0.0.0.0:8080",
PublicURL: strings.TrimRight(publicURL, "/"),
},
Media: config.MediaConfig{
Movies: nonEmpty(moviesDir),
Series: nonEmpty(seriesDir),
},
Scanner: config.ScannerConfig{FullScanInterval: "off", FullScanAt: "03:30"},
Security: config.SecurityConfig{SessionIdleTimeout: "20m", MaxConcurrentPerIP: 64},
}
configDir := "/etc/media-node"
dataDir := "/var/lib/media-node"
credsDir := filepath.Join(configDir, "credentials")
if err := os.MkdirAll(credsDir, 0o750); err != nil {
fmt.Printf("Warning: system paths unavailable (%v), using ./config\n", err)
configDir = "./config"
dataDir = "./data"
credsDir = filepath.Join(configDir, "credentials")
_ = os.MkdirAll(credsDir, 0o750)
_ = os.MkdirAll(dataDir, 0o750)
} else {
_ = os.MkdirAll(dataDir, 0o750)
}
configPath := filepath.Join(configDir, "config.yaml")
if err := config.Save(configPath, cfg); err != nil {
log.Fatalf("Failed to save config: %v", err)
}
fmt.Printf("Config saved to %s\n", configPath)
creds, err := enroll(cfg.Master.URL, token, nodeName, cfg.Stream.PublicURL)
if err != nil {
log.Fatalf("Enrollment failed: %v", err)
}
credsPath := filepath.Join(credsDir, "credentials.json")
data, _ := json.MarshalIndent(creds, "", " ")
if err := os.WriteFile(credsPath, data, 0o600); err != nil {
log.Fatalf("Failed to save credentials: %v", err)
}
fmt.Printf("Credentials saved to %s\n", credsPath)
if runtime.GOOS == "linux" && os.Geteuid() == 0 {
installSystemd()
fmt.Println("\nInstallation complete. Service started.")
fmt.Println(" systemctl status media-node")
} else {
fmt.Println("\nInstallation complete! Run: media-node run")
}
}
func installSystemd() {
exe, err := os.Executable()
if err != nil {
log.Printf("Could not determine binary path: %v", err)
return
}
target := "/usr/local/bin/media-node"
if exe != target {
in, err := os.ReadFile(exe)
if err == nil {
_ = os.WriteFile(target, in, 0o755)
}
}
_ = exec.Command("useradd", "--system", "--no-create-home", "--shell", "/usr/sbin/nologin", "media-node").Run()
_ = exec.Command("chown", "-R", "media-node:media-node", "/var/lib/media-node").Run()
_ = exec.Command("chown", "-R", "root:media-node", "/etc/media-node").Run()
unit := `[Unit]
Description=Media Node
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User=media-node
Group=media-node
ExecStart=/usr/local/bin/media-node run
Restart=always
RestartSec=5
Environment=MEDIA_NODE_CONFIG=/etc/media-node/config.yaml
Environment=MEDIA_NODE_DATA=/var/lib/media-node
Environment=MEDIA_NODE_CREDENTIALS=/etc/media-node/credentials/credentials.json
NoNewPrivileges=true
ProtectSystem=strict
ProtectHome=true
ReadWritePaths=/var/lib/media-node
PrivateTmp=true
[Install]
WantedBy=multi-user.target
`
_ = os.WriteFile("/etc/systemd/system/media-node.service", []byte(unit), 0o644)
_ = exec.Command("systemctl", "daemon-reload").Run()
_ = exec.Command("systemctl", "enable", "--now", "media-node").Run()
}
func runSystemctl(action string) {
if runtime.GOOS != "linux" {
fmt.Println("systemctl commands are only available on Linux")
return
}
cmd := exec.Command("systemctl", action, "media-node")
cmd.Stdout = os.Stdout
cmd.Stderr = os.Stderr
_ = cmd.Run()
}
func runRegister() {
configPath := envOr("MEDIA_NODE_CONFIG", config.DefaultConfigPath)
credsDir := envOr("MEDIA_NODE_CREDENTIALS_DIR", "/etc/media-node/credentials")
cfg, err := config.Load(configPath)
if err != nil {
log.Fatalf("Failed to load config: %v", err)
}
token := prompt("Enrollment token: ")
creds, err := enroll(cfg.Master.URL, token, cfg.Node.Name, cfg.Stream.PublicURL)
if err != nil {
log.Fatalf("Enrollment failed: %v", err)
}
_ = os.MkdirAll(credsDir, 0o750)
credsPath := filepath.Join(credsDir, "credentials.json")
data, _ := json.MarshalIndent(creds, "", " ")
if err := os.WriteFile(credsPath, data, 0o600); err != nil {
log.Fatalf("Failed to save credentials: %v", err)
}
fmt.Println("Registration successful!")
}
func runScan() {
configPath := envOr("MEDIA_NODE_CONFIG", config.DefaultConfigPath)
dataDir := envOr("MEDIA_NODE_DATA", config.DefaultDataDir)
cfg, err := config.Load(configPath)
if err != nil {
log.Fatalf("Failed to load config: %v", err)
}
store, err := database.Open(dataDir)
if err != nil {
log.Fatalf("Failed to open database: %v", err)
}
defer store.Close()
sc := scanner.New(cfg, store, func(events []scanner.LibraryEvent) {
fmt.Printf("Event batch: %d\n", len(events))
})
sc.FullScan()
fmt.Println("Scan complete")
}
func runConfig() {
configPath := envOr("MEDIA_NODE_CONFIG", config.DefaultConfigPath)
cfg, err := config.Load(configPath)
if err != nil {
log.Fatalf("Failed to load config: %v", err)
}
data, _ := json.MarshalIndent(cfg, "", " ")
fmt.Println(string(data))
}
func runDiagnostics() {
configPath := envOr("MEDIA_NODE_CONFIG", config.DefaultConfigPath)
cfg, err := config.Load(configPath)
if err != nil {
fmt.Printf("✗ Config: %v\n", err)
return
}
fmt.Printf("✓ Config loaded from %s\n", configPath)
u := strings.TrimPrefix(strings.TrimPrefix(cfg.Master.URL, "https://"), "http://")
host := strings.Split(u, "/")[0]
if _, err := net.LookupHost(strings.Split(host, ":")[0]); err != nil {
fmt.Printf("✗ Master DNS: %v\n", err)
} else {
fmt.Printf("✓ Master DNS: %s\n", host)
}
client := &http.Client{Timeout: 10 * time.Second}
resp, err := client.Get(cfg.Master.URL + "/health")
if err != nil {
fmt.Printf("✗ Master HTTPS: %v\n", err)
} else {
resp.Body.Close()
fmt.Printf("✓ Master HTTPS: %s (HTTP %d)\n", cfg.Master.URL, resp.StatusCode)
}
for _, dir := range append(cfg.Media.Movies, cfg.Media.Series...) {
if dir == "" {
continue
}
info, err := os.Stat(dir)
if err != nil {
fmt.Printf("✗ Media directory %s: %v\n", dir, err)
} else if !info.IsDir() {
fmt.Printf("✗ Media path not a directory: %s\n", dir)
} else {
fmt.Printf("✓ Media directory accessible: %s\n", dir)
}
}
ln, err := net.Listen("tcp", cfg.Stream.Listen)
if err != nil {
fmt.Printf("✗ Stream listen %s: %v\n", cfg.Stream.Listen, err)
} else {
_ = ln.Close()
fmt.Printf("✓ Stream listen port available: %s\n", cfg.Stream.Listen)
}
dataDir := envOr("MEDIA_NODE_DATA", config.DefaultDataDir)
if store, err := database.Open(dataDir); err != nil {
fmt.Printf("✗ Node database: %v\n", err)
} else {
_ = store.Close()
fmt.Printf("✓ Node database: %s\n", filepath.Join(dataDir, "node.db"))
}
fmt.Printf("✓ System clock: %s\n", time.Now().UTC().Format(time.RFC3339))
fmt.Printf("✓ Version: %s (%s/%s)\n", version, runtime.GOOS, runtime.GOARCH)
}
func runUninstall() {
if runtime.GOOS == "linux" && os.Geteuid() == 0 {
_ = exec.Command("systemctl", "stop", "media-node").Run()
_ = exec.Command("systemctl", "disable", "media-node").Run()
_ = os.Remove("/etc/systemd/system/media-node.service")
_ = exec.Command("systemctl", "daemon-reload").Run()
fmt.Println("Service removed. Config/data left at /etc/media-node and /var/lib/media-node")
return
}
fmt.Println("To uninstall:")
fmt.Println(" sudo systemctl stop media-node")
fmt.Println(" sudo systemctl disable media-node")
fmt.Println(" sudo rm /etc/systemd/system/media-node.service")
fmt.Println(" sudo rm -rf /etc/media-node /var/lib/media-node")
}
func enroll(masterURL, token, name, publicURL string) (*config.Credentials, error) {
hostname, _ := os.Hostname()
body := map[string]string{
"token": token,
"name": name,
"hostname": hostname,
"architecture": runtime.GOARCH,
"version": version,
"publicStreamUrl": publicURL,
}
data, _ := json.Marshal(body)
resp, err := http.Post(masterURL+"/api/v1/nodes/enroll", "application/json", bytes.NewReader(data))
if err != nil {
return nil, err
}
defer resp.Body.Close()
respData, _ := io.ReadAll(resp.Body)
if resp.StatusCode != 200 {
return nil, fmt.Errorf("enrollment failed (HTTP %d): %s", resp.StatusCode, string(respData))
}
var result config.Credentials
if err := json.Unmarshal(respData, &result); err != nil {
return nil, err
}
return &result, nil
}
func loadCredentials(path string) (*config.Credentials, error) {
data, err := os.ReadFile(path)
if err != nil {
return nil, err
}
var creds config.Credentials
if err := json.Unmarshal(data, &creds); err != nil {
return nil, err
}
return &creds, nil
}
func prompt(label string) string {
fmt.Print(label)
var val string
_, _ = fmt.Scanln(&val)
return strings.TrimSpace(val)
}
func envOr(key, fallback string) string {
if v := os.Getenv(key); v != "" {
return v
}
return fallback
}
func nonEmpty(s string) []string {
s = strings.TrimSpace(s)
if s == "" {
return nil
}
return []string{s}
}
func ensureDirs(paths ...string) {
for _, p := range paths {
if p == "" || p == "." {
continue
}
_ = os.MkdirAll(p, 0o750)
}
}