diff --git a/.gitignore b/.gitignore index 4952588..5c08c6f 100644 --- a/.gitignore +++ b/.gitignore @@ -12,3 +12,4 @@ de.json agent.md skill.md node_modules/ +.claude/ diff --git a/src/buffer.go b/src/buffer.go index 5331362..40dd0cb 100644 --- a/src/buffer.go +++ b/src/buffer.go @@ -7,7 +7,7 @@ package src import ( "bufio" - "bytes" + "context" "errors" "fmt" "io" @@ -19,10 +19,11 @@ import ( "sort" "strconv" "strings" + "sync/atomic" "time" ) -func createStreamID(stream map[int]ThisStream) (streamID int) { +func createStreamID(stream map[int]*bufferStream) (streamID int) { var debug string @@ -46,363 +47,170 @@ func bufferingStream(playlistID, streamingURL, channelName string, w http.Respon time.Sleep(time.Duration(Settings.BufferTimeout) * time.Millisecond) - var playlist Playlist - var client ThisClient - var stream ThisStream - var streaming = false - var streamID int - var debug string - var timeOut = 0 - var newStream = true - - //w.Header().Set("Connection", "keep-alive") w.Header().Set("Connection", "close") - // Überprüfen ob die Playlist schon verwendet wird - if p, ok := BufferInformation.Load(playlistID); !ok { - - var playlistType string - // Playlist wird noch nicht verwendet, Default-Werte für die Playlist erstellen - playlist.Folder = System.Folder.Temp + playlistID + string(os.PathSeparator) - playlist.PlaylistID = playlistID - playlist.Streams = make(map[int]ThisStream) - playlist.Clients = make(map[int]ThisClient) - - err := checkFolder(playlist.Folder) - if err != nil { - ShowError(err, 000) - httpStatusError(w, r, 404) + streamID, isNew, err := bufferAcquireStream(playlistID, streamingURL, channelName) + if err != nil { + if errors.Is(err, errTunerLimit) { + serveTunerLimit(w, r) return } - - switch playlist.PlaylistID[0:1] { - - case "M": - playlistType = "m3u" - - case "H": - playlistType = "hdhr" - - } - - playlist.Tuner = getTuner(playlistID, playlistType) - - playlist.PlaylistName = getProviderParameter(playlist.PlaylistID, playlistType, "name") - - // Default-Werte für den Stream erstellen - streamID = createStreamID(playlist.Streams) - - client.Connection = 1 - stream.URL = streamingURL - stream.ChannelName = channelName - stream.Status = false - - playlist.Streams[streamID] = stream - playlist.Clients[streamID] = client - - BufferInformation.Store(playlistID, playlist) - - } else { - - // Playlist wird bereits zum streamen verwendet - // Überprüfen ob die URL bereit von einem anderen Client gestreamt wird. - - playlist = p.(Playlist) - - for id := range playlist.Streams { - - stream = playlist.Streams[id] - client = playlist.Clients[id] - - if streamingURL == stream.URL { - - streamID = id - newStream = false - client.Connection++ - - //playlist.Streams[streamID] = stream - playlist.Clients[streamID] = client - - BufferInformation.Store(playlistID, playlist) - - debug = fmt.Sprintf("Restream Status:Playlist: %s - Channel: %s - Connections: %d", playlist.PlaylistName, stream.ChannelName, client.Connection) - - showDebug(debug, 1) - - if c, ok := BufferClients.Load(playlistID + stream.MD5); ok { - - var clients = c.(ClientConnection) - clients.Connection = clients.Connection + 1 - showInfo(fmt.Sprintf("Streaming Status:Channel: %s (Clients: %d)", stream.ChannelName, clients.Connection)) - - BufferClients.Store(playlistID+stream.MD5, clients) - - } - - break - } - - } - - // Neuer Stream bei einer bereits aktiven Playlist - if newStream { - - // Prüfen ob die Playlist noch einen weiteren Stream erlaubt (Tuner) - if len(playlist.Streams) >= playlist.Tuner { - - showInfo(fmt.Sprintf("Streaming Status:Playlist: %s - No new connections available. Tuner = %d", playlist.PlaylistName, playlist.Tuner)) - - if content, err := readWebFile("video/stream-limit.ts"); err == nil { - - w.WriteHeader(200) - w.Header().Set("Content-type", "video/mpeg") - w.Header().Set("Content-Length:", "0") - - for i := 1; i < 60; i++ { - _ = i - w.Write([]byte(content)) - time.Sleep(time.Duration(500) * time.Millisecond) - } - - return - } - - return - } - - // Playlist erlaubt einen weiterern Stream (Das Limit des Tuners ist noch nicht erreicht) - // Default-Werte für den Stream erstellen - stream = ThisStream{} - client = ThisClient{} - - streamID = createStreamID(playlist.Streams) - - client.Connection = 1 - stream.URL = streamingURL - stream.ChannelName = channelName - stream.Status = false - - playlist.Streams[streamID] = stream - playlist.Clients[streamID] = client - - BufferInformation.Store(playlistID, playlist) - - } - + ShowError(err, 000) + httpStatusError(w, r, 404) + return } + defer bufferReleaseClient(playlistID, streamID) - // Überprüfen ob der Stream breits von einem anderen Client abgespielt wird - if !playlist.Streams[streamID].Status && newStream { - - // Neuer Buffer wird benötigt - stream = playlist.Streams[streamID] - stream.MD5 = getMD5(streamingURL) - stream.Folder = playlist.Folder + stream.MD5 + string(os.PathSeparator) - stream.PlaylistID = playlistID - stream.PlaylistName = playlist.PlaylistName - - playlist.Streams[streamID] = stream - BufferInformation.Store(playlistID, playlist) + if isNew { switch Settings.Buffer { case "xteve": go connectToStreamingServer(streamID, playlistID) + case "ffmpeg", "vlc": go thirdPartyBuffer(streamID, playlistID) - default: - break - } - showInfo(fmt.Sprintf("Streaming Status:Playlist: %s - Tuner: %d / %d", playlist.PlaylistName, len(playlist.Streams), playlist.Tuner)) - - var clients ClientConnection - clients.Connection = 1 - BufferClients.Store(playlistID+stream.MD5, clients) - } - w.WriteHeader(200) + var ctx = r.Context() + var stream bufferStream - for { // Loop 1: Warten bis das erste Segment durch den Buffer heruntergeladen wurde + // Loop 1: wait until the downloader has written the first segment. + for waited := 0; ; waited++ { - if p, ok := BufferInformation.Load(playlistID); ok { + s, ok := bufferStreamSnapshot(playlistID, streamID) + if !ok { + return + } + if s.Err != nil { + ShowError(s.Err, 0) + return + } + if s.Status { + stream = s + break + } + if waited > 200 { + showInfo("Streaming Status:No data received from the buffer, giving up") + return + } - var playlist = p.(Playlist) + select { + case <-ctx.Done(): + return + case <-time.After(100 * time.Millisecond): + } - if stream, ok := playlist.Streams[streamID]; ok { + } - if !stream.Status { + w.Header().Set("Content-Type", "video/mpeg") + w.WriteHeader(http.StatusOK) + flusher, _ := w.(http.Flusher) - timeOut++ + var seen []string // segment files already handed to this client + var oldSegments []string // kept on disk for a while, oldest first - time.Sleep(time.Duration(100) * time.Millisecond) + // Loop 2: send finished segments as they appear. + for { - if c, ok := BufferClients.Load(playlistID + stream.MD5); ok { + select { + case <-ctx.Done(): + return + default: + } - var clients = c.(ClientConnection) + s, ok := bufferStreamSnapshot(playlistID, streamID) + if !ok { + return + } + if s.Err != nil { + ShowError(s.Err, 0) + return + } - if clients.Error != nil || timeOut > 200 { - killClientConnection(streamID, stream.PlaylistID, false) - return - } + if _, err := os.Stat(stream.Folder); os.IsNotExist(err) { + return + } - } + var tmpFiles = getTmpFiles(stream.Folder, &seen) - continue - } + for _, f := range tmpFiles { - var oldSegments []string - - for { // Loop 2: Temporäre Datein sind vorhanden, Daten können zum Client gesendet werden - - // HTTP Clientverbindung überwachen - - //lint:ignore SA1019 replaced by request context in Phase 2 - cn, ok := w.(http.CloseNotifier) - if ok { - - select { - - case <-cn.CloseNotify(): - killClientConnection(streamID, playlistID, false) - return - - default: - if c, ok := BufferClients.Load(playlistID + stream.MD5); ok { - - var clients = c.(ClientConnection) - if clients.Error != nil { - ShowError(clients.Error, 0) - killClientConnection(streamID, playlistID, false) - return - } - - } else { - - return - - } - - } - - } - - if _, err := os.Stat(stream.Folder); os.IsNotExist(err) { - killClientConnection(streamID, playlistID, false) - return - } - - var tmpFiles = getTmpFiles(&stream) - //fmt.Println("Buffer Loop:", stream.Connection) - - for _, f := range tmpFiles { - - if _, err := os.Stat(stream.Folder); os.IsNotExist(err) { - killClientConnection(streamID, playlistID, false) - return - } - - oldSegments = append(oldSegments, f) - - var fileName = stream.Folder + f - - file, err := os.Open(fileName) - - if err == nil { - - defer file.Close() - - l, err := file.Stat() - if err == nil { - - debug = fmt.Sprintf("Buffer Status:Send to client (%s)", fileName) - showDebug(debug, 2) - - var buffer = make([]byte, int(l.Size())) - _, err = file.Read(buffer) - - if err == nil { - - file.Seek(0, 0) - - if !streaming { - - contentType := http.DetectContentType(buffer) - _ = contentType - //w.Header().Set("Content-type", "video/mpeg") - w.Header().Set("Content-type", contentType) - w.Header().Set("Content-Length", "0") - w.Header().Set("Connection", "close") - - } - - /* - // HDHR Header - w.Header().Set("Cache-Control", "no-cache") - w.Header().Set("Pragma", "no-cache") - w.Header().Set("transferMode.dlna.org", "Streaming") - */ - - _, err := w.Write(buffer) - - if err != nil { - file.Close() - killClientConnection(streamID, playlistID, false) - return - } - - file.Close() - streaming = true - - } - - file.Close() - - } - - var n = indexOfString(f, oldSegments) - - if n > 20 { - - var fileToRemove = stream.Folder + oldSegments[0] - os.RemoveAll(getPlatformFile(fileToRemove)) - oldSegments = append(oldSegments[:0], oldSegments[0+1:]...) - - } - - } - - file.Close() - - } - - if len(tmpFiles) == 0 { - time.Sleep(time.Duration(100) * time.Millisecond) - } - - } // Ende Loop 2 - - } else { - - // Stream nicht vorhanden - killClientConnection(streamID, stream.PlaylistID, false) - showInfo(fmt.Sprintf("Streaming Status:Playlist: %s - Tuner: %d / %d", playlist.PlaylistName, len(playlist.Streams), playlist.Tuner)) - return + var fileName = stream.Folder + f + content, err := os.ReadFile(fileName) + if err != nil { + continue } - } // Ende BufferInformation + showDebug(fmt.Sprintf("Buffer Status:Send to client (%s)", fileName), 2) - } // Ende Loop 1 + if _, err := w.Write(content); err != nil { + return + } + if flusher != nil { + flusher.Flush() + } + + oldSegments = append(oldSegments, f) + if len(oldSegments) > 20 { + os.RemoveAll(getPlatformFile(stream.Folder + oldSegments[0])) + oldSegments = oldSegments[1:] + } + + } + + if len(tmpFiles) == 0 { + select { + case <-ctx.Done(): + return + case <-time.After(100 * time.Millisecond): + } + } + + } } -func getTmpFiles(stream *ThisStream) (tmpFiles []string) { +// serveTunerLimit : plays the "tuner limit reached" clip for up to 30 seconds. +func serveTunerLimit(w http.ResponseWriter, r *http.Request) { + + content, err := readWebFile("video/stream-limit.ts") + if err != nil { + httpStatusError(w, r, 503) + return + } + + w.Header().Set("Content-Type", "video/mpeg") + w.WriteHeader(http.StatusOK) + flusher, _ := w.(http.Flusher) + + for i := 0; i < 60; i++ { + + if _, err := w.Write(content); err != nil { + return + } + if flusher != nil { + flusher.Flush() + } + + select { + case <-r.Context().Done(): + return + case <-time.After(500 * time.Millisecond): + } + + } + +} + +// getTmpFiles : finished segment files in tmpFolder not yet listed in *seen, +// oldest first. The newest file is skipped because it is still being written. +func getTmpFiles(tmpFolder string, seen *[]string) (tmpFiles []string) { - var tmpFolder = stream.Folder var fileIDs []float64 if _, err := os.Stat(tmpFolder); !os.IsNotExist(err) { @@ -433,9 +241,9 @@ func getTmpFiles(stream *ThisStream) (tmpFiles []string) { var fileName = fmt.Sprintf("%d.ts", int64(file)) - if indexOfString(fileName, stream.OldSegments) == -1 { + if indexOfString(fileName, *seen) == -1 { tmpFiles = append(tmpFiles, fileName) - stream.OldSegments = append(stream.OldSegments, fileName) + *seen = append(*seen, fileName) } } @@ -447,116 +255,27 @@ func getTmpFiles(stream *ThisStream) (tmpFiles []string) { return } -func killClientConnection(streamID int, playlistID string, force bool) { - - Lock.Lock() - defer Lock.Unlock() - - if p, ok := BufferInformation.Load(playlistID); ok { - - var playlist = p.(Playlist) - - if force { - delete(playlist.Streams, streamID) - showInfo(fmt.Sprintf("Streaming Status:Playlist: %s - Tuner: %d / %d", playlist.PlaylistName, len(playlist.Streams), playlist.Tuner)) - return - } - - if stream, ok := playlist.Streams[streamID]; ok { - - if c, ok := BufferClients.Load(playlistID + stream.MD5); ok { - - var clients = c.(ClientConnection) - clients.Connection = clients.Connection - 1 - BufferClients.Store(playlistID+stream.MD5, clients) - - showInfo("Streaming Status:Client has terminated the connection") - showInfo(fmt.Sprintf("Streaming Status:Channel: %s (Clients: %d)", stream.ChannelName, clients.Connection)) - - if clients.Connection <= 0 { - BufferClients.Delete(playlistID + stream.MD5) - delete(playlist.Streams, streamID) - delete(playlist.Clients, streamID) - } - - } - - BufferInformation.Store(playlistID, playlist) - - if len(playlist.Streams) > 0 { - showInfo(fmt.Sprintf("Streaming Status:Playlist: %s - Tuner: %d / %d", playlist.PlaylistName, len(playlist.Streams), playlist.Tuner)) - } - - } - - } - -} - -func clientConnection(stream ThisStream) (status bool) { - - status = true - Lock.Lock() - defer Lock.Unlock() - - if _, ok := BufferClients.Load(stream.PlaylistID + stream.MD5); !ok { - - var debug = fmt.Sprintf("Streaming Status:Remove temporary files (%s)", stream.Folder) - showDebug(debug, 1) - - status = false - - debug = fmt.Sprintf("Remove tmp folder:%s", stream.Folder) - showDebug(debug, 1) - - os.RemoveAll(stream.Folder) - - if p, ok := BufferInformation.Load(stream.PlaylistID); ok { - - showInfo(fmt.Sprintf("Streaming Status:Channel: %s - No client is using this channel anymore. Streaming Server connection has ended", stream.ChannelName)) - - var playlist = p.(Playlist) - - showInfo(fmt.Sprintf("Streaming Status:Playlist: %s - Tuner: %d / %d", playlist.PlaylistName, len(playlist.Streams), playlist.Tuner)) - - if len(playlist.Streams) <= 0 { - BufferInformation.Delete(stream.PlaylistID) - } - - } - - status = false - - } - - return -} - -// addErrorToStream : Fehler an die Clientverbindung des Streams weitergeben -func addErrorToStream(playlist Playlist, streamID int, playlistID string, err error) { - - var stream = playlist.Streams[streamID] - - if c, ok := BufferClients.Load(playlistID + stream.MD5); ok { - - var clients = c.(ClientConnection) - clients.Error = err - BufferClients.Store(playlistID+stream.MD5, clients) - - } - -} - func connectToStreamingServer(streamID int, playlistID string) { - if p, ok := BufferInformation.Load(playlistID); ok { + snapshot, ok := bufferStreamSnapshot(playlistID, streamID) + if !ok { + return + } - var playlist = p.(Playlist) + // Private working copy; shared state is published via buffer_state.go. + var stream ThisStream + stream.URL = snapshot.URL + stream.ChannelName = snapshot.ChannelName + stream.Folder = snapshot.Folder + stream.MD5 = snapshot.MD5 + stream.PlaylistID = snapshot.PlaylistID + stream.PlaylistName = snapshot.PlaylistName + { var timeOut = 0 var debug string var tmpSegment = 1 - var tmpFolder = playlist.Streams[streamID].Folder + var tmpFolder = stream.Folder var m3u8Segments []string var bandwidth BandwidthCalculation var networkBandwidth = Settings.M3U8AdaptiveBandwidthMBPS * 1e+6 @@ -568,15 +287,14 @@ func connectToStreamingServer(streamID int, playlistID string) { var segment Segment - if len(playlist.Streams[streamID].Location) > 0 { - segment.URL = playlist.Streams[streamID].Location + if len(stream.Location) > 0 { + segment.URL = stream.Location } else { - segment.URL = playlist.Streams[streamID].URL + segment.URL = stream.URL } segment.Duration = 0 - var stream = playlist.Streams[streamID] stream.Segment = []Segment{} stream.Segment = append(stream.Segment, segment) @@ -585,8 +303,6 @@ func connectToStreamingServer(streamID int, playlistID string) { stream.Wait = 0 stream.NetworkBandwidth = networkBandwidth - playlist.Streams[streamID] = stream - timeOut++ } @@ -596,7 +312,7 @@ func connectToStreamingServer(streamID int, playlistID string) { err := checkFolder(tmpFolder) if err != nil { ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) + bufferSetError(playlistID, streamID, err) return } @@ -611,8 +327,6 @@ func connectToStreamingServer(streamID int, playlistID string) { return } - var stream = playlist.Streams[streamID] - if !stream.Status { if strings.Contains(stream.URL, ".m3u8") { @@ -633,7 +347,7 @@ func connectToStreamingServer(streamID int, playlistID string) { for { - if !clientConnection(stream) { + if !bufferStreamActive(playlistID, streamID) { return } @@ -658,9 +372,7 @@ func connectToStreamingServer(streamID int, playlistID string) { req, err := http.NewRequest("GET", currentURL, nil) if err != nil { ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) - killClientConnection(streamID, playlistID, true) - clientConnection(stream) + bufferSetError(playlistID, streamID, err) return } @@ -670,12 +382,7 @@ func connectToStreamingServer(streamID int, playlistID string) { req.Header.Set("Accept", "*/*") debugRequest(req) - client := &http.Client{} - client.CheckRedirect = func(req *http.Request, via []*http.Request) error { - return errors.New("Redirect") - } - - resp, err := client.Do(req) + resp, err := streamHTTPClient.Do(req) if resp != nil && err != nil { debugResponse(resp) @@ -689,10 +396,7 @@ func connectToStreamingServer(streamID int, playlistID string) { fmt.Println("Current URL:", currentURL) ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) - - killClientConnection(streamID, playlistID, true) - clientConnection(stream) + bufferSetError(playlistID, streamID, err) return } @@ -711,17 +415,15 @@ func connectToStreamingServer(streamID int, playlistID string) { debug = fmt.Sprintf("HTTP Redirect:%s", stream.Location) showDebug(debug, 2) - defer resp.Body.Close() + resp.Body.Close() goto Redirect } else { err = errors.New("Streaming server") ShowError(err, 4002) - addErrorToStream(playlist, streamID, playlistID, err) - - defer resp.Body.Close() - + bufferSetError(playlistID, streamID, err) + resp.Body.Close() return } @@ -729,18 +431,14 @@ func connectToStreamingServer(streamID int, playlistID string) { } else { ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) - - defer resp.Body.Close() - + bufferSetError(playlistID, streamID, err) + resp.Body.Close() return } } - defer resp.Body.Close() - // HTTP Status überprüfen, bei Fehlern wird der Stream beendet var contentType = resp.Header.Get("Content-Type") var httpStatusCode = resp.StatusCode @@ -755,14 +453,11 @@ func connectToStreamingServer(streamID int, playlistID string) { var err = errors.New(http.StatusText(resp.StatusCode)) ShowError(err, 4004) - debug = fmt.Sprintf("Streaming Status:Playlist: %s - Tuner: %d / %d", playlist.PlaylistName, len(playlist.Streams), playlist.Tuner) + debug = "Streaming Status:" + bufferTunerStatus(playlistID) showDebug(debug, 1) - BufferInformation.Store(playlist.PlaylistID, playlist) - addErrorToStream(playlist, streamID, playlistID, err) + bufferSetError(playlistID, streamID, err) - killClientConnection(streamID, playlistID, true) - clientConnection(stream) resp.Body.Close() return @@ -815,7 +510,7 @@ func connectToStreamingServer(streamID int, playlistID string) { body, err := io.ReadAll(resp.Body) if err != nil { ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) + bufferSetError(playlistID, streamID, err) } stream.Body = string(body) @@ -825,7 +520,7 @@ func connectToStreamingServer(streamID int, playlistID string) { err = parseM3U8(&stream) if err != nil { ShowError(err, 4050) - addErrorToStream(playlist, streamID, playlistID, err) + bufferSetError(playlistID, streamID, err) } // Video Stream (TS) @@ -845,19 +540,16 @@ func connectToStreamingServer(streamID int, playlistID string) { var tmpFile = fmt.Sprintf("%s%d.ts", tmpFolder, tmpSegment) - if !clientConnection(stream) { + if !bufferStreamActive(playlistID, streamID) { resp.Body.Close() return } bufferFile, err := os.Create(tmpFile) if err != nil { - - addErrorToStream(playlist, streamID, playlistID, err) - bufferFile.Close() + bufferSetError(playlistID, streamID, err) resp.Body.Close() return - } for { @@ -874,30 +566,24 @@ func connectToStreamingServer(streamID int, playlistID string) { n, err := resp.Body.Read(buffer) if err != nil && err != io.EOF { - ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) + bufferSetError(playlistID, streamID, err) + bufferFile.Close() resp.Body.Close() return - } - defer resp.Body.Close() - if _, err := bufferFile.Write(buffer[:n]); err != nil { - ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) + bufferSetError(playlistID, streamID, err) + bufferFile.Close() resp.Body.Close() return - } - defer bufferFile.Close() - fileSize = fileSize + n - if !clientConnection(stream) { + if !bufferStreamActive(playlistID, streamID) { resp.Body.Close() bufferFile.Close() @@ -913,8 +599,6 @@ func connectToStreamingServer(streamID int, playlistID string) { // Buffer auf die Festplatte speichern if fileSize >= tmpFileSize/2 || n == 0 { - Lock.Lock() - bandwidth.Stop = time.Now() bandwidth.Size += fileSize @@ -931,15 +615,13 @@ func connectToStreamingServer(streamID int, playlistID string) { bufferFile.Close() stream.Status = true - playlist.Streams[streamID] = stream - BufferInformation.Store(playlistID, playlist) - Lock.Unlock() + bufferSetStatus(playlistID, streamID, true) tmpSegment++ tmpFile = fmt.Sprintf("%s%d.ts", tmpFolder, tmpSegment) - if !clientConnection(stream) { + if !bufferStreamActive(playlistID, streamID) { bufferFile.Close() resp.Body.Close() @@ -954,7 +636,7 @@ func connectToStreamingServer(streamID int, playlistID string) { bufferFile, err = os.Create(tmpFile) if err != nil { - addErrorToStream(playlist, streamID, playlistID, err) + bufferSetError(playlistID, streamID, err) resp.Body.Close() return } @@ -980,7 +662,7 @@ func connectToStreamingServer(streamID int, playlistID string) { err = errors.New("Streaming error") ShowError(err, 4003) - addErrorToStream(playlist, streamID, playlistID, err) + bufferSetError(playlistID, streamID, err) resp.Body.Close() return } @@ -1031,7 +713,7 @@ func connectToStreamingServer(streamID int, playlistID string) { } // Ende for loop - } // Ende BufferInformation + } } @@ -1346,293 +1028,233 @@ func switchBandwidth(stream *ThisStream) (err error) { } // Buffer mit FFMPEG +// buildFFmpegArgs : ffmpeg command line from the configured options. +func buildFFmpegArgs(options, url, userAgent string) (args []string) { + + if len(userAgent) != 0 { + args = []string{"-user_agent", userAgent} + } + + for _, a := range strings.Fields(options) { + args = append(args, strings.Replace(a, "[URL]", url, -1)) + } + + return +} + +// buildVLCArgs : VLC / cvlc command line from the configured options. +func buildVLCArgs(options, url, userAgent string) (args []string) { + + for _, a := range strings.Fields(options) { + + if a == "[URL]" { + args = append(args, url) + if len(userAgent) != 0 { + args = append(args, ":http-user-agent="+userAgent) + } + continue + } + + args = append(args, a) + } + + return +} + +// thirdPartyBuffer : runs ffmpeg or VLC and writes its output into segment +// files. The process lives in a context that is cancelled when the last +// client leaves (bufferReleaseClient) or when no usable data arrives within +// 20 seconds of starting. func thirdPartyBuffer(streamID int, playlistID string) { - if p, ok := BufferInformation.Load(playlistID); ok { + snapshot, ok := bufferStreamSnapshot(playlistID, streamID) + if !ok { + return + } - var playlist = p.(Playlist) - var debug, path, options, bufferType string - var tmpSegment = 1 - var bufferSize = Settings.BufferSize * 1024 - var stream = playlist.Streams[streamID] - var buf bytes.Buffer - var fileSize = 0 - var streamStatus = make(chan bool) + var path, options, bufferType string - var tmpFolder = playlist.Streams[streamID].Folder - var url = playlist.Streams[streamID].URL + switch Settings.Buffer { - stream.Status = false + case "ffmpeg": + path, options, bufferType = Settings.FFmpegPath, Settings.FFmpegOptions, "FFMPEG" - bufferType = strings.ToUpper(Settings.Buffer) - - switch Settings.Buffer { - - case "ffmpeg": - path = Settings.FFmpegPath - options = Settings.FFmpegOptions - - case "vlc": - path = Settings.VLCPath - options = Settings.VLCOptions - - default: - return - } - - os.RemoveAll(getPlatformPath(tmpFolder)) - - err := checkFolder(tmpFolder) - if err != nil { - ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) - return - } - - err = checkStreamingBinary(path) - if err != nil { - ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) - return - } - - err = checkStreamURL(url) - if err != nil { - ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) - return - } - - showInfo(fmt.Sprintf("%s path:%s", bufferType, path)) - showInfo("Streaming URL:" + stream.URL) - - var tmpFile = fmt.Sprintf("%s%d.ts", tmpFolder, tmpSegment) - - f, err := os.Create(tmpFile) - f.Close() - if err != nil { - addErrorToStream(playlist, streamID, playlistID, err) - return - } - - //args = strings.Replace(args, "[USER-AGENT]", Settings.UserAgent, -1) - - // User-Agent setzen - var args []string - - var userAgent = getUserAgent() - - for i, a := range strings.Split(options, " ") { - - switch bufferType { - case "FFMPEG": - a = strings.Replace(a, "[URL]", url, -1) - if i == 0 { - if len(userAgent) != 0 { - args = []string{"-user_agent", userAgent} - } - } - - args = append(args, a) - - case "VLC": - if a == "[URL]" { - a = strings.Replace(a, "[URL]", url, -1) - args = append(args, a) - - if len(userAgent) != 0 { - args = append(args, fmt.Sprintf(":http-user-agent=%s", userAgent)) - } - - } else { - args = append(args, a) - } - - } - - } - - var cmd = exec.Command(path, args...) - - debug = fmt.Sprintf("%s:%s %s", bufferType, path, args) - showDebug(debug, 1) - - // Byte-Daten vom Prozess - stdOut, err := cmd.StdoutPipe() - if err != nil { - ShowError(err, 0) - cmd.Process.Kill() - cmd.Wait() - addErrorToStream(playlist, streamID, playlistID, err) - return - } - - // Log-Daten vom Prozess - logOut, err := cmd.StderrPipe() - if err != nil { - ShowError(err, 0) - cmd.Process.Kill() - cmd.Wait() - addErrorToStream(playlist, streamID, playlistID, err) - return - } - - if len(buf.Bytes()) == 0 && !stream.Status { - showInfo(bufferType + ":Processing data") - } - - cmd.Start() - defer cmd.Wait() - - go func() { - - // Log Daten vom Prozess im Dubug Mode 1 anzeigen. - scanner := bufio.NewScanner(logOut) - scanner.Split(bufio.ScanLines) - - for scanner.Scan() { - - debug = fmt.Sprintf("%s log:%s", bufferType, strings.TrimSpace(scanner.Text())) - - select { - case <-streamStatus: - showDebug(debug, 1) - default: - showInfo(debug) - } - - time.Sleep(time.Duration(10) * time.Millisecond) - - } - - }() - - f, err = os.OpenFile(tmpFile, os.O_APPEND|os.O_WRONLY, 0600) - if err != nil { - panic(err) - } - defer f.Close() - - buffer := make([]byte, 1024*4) - - reader := bufio.NewReader(stdOut) - - t := make(chan int) - - go func() { - - var timeout = 0 - for { - time.Sleep(time.Duration(1000) * time.Millisecond) - timeout++ - - select { - case <-t: - return - default: - t <- timeout - } - - } - - }() - - for { - - select { - case timeout := <-t: - if timeout >= 20 && tmpSegment == 1 { - cmd.Process.Kill() - err = errors.New("Timout") - ShowError(err, 4006) - addErrorToStream(playlist, streamID, playlistID, err) - cmd.Wait() - f.Close() - return - } - - default: - - } - - if fileSize == 0 && !stream.Status { - showInfo("Streaming Status:Receive data from " + bufferType) - } - - if !clientConnection(stream) { - cmd.Process.Kill() - f.Close() - cmd.Wait() - return - } - - n, err := reader.Read(buffer) - if err == io.EOF { - break - } - - fileSize = fileSize + len(buffer[:n]) - - if _, err := f.Write(buffer[:n]); err != nil { - cmd.Process.Kill() - ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) - cmd.Wait() - f.Close() - return - } - - if fileSize >= bufferSize/2 { - - if tmpSegment == 1 && !stream.Status { - close(t) - close(streamStatus) - showInfo(fmt.Sprintf("Streaming Status:Buffering data from %s", bufferType)) - } - - f.Close() - tmpSegment++ - - if !stream.Status { - Lock.Lock() - stream.Status = true - playlist.Streams[streamID] = stream - BufferInformation.Store(playlistID, playlist) - Lock.Unlock() - } - - tmpFile = fmt.Sprintf("%s%d.ts", tmpFolder, tmpSegment) - - fileSize = 0 - - f, err = os.OpenFile(tmpFile, os.O_CREATE|os.O_TRUNC|os.O_APPEND|os.O_WRONLY, 0666) - if err != nil { - cmd.Process.Kill() - ShowError(err, 0) - addErrorToStream(playlist, streamID, playlistID, err) - cmd.Wait() - f.Close() - return - } - - } - - } - - cmd.Process.Kill() - cmd.Wait() - - err = errors.New(bufferType + " error") - addErrorToStream(playlist, streamID, playlistID, err) - ShowError(err, 1204) - - time.Sleep(time.Duration(500) * time.Millisecond) - clientConnection(stream) + case "vlc": + path, options, bufferType = Settings.VLCPath, Settings.VLCOptions, "VLC" + default: return } + var fail = func(err error, code int) { + ShowError(err, code) + bufferSetError(playlistID, streamID, err) + } + + var tmpFolder = snapshot.Folder + + os.RemoveAll(getPlatformPath(tmpFolder)) + + if err := checkFolder(tmpFolder); err != nil { + fail(err, 0) + return + } + + if err := checkStreamingBinary(path); err != nil { + fail(err, 0) + return + } + + if err := checkStreamURL(snapshot.URL); err != nil { + fail(err, 0) + return + } + + showInfo(fmt.Sprintf("%s path:%s", bufferType, path)) + showInfo("Streaming URL:" + snapshot.URL) + + var args []string + switch bufferType { + case "FFMPEG": + args = buildFFmpegArgs(options, snapshot.URL, getUserAgent()) + case "VLC": + args = buildVLCArgs(options, snapshot.URL, getUserAgent()) + } + + ctx, cancel := context.WithCancel(context.Background()) + bufferSetCancel(playlistID, streamID, cancel) + + var cmd = exec.CommandContext(ctx, path, args...) + + showDebug(fmt.Sprintf("%s:%s %s", bufferType, path, args), 1) + + stdOut, err := cmd.StdoutPipe() + if err != nil { + cancel() + fail(err, 0) + return + } + + logOut, err := cmd.StderrPipe() + if err != nil { + cancel() + fail(err, 0) + return + } + + showInfo(bufferType + ":Processing data") + + if err := cmd.Start(); err != nil { + cancel() + fail(err, 0) + return + } + + // Cancel first so Wait cannot block on a process that is still running. + defer func() { + cancel() + cmd.Wait() + }() + + var started atomic.Bool // first segment written, clients may read + + go func() { + scanner := bufio.NewScanner(logOut) + for scanner.Scan() { + var line = fmt.Sprintf("%s log:%s", bufferType, strings.TrimSpace(scanner.Text())) + if started.Load() { + showDebug(line, 1) + } else { + showInfo(line) + } + } + }() + + // Startup watchdog: nothing usable within 20 s means the process is stuck. + var watchdog = time.AfterFunc(20*time.Second, func() { + if !started.Load() { + fail(errors.New("Timeout"), 4006) + cancel() + } + }) + defer watchdog.Stop() + + var segmentSize = Settings.BufferSize * 1024 / 2 + var tmpSegment = 1 + var fileSize = 0 + + var openSegment = func() (*os.File, error) { + return os.OpenFile(fmt.Sprintf("%s%d.ts", tmpFolder, tmpSegment), os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0600) + } + + f, err := openSegment() + if err != nil { + fail(err, 0) + return + } + + showInfo("Streaming Status:Receive data from " + bufferType) + + var buffer = make([]byte, 1024*4) + var reader = bufio.NewReader(stdOut) + + for { + + if !bufferStreamActive(playlistID, streamID) { + f.Close() + return + } + + n, readErr := reader.Read(buffer) + + if n > 0 { + + if _, err := f.Write(buffer[:n]); err != nil { + f.Close() + fail(err, 0) + return + } + + fileSize += n + + if fileSize >= segmentSize { + + f.Close() + + if !started.Load() { + started.Store(true) + watchdog.Stop() + showInfo(fmt.Sprintf("Streaming Status:Buffering data from %s", bufferType)) + bufferSetStatus(playlistID, streamID, true) + } + + tmpSegment++ + fileSize = 0 + + if f, err = openSegment(); err != nil { + fail(err, 0) + return + } + + } + + } + + if readErr != nil { + break + } + + } + + f.Close() + + // Cancelled by us (client gone or watchdog): nothing more to report. + if ctx.Err() != nil { + return + } + + fail(errors.New(bufferType+" error"), 1204) } func getTuner(id, playlistType string) (tuner int) { diff --git a/src/buffer_state.go b/src/buffer_state.go new file mode 100644 index 0000000..1a307ba --- /dev/null +++ b/src/buffer_state.go @@ -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 +} diff --git a/src/buffer_test.go b/src/buffer_test.go new file mode 100644 index 0000000..48d6703 --- /dev/null +++ b/src/buffer_test.go @@ -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 diff --git a/src/config.go b/src/config.go index 6c819fe..6073c76 100644 --- a/src/config.go +++ b/src/config.go @@ -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) { diff --git a/src/hdhr.go b/src/hdhr.go index 2cce878..830ba5f 100644 --- a/src/hdhr.go +++ b/src/hdhr.go @@ -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 } diff --git a/src/struct-buffer.go b/src/struct-buffer.go index 68332d6..ca16dad 100644 --- a/src/struct-buffer.go +++ b/src/struct-buffer.go @@ -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 diff --git a/src/system.go b/src/system.go index d039333..241c8b0 100644 --- a/src/system.go +++ b/src/system.go @@ -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) diff --git a/src/xepg.go b/src/xepg.go index b3eb318..ccdffeb 100644 --- a/src/xepg.go +++ b/src/xepg.go @@ -1166,7 +1166,7 @@ func createM3UFile() { ShowError(err, 000) } - saveMapToJSONFile(System.File.URLS, Data.Cache.StreamingURLS) + saveMapToJSONFile(System.File.URLS, streamingURLsSnapshot()) } diff --git a/tasks/todo.md b/tasks/todo.md index d5590a7..745b63c 100644 --- a/tasks/todo.md +++ b/tasks/todo.md @@ -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