From 1c69d6787f61da5da498526e83e957aa3f25d108 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Tue, 21 Jul 2026 17:19:18 -0700 Subject: [PATCH] s3: bound multipart ops; rtcrec: log failed recording commits MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Addresses the two Greptile P1s on #1209: - MultipartWriter's S3 calls had no deadline (the SDK's default HTTP client has none), so a stalled connection could block its caller forever — notably a debug-recording commit, whose writer deliberately runs on a non-cancellable ctx so session teardown can't abort it. Every operation now carries a generous per-op timeout (10m per part, 2m for create/complete/abort), so a genuine stall errors out and surfaces instead of leaking a wedged goroutine + dangling multipart upload. - rtcrec's finishRecording discarded pc.file.Close()'s error, so a failed S3 commit looked identical to success. The outcome is now logged either way; recDone still means "attempt finished" — with the session over there's nothing better to do with a failure than say so loudly. Co-Authored-By: Claude Fable 5 --- pkg/rtcrec/recording_peerconnection.go | 24 ++++++++++++++++-------- pkg/s3/multipart_writer.go | 21 ++++++++++++++++++++- 2 files changed, 36 insertions(+), 9 deletions(-) diff --git a/pkg/rtcrec/recording_peerconnection.go b/pkg/rtcrec/recording_peerconnection.go index 34a8347d2..3e43f0dc5 100644 --- a/pkg/rtcrec/recording_peerconnection.go +++ b/pkg/rtcrec/recording_peerconnection.go @@ -18,8 +18,9 @@ type RecordingPeerConnection struct { pionpc *webrtc.PeerConnection file config.DebugRecordingFile stream *RecorderStream + logCtx context.Context // for finishRecording's logs (it outlives the session) closeOnce sync.Once - recDone chan struct{} // closed once the recording file/upload is committed + recDone chan struct{} // closed once the finalize ATTEMPT is over — check the logs for commit failures } func NewRecordingPeerConnection(ctx context.Context, cli config.CLI, user string, pionpc *webrtc.PeerConnection, enabled bool) (PeerConnection, error) { @@ -46,6 +47,7 @@ func NewRecordingPeerConnection(ctx context.Context, cli config.CLI, user string file: f, stream: stream, enabled: enabled, + logCtx: context.WithoutCancel(ctx), recDone: make(chan struct{}), }, nil } @@ -63,21 +65,27 @@ func (pc *RecordingPeerConnection) Close() error { // finishRecording drains stragglers, commits the recording (for S3, Close IS // the commit), and signals recDone. Idempotent — Close on the disconnect path -// and FinalizeRecording at worker exit can both trigger it. +// and FinalizeRecording at worker exit can both trigger it. recDone means the +// attempt finished, not that it succeeded: a failed commit is logged loudly +// (there is nothing better to do with it at this point — the session is over). func (pc *RecordingPeerConnection) finishRecording() { pc.closeOnce.Do(func() { // This is sloppy but there might be other goroutines still writing so let's chill for a sec time.Sleep(10 * time.Second) - pc.file.Close() + if err := pc.file.Close(); err != nil { + log.Error(pc.logCtx, "debug recording commit FAILED; the recording is lost", "file", pc.file.Name(), "error", err) + } else { + log.Log(pc.logCtx, "debug recording committed", "file", pc.file.Name()) + } close(pc.recDone) }) } -// FinalizeRecording blocks until the debug recording is committed (bounded). -// Call it before process exit on paths like the WHIP ingest worker: Close only -// *starts* the drain+commit on a goroutine, and a process that exits first -// strands an uncommitted S3 upload — the object never appears. No-op when not -// recording. +// FinalizeRecording blocks until the debug recording's commit attempt finishes +// (bounded; failures are logged by finishRecording). Call it before process +// exit on paths like the WHIP ingest worker: Close only *starts* the +// drain+commit on a goroutine, and a process that exits first strands an +// uncommitted S3 upload — the object never appears. No-op when not recording. func (pc *RecordingPeerConnection) FinalizeRecording(ctx context.Context) { if !pc.enabled { return diff --git a/pkg/s3/multipart_writer.go b/pkg/s3/multipart_writer.go index 714f07c09..9891030c3 100644 --- a/pkg/s3/multipart_writer.go +++ b/pkg/s3/multipart_writer.go @@ -37,6 +37,17 @@ const MultipartPartSize = 16 * 1024 * 1024 // so a 1.5 GB upload dragged on for tens of minutes. const multipartUploadConcurrency = 8 +// Per-operation deadlines for MultipartWriter's S3 calls. The SDK's default +// HTTP client has no request timeout, so without these a stalled connection +// blocks its caller forever — e.g. a debug-recording commit, whose writer +// deliberately runs on a non-cancellable ctx (config.DebugRecordingCreate) so +// session teardown can't abort it. Values are far above healthy operation +// times; only genuine stalls hit them. +const ( + s3PartOpTimeout = 10 * time.Minute // one ≤MultipartPartSize UploadPart + s3ControlOpTimeout = 2 * time.Minute // create/complete/abort/empty-put +) + // multipartAPI is the subset of *s3.Client that MultipartWriter calls. // Pulled out so tests can inject a fake; *s3.Client satisfies it. type multipartAPI interface { @@ -105,7 +116,9 @@ func newMultipartWriter(ctx context.Context, client multipartAPI, bucket, key, c if contentType != "" { in.ContentType = aws.String(contentType) } - resp, err := client.CreateMultipartUpload(ctx, in) + cctx, cancel := context.WithTimeout(ctx, s3ControlOpTimeout) + defer cancel() + resp, err := client.CreateMultipartUpload(cctx, in) if err != nil { span.RecordError(err) return nil, fmt.Errorf("create multipart upload s3://%s/%s: %w", bucket, key, err) @@ -175,6 +188,8 @@ func (w *MultipartWriter) uploadPart(num int32, body []byte) { attribute.Int("part_size_bytes", len(body)), )) defer span.End() + ctx, cancel := context.WithTimeout(ctx, s3PartOpTimeout) + defer cancel() resp, err := w.client.UploadPart(ctx, &s3.UploadPartInput{ Bucket: aws.String(w.bucket), Key: aws.String(w.key), @@ -244,6 +259,8 @@ func (w *MultipartWriter) Complete() error { return err } span.SetAttributes(attribute.Int("part_count", len(w.parts))) + ctx, cancel := context.WithTimeout(ctx, s3ControlOpTimeout) + defer cancel() if len(w.parts) == 0 { // Zero-byte upload: S3 won't accept an empty CompletedMultipartUpload, // so abort and create an empty object via PutObject. @@ -308,6 +325,8 @@ func (w *MultipartWriter) Abort() error { attribute.Int("parts_pending", len(w.parts)), )) defer span.End() + ctx, cancel := context.WithTimeout(ctx, s3ControlOpTimeout) + defer cancel() _, err := w.client.AbortMultipartUpload(ctx, &s3.AbortMultipartUploadInput{ Bucket: aws.String(w.bucket), Key: aws.String(w.key), -- 2.51.2