diff --git a/pkg/config/config.go b/pkg/config/config.go index bafd2086e..df12ac330 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -32,6 +32,7 @@ import ( "stream.place/streamplace/pkg/integrations/discord/discordtypes" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/moderation" + "stream.place/streamplace/pkg/s3" placestream "stream.place/streamplace/pkg/streamplace" ) @@ -1428,6 +1429,40 @@ func (cli *CLI) S3Configured() bool { return cli.S3Endpoint != "" && cli.S3Bucket != "" && cli.S3AccessKeyID != "" && cli.S3SecretAccessKey != "" } +// S3Config assembles an s3.Config from the CLI's S3 flags. +func (cli *CLI) S3Config() s3.Config { + return s3.Config{ + Endpoint: cli.S3Endpoint, + Bucket: cli.S3Bucket, + AccessKeyID: cli.S3AccessKeyID, + SecretAccessKey: cli.S3SecretAccessKey, + Region: cli.S3Region, + } +} + +// DebugRecordingFile is the write target returned by DebugRecordingCreate: an +// *os.File on local disk, or an S3 upload that commits on Close. Name() reports +// the destination (path or object key) for logging. +type DebugRecordingFile interface { + io.WriteCloser + Name() string +} + +// DebugRecordingCreate opens a write target for a debug recording (RTMP/MKV +// dumps, WHIP rtcrec sessions). When S3 is configured the recording streams to +// an S3 object at the key formed by joining fpath with "/" (so the bucket +// mirrors the on-disk debug-recordings// layout); otherwise it falls +// back to a local file under DataDir — the dev default. The returned value must +// be Closed to finalize (Close commits the S3 upload). overwrite only affects +// the local-disk path (S3 puts always overwrite). +func (cli *CLI) DebugRecordingCreate(ctx context.Context, fpath []string, contentType string, overwrite bool) (DebugRecordingFile, error) { + if cli.S3Configured() { + key := strings.Join(fpath, "/") + return s3.NewUploadWriter(ctx, s3.NewClient(cli.S3Config()), cli.S3Bucket, key, contentType) + } + return cli.DataFileCreate(fpath, overwrite) +} + func (cli *CLI) ShouldSyndicate(did string) bool { if cli.DisableSyndication { return false diff --git a/pkg/director/s3_upload.go b/pkg/director/s3_upload.go index 6030fc972..7c929666e 100644 --- a/pkg/director/s3_upload.go +++ b/pkg/director/s3_upload.go @@ -74,13 +74,7 @@ func (ss *StreamSession) maybeStartS3Upload(ctx context.Context, repoDID string) log.Debug(ctx, "live recording disabled for streamer (not in VOD beta or recording not enabled)", "repoDID", repoDID) return } - cfg := s3.Config{ - Endpoint: ss.cli.S3Endpoint, - Bucket: ss.cli.S3Bucket, - AccessKeyID: ss.cli.S3AccessKeyID, - SecretAccessKey: ss.cli.S3SecretAccessKey, - Region: ss.cli.S3Region, - } + cfg := ss.cli.S3Config() // live-rec/ namespaces the in-progress livestream recordings away from the // finalized VOD blobs (blobs/) and anything else in the bucket. keyPrefix := liveRecPrefix + repoDID + "/" diff --git a/pkg/media/mkv_ingest.go b/pkg/media/mkv_ingest.go index ca8bf5626..aac3a36cf 100644 --- a/pkg/media/mkv_ingest.go +++ b/pkg/media/mkv_ingest.go @@ -109,14 +109,18 @@ func buildMKVIngestPipeline(ctx context.Context, input io.Reader, signerElem *gs 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) + // Streams to S3 when configured (production), else a local file under DataDir + // (dev). Close finalizes either target — for S3 it commits the upload. + f, err := mm.cli.DebugRecordingCreate(ctx, []string{"debug-recordings", user, filename}, "video/x-matroska", false) if err != nil { - return fmt.Errorf("failed to create data file: %w", err) + return fmt.Errorf("failed to create debug recording: %w", err) } - defer f.Close() - _, err = io.Copy(f, r) - if err != nil { - return fmt.Errorf("failed to copy to file: %w", err) + if _, err = io.Copy(f, r); err != nil { + f.Close() + return fmt.Errorf("failed to copy to debug recording: %w", err) + } + if err = f.Close(); err != nil { + return fmt.Errorf("failed to finalize debug recording: %w", err) } return nil } diff --git a/pkg/rtcrec/recording_peerconnection.go b/pkg/rtcrec/recording_peerconnection.go index 58f6f16dd..29de5eaca 100644 --- a/pkg/rtcrec/recording_peerconnection.go +++ b/pkg/rtcrec/recording_peerconnection.go @@ -3,7 +3,6 @@ package rtcrec import ( "context" "fmt" - "os" "time" "github.com/pion/rtcp" @@ -16,7 +15,7 @@ import ( type RecordingPeerConnection struct { enabled bool pionpc *webrtc.PeerConnection - file *os.File + file config.DebugRecordingFile stream *RecorderStream } @@ -28,9 +27,11 @@ func NewRecordingPeerConnection(ctx context.Context, cli config.CLI, user string }, nil } aqt := aqtime.FromTime(time.Now()) - f, err := cli.DataFileCreate([]string{"debug-recordings", user, fmt.Sprintf("%s.rtcrec.cbor", aqt.FileSafeString())}, true) + // Streams to S3 when configured (production), else a local file under DataDir + // (dev). Close (after the drain delay below) finalizes either target. + f, err := cli.DebugRecordingCreate(ctx, []string{"debug-recordings", user, fmt.Sprintf("%s.rtcrec.cbor", aqt.FileSafeString())}, "application/cbor", true) if err != nil { - return nil, fmt.Errorf("failed to create data file: %w", err) + return nil, fmt.Errorf("failed to create debug recording: %w", err) } log.Log(ctx, "logging webrtc session to file", "file", f.Name()) stream, err := MakeWebRTCEncoder(f) diff --git a/pkg/rtcrec/webrtc_recording.go b/pkg/rtcrec/webrtc_recording.go index c461b854a..ded0dfc19 100644 --- a/pkg/rtcrec/webrtc_recording.go +++ b/pkg/rtcrec/webrtc_recording.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "io" + "sync" "time" "github.com/fxamacker/cbor/v2" @@ -113,6 +114,10 @@ type TrackSSRC struct { } type RecorderStream struct { + // mu serializes Event: RecordingPeerConnection.Do fires each event on its own + // goroutine, so without this the concurrent Encode calls would race the CBOR + // encoder and (for the S3 upload target) corrupt the multipart writer's buffer. + mu sync.Mutex encoder *cbor.Encoder } @@ -131,6 +136,8 @@ func MakeWebRTCEncoder(w io.Writer) (*RecorderStream, error) { } func (s *RecorderStream) Event(event WebRTCEvent) { + s.mu.Lock() + defer s.mu.Unlock() err := s.encoder.Encode(event) if err != nil { log.Log(context.Background(), "error encoding event", "error", err) diff --git a/pkg/s3/multipart_writer.go b/pkg/s3/multipart_writer.go index 6f675e8a6..714f07c09 100644 --- a/pkg/s3/multipart_writer.go +++ b/pkg/s3/multipart_writer.go @@ -325,3 +325,43 @@ func (w *MultipartWriter) Abort() error { func (w *MultipartWriter) Close() error { return w.Abort() } var _ io.WriteCloser = (*MultipartWriter)(nil) + +// UploadWriter streams a single object to S3 and finalizes it on Close. It is a +// thin adapter over MultipartWriter for callers that just want a plain +// io.WriteCloser whose Close() *commits* the object (MultipartWriter.Close() +// aborts, which is the wrong default for a fire-and-forget upload). If any Write +// failed, Close surfaces that error. Like MultipartWriter, Write must be called +// from a single goroutine. +type UploadWriter struct { + mw *MultipartWriter + key string +} + +// NewUploadWriter starts a multipart upload at key and returns a writer that +// commits it when closed. +func NewUploadWriter(ctx context.Context, client *s3.Client, bucket, key, contentType string) (*UploadWriter, error) { + return newUploadWriter(ctx, client, bucket, key, contentType) +} + +// newUploadWriter is the client-injectable constructor behind NewUploadWriter; +// the fake-client tests use it directly. +func newUploadWriter(ctx context.Context, client multipartAPI, bucket, key, contentType string) (*UploadWriter, error) { + mw, err := newMultipartWriter(ctx, client, bucket, key, contentType) + if err != nil { + return nil, err + } + return &UploadWriter{mw: mw, key: key}, nil +} + +func (w *UploadWriter) Write(p []byte) (int, error) { return w.mw.Write(p) } + +// Close completes the multipart upload, flushing any buffered bytes. Idempotent +// on the underlying writer (a second Complete returns an error, so callers +// should Close exactly once). +func (w *UploadWriter) Close() error { return w.mw.Complete() } + +// Name reports the object key, mirroring *os.File.Name() so callers can log a +// destination uniformly whether they got a file or an S3 upload. +func (w *UploadWriter) Name() string { return w.key } + +var _ io.WriteCloser = (*UploadWriter)(nil) diff --git a/pkg/s3/multipart_writer_test.go b/pkg/s3/multipart_writer_test.go index c478350b3..74a03dad5 100644 --- a/pkg/s3/multipart_writer_test.go +++ b/pkg/s3/multipart_writer_test.go @@ -123,3 +123,35 @@ func TestMultipartWriterConcurrent(t *testing.T) { } require.Equal(t, want, got) } + +// TestUploadWriterCommitsOnClose verifies UploadWriter.Close finalizes the +// object (CompleteMultipartUpload) rather than aborting it, and that the +// written bytes reassemble into a single object. This is the property that +// makes UploadWriter safe as a fire-and-forget io.WriteCloser for debug +// recordings, where the raw MultipartWriter.Close() would silently abort. +func TestUploadWriterCommitsOnClose(t *testing.T) { + fake := newFakeMultipartClient(0) + w, err := newUploadWriter(context.Background(), fake, "bucket", "debug-recordings/did/x.rtcrec.cbor", "application/cbor") + require.NoError(t, err) + require.Equal(t, "debug-recordings/did/x.rtcrec.cbor", w.Name()) + + want := make([]byte, MultipartPartSize+4096) + for i := range want { + want[i] = byte(i * 3 % 251) + } + n, err := w.Write(want) + require.NoError(t, err) + require.Equal(t, len(want), n) + + require.NoError(t, w.Close()) + + // Close must have completed (not aborted) the upload. + require.NotEmpty(t, fake.completed, "Close should CompleteMultipartUpload") + + // Reassemble parts in order and compare to what we wrote. + got := make([]byte, 0, len(want)) + for i := int32(1); i <= int32(len(fake.parts)); i++ { + got = append(got, fake.parts[i]...) + } + require.Equal(t, want, got) +} diff --git a/pkg/s3/s3.go b/pkg/s3/s3.go index eb475578c..1eb340f8e 100644 --- a/pkg/s3/s3.go +++ b/pkg/s3/s3.go @@ -111,7 +111,14 @@ var DefaultCutoverEvery = 10 * time.Minute // nil to disable persistence. Starts the muxl Concatenator and a background // goroutine that reads processed segments and uploads them. func NewS3Uploader(cfg Config, userDID, keyPrefix string, cutoverEvery time.Duration, recorder Recorder) *S3Uploader { - client := s3.New(s3.Options{ + return newS3Uploader(NewClient(cfg), cfg.Bucket, userDID, keyPrefix, cutoverEvery, recorder) +} + +// NewClient builds an *s3.Client for an S3-compatible endpoint from cfg. Uses +// static credentials, an explicit BaseEndpoint, and path-style addressing (so it +// works against MinIO / R2 / plain-IP endpoints, not just AWS virtual-host style). +func NewClient(cfg Config) *s3.Client { + return s3.New(s3.Options{ Region: cfg.Region, Credentials: credentials.NewStaticCredentialsProvider( cfg.AccessKeyID, @@ -121,7 +128,6 @@ func NewS3Uploader(cfg Config, userDID, keyPrefix string, cutoverEvery time.Dura BaseEndpoint: aws.String(cfg.Endpoint), UsePathStyle: true, }) - return newS3Uploader(client, cfg.Bucket, userDID, keyPrefix, cutoverEvery, recorder) } // newS3Uploader is the client-injectable constructor behind NewS3Uploader; the