120 lines
3.2 KiB
Go
120 lines
3.2 KiB
Go
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
|
|
}
|
|
}
|
|
}
|