package main import ( "context" "flag" "fmt" "log" "net/http" "os" "os/signal" "path/filepath" "strings" "sync" "syscall" "time" "github.com/Sendspin/sendspin-go/pkg/discovery" "github.com/Sendspin/sendspin-go/pkg/sendspin" "github.com/gorilla/websocket" "rpi-sendspin/internal/pipewire" ) func main() { server := flag.String("server", "", "server address for client-initiated mode (dial out). Leave empty for server-initiated mode: advertise via mDNS and wait for the server to connect.") name := flag.String("name", hostname(), "player name") device := flag.String("device", "", "PipeWire sink node name (default: system default)") listDevices := flag.Bool("list-devices", false, "list available PipeWire sinks and exit") volume := flag.Int("volume", 80, "initial volume (0-100)") mdnsPort := flag.Int("mdns-port", 8928, "port to advertise in mDNS (_sendspin._tcp)") noMDNS := flag.Bool("no-mdns", false, "disable mDNS advertisement") muteLoopback := flag.String("mute-loopback", "", "PipeWire node.name of a loopback to mute while streaming from --mute-server (requires --mute-server)") muteServer := flag.String("mute-server", "", "server address that arms loopback muting; must match --server (requires --mute-loopback)") flag.Parse() if *listDevices { sinks, err := pipewire.ListSinks() if err != nil { log.Fatalf("list sinks: %v", err) } if len(sinks) == 0 { fmt.Println("no PipeWire sinks found") return } fmt.Println("Available PipeWire sinks:") for _, s := range sinks { fmt.Printf(" %s\n", s) } return } idFile := clientIDFile() savedID, _ := os.ReadFile(idFile) clientID, err := sendspin.ResolveClientID("", string(savedID), func(id string) error { if err := os.MkdirAll(filepath.Dir(idFile), 0700); err != nil { return err } return os.WriteFile(idFile, []byte(id), 0600) }) if err != nil { log.Fatalf("resolve client ID: %v", err) } if !*noMDNS { disc := discovery.NewManager(discovery.Config{ ServiceName: *name, Port: *mdnsPort, }) if err := disc.Advertise(); err != nil { log.Printf("mdns advertise: %v", err) } else { defer disc.Stop() } } output := pipewire.NewOutput(*device, *name) muter := buildMuter(*server, *muteServer, *muteLoopback) player, err := sendspin.NewPlayer(sendspin.PlayerConfig{ ServerAddr: *server, PlayerName: *name, ClientID: clientID, Volume: *volume, Output: output, Reconnect: sendspin.ReconnectConfig{ Enabled: true, InitialDelay: 2 * time.Second, MaxDelay: 30 * time.Second, Multiplier: 2.0, }, OnMetadata: func(meta sendspin.Metadata) { if meta.Artist != "" || meta.Title != "" { log.Printf("now playing: %s — %s", meta.Artist, meta.Title) } }, OnStateChange: func(s sendspin.PlayerState) { if muter != nil { muter.SetDesiredMute(s.State == "playing") } }, OnError: func(err error) { log.Printf("error: %v", err) }, }) if err != nil { log.Fatalf("create player: %v", err) } // LIFO: muter.Stop runs AFTER player.Close, so the player has stopped // emitting "playing" state by the time we force the final unmute. if muter != nil { defer func() { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() muter.Stop(ctx) }() } defer player.Close() if *device != "" { log.Printf("output: %s", *device) } else { log.Printf("output: PipeWire default sink") } sig := make(chan os.Signal, 1) signal.Notify(sig, os.Interrupt, syscall.SIGTERM) if *server == "" { // Server-initiated mode: the server discovers us via mDNS and dials // our /sendspin endpoint. We just listen and let player.Accept drive // each connection. if *noMDNS { log.Printf("warning: --no-mdns with no --server — the server cannot discover this player") } runServerInitiated(player, *mdnsPort, sig) } else { // Client-initiated mode: dial the configured server. log.Printf("connecting to %s as %q...", *server, *name) if err := player.Connect(); err != nil { log.Fatalf("connect: %v", err) } if err := player.Play(); err != nil { log.Fatalf("play: %v", err) } <-sig } log.Println("stopping") player.Stop() } // runServerInitiated runs an HTTP/WebSocket server on the mDNS-advertised // port and hands each accepted /sendspin connection to player.Accept. Only // one session runs at a time; a newer connection evicts the current one // (newest-wins) so a server that reconnects isn't blocked behind a stale // socket. Blocks until a signal arrives on sig. func runServerInitiated(player *sendspin.Player, port int, sig <-chan os.Signal) { upgrader := websocket.Upgrader{ CheckOrigin: func(*http.Request) bool { return true }, } var ( mu sync.Mutex // serializes sessions; held for a connection's lifetime connMu sync.Mutex // guards curConn curConn *websocket.Conn ) handler := func(w http.ResponseWriter, r *http.Request) { conn, err := upgrader.Upgrade(w, r, nil) if err != nil { log.Printf("websocket upgrade failed: %v", err) return } // Evict any in-flight session so its Accept returns and releases mu. connMu.Lock() if curConn != nil { curConn.Close() } curConn = conn connMu.Unlock() mu.Lock() defer mu.Unlock() // We may have been superseded while waiting for the lock. connMu.Lock() superseded := curConn != conn connMu.Unlock() if superseded { conn.Close() return } if err := player.Accept(conn); err != nil { log.Printf("session ended: %v", err) } connMu.Lock() if curConn == conn { curConn = nil } connMu.Unlock() } mux := http.NewServeMux() mux.HandleFunc("/sendspin", handler) srv := &http.Server{Addr: fmt.Sprintf(":%d", port), Handler: mux} go func() { log.Printf("server-initiated mode: listening on :%d/sendspin, waiting for a server to connect", port) if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed { log.Fatalf("listen: %v", err) } }() <-sig // Drop the active session first so the blocked handler returns, then let // Shutdown drain the (now-unblocked) server. player.CloseConnection() ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() srv.Shutdown(ctx) } func hostname() string { h, _ := os.Hostname() if h == "" { return "sendspin-player" } return h } // buildMuter wires up the loopback muter when both --mute-loopback and // --mute-server are set AND --mute-server matches the configured --server. // Mismatches are intentional silent no-ops so the same config file can ship // to machines where the server happens to be different — only logged so // it's visible. Returns nil when muting is disabled. // // The address check is a startup-time string compare; mDNS rediscovery // during reconnect could in theory move the player to a different server, // but for the same-machine echo-cancel use case that's acceptable. func buildMuter(server, muteServer, muteLoopback string) *pipewire.Muter { if muteLoopback == "" { if muteServer != "" { log.Printf("muter: --mute-server set without --mute-loopback; loopback mute disabled") } return nil } // Server-initiated (listen) mode has no fixed --server to compare // against, so arm the muter purely on --mute-loopback. if server == "" { log.Printf("muter: will mute PipeWire node %q while streaming (server-initiated mode)", muteLoopback) return pipewire.NewMuter(muteLoopback) } if muteServer == "" { log.Printf("muter: --mute-loopback and --mute-server must be set together in client-initiated mode; loopback mute disabled") return nil } if !sameServerAddr(server, muteServer) { log.Printf("muter: --mute-server %q != --server %q; loopback mute disabled", muteServer, server) return nil } log.Printf("muter: will mute PipeWire node %q while streaming from %s", muteLoopback, server) return pipewire.NewMuter(muteLoopback) } // sameServerAddr normalizes two host[:port] strings for equality so // "localhost" and "localhost:8927" compare equal when 8927 is the default. func sameServerAddr(a, b string) bool { const defaultPort = "8927" norm := func(s string) string { s = strings.TrimSpace(strings.ToLower(s)) if !strings.Contains(s, ":") { s += ":" + defaultPort } return s } return norm(a) == norm(b) } func clientIDFile() string { dir, err := os.UserConfigDir() if err != nil { dir = os.TempDir() } return filepath.Join(dir, "rpi-sendspin", "client_id") }