From ad822138b69bd06620c51b43acd0ab21bd7bf6e4 Mon Sep 17 00:00:00 2001 From: hermes Date: Wed, 30 Sep 2026 21:59:30 +0200 Subject: [PATCH] add cmd/streamer/main.go --- cmd/streamer/main.go | 186 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 186 insertions(+) create mode 100644 cmd/streamer/main.go diff --git a/cmd/streamer/main.go b/cmd/streamer/main.go new file mode 100644 index 0000000..dc9182b --- /dev/null +++ b/cmd/streamer/main.go @@ -0,0 +1,186 @@ +// Command streamer runs on the streaming VPS. +// +// It reads the shipped playlist.m3u8 from the local web root, streams each +// MP3 into the icecast mount as one endless live stream, and publishes a +// now-playing JSON message to MQTT for every track (title display, logging). +// +// The generator machine can be asleep — the stream keeps rolling on whatever +// was last shipped, re-reading the playlist each lap so new shipments appear. +// +// Config: environment variables (see .env.example). +package main + +import ( + "bufio" + "encoding/json" + "log" + "os" + "path/filepath" + "strings" + "time" + + "radio-go/internal/infra" +) +func env(k, def string) string { + if v := os.Getenv(k); v != "" { + return v + } + return def +} + +// nowPlaying is the MQTT payload for the current track. +type nowPlaying struct { + Station string `json:"station"` + Title string `json:"title"` + Artist string `json:"artist"` + Track string `json:"track"` + Filename string `json:"file"` + StreamURL string `json:"stream_url"` + Timestamp int64 `json:"ts"` +} + +func main() { + webroot := env("STREAM_WEBROOT", "/usr/local/share/icecast/web/ai-radio") + iceAddr := env("ICECAST_ADDR", "icecast:8000") + icePass := env("ICECAST_SOURCE_PASSWORD", "hackme") + mount := env("ICECAST_MOUNT", "/ai-radio.mp3") + station := env("STATION_NAME", "AI Radio") + bitrate := parseInt(env("MP3_BITRATE", "128")) + streamURL := env("STREAM_PUBLIC_URL", "") + + var mqtt *infra.MQTTNowPlaying + if addr := env("MQTT_ADDR", ""); addr != "" { + mqtt = &infra.MQTTNowPlaying{ + Addr: addr, + User: os.Getenv("MQTT_USER"), + Pass: os.Getenv("MQTT_PASS"), + ClientID: env("MQTT_CLIENT_ID", "radio-streamer"), + Topic: env("MQTT_TOPIC", "airstudio/nowplaying"), + } + log.Printf("mqtt now-playing -> %s topic %s", addr, mqtt.Topic) + } + + src := &infra.IcecastSource{ + Addr: iceAddr, + Password: icePass, + Meta: infra.SourceMeta{ + Mount: mount, + Name: station, + Genre: "AI Generated", + Bitrate: env("MP3_BITRATE", "128"), + URL: streamURL, + Description: env("STATION_DESCRIPTION", "Endless AI-generated radio"), + }, + } + + for { + conn, err := src.Connect() + if err != nil { + log.Printf("icecast connect failed: %v — retry in 10s", err) + time.Sleep(10 * time.Second) + continue + } + log.Printf("live on %s%s", iceAddr, mount) + lap(webroot, conn, mqtt, station, bitrate, streamURL) + conn.Close() + } +} + +// loadTitles reads the titles.json sidecar shipped by the shipper. +func loadTitles(webroot string) map[string]string { + raw, err := os.ReadFile(filepath.Join(webroot, "titles.json")) + if err != nil { + return nil + } + var m map[string]string + if json.Unmarshal(raw, &m) != nil { + return nil + } + return m +} + +// lap plays one pass through the shipped playlist. +func lap(webroot string, conn interface{ Write([]byte) (int, error) }, mqtt *infra.MQTTNowPlaying, station string, bitrate int, streamURL string) { + files := readPlaylist(filepath.Join(webroot, "playlist.m3u8")) + if len(files) == 0 { + log.Printf("no playable files in %s — waiting 30s", webroot) + time.Sleep(30 * time.Second) + return + } + titles := loadTitles(webroot) + for _, f := range files { + fh, err := os.Open(filepath.Join(webroot, f)) + if err != nil { + log.Printf("skip %s: %v", f, err) + continue + } + title := titleFromName(f) + if t, ok := titles[strings.TrimSuffix(f, ".mp3")]; ok && t != "" { + title = t + } + publishNowPlaying(mqtt, station, title, f, streamURL) + sent, err := infra.StreamFile(conn, fh, bitrate) + fh.Close() + if err != nil { + log.Printf("stream dropped at %s: %v", f, err) + return + } + log.Printf("played %s (%.1f MB)", f, float64(sent)/1e6) + } +} + +// readPlaylist returns the mp3 filenames listed in an m3u8, in order. +func readPlaylist(path string) []string { + f, err := os.Open(path) + if err != nil { + return nil + } + defer f.Close() + var out []string + sc := bufio.NewScanner(f) + for sc.Scan() { + line := strings.TrimSpace(sc.Text()) + if line == "" || strings.HasPrefix(line, "#") { + continue + } + if strings.HasSuffix(line, ".mp3") { + out = append(out, filepath.Base(line)) // never trust paths in the list + } + } + return out +} + +// titleFromName pulls the ID3-free display name; the mp3 itself carries the +// real title for the stereo, this is for MQTT/logging. +func titleFromName(name string) string { + return strings.TrimSuffix(name, ".mp3") +} + +func publishNowPlaying(m *infra.MQTTNowPlaying, station, title, file, streamURL string) { + if m == nil { + return + } + payload, _ := json.Marshal(nowPlaying{ + Station: station, + Title: title, + Artist: station, + Track: title, + Filename: file, + StreamURL: streamURL, + Timestamp: time.Now().Unix(), + }) + if err := m.Publish(payload); err != nil { + log.Printf("mqtt publish failed (continuing): %v", err) + } +} + +func parseInt(s string) int { + var n int + for _, c := range s { + if c < '0' || c > '9' { + return 0 + } + n = n*10 + int(c-'0') + } + return n +}