Phase 2a: buffer state under one lock, atomic tuner limit, context-driven external buffer
continuous-integration/drone/push Build encountered an error
continuous-integration/drone/push Build encountered an error
- src/buffer_state.go replaces the two sync.Maps plus a global RWMutex with one bufferMu guarding a map of *Playlist. Each stream has one shared bufferStream (URL, folder, status, client count, error, cancel hook); downloaders keep a private ThisStream and publish through helpers. Playlist.Clients / ThisClient / ClientConnection are gone: there was one client counter per stream in two places that could disagree. - bufferAcquireStream does the tuner check and the registration under the same lock, so two clients tuning at once cannot both pass the limit. bufferReleaseClient removes the stream when the last client leaves, cancels its process and deletes its segment folder. - bufferingStream rewritten: waits on r.Context() instead of the deprecated CloseNotifier, sets Content-Type before WriteHeader (the old code set headers after and one was literally named "Content-Length:"), flushes each segment to the client, no defer inside the segment loop. - connectToStreamingServer: shared streamHTTPClient with dial, TLS and response-header timeouts (no overall timeout, bodies are endless); the deferred Body/segment closes inside the redirect and read loops are now explicit closes, so a multi-hour stream no longer accumulates them. - thirdPartyBuffer rewritten around exec.CommandContext: the process is killed when the last client leaves or when no usable data arrives within 20 s (time.AfterFunc watchdog, no leaked goroutine); cmd.Start error is checked; no panic; one file handle per segment. ffmpeg and VLC command lines come from buildFFmpegArgs / buildVLCArgs, which are unit tested. - Data.Cache.StreamingURLS is guarded by streamingURLsMu and persisted from a snapshot. - Tests (src/buffer_test.go, run with -race): restream shares one provider connection, six concurrent tunes against a tuner limit of two, 50-way acquire/release contention, cleanup and cancel on last release.
This commit is contained in:
@@ -12,3 +12,4 @@ de.json
|
||||
agent.md
|
||||
skill.md
|
||||
node_modules/
|
||||
.claude/
|
||||
|
||||
+303
-681
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,246 @@
|
||||
package src
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net"
|
||||
"net/http"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Shared buffer state.
|
||||
//
|
||||
// One mutex guards the playlist table and every bufferStream in it. The
|
||||
// download goroutines keep their own private ThisStream (HLS sequence,
|
||||
// timing, bandwidth) and only publish status, errors and the cancel hook
|
||||
// through the helpers below. HTTP client goroutines only read snapshots.
|
||||
|
||||
var bufferMu sync.Mutex
|
||||
var bufferPlaylists = make(map[string]*Playlist)
|
||||
|
||||
// errTunerLimit : the playlist already serves as many streams as it has tuners.
|
||||
var errTunerLimit = errors.New("tuner limit reached")
|
||||
|
||||
// bufferStream : the shared view of one buffered stream.
|
||||
type bufferStream struct {
|
||||
URL string
|
||||
ChannelName string
|
||||
MD5 string
|
||||
Folder string
|
||||
PlaylistID string
|
||||
PlaylistName string
|
||||
|
||||
Status bool // first segment is on disk, clients may read
|
||||
Clients int // HTTP clients currently attached
|
||||
Err error // fatal downloader error; clients disconnect when set
|
||||
|
||||
cancel context.CancelFunc // stops an external process buffer
|
||||
}
|
||||
|
||||
// streamHTTPClient : used by the xTeVe buffer to fetch provider streams.
|
||||
// No overall timeout (bodies are endless), but dial and response-header
|
||||
// timeouts so a dead provider cannot hang a tuner forever. Redirects are
|
||||
// surfaced to the caller, which follows them itself.
|
||||
var streamHTTPClient = &http.Client{
|
||||
Transport: &http.Transport{
|
||||
Proxy: http.ProxyFromEnvironment,
|
||||
DialContext: (&net.Dialer{Timeout: 10 * time.Second, KeepAlive: 30 * time.Second}).DialContext,
|
||||
TLSHandshakeTimeout: 10 * time.Second,
|
||||
ResponseHeaderTimeout: 20 * time.Second,
|
||||
},
|
||||
CheckRedirect: func(req *http.Request, via []*http.Request) error {
|
||||
return errors.New("Redirect")
|
||||
},
|
||||
}
|
||||
|
||||
// bufferAcquireStream : registers a client for streamingURL on the playlist.
|
||||
// Returns the stream ID and whether a downloader has to be started. The
|
||||
// tuner check and the registration happen under one lock, so two clients
|
||||
// tuning at the same moment cannot both slip past the limit.
|
||||
func bufferAcquireStream(playlistID, streamingURL, channelName string) (streamID int, isNew bool, err error) {
|
||||
|
||||
bufferMu.Lock()
|
||||
defer bufferMu.Unlock()
|
||||
|
||||
playlist, ok := bufferPlaylists[playlistID]
|
||||
if !ok {
|
||||
|
||||
var playlistType string
|
||||
switch playlistID[0:1] {
|
||||
case "M":
|
||||
playlistType = "m3u"
|
||||
case "H":
|
||||
playlistType = "hdhr"
|
||||
}
|
||||
|
||||
playlist = &Playlist{
|
||||
Folder: System.Folder.Temp + playlistID + string(os.PathSeparator),
|
||||
PlaylistID: playlistID,
|
||||
Streams: make(map[int]*bufferStream),
|
||||
}
|
||||
|
||||
if err = checkFolder(playlist.Folder); err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
playlist.Tuner = getTuner(playlistID, playlistType)
|
||||
playlist.PlaylistName = getProviderParameter(playlistID, playlistType, "name")
|
||||
|
||||
bufferPlaylists[playlistID] = playlist
|
||||
}
|
||||
|
||||
// Same URL already buffered: attach to it (restream).
|
||||
for id, s := range playlist.Streams {
|
||||
if s.URL == streamingURL {
|
||||
s.Clients++
|
||||
showDebug(fmt.Sprintf("Restream Status:Playlist: %s - Channel: %s - Connections: %d", playlist.PlaylistName, s.ChannelName, s.Clients), 1)
|
||||
showInfo(fmt.Sprintf("Streaming Status:Channel: %s (Clients: %d)", s.ChannelName, s.Clients))
|
||||
return id, false, nil
|
||||
}
|
||||
}
|
||||
|
||||
if len(playlist.Streams) >= playlist.Tuner {
|
||||
showInfo(fmt.Sprintf("Streaming Status:Playlist: %s - No new connections available. Tuner = %d", playlist.PlaylistName, playlist.Tuner))
|
||||
return 0, false, errTunerLimit
|
||||
}
|
||||
|
||||
streamID = createStreamID(playlist.Streams)
|
||||
|
||||
var md5 = getMD5(streamingURL)
|
||||
playlist.Streams[streamID] = &bufferStream{
|
||||
URL: streamingURL,
|
||||
ChannelName: channelName,
|
||||
MD5: md5,
|
||||
Folder: playlist.Folder + md5 + string(os.PathSeparator),
|
||||
PlaylistID: playlistID,
|
||||
PlaylistName: playlist.PlaylistName,
|
||||
Clients: 1,
|
||||
}
|
||||
|
||||
showInfo("Streaming Status:" + tunerStatusLocked(playlist))
|
||||
|
||||
return streamID, true, nil
|
||||
}
|
||||
|
||||
func tunerStatusLocked(playlist *Playlist) string {
|
||||
return fmt.Sprintf("Playlist: %s - Tuner: %d / %d", playlist.PlaylistName, len(playlist.Streams), playlist.Tuner)
|
||||
}
|
||||
|
||||
// bufferTunerStatus : "Playlist: x - Tuner: n / m" for log lines.
|
||||
func bufferTunerStatus(playlistID string) string {
|
||||
bufferMu.Lock()
|
||||
defer bufferMu.Unlock()
|
||||
if playlist, ok := bufferPlaylists[playlistID]; ok {
|
||||
return tunerStatusLocked(playlist)
|
||||
}
|
||||
return "Playlist: " + playlistID + " - Tuner: 0 / 0"
|
||||
}
|
||||
|
||||
func lookupStreamLocked(playlistID string, streamID int) (*bufferStream, bool) {
|
||||
if playlist, ok := bufferPlaylists[playlistID]; ok {
|
||||
s, ok := playlist.Streams[streamID]
|
||||
return s, ok
|
||||
}
|
||||
return nil, false
|
||||
}
|
||||
|
||||
// bufferStreamSnapshot : copy of the shared stream state.
|
||||
func bufferStreamSnapshot(playlistID string, streamID int) (bufferStream, bool) {
|
||||
bufferMu.Lock()
|
||||
defer bufferMu.Unlock()
|
||||
if s, ok := lookupStreamLocked(playlistID, streamID); ok {
|
||||
var copy = *s
|
||||
copy.cancel = nil
|
||||
return copy, true
|
||||
}
|
||||
return bufferStream{}, false
|
||||
}
|
||||
|
||||
// bufferStreamActive : false once the last client left and the stream was removed.
|
||||
func bufferStreamActive(playlistID string, streamID int) bool {
|
||||
bufferMu.Lock()
|
||||
defer bufferMu.Unlock()
|
||||
_, ok := lookupStreamLocked(playlistID, streamID)
|
||||
return ok
|
||||
}
|
||||
|
||||
// bufferSetStatus : downloader reports that segments are available.
|
||||
func bufferSetStatus(playlistID string, streamID int, status bool) {
|
||||
bufferMu.Lock()
|
||||
defer bufferMu.Unlock()
|
||||
if s, ok := lookupStreamLocked(playlistID, streamID); ok {
|
||||
s.Status = status
|
||||
}
|
||||
}
|
||||
|
||||
// bufferSetError : downloader reports a fatal error; attached clients disconnect.
|
||||
func bufferSetError(playlistID string, streamID int, err error) {
|
||||
bufferMu.Lock()
|
||||
defer bufferMu.Unlock()
|
||||
if s, ok := lookupStreamLocked(playlistID, streamID); ok && s.Err == nil {
|
||||
s.Err = err
|
||||
}
|
||||
}
|
||||
|
||||
// bufferSetCancel : registers the function that stops an external buffer process.
|
||||
func bufferSetCancel(playlistID string, streamID int, cancel context.CancelFunc) {
|
||||
bufferMu.Lock()
|
||||
defer bufferMu.Unlock()
|
||||
if s, ok := lookupStreamLocked(playlistID, streamID); ok {
|
||||
s.cancel = cancel
|
||||
}
|
||||
}
|
||||
|
||||
// bufferReleaseClient : detaches one client. When the last one leaves the
|
||||
// stream is removed, its process (if any) cancelled and its files deleted.
|
||||
func bufferReleaseClient(playlistID string, streamID int) {
|
||||
|
||||
bufferMu.Lock()
|
||||
defer bufferMu.Unlock()
|
||||
|
||||
playlist, ok := bufferPlaylists[playlistID]
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
s, ok := playlist.Streams[streamID]
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
s.Clients--
|
||||
showInfo("Streaming Status:Client has terminated the connection")
|
||||
showInfo(fmt.Sprintf("Streaming Status:Channel: %s (Clients: %d)", s.ChannelName, s.Clients))
|
||||
|
||||
if s.Clients > 0 {
|
||||
return
|
||||
}
|
||||
|
||||
if s.cancel != nil {
|
||||
s.cancel()
|
||||
}
|
||||
|
||||
delete(playlist.Streams, streamID)
|
||||
|
||||
showDebug(fmt.Sprintf("Remove tmp folder:%s", s.Folder), 1)
|
||||
os.RemoveAll(s.Folder)
|
||||
|
||||
showInfo(fmt.Sprintf("Streaming Status:Channel: %s - No client is using this channel anymore. Streaming Server connection has ended", s.ChannelName))
|
||||
showInfo("Streaming Status:" + tunerStatusLocked(playlist))
|
||||
|
||||
if len(playlist.Streams) == 0 {
|
||||
delete(bufferPlaylists, playlistID)
|
||||
}
|
||||
}
|
||||
|
||||
// bufferActiveStreams : number of buffered streams for a playlist (0 when idle).
|
||||
func bufferActiveStreams(playlistID string) int {
|
||||
bufferMu.Lock()
|
||||
defer bufferMu.Unlock()
|
||||
if playlist, ok := bufferPlaylists[playlistID]; ok {
|
||||
return len(playlist.Streams)
|
||||
}
|
||||
return 0
|
||||
}
|
||||
@@ -0,0 +1,285 @@
|
||||
package src
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"reflect"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// fakeProvider streams endless TS-looking bytes and counts concurrent connections.
|
||||
type fakeProvider struct {
|
||||
srv *httptest.Server
|
||||
active atomic.Int32
|
||||
max atomic.Int32
|
||||
}
|
||||
|
||||
func newFakeProvider(t *testing.T) *fakeProvider {
|
||||
t.Helper()
|
||||
p := &fakeProvider{}
|
||||
p.srv = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
n := p.active.Add(1)
|
||||
defer p.active.Add(-1)
|
||||
for {
|
||||
m := p.max.Load()
|
||||
if n <= m || p.max.CompareAndSwap(m, n) {
|
||||
break
|
||||
}
|
||||
}
|
||||
w.Header().Set("Content-Type", "video/mp2t")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
flusher, _ := w.(http.Flusher)
|
||||
chunk := bytes.Repeat([]byte{0x47}, 2048)
|
||||
for {
|
||||
select {
|
||||
case <-r.Context().Done():
|
||||
return
|
||||
default:
|
||||
}
|
||||
if _, err := w.Write(chunk); err != nil {
|
||||
return
|
||||
}
|
||||
if flusher != nil {
|
||||
flusher.Flush()
|
||||
}
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
}))
|
||||
t.Cleanup(p.srv.Close)
|
||||
return p
|
||||
}
|
||||
|
||||
func setupBuffer(t *testing.T, tuner string) (playlistID string, clientSrv *httptest.Server) {
|
||||
t.Helper()
|
||||
playlistID = "M-test"
|
||||
|
||||
System.Folder.Temp = t.TempDir() + string(os.PathSeparator)
|
||||
Settings.Buffer = "xteve"
|
||||
Settings.BufferSize = 1
|
||||
Settings.BufferTimeout = 0
|
||||
Settings.Files.M3U = map[string]any{playlistID: map[string]any{"name": "Test", "tuner": tuner}}
|
||||
|
||||
bufferMu.Lock()
|
||||
bufferPlaylists = make(map[string]*Playlist)
|
||||
bufferMu.Unlock()
|
||||
|
||||
clientSrv = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
bufferingStream(playlistID, r.URL.Query().Get("u"), "chan", w, r)
|
||||
}))
|
||||
t.Cleanup(clientSrv.Close)
|
||||
return
|
||||
}
|
||||
|
||||
// tune requests a stream, reads at least want bytes (or gives up after 10 s)
|
||||
// and returns what it got.
|
||||
func tune(t *testing.T, clientSrv *httptest.Server, streamURL string, want int) []byte {
|
||||
t.Helper()
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
req, _ := http.NewRequestWithContext(ctx, "GET", clientSrv.URL+"/?u="+streamURL, nil)
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
t.Errorf("tune %s: %v", streamURL, err)
|
||||
return nil
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
var got []byte
|
||||
buf := make([]byte, 4096)
|
||||
for len(got) < want {
|
||||
n, err := resp.Body.Read(buf)
|
||||
got = append(got, buf[:n]...)
|
||||
if err != nil {
|
||||
break
|
||||
}
|
||||
}
|
||||
return got
|
||||
}
|
||||
|
||||
func waitIdle(t *testing.T, playlistID string) {
|
||||
t.Helper()
|
||||
deadline := time.Now().Add(5 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
if bufferActiveStreams(playlistID) == 0 {
|
||||
return
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
}
|
||||
t.Fatalf("streams still active: %d", bufferActiveStreams(playlistID))
|
||||
}
|
||||
|
||||
func TestBufferRestreamSharesOneProviderConnection(t *testing.T) {
|
||||
provider := newFakeProvider(t)
|
||||
playlistID, clientSrv := setupBuffer(t, "1")
|
||||
streamURL := provider.srv.URL + "/live.ts"
|
||||
|
||||
var wg sync.WaitGroup
|
||||
results := make([][]byte, 2)
|
||||
for i := range results {
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
results[i] = tune(t, clientSrv, streamURL, 4096)
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
for i, got := range results {
|
||||
if len(got) == 0 {
|
||||
t.Errorf("client %d received no data", i)
|
||||
} else if got[0] != 0x47 {
|
||||
t.Errorf("client %d received something other than the provider stream", i)
|
||||
}
|
||||
}
|
||||
if m := provider.max.Load(); m != 1 {
|
||||
t.Errorf("provider saw %d concurrent connections, want 1 (restream)", m)
|
||||
}
|
||||
|
||||
waitIdle(t, playlistID)
|
||||
|
||||
entries, _ := os.ReadDir(System.Folder.Temp + playlistID)
|
||||
if len(entries) != 0 {
|
||||
t.Errorf("stream folders left behind: %v", entries)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBufferTunerLimitIsAtomic(t *testing.T) {
|
||||
provider := newFakeProvider(t)
|
||||
playlistID, clientSrv := setupBuffer(t, "2")
|
||||
|
||||
limitClip, err := readWebFile("video/stream-limit.ts")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
const clients = 6
|
||||
var wg sync.WaitGroup
|
||||
results := make([][]byte, clients)
|
||||
for i := 0; i < clients; i++ {
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
results[i] = tune(t, clientSrv, provider.srv.URL+"/"+string(rune('a'+i))+".ts", 2048)
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
// Both the provider stream and the tuner-limit clip are MPEG-TS, so both
|
||||
// start with 0x47. The fake provider sends nothing but 0x47 bytes; the clip
|
||||
// does not.
|
||||
providerPrefix := bytes.Repeat([]byte{0x47}, 188)
|
||||
var real, limited int
|
||||
for i, got := range results {
|
||||
switch {
|
||||
case len(got) < 188:
|
||||
t.Errorf("client %d received only %d bytes", i, len(got))
|
||||
case bytes.HasPrefix(limitClip, got[:188]):
|
||||
limited++
|
||||
case bytes.Equal(got[:188], providerPrefix):
|
||||
real++
|
||||
default:
|
||||
t.Errorf("client %d received unexpected bytes", i)
|
||||
}
|
||||
}
|
||||
|
||||
if real != 2 || limited != 4 {
|
||||
t.Errorf("real=%d limited=%d, want 2 real streams and 4 tuner-limit responses", real, limited)
|
||||
}
|
||||
if m := provider.max.Load(); m > 2 {
|
||||
t.Errorf("provider saw %d concurrent connections, tuner limit is 2", m)
|
||||
}
|
||||
|
||||
waitIdle(t, playlistID)
|
||||
}
|
||||
|
||||
func TestBufferAcquireReleaseUnderContention(t *testing.T) {
|
||||
playlistID, _ := setupBuffer(t, "3")
|
||||
Settings.Buffer = "-" // no downloader needed for this test
|
||||
Settings.Tuner = 3
|
||||
|
||||
var wg sync.WaitGroup
|
||||
var over atomic.Int32
|
||||
for i := 0; i < 50; i++ {
|
||||
wg.Add(1)
|
||||
go func(i int) {
|
||||
defer wg.Done()
|
||||
id, isNew, err := bufferAcquireStream(playlistID, "http://x/"+string(rune('A'+i%26))+string(rune('a'+i/26)), "c")
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
if n := bufferActiveStreams(playlistID); n > 3 {
|
||||
over.Store(int32(n))
|
||||
}
|
||||
_ = isNew
|
||||
time.Sleep(time.Millisecond)
|
||||
bufferReleaseClient(playlistID, id)
|
||||
}(i)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
if n := over.Load(); n != 0 {
|
||||
t.Errorf("tuner limit exceeded: saw %d streams", n)
|
||||
}
|
||||
if n := bufferActiveStreams(playlistID); n != 0 {
|
||||
t.Errorf("%d streams left after release", n)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBufferReleaseCancelsProcessAndRemovesFolder(t *testing.T) {
|
||||
playlistID, _ := setupBuffer(t, "1")
|
||||
Settings.Buffer = "-"
|
||||
|
||||
id, isNew, err := bufferAcquireStream(playlistID, "http://x/one", "c")
|
||||
if err != nil || !isNew {
|
||||
t.Fatalf("acquire: new=%v err=%v", isNew, err)
|
||||
}
|
||||
|
||||
snap, _ := bufferStreamSnapshot(playlistID, id)
|
||||
if err := os.MkdirAll(snap.Folder, 0755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
var cancelled atomic.Bool
|
||||
bufferSetCancel(playlistID, id, func() { cancelled.Store(true) })
|
||||
|
||||
bufferReleaseClient(playlistID, id)
|
||||
|
||||
if !cancelled.Load() {
|
||||
t.Error("cancel hook not called when the last client left")
|
||||
}
|
||||
if _, err := os.Stat(snap.Folder); !os.IsNotExist(err) {
|
||||
t.Error("stream folder not removed")
|
||||
}
|
||||
if bufferStreamActive(playlistID, id) {
|
||||
t.Error("stream still active")
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildFFmpegArgs(t *testing.T) {
|
||||
got := buildFFmpegArgs("-hide_banner -i [URL] -c copy -f mpegts pipe:1", "http://h/s.ts", "UA/1")
|
||||
want := []string{"-user_agent", "UA/1", "-hide_banner", "-i", "http://h/s.ts", "-c", "copy", "-f", "mpegts", "pipe:1"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("got %v\nwant %v", got, want)
|
||||
}
|
||||
if got := buildFFmpegArgs("-i [URL]", "u", ""); !reflect.DeepEqual(got, []string{"-i", "u"}) {
|
||||
t.Errorf("no user agent: got %v", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildVLCArgs(t *testing.T) {
|
||||
got := buildVLCArgs("-I dummy [URL] --sout #std{mux=ts,access=file,dst=-}", "http://h/s.ts", "UA/1")
|
||||
want := []string{"-I", "dummy", "http://h/s.ts", ":http-user-agent=UA/1", "--sout", "#std{mux=ts,access=file,dst=-}"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Errorf("got %v\nwant %v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
var _ = io.EOF
|
||||
@@ -5,7 +5,6 @@ import (
|
||||
"os"
|
||||
"runtime"
|
||||
"strings"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// System : Beinhaltet alle Systeminformationen
|
||||
@@ -23,15 +22,6 @@ var Data DataStruct
|
||||
// SystemFiles : Alle Systemdateien
|
||||
var SystemFiles = []string{"authentication.json", "pms.json", "settings.json", "xepg.json", "urls.json"}
|
||||
|
||||
// BufferInformation : Informationen über den Buffer (aktive Streams, maximale Streams)
|
||||
var BufferInformation sync.Map
|
||||
|
||||
// BufferClients : Anzahl der Clients die einen Stream über den Buffer abspielen
|
||||
var BufferClients sync.Map
|
||||
|
||||
// Lock : Lock Map
|
||||
var Lock = sync.RWMutex{}
|
||||
|
||||
// Init : Systeminitialisierung
|
||||
func Init() (err error) {
|
||||
|
||||
|
||||
+1
-1
@@ -176,7 +176,7 @@ func getLineup() (jsonContent []byte, err error) {
|
||||
|
||||
Data.Cache.PMS = nil
|
||||
|
||||
saveMapToJSONFile(System.File.URLS, Data.Cache.StreamingURLS)
|
||||
saveMapToJSONFile(System.File.URLS, streamingURLsSnapshot())
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
+2
-14
@@ -2,20 +2,14 @@ package src
|
||||
|
||||
import "time"
|
||||
|
||||
// Playlist : Enthält allen Playlistinformationen, die der Buffer benötigr
|
||||
// Playlist : per-provider buffer state (guarded by bufferMu in buffer_state.go)
|
||||
type Playlist struct {
|
||||
Folder string
|
||||
PlaylistID string
|
||||
PlaylistName string
|
||||
Tuner int
|
||||
|
||||
Clients map[int]ThisClient
|
||||
Streams map[int]ThisStream
|
||||
}
|
||||
|
||||
// ThisClient : Clientinfos
|
||||
type ThisClient struct {
|
||||
Connection int
|
||||
Streams map[int]*bufferStream
|
||||
}
|
||||
|
||||
// ThisStream : Enthält Informationen zu dem abzuspielenden Stream einer Playlist
|
||||
@@ -94,12 +88,6 @@ type DynamicStream struct {
|
||||
URL string
|
||||
}
|
||||
|
||||
// ClientConnection : Client Verbindungen
|
||||
type ClientConnection struct {
|
||||
Connection int
|
||||
Error error
|
||||
}
|
||||
|
||||
// BandwidthCalculation : Bandbreitenberechnung für den Stream
|
||||
type BandwidthCalculation struct {
|
||||
NetworkBandwidth int
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
"os"
|
||||
"reflect"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
@@ -270,6 +271,21 @@ func saveSettings(settings SettingsStruct) (err error) {
|
||||
}
|
||||
|
||||
// Zugriff über die Domain ermöglichen
|
||||
// streamingURLsMu guards Data.Cache.StreamingURLS, which is written by
|
||||
// /lineup.json and /m3u/ requests while /stream/ requests read it.
|
||||
var streamingURLsMu sync.Mutex
|
||||
|
||||
// streamingURLsSnapshot : copy of the URL cache for persisting to disk.
|
||||
func streamingURLsSnapshot() map[string]StreamInfo {
|
||||
streamingURLsMu.Lock()
|
||||
defer streamingURLsMu.Unlock()
|
||||
var out = make(map[string]StreamInfo, len(Data.Cache.StreamingURLS))
|
||||
for k, v := range Data.Cache.StreamingURLS {
|
||||
out[k] = v
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func setGlobalDomain(domain string) {
|
||||
|
||||
System.Domain = domain
|
||||
@@ -329,6 +345,9 @@ func createStreamingURL(streamingType, playlistID, channelNumber, channelName, u
|
||||
var streamInfo StreamInfo
|
||||
var serverProtocol string
|
||||
|
||||
streamingURLsMu.Lock()
|
||||
defer streamingURLsMu.Unlock()
|
||||
|
||||
if len(Data.Cache.StreamingURLS) == 0 {
|
||||
Data.Cache.StreamingURLS = make(map[string]StreamInfo)
|
||||
}
|
||||
@@ -368,6 +387,9 @@ func createStreamingURL(streamingType, playlistID, channelNumber, channelName, u
|
||||
|
||||
func getStreamInfo(urlID string) (streamInfo StreamInfo, err error) {
|
||||
|
||||
streamingURLsMu.Lock()
|
||||
defer streamingURLsMu.Unlock()
|
||||
|
||||
if len(Data.Cache.StreamingURLS) == 0 {
|
||||
|
||||
tmp, err := loadJSONFileToMap(System.File.URLS)
|
||||
|
||||
+1
-1
@@ -1166,7 +1166,7 @@ func createM3UFile() {
|
||||
ShowError(err, 000)
|
||||
}
|
||||
|
||||
saveMapToJSONFile(System.File.URLS, Data.Cache.StreamingURLS)
|
||||
saveMapToJSONFile(System.File.URLS, streamingURLsSnapshot())
|
||||
|
||||
}
|
||||
|
||||
|
||||
+7
-7
@@ -29,17 +29,17 @@ Status: Phase 0 committed on branch `improvements` 2026-09-26.
|
||||
- [x] Decide default for web auth on fresh installs (keep off; LAN only, decided 2026-09-25)
|
||||
|
||||
## Phase 2: Streaming stability
|
||||
- [ ] Race tests: two concurrent tuners against a fake TS server, run with `-race`
|
||||
- [ ] `*Playlist` with own mutex; remove stale store-backs
|
||||
- [ ] Atomic tuner reservation; clean `BufferClients` on force kill
|
||||
- [ ] RWMutex around `StreamingURLS`
|
||||
- [ ] Remove `defer` inside read loops (buffer, compression, imgcache)
|
||||
- [ ] Single external-process buffer (`exec.CommandContext` + request context) with ffmpeg and VLC argument builders; drop `CloseNotifier`
|
||||
- [x] Race tests in `src/buffer_test.go`: restream shares one provider connection, tuner limit atomic under 6 concurrent tunes, acquire/release contention, cleanup on last client; run with `-race`
|
||||
- [x] New `src/buffer_state.go`: one `bufferMu`, playlists by pointer, one `bufferStream` per stream (status, client count, error, cancel hook); downloaders keep a private working copy and publish through helpers. `BufferInformation`/`BufferClients`/`Lock` removed
|
||||
- [x] Tuner check and registration under one lock (`bufferAcquireStream`); last client out removes the stream, cancels its process and deletes its folder
|
||||
- [x] `streamingURLsMu` around `Data.Cache.StreamingURLS`; persisted from a snapshot
|
||||
- [x] No more `defer` inside the buffer read/redirect loops (compression done in Phase 1; imgcache in 2b)
|
||||
- [x] `thirdPartyBuffer` rewritten: `exec.CommandContext` cancelled when the last client leaves or the 20 s startup watchdog fires; `buildFFmpegArgs`/`buildVLCArgs` tested; no panic, `cmd.Start` checked, one file handle per segment. Client loop uses `r.Context()` instead of `CloseNotifier`
|
||||
- [ ] Real mutex for screen log
|
||||
- [ ] `http.Server` with timeouts; timeouts on all outbound clients
|
||||
- [ ] imgcache: download outside the lock
|
||||
- [ ] Atomic write helper (temp + fsync + rename)
|
||||
- [ ] Fix listed concrete bugs (xepg append, range mutation, log overflow, headers after WriteHeader, random notification eviction, unchecked assertions, API double write, WS struct reuse)
|
||||
- [x] Buffer: headers set before `WriteHeader`, bogus `Content-Length:` header gone, segments flushed to the client as they arrive. Others in 2b
|
||||
- [ ] Error hygiene: no-op closures, `os.Exit`/`panic` outside main, ignored results
|
||||
|
||||
## Phase 3: Build, embed, CI, Docker
|
||||
|
||||
Reference in New Issue
Block a user