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