add cmd/streamer/main.go
This commit is contained in:
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user