// 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 }