diff --git a/pkg/media/mkv_ingest.go b/pkg/media/mkv_ingest.go index 66f3ce9e3..2f5c1a974 100644 --- a/pkg/media/mkv_ingest.go +++ b/pkg/media/mkv_ingest.go @@ -5,14 +5,33 @@ import ( "fmt" "io" "strings" + "time" "github.com/go-gst/go-gst/gst" "github.com/go-gst/go-gst/gst/app" + "stream.place/streamplace/pkg/aqtime" "stream.place/streamplace/pkg/log" ) // ingest a H264+AAC MKV stream (prolly from an RTMP server) func (mm *MediaManager) MKVIngest(ctx context.Context, input io.Reader, ms MediaSigner) error { + shouldRecord, err := mm.shouldRecord(ctx, ms.Streamer()) + if err != nil { + return err + } + if shouldRecord { + log.Log(ctx, "recording RTMP stream to file", "streamer", ms.Streamer()) + pr, pw := io.Pipe() + input = io.TeeReader(input, pw) + go func() { + err := mm.dumpToFile(ctx, pr, ms.Streamer(), ".rtmp.mkv") + if err != nil { + log.Error(ctx, "error dumping to file", "error", err) + } + }() + } else { + log.Log(ctx, "not recording RTMP stream to file", "streamer", ms.Streamer()) + } ctx, cancel := context.WithCancel(ctx) defer cancel() pipelineSlice := []string{ @@ -86,3 +105,18 @@ func (mm *MediaManager) MKVIngest(ctx context.Context, input io.Reader, ms Media return nil } + +func (mm *MediaManager) dumpToFile(ctx context.Context, r io.Reader, user string, filesuffix string) error { + now := aqtime.FromTime(time.Now()) + filename := fmt.Sprintf("%s%s", now.FileSafeString(), filesuffix) + f, err := mm.cli.DataFileCreate([]string{"debug-recordings", user, filename}, false) + if err != nil { + return fmt.Errorf("failed to create data file: %w", err) + } + defer f.Close() + _, err = io.Copy(f, r) + if err != nil { + return fmt.Errorf("failed to copy to file: %w", err) + } + return nil +} diff --git a/pkg/media/peer_connection.go b/pkg/media/peer_connection.go index 6b2b3fefd..de68767d8 100644 --- a/pkg/media/peer_connection.go +++ b/pkg/media/peer_connection.go @@ -7,23 +7,29 @@ import ( "stream.place/streamplace/pkg/rtcrec" ) -func (mm *MediaManager) NewPeerConnection(ctx context.Context, user string) (rtcrec.PeerConnection, error) { +func (mm *MediaManager) shouldRecord(ctx context.Context, user string) (bool, error) { shouldRecord := false settings, err := mm.model.GetServerSettings(ctx, mm.cli.BroadcasterHost, user) if err != nil { - return nil, err + return false, err } if settings != nil { spsettings, err := settings.ToStreamplaceServerSettings() if err != nil { - return nil, err + return false, err } if spsettings.DebugRecording != nil { shouldRecord = *spsettings.DebugRecording } } - if !shouldRecord { - log.Warn(ctx, "no server settings found, will not record") + return shouldRecord, nil +} + +func (mm *MediaManager) NewPeerConnection(ctx context.Context, user string) (rtcrec.PeerConnection, error) { + ctx = log.WithLogValues(ctx, "func", "NewPeerConnection", "streamer", user) + shouldRecord, err := mm.shouldRecord(ctx, user) + if err != nil { + return nil, err } pionpc, err := mm.webrtcAPI.NewPeerConnection(mm.webrtcConfig) if err != nil { diff --git a/pkg/rtcrec/recording_peerconnection.go b/pkg/rtcrec/recording_peerconnection.go index 49963bb4d..58f6f16dd 100644 --- a/pkg/rtcrec/recording_peerconnection.go +++ b/pkg/rtcrec/recording_peerconnection.go @@ -28,7 +28,7 @@ func NewRecordingPeerConnection(ctx context.Context, cli config.CLI, user string }, nil } aqt := aqtime.FromTime(time.Now()) - f, err := cli.DataFileCreate([]string{user, "rtcrec", fmt.Sprintf("%s.cbor", aqt.FileSafeString())}, true) + f, err := cli.DataFileCreate([]string{"debug-recordings", user, fmt.Sprintf("%s.rtcrec.cbor", aqt.FileSafeString())}, true) if err != nil { return nil, fmt.Errorf("failed to create data file: %w", err) }