From 248c489a368754ac217613f742b1dc283f4f3d6f Mon Sep 17 00:00:00 2001 From: hermes Date: Wed, 30 Sep 2026 21:59:33 +0200 Subject: [PATCH] add internal/infra/icecast.go --- internal/infra/icecast.go | 119 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 119 insertions(+) create mode 100644 internal/infra/icecast.go diff --git a/internal/infra/icecast.go b/internal/infra/icecast.go new file mode 100644 index 0000000..d09bfda --- /dev/null +++ b/internal/infra/icecast.go @@ -0,0 +1,119 @@ +package infra + +import ( + "bufio" + "fmt" + "io" + "net" + "strings" + "time" +) + +// Icecast source-protocol client (SOURCE command, HTTP/1.0 style headers). +// Spec: Icecast 2 docs, "Source Protocol". The source authenticates with the +// source password, then writes raw MP3 bytes forever. + +// SourceMeta describes the stream to the icecast server. +type SourceMeta struct { + Mount string // mount point, e.g. "/ai-radio.mp3" + Name string // stream name shown in the directory + Genre string + Bitrate string // kbps, informational + URL string + Description string +} + +// buildSourceRequest renders the SOURCE handshake request. +func buildSourceRequest(password string, m SourceMeta) string { + h := []string{ + "SOURCE " + password, + "ice-name: " + m.Name, + "ice-genre: " + m.Genre, + "ice-bitrate: " + m.Bitrate, + "ice-public: 1", + "ice-description: " + m.Description, + "ice-url: " + m.URL, + "content-type: audio/mpeg", + "mount: " + m.Mount, + } + return strings.Join(h, "\r\n") + "\r\n\r\n" +} + +// parseSourceResponse inspects the first status line of the handshake reply. +func parseSourceResponse(resp string) (bool, string) { + line := strings.SplitN(strings.TrimSpace(resp), "\r\n", 2)[0] + if strings.Contains(line, " 200") { + return true, line + } + return false, line +} + +// IcecastSource connects and streams. Reconnection is the caller's job +// (the streamer loop handles it). +type IcecastSource struct { + Addr string // host:port + Password string + Meta SourceMeta +} + +// Connect performs the SOURCE handshake and returns the live connection. +// Every subsequent Write goes straight to listeners. +func (s *IcecastSource) Connect() (net.Conn, error) { + conn, err := net.DialTimeout("tcp", s.Addr, 10*time.Second) + if err != nil { + return nil, fmt.Errorf("icecast dial %s: %w", s.Addr, err) + } + if _, err := conn.Write([]byte(buildSourceRequest(s.Password, s.Meta))); err != nil { + conn.Close() + return nil, fmt.Errorf("icecast handshake write: %w", err) + } + _ = conn.SetReadDeadline(time.Now().Add(10 * time.Second)) + br := bufio.NewReader(conn) + status, err := br.ReadString('\n') + if err != nil && status == "" { + conn.Close() + return nil, fmt.Errorf("icecast handshake read: %w", err) + } + // drain remaining headers + for { + l, err := br.ReadString('\n') + if err != nil || strings.TrimSpace(l) == "" { + break + } + } + ok, line := parseSourceResponse(status) + if !ok { + conn.Close() + return nil, fmt.Errorf("icecast rejected source: %s", line) + } + return conn, nil +} + +// StreamFile sends one MP3 file paced in real time (bytes/sec ≈ bitrate). +// Pacing keeps icecast's buffer healthy; too fast overflows, too slow starves. +func StreamFile(w io.Writer, r io.Reader, kbps int) (int64, error) { + const chunk = 4096 + bytesPerSec := int64(kbps) * 1000 / 8 + buf := make([]byte, chunk) + start := time.Now() + var sent int64 + for { + n, rerr := r.Read(buf) + if n > 0 { + if _, werr := w.Write(buf[:n]); werr != nil { + return sent, werr + } + sent += int64(n) + lead := float64(sent)/float64(bytesPerSec) - time.Since(start).Seconds() + if lead > 0.05 { + time.Sleep(time.Duration(lead*float64(time.Second)) / 2) + } + } + if rerr != nil { + if rerr == io.EOF { + return sent, nil + } + return sent, rerr + } + } +}