transcode media particles to mp4
This commit is contained in:
@@ -14,6 +14,8 @@ RUN ls
|
|||||||
|
|
||||||
FROM alpine:latest AS second-stage
|
FROM alpine:latest AS second-stage
|
||||||
|
|
||||||
|
RUN apk add --no-cache ffmpeg ca-certificates
|
||||||
|
|
||||||
WORKDIR /app
|
WORKDIR /app
|
||||||
COPY --from=first-stage /app/cmd/particleprocessorworker .
|
COPY --from=first-stage /app/cmd/particleprocessorworker .
|
||||||
RUN echo "copied over binary to production stage"
|
RUN echo "copied over binary to production stage"
|
||||||
|
|||||||
@@ -1,11 +1,13 @@
|
|||||||
package main
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"os"
|
"os"
|
||||||
|
"os/exec"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -107,6 +109,7 @@ func main() {
|
|||||||
|
|
||||||
updateParentLastChildCreatedAt(ctx, change.Doc)
|
updateParentLastChildCreatedAt(ctx, change.Doc)
|
||||||
transcribeMediaParticle(ctx, depotSvc, speechSvc, change.Doc)
|
transcribeMediaParticle(ctx, depotSvc, speechSvc, change.Doc)
|
||||||
|
transcodeMediaParticle(ctx, depotSvc, change.Doc)
|
||||||
recordFreemiumUsage(ctx, billingSvc, change.Doc)
|
recordFreemiumUsage(ctx, billingSvc, change.Doc)
|
||||||
|
|
||||||
if err := processingRepo.MarkProcessed(ctx, particleID); err != nil {
|
if err := processingRepo.MarkProcessed(ctx, particleID); err != nil {
|
||||||
@@ -162,6 +165,142 @@ func transcribeMediaParticle(ctx context.Context, depotSvc depot.Service, speech
|
|||||||
slog.Info("transcribed media particle", "particleID", doc.Ref.ID)
|
slog.Info("transcribed media particle", "particleID", doc.Ref.ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// transcodeMediaParticle produces an iOS-playable MP4/m4a derivative for media
|
||||||
|
// particles whose original mime type AVPlayer can't decode (notably the WebM
|
||||||
|
// the desktop recorder emits today). Skips when the source is already in an
|
||||||
|
// iOS-playable family or when a transcoded variant has already been written.
|
||||||
|
//
|
||||||
|
// ffmpeg reads directly from the GCS signed URL and writes to a local temp
|
||||||
|
// file — `+faststart` requires seekable output, so a stdout pipe wouldn't work.
|
||||||
|
// The temp file is then streamed to GCS via depotSvc.CreateFromReader (no
|
||||||
|
// presigned-PUT round-trip — the worker has direct SDK access).
|
||||||
|
func transcodeMediaParticle(ctx context.Context, depotSvc depot.Service, doc *firestore.DocumentSnapshot) {
|
||||||
|
var mediaParticle particle.FirestoreMediaParticle
|
||||||
|
if err := doc.DataTo(&mediaParticle); err != nil {
|
||||||
|
slog.Error("transcode: unable to marshal particle data", "error", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
particleType, err := particle.ParseParticleType(mediaParticle.Type)
|
||||||
|
if err != nil {
|
||||||
|
slog.Error("transcode: invalid particle type", "error", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if particleType != particle.TypeMedia {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// Already transcoded — re-delivery within the 5-min Firestore window.
|
||||||
|
if mediaParticle.Properties.TranscodedObjectId != "" {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// Source is already iOS-playable; nothing to do.
|
||||||
|
if isIOSPlayableMime(mediaParticle.Properties.MimeType) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
sourceURL, err := depotSvc.GetDownloadURL(ctx, mediaParticle.Properties.ObjectId)
|
||||||
|
if err != nil {
|
||||||
|
slog.Error("transcode: failed to get download URL", "error", err, "object_id", mediaParticle.Properties.ObjectId)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
isAudio := strings.HasPrefix(mediaParticle.Properties.MimeType, "audio/")
|
||||||
|
var outputExt, outputMime string
|
||||||
|
if isAudio {
|
||||||
|
outputExt = ".m4a"
|
||||||
|
outputMime = "audio/mp4"
|
||||||
|
} else {
|
||||||
|
outputExt = ".mp4"
|
||||||
|
outputMime = "video/mp4"
|
||||||
|
}
|
||||||
|
|
||||||
|
tmp, err := os.CreateTemp("", "transcode-*"+outputExt)
|
||||||
|
if err != nil {
|
||||||
|
slog.Error("transcode: failed to create temp file", "error", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
tmpPath := tmp.Name()
|
||||||
|
tmp.Close()
|
||||||
|
defer os.Remove(tmpPath)
|
||||||
|
|
||||||
|
transcodeCtx, cancel := context.WithTimeout(ctx, 5*time.Minute)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
var args []string
|
||||||
|
if isAudio {
|
||||||
|
args = []string{
|
||||||
|
"-y", "-i", sourceURL,
|
||||||
|
"-vn",
|
||||||
|
"-c:a", "aac", "-b:a", "128k",
|
||||||
|
"-movflags", "+faststart",
|
||||||
|
tmpPath,
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
args = []string{
|
||||||
|
"-y", "-i", sourceURL,
|
||||||
|
"-c:v", "libx264", "-preset", "veryfast", "-crf", "23",
|
||||||
|
"-pix_fmt", "yuv420p", "-profile:v", "baseline", "-level", "3.1",
|
||||||
|
"-c:a", "aac", "-b:a", "128k",
|
||||||
|
"-movflags", "+faststart",
|
||||||
|
tmpPath,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
cmd := exec.CommandContext(transcodeCtx, "ffmpeg", args...)
|
||||||
|
var stderr bytes.Buffer
|
||||||
|
cmd.Stderr = &stderr
|
||||||
|
if err := cmd.Run(); err != nil {
|
||||||
|
slog.Error("transcode: ffmpeg failed", "error", err, "stderr", stderr.String(), "particleID", doc.Ref.ID)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
networkID, err := networkIDFromParticlePath(doc.Ref.Path)
|
||||||
|
if err != nil {
|
||||||
|
slog.Error("transcode: failed to derive network id", "error", err, "path", doc.Ref.Path)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
f, err := os.Open(tmpPath)
|
||||||
|
if err != nil {
|
||||||
|
slog.Error("transcode: failed to open transcoded file", "error", err, "path", tmpPath)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
defer f.Close()
|
||||||
|
|
||||||
|
newObj, err := depotSvc.CreateFromReader(ctx, depot.CreateFromReaderInput{
|
||||||
|
Prefix: networkID,
|
||||||
|
Name: "transcoded" + outputExt,
|
||||||
|
ContentType: outputMime,
|
||||||
|
}, f)
|
||||||
|
if err != nil {
|
||||||
|
slog.Error("transcode: failed to upload transcoded object", "error", err, "particleID", doc.Ref.ID)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = doc.Ref.Set(ctx, map[string]interface{}{
|
||||||
|
"properties": map[string]interface{}{
|
||||||
|
"transcoded_object_id": newObj.ID,
|
||||||
|
"transcoded_mime_type": outputMime,
|
||||||
|
},
|
||||||
|
}, firestore.MergeAll)
|
||||||
|
if err != nil {
|
||||||
|
slog.Error("transcode: failed to update particle in firestore", "error", err, "particleID", doc.Ref.ID)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
slog.Info("transcoded media particle", "particleID", doc.Ref.ID, "transcoded_object_id", newObj.ID, "mime", outputMime)
|
||||||
|
}
|
||||||
|
|
||||||
|
func isIOSPlayableMime(mime string) bool {
|
||||||
|
switch mime {
|
||||||
|
case "video/mp4", "video/quicktime", "audio/mp4", "audio/aac", "audio/x-m4a", "audio/mpeg":
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
func toFirestoreTranscript(result *speech.TranscriptResult) particle.FirestoreTranscript {
|
func toFirestoreTranscript(result *speech.TranscriptResult) particle.FirestoreTranscript {
|
||||||
words := make([]particle.FirestoreTranscriptWord, len(result.Words))
|
words := make([]particle.FirestoreTranscriptWord, len(result.Words))
|
||||||
for i, w := range result.Words {
|
for i, w := range result.Words {
|
||||||
|
|||||||
@@ -29,6 +29,15 @@ type PrepareUploadResult struct {
|
|||||||
UploadHeaders map[string]string
|
UploadHeaders map[string]string
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// CreateFromReaderInput is for server-side direct uploads (no presigned URL).
|
||||||
|
// Used by background workers that already have the bytes on hand and don't
|
||||||
|
// need a client round-trip.
|
||||||
|
type CreateFromReaderInput struct {
|
||||||
|
Prefix string // Optional prefix for organizing objects (e.g., network_id)
|
||||||
|
Name string
|
||||||
|
ContentType string
|
||||||
|
}
|
||||||
|
|
||||||
// Config holds configuration for the depot service
|
// Config holds configuration for the depot service
|
||||||
type Config struct {
|
type Config struct {
|
||||||
GoogleServiceAccountEmail string
|
GoogleServiceAccountEmail string
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -20,6 +21,7 @@ const (
|
|||||||
type Service interface {
|
type Service interface {
|
||||||
PrepareUpload(ctx context.Context, input PrepareUploadInput) (*PrepareUploadResult, error)
|
PrepareUpload(ctx context.Context, input PrepareUploadInput) (*PrepareUploadResult, error)
|
||||||
ConfirmUpload(ctx context.Context, objectID string) (*Object, error)
|
ConfirmUpload(ctx context.Context, objectID string) (*Object, error)
|
||||||
|
CreateFromReader(ctx context.Context, input CreateFromReaderInput, body io.Reader) (*Object, error)
|
||||||
GetByID(ctx context.Context, objectID string) (*Object, error)
|
GetByID(ctx context.Context, objectID string) (*Object, error)
|
||||||
GetDownloadURL(ctx context.Context, objectID string) (string, error)
|
GetDownloadURL(ctx context.Context, objectID string) (string, error)
|
||||||
Delete(ctx context.Context, objectID string) error
|
Delete(ctx context.Context, objectID string) error
|
||||||
@@ -155,6 +157,56 @@ func (s *serviceImpl) ConfirmUpload(ctx context.Context, objectID string) (*Obje
|
|||||||
return s.repo.getByID(ctx, objectID)
|
return s.repo.getByID(ctx, objectID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// CreateFromReader streams bytes directly to GCS using the storage client and
|
||||||
|
// records the depot_objects row in one shot. Unlike PrepareUpload, there is no
|
||||||
|
// signed URL or client round-trip — the caller already has the bytes. Intended
|
||||||
|
// for worker-side flows (e.g. transcoded media variants).
|
||||||
|
func (s *serviceImpl) CreateFromReader(ctx context.Context, input CreateFromReaderInput, body io.Reader) (*Object, error) {
|
||||||
|
if input.Name == "" {
|
||||||
|
return nil, errors.Join(ErrInvalidInput, errors.New("name is required"))
|
||||||
|
}
|
||||||
|
if input.ContentType == "" {
|
||||||
|
return nil, errors.Join(ErrInvalidInput, errors.New("content_type is required"))
|
||||||
|
}
|
||||||
|
|
||||||
|
objectKey := fmt.Sprintf("%s/%s/%s", input.Prefix, uuid.New().String(), input.Name)
|
||||||
|
|
||||||
|
w := s.storageClient.Bucket(s.bucketName).Object(objectKey).NewWriter(ctx)
|
||||||
|
w.ContentType = input.ContentType
|
||||||
|
if _, err := io.Copy(w, body); err != nil {
|
||||||
|
// Close to release resources, then surface the original copy error.
|
||||||
|
if cerr := w.Close(); cerr != nil {
|
||||||
|
slog.Warn("failed to close GCS writer after copy failure", "error", cerr, "object_key", objectKey)
|
||||||
|
}
|
||||||
|
slog.Error("failed to stream object to GCS", "error", err, "bucket", s.bucketName, "object_key", objectKey)
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if err := w.Close(); err != nil {
|
||||||
|
slog.Error("failed to close GCS writer", "error", err, "bucket", s.bucketName, "object_key", objectKey)
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
obj := &Object{
|
||||||
|
Name: input.Name,
|
||||||
|
ContentType: input.ContentType,
|
||||||
|
ContentLength: w.Attrs().Size,
|
||||||
|
BucketName: s.bucketName,
|
||||||
|
ObjectKey: objectKey,
|
||||||
|
ContainsContent: true,
|
||||||
|
}
|
||||||
|
|
||||||
|
created, err := s.repo.create(ctx, obj)
|
||||||
|
if err != nil {
|
||||||
|
// Best-effort: clean up the GCS object since we can't track it in the DB.
|
||||||
|
if delErr := s.storageClient.Bucket(s.bucketName).Object(objectKey).Delete(ctx); delErr != nil {
|
||||||
|
slog.Warn("failed to clean up GCS object after db create failure", "error", delErr, "object_key", objectKey)
|
||||||
|
}
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return created, nil
|
||||||
|
}
|
||||||
|
|
||||||
func (s *serviceImpl) GetByID(ctx context.Context, objectID string) (*Object, error) {
|
func (s *serviceImpl) GetByID(ctx context.Context, objectID string) (*Object, error) {
|
||||||
obj, err := s.repo.getByID(ctx, objectID)
|
obj, err := s.repo.getByID(ctx, objectID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -35,11 +35,13 @@ type FirestoreTranscript struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type FirestoreMediaParticleProperties struct {
|
type FirestoreMediaParticleProperties struct {
|
||||||
ObjectId string `firestore:"object_id"`
|
ObjectId string `firestore:"object_id"`
|
||||||
MimeType string `firestore:"mime_type"`
|
MimeType string `firestore:"mime_type"`
|
||||||
DurationMs int `firestore:"duration_ms"`
|
DurationMs int `firestore:"duration_ms"`
|
||||||
SizeBytes int `firestore:"size_bytes"`
|
SizeBytes int `firestore:"size_bytes"`
|
||||||
Transcript *FirestoreTranscript `firestore:"transcript,omitempty"`
|
Transcript *FirestoreTranscript `firestore:"transcript,omitempty"`
|
||||||
|
TranscodedObjectId string `firestore:"transcoded_object_id,omitempty"`
|
||||||
|
TranscodedMimeType string `firestore:"transcoded_mime_type,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type FirestoreStreamParticle struct {
|
type FirestoreStreamParticle struct {
|
||||||
|
|||||||
Reference in New Issue
Block a user