Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
27 kB · 707 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708// Package vod implements the post-upload processing pipeline for a VOD// upload: read the user's raw file from wherever the upload manager// stored it, run it through gstreamer (parsebin -> mp4mux -> muxl// concatenator), and write the resulting fMP4 to a content-addressed// key under a blob.Store.//// The pipeline is streaming end-to-end — bytes flow from the source// (file or ranged S3 GETs), through gstreamer's appsrc -> parsebin ->// fdkaacenc/h264parse -> mp4mux -> appsink, into the muxl wasm// concatenator, and out to a blob.Store writer. A bdasl.Writer runs// as a tee on the way out; the final CID is known only when the last// byte is written, at which point we Move the staging blob to the// content-addressed key.package vod
import ( "context" "errors" "fmt" "io" "os" "sync/atomic" "time"
"go.opentelemetry.io/otel" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" "go.opentelemetry.io/otel/trace"
"stream.place/streamplace/pkg/bdasl" "stream.place/streamplace/pkg/blob" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" "stream.place/streamplace/pkg/muxl" "stream.place/streamplace/pkg/spmetrics" "stream.place/streamplace/pkg/statedb")
var vodTracer = otel.Tracer("vod")
// vodProgressLogInterval is how often consumeConcatTraced emits an// info-level heartbeat while segments flow, so a long-running mux is// visible in the logs rather than silent between start and finish.const vodProgressLogInterval = 5 * time.Second
const ( // vodStageHeartbeat is how often runVODStage logs a "still working" // line for a post-pipeline stage (S3 upload/move, metafile, thumbnail, // publish), so a slow or wedged backend call is visible rather than // silent. It keeps ticking until the stage actually returns. vodStageHeartbeat = 5 * time.Second // vodStageTimeout bounds any single post-pipeline stage. These call // out to S3 and the user's PDS; a wedged connection should fail the // task — freeing the worker — well before the 30-minute task lock // expires and a second worker picks the same upload up. NOTE: it can // only interrupt stages that honor the context; completeStaging's // underlying Writer.Complete() takes none, so there the timeout bounds // only the logging, not the call itself. vodStageTimeout = 10 * time.Minute)
// stage labels for spmetrics.VODProcessErrorsTotal — keep these in sync// with whatever the dashboard / alert routing keys off.const ( stageOpenSource = "open_source" stageStaging = "start_staging" stageSigner = "create_signer" stagePipeline = "gstreamer_pipeline" stageEmptyOutput = "empty_muxl_output" stageStagingComplete = "store_complete" stageContentAddressCopy = "content_address_move" stageMetafile = "write_metafile" stagePublish = "publish_records")
// Input is the per-upload state needed to drive ProcessVOD. It's a// flattened copy of statedb.VODProcessTask so this package doesn't// import statedb (which would invite a dependency cycle with the queue// processor).type Input struct { UploadID string RepoDID string MimeType string Filename string Size int64 // Backend describes the upload source backend ("file" or "s3"), // matching what the upload manager stored on the row. Used only // for metrics labeling; the actual reading goes through the // blob.Store passed to ProcessVOD. Backend string // Location is the backend-specific URL/path (a file path or // "s3://bucket/key" URL) the upload landed at. Translated to a // Store key via Store.ParseLocation. Location string}
// StagingPrefix is where in-progress VOD outputs land before being// renamed to their content-addressed key. We park them under a// dedicated prefix so a periodic janitor can sweep abandoned uploads.const StagingPrefix = "vod-staging/"
// BlobsPrefix is where finalized, content-addressed blobs land in the// Store: `<BlobsPrefix><cid>.mp4` for content + init segments,// `<BlobsPrefix><cid>.json` for metafiles. The prefix is intentionally// content-agnostic — the blob doesn't know what kind of video it's for// or whether it's a primary fmp4, a per-track init, or a metafile.// CDN deployments map `<vod-cdn-url>/blobs/...` straight at this layout.const BlobsPrefix = "blobs/"
// ProcessVOD runs the streaming pipeline for one VOD upload and returns// the BDASL CID of the resulting fMP4. Reads come from `in` via the// supplied Store; the output blob lands at BlobsPrefix+<cid>.mp4 in// the same Store; staging blobs are cleaned up on success or failure.//// The Store is the storage layer; it can be either a FileStore (single-// node deployments) or an S3Store (production). Future Stores can mix// caches and archives behind the same interface.func ProcessVOD(ctx context.Context, cli *config.CLI, state *statedb.StatefulDB, store blob.Store, in Input) (string, error) { ctx = log.WithLogValues(ctx, "func", "ProcessVOD", "uploadId", in.UploadID, "did", in.RepoDID) ctx, span := vodTracer.Start(ctx, "vod.ProcessVOD", trace.WithAttributes( attribute.String("upload_id", in.UploadID), attribute.String("did", in.RepoDID), attribute.String("backend", in.Backend), attribute.String("mime_type", in.MimeType), attribute.Int64("input_size_bytes", in.Size), )) defer span.End()
startTime := time.Now() spmetrics.VODProcessAttemptsTotal.WithLabelValues(in.Backend).Inc() defer func() { spmetrics.VODProcessDurationMS.Observe(float64(time.Since(startTime).Milliseconds())) }()
log.Log(ctx, "starting VOD processing", "backend", in.Backend, "size", in.Size, "mimeType", in.MimeType, "location", in.Location, )
src, size, closer, err := openSource(ctx, store, in) if err != nil { recordErr(span, stageOpenSource, err) return "", fmt.Errorf("open source: %w", err) } defer closer() span.SetAttributes(attribute.Int64("source_size_bytes", size)) spmetrics.VODInputBytes.Observe(float64(size))
stagingKey := StagingPrefix + in.UploadID + ".mp4" span.SetAttributes(attribute.String("staging_key", stagingKey)) staging, err := store.NewWriter(ctx, stagingKey, "video/mp4") if err != nil { recordErr(span, stageStaging, err) return "", fmt.Errorf("start staging upload: %w", err) } defer staging.Close() // Abort is idempotent; no-op once Complete ran
hasher := bdasl.NewWriter() counter := &countingWriter{} final := io.MultiWriter(hasher, counter, staging)
// Start a progress reporter that polls the byte counter and writes // percentage updates to the DB every 2 seconds. It stops when the // pipeline completes (progressDone closed) or the context is cancelled. progressDone := make(chan struct{}) defer close(progressDone) go reportProgress(ctx, state, in.UploadID, counter, size, progressDone, 2*time.Second)
// Ephemeral per-upload signing key. Generated here as muxing starts, // used to C2PA-sign every segment, and dropped when this function // returns — so once the upload is processed nobody can mint more // signatures under its did:key. signer, err := newUploadSigner(startTime) if err != nil { recordErr(span, stageSigner, err) return "", fmt.Errorf("create upload signer: %w", err) } span.SetAttributes(attribute.String("signing_did", signer.DIDKey))
// Fragment metafile builder: dst (the staging blob) is the bare canonical // fragments alone — no leading init — so its hash is the MUXL CID and the // metafile offsets are fragment-relative. The init is re-synthesized as a // flat-MP4 faststart header below and prepended at assembly. metaBuilder := newFragmentMetafileBuilder(ctx, store) probe, err := streamThroughMuxl(ctx, src, size, final, metaBuilder, signer.SignerInput) if err != nil { recordErr(span, stagePipeline, err) return "", err } _ = state.SetUploadProgress(ctx, in.UploadID, 85) span.SetAttributes( attribute.Int64("output_size_bytes", counter.load()), attribute.Int64("probe_duration_ms", probe.DurationMS), ) if probe.Video != nil { span.SetAttributes( attribute.String("probe_video_codec", probe.Video.Codec), attribute.Int("probe_video_width", probe.Video.Width), attribute.Int("probe_video_height", probe.Video.Height), ) } if probe.Audio != nil { span.SetAttributes( attribute.String("probe_audio_codec", probe.Audio.Codec), attribute.Int("probe_audio_rate", probe.Audio.Rate), attribute.Int("probe_audio_channels", probe.Audio.Channels), ) } spmetrics.VODOutputBytes.Observe(float64(counter.load()))
if counter.load() == 0 { err := errors.New("muxl produced zero bytes") recordErr(span, stageEmptyOutput, err) return "", err }
if err := runVODStage(ctx, stageStagingComplete, func(ctx context.Context) error { return completeStaging(ctx, staging, stagingKey) }); err != nil { recordErr(span, stageStagingComplete, err) return "", fmt.Errorf("complete staging upload: %w", err) } _ = state.SetUploadProgress(ctx, in.UploadID, 90)
// The staging blob is the bare canonical fragments; its BDASL hash is the // MUXL CID — the stable, header-independent VOD id. Keying the content blob // off it means a better flat header can later be re-synthesized and // rewritten in place without changing the VOD's id or any published record. muxlCID := hasher.CID() fragSize := counter.load() contentKey := BlobsPrefix + muxlCID + ".mp4"
// Synthesize the flat-MP4 faststart header from the fragments' metafile, // then assemble [flat-header][fragments] at the content key and drop the // staging blob. muxl owns all co64 offsets over the [header][fragments] // layout, so the result is a real flat MP4 served whole, with HLS reading // byte-ranges at flatHeaderSize+offset into the same blob. flatHeader, err := synthFlatHeaderForObjects(ctx, store, []string{stagingKey}) if err != nil { recordErr(span, stageContentAddressCopy, err) return "", fmt.Errorf("synthesize flat header: %w", err) } blobSize := int64(len(flatHeader)) + fragSize span.SetAttributes( attribute.String("cid", muxlCID), attribute.String("content_key", contentKey), attribute.String("content_url", store.URL(contentKey)), attribute.Int("flat_header_bytes", len(flatHeader)), attribute.Int64("blob_size", blobSize), )
if err := runVODStage(ctx, stageContentAddressCopy, func(ctx context.Context) error { if err := assembleContentBlob(ctx, store, flatHeader, []string{stagingKey}, contentKey); err != nil { return err } return store.Delete(ctx, stagingKey) }); err != nil { recordErr(span, stageContentAddressCopy, err) return "", fmt.Errorf("finalize: %w", err) } _ = state.SetUploadProgress(ctx, in.UploadID, 95)
metafile := metaBuilder.Finalize(muxlCID, fragSize) metafile.FlatHeaderSize = int64(len(flatHeader)) if err := runVODStage(ctx, stageMetafile, func(ctx context.Context) error { return writeMetafile(ctx, store, muxlCID, metafile) }); err != nil { recordErr(span, stageMetafile, err) return "", fmt.Errorf("write metafile: %w", err) }
// Thumbnail generation is deferred: the place.stream.video record is // created client-side now, so there's no server-side record to attach // one to. generateThumbnail (and its fmp4 fix) stay in the tree, // exercised by vod-test, for when thumbnails are wired to the client. if err := runVODStage(ctx, stagePublish, func(ctx context.Context) error { return publishRecords(ctx, publishParams{ cli: cli, state: state, in: in, cid: muxlCID, size: blobSize, mimeType: "video/mp4", probe: probe, signingKey: signer.DIDKey, }) }); err != nil { recordErr(span, stagePublish, err) return "", fmt.Errorf("publish records: %w", err) }
spmetrics.VODProcessSuccessesTotal.WithLabelValues(in.Backend).Inc() span.SetStatus(codes.Ok, "") log.Log(ctx, "VOD processed", "cid", muxlCID, "url", store.URL(contentKey), "input_size", size, "frag_size", fragSize, "output_size", blobSize, "duration_ms", time.Since(startTime).Milliseconds(), ) return muxlCID, nil}
// runVODStage runs a post-pipeline stage with the default stage timeout.func runVODStage(ctx context.Context, name string, fn func(context.Context) error) error { return runVODStageWithin(ctx, name, vodStageTimeout, fn)}
// runVODStageWithin runs one post-pipeline stage with an entry/exit log, a// heartbeat while it's in flight, and a hard timeout. The heartbeat turns// a silent hang into a visible "stage X still running after Ns" trail and// keeps ticking until fn returns — so a stage that ignores the deadline// (e.g. completeStaging) is still observable past the timeout. The timeout// frees the worker when a context-aware stage wedges. The stage error (or// the context error on timeout) is returned unwrapped; the caller adds the// stage label and the upload ID is attached upstream.func runVODStageWithin(ctx context.Context, name string, timeout time.Duration, fn func(context.Context) error) error { ctx, cancel := context.WithTimeout(ctx, timeout) defer cancel() start := time.Now() log.Log(ctx, "vod stage starting", "stage", name)
done := make(chan struct{}) go func() { ticker := time.NewTicker(vodStageHeartbeat) defer ticker.Stop() for { select { case <-done: return case <-ticker.C: log.Log(ctx, "vod stage in progress", "stage", name, "elapsed_seconds", time.Since(start).Seconds()) } } }()
err := fn(ctx) close(done) log.Log(ctx, "vod stage finished", "stage", name, "elapsed_seconds", time.Since(start).Seconds(), "ok", err == nil) return err}
// completeStaging wraps the writer Complete call in its own span so we// can see how much of total processing time is spent waiting for the// storage layer to acknowledge the upload.func completeStaging(ctx context.Context, staging blob.Writer, stagingKey string) error { _, span := vodTracer.Start(ctx, "vod.completeStaging", trace.WithAttributes( attribute.String("staging_key", stagingKey), )) defer span.End() if err := staging.Complete(); err != nil { span.RecordError(err) span.SetStatus(codes.Error, err.Error()) return err } return nil}
// countingWriter is an io.Writer that tallies bytes written via an// atomic counter so a progress goroutine can read it concurrently.type countingWriter struct{ n atomic.Int64 }
func (c *countingWriter) Write(p []byte) (int, error) { c.n.Add(int64(len(p))) return len(p), nil}
func (c *countingWriter) load() int64 { return c.n.Load() }
// reportProgress periodically reads the byte counter and writes a// percentage to the DB until doneCh closes. The estimate is// counter.bytes / inputSize, capped at 90 so the remaining stages// (staging complete, finalize, publish) each have visible increments.func reportProgress(ctx context.Context, state *statedb.StatefulDB, uploadID string, counter *countingWriter, inputSize int64, doneCh <-chan struct{}, interval time.Duration) { ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-doneCh: return case <-ticker.C: if inputSize <= 0 { continue } pct := int(float64(counter.load()) / float64(inputSize) * 100) if pct > 90 { pct = 90 } if pct < 5 { pct = 5 } _ = state.SetUploadProgress(ctx, uploadID, pct) } }}
// ProcessToDiscard runs the full VOD media path against a local file —// gstreamer remux, muxl sign-segment, metafile assembly, AND thumbnail// generation — exactly as ProcessVOD does, into a throwaway temp store// that's deleted on return. The only things it skips vs production are the// S3 staging upload and the record publish (no statedb, CLI, or PDS// needed), which makes it the engine behind `streamplace vod-test`: a way// to reproduce pipeline crashes/stalls — including thumbnail-decode hangs —// against an arbitrary file. Returns the probe metadata and output bytes.func ProcessToDiscard(ctx context.Context, src io.ReaderAt, size int64) (media.VODResult, int64, error) { signer, err := newUploadSigner(time.Now()) if err != nil { return media.VODResult{}, 0, fmt.Errorf("create upload signer: %w", err) } // A throwaway temp dir stands in for the production blob store, so the // metafile init blobs, the content blob, and the thumbnail's reads all // exercise the same code paths. dir, err := os.MkdirTemp("", "vod-test-") if err != nil { return media.VODResult{}, 0, fmt.Errorf("temp dir: %w", err) } defer os.RemoveAll(dir) store, err := blob.NewFileStore(dir) if err != nil { return media.VODResult{}, 0, fmt.Errorf("file store: %w", err) }
mb := newMetafileBuilder(ctx, store) // Write the muxl output into the store (hashed + counted) so it can be // finalized to a content-addressed blob the thumbnail stage can read, // mirroring ProcessVOD instead of writing straight to /dev/null. stagingKey := StagingPrefix + "vod-test.mp4" staging, err := store.NewWriter(ctx, stagingKey, "video/mp4") if err != nil { return media.VODResult{}, 0, fmt.Errorf("staging writer: %w", err) } defer staging.Close()
hasher := bdasl.NewWriter() counter := &countingWriter{} final := io.MultiWriter(hasher, counter, staging) result, err := streamThroughMuxl(ctx, src, size, final, mb, signer.SignerInput) if err != nil { return result, counter.load(), err } if err := staging.Complete(); err != nil { return result, counter.load(), fmt.Errorf("complete staging: %w", err) }
cid := hasher.CID() if err := store.Move(ctx, stagingKey, BlobsPrefix+cid+".mp4"); err != nil { return result, counter.load(), fmt.Errorf("finalize: %w", err) } metafile := mb.Finalize(cid, counter.load())
// Thumbnail is non-fatal in production; here we run it unbounded so a // hang is observable (the per-step timing logs inside generateThumbnail // show where it wedges). Failures are logged, not returned. if thumb, terr := generateThumbnail(ctx, store, cid, metafile); terr != nil { log.Warn(ctx, "vod-test: thumbnail generation failed", "error", terr) } else { log.Log(ctx, "vod-test: thumbnail generated", "bytes", len(thumb)) }
return result, counter.load(), nil}
// recordErr is a small helper to attach the standard error attributes// to a span + bump the per-stage counter. Both spans and counters need// to be kept in sync for dashboards to work.func recordErr(span trace.Span, stage string, err error) { span.RecordError(err) span.SetStatus(codes.Error, err.Error()) span.SetAttributes(attribute.String("error_stage", stage)) spmetrics.VODProcessErrorsTotal.WithLabelValues(stage).Inc()}
// streamThroughMuxl wires up the goroutines that connect://// [gstreamer pipeline] -> mp4mux output (io.Pipe)// -> muxl sign-segment -> dst//// gstreamer writes the fMP4 stream to one end of an io.Pipe; a goroutine// forwards those bytes into muxl-sign's sign-segment subcommand, which// segments per-GoP and C2PA-signs each canonical segment; another// goroutine drains the init+seg channels and writes them to dst. The// segment bytes that land in dst are therefore signed// ([c2pa-uuid][muxl-uuid][moof][mdat] per track).//// If metaBuilder is non-nil, a third goroutine consumes the rich event// channel and feeds events to it — used to build the per-blob metafile// sidecar (segment offsets/sizes/durations and per-track init blobs)// without re-parsing the bytes.//// The returned media.VODResult is the probe metadata gathered during// the pipeline run (codec/dimensions/duration). The CID is computed// by the caller from the bdasl.Writer tee'd into dst, since the hash// is only final once the last byte lands.func streamThroughMuxl(ctx context.Context, src io.ReaderAt, size int64, dst io.Writer, metaBuilder *metafileBuilder, signerInput muxl.SignerInput) (media.VODResult, error) { ctx, span := vodTracer.Start(ctx, "vod.streamThroughMuxl", trace.WithAttributes( attribute.Int64("source_size_bytes", size), )) defer span.End()
ctx, cancel := context.WithCancel(ctx) defer cancel()
pr, pw := io.Pipe()
// muxStats: bytes counted at the muxer output (input to muxl), and at // the concatenator output (final fMP4). Reported in span attributes // at the end so we can see the per-stage byte amplification. var mp4muxBytes int64 var initBytes int64 var segBytes int64 var initEmits int64 var segEmits int64
// Forwarder: gstreamer mp4mux output -> muxl sign-segment stdin. // Run as a goroutine because both ends are pipes; both sides need to // be live for either to make progress. concat := muxl.NewSigningSegmenter(ctx, signerInput) feedDone := make(chan error, 1) go func() { defer close(feedDone) buf := make([]byte, 64*1024) for { n, err := pr.Read(buf) if n > 0 { mp4muxBytes += int64(n) if werr := concat.Write(buf[:n]); werr != nil { feedDone <- werr return } } if err != nil { if errors.Is(err, io.EOF) { log.Debug(ctx, "muxl feed: EOF, closing concatenator", "mp4mux_bytes", mp4muxBytes) feedDone <- concat.Close() } else { feedDone <- err } return } } }()
// Consumer: muxl concatenator output channels -> dst. The metafile // builder's offset model dictates the blob layout: leadingInitInBlob // means dst is [init][segments] (the legacy/transfer shape), otherwise // dst is the bare fragments alone (the flat-MP4 shape, init synthesized // separately). A nil builder keeps the legacy behavior. writeInit := metaBuilder == nil || metaBuilder.leadingInitInBlob consumeDone := make(chan error, 1) go func() { consumeDone <- consumeConcatTraced(ctx, concat, dst, writeInit, &initBytes, &segBytes, &initEmits, &segEmits) }()
// Optional metafile builder: parallel consumer on the event channel. // Buffered enough that a stall in the metafile path doesn't backpressure // the byte consumer. EventCh always exists on the concatenator; if no // builder was provided we still need to drain it so the wasm parser // doesn't block. metaDone := make(chan error, 1) go func() { var err error for ev := range concat.EventCh { if metaBuilder == nil { continue } if oerr := metaBuilder.Observe(ev); oerr != nil { err = oerr // Keep draining the channel so upstream doesn't block; // the first error wins. } } metaDone <- err }()
// Run gstreamer on this goroutine; mp4mux output bytes land at pw. result, pipelineErr := media.RunVODPipeline(ctx, src, size, pw) // Close pw so the forwarder sees EOF and can close the concat. _ = pw.Close()
// Wait for all downstream goroutines. We collect any error but // prefer surfacing the pipeline error since downstream errors are // often a consequence of it. feedErr := <-feedDone consumeErr := <-consumeDone metaErr := <-metaDone
span.SetAttributes( attribute.Int64("mp4mux_output_bytes", mp4muxBytes), attribute.Int64("muxl_init_bytes", initBytes), attribute.Int64("muxl_seg_bytes", segBytes), attribute.Int64("muxl_init_emits", initEmits), attribute.Int64("muxl_seg_emits", segEmits), )
if pipelineErr != nil { span.RecordError(pipelineErr) span.SetStatus(codes.Error, "pipeline") return result, fmt.Errorf("gstreamer pipeline: %w", pipelineErr) } if feedErr != nil { span.RecordError(feedErr) span.SetStatus(codes.Error, "feed") return result, fmt.Errorf("muxl feed: %w", feedErr) } if consumeErr != nil { span.RecordError(consumeErr) span.SetStatus(codes.Error, "consume") return result, fmt.Errorf("muxl consume: %w", consumeErr) } if metaErr != nil { span.RecordError(metaErr) span.SetStatus(codes.Error, "metafile") return result, fmt.Errorf("muxl metafile builder: %w", metaErr) } return result, nil}
// consumeConcatTraced drains the concatenator's two output channels and// writes their contents (init, then segments) to dst in order of// arrival. Increments the supplied counters so the parent span can// report them. If the init segment changes mid-stream (multi-input// concatenation), the new init is written too — for VOD with a single// input that doesn't happen, but the loop handles it for free.//// writeInit controls whether the init segment(s) land in dst. The flat-MP4// model wants dst to be the bare canonical fragments alone (the init is// re-synthesized as a faststart header at finalize time and prepended// separately); writeInit=false drains and counts the init for tracing but// keeps it out of the blob — so dst hashes to the MUXL CID of the fragments.func consumeConcatTraced(ctx context.Context, c *muxl.Concatenator, dst io.Writer, writeInit bool, initBytes, segBytes, initEmits, segEmits *int64) error { initCh, segCh := c.InitCh, c.SegCh start := time.Now() ticker := time.NewTicker(vodProgressLogInterval) defer ticker.Stop() for initCh != nil || segCh != nil { select { case <-ticker.C: log.Log(ctx, "muxl progress", "seg_emits", *segEmits, "seg_bytes", *segBytes, "elapsed_seconds", time.Since(start).Seconds()) case init, ok := <-initCh: if !ok { initCh = nil continue } *initEmits++ *initBytes += int64(len(init)) log.Debug(ctx, "muxl init segment", "size", len(init), "emit_n", *initEmits, "in_blob", writeInit) if writeInit { if _, err := dst.Write(init); err != nil { return err } } case seg, ok := <-segCh: if !ok { segCh = nil continue } *segEmits++ *segBytes += int64(len(seg)) if _, err := dst.Write(seg); err != nil { return err } } } log.Debug(ctx, "muxl drain complete", "init_emits", *initEmits, "seg_emits", *segEmits, "init_bytes", *initBytes, "seg_bytes", *segBytes) return nil}
// openSource returns a Reader + size for the upload via the supplied// Store. The returned closer must be called when the caller is done// with the source.//// The upload's Location field (a backend-specific URL/path) is// translated to a Store-relative key via store.ParseLocation. If the// configured Store can't claim the Location — e.g. the upload was// stored in S3 but the deployment is now serving file-only — the open// fails with a clear error.func openSource(ctx context.Context, store blob.Store, in Input) (io.ReaderAt, int64, func(), error) { ctx, span := vodTracer.Start(ctx, "vod.openSource", trace.WithAttributes( attribute.String("backend", in.Backend), attribute.String("location", in.Location), )) defer span.End() key, ok := store.ParseLocation(in.Location) if !ok { err := fmt.Errorf("store does not own upload location %q (backend %s)", in.Location, in.Backend) span.RecordError(err) return nil, 0, nil, err } span.SetAttributes(attribute.String("key", key)) rdr, err := store.Open(ctx, key) if err != nil { span.RecordError(err) return nil, 0, nil, fmt.Errorf("open upload: %w", err) } size := rdr.Size() span.SetAttributes(attribute.Int64("size_bytes", size)) log.Debug(ctx, "opened upload via Store", "url", store.URL(key), "size", size) return rdr, size, func() { _ = rdr.Close() }, nil}