diff --git a/go/Dockerfile.particleprocessorworker b/go/Dockerfile.particleprocessorworker index f68a2ac..385a93c 100644 --- a/go/Dockerfile.particleprocessorworker +++ b/go/Dockerfile.particleprocessorworker @@ -14,6 +14,8 @@ RUN ls FROM alpine:latest AS second-stage +RUN apk add --no-cache ffmpeg ca-certificates + WORKDIR /app COPY --from=first-stage /app/cmd/particleprocessorworker . RUN echo "copied over binary to production stage" diff --git a/go/cmd/particleprocessorworker/main.go b/go/cmd/particleprocessorworker/main.go index 6642d42..35e482d 100644 --- a/go/cmd/particleprocessorworker/main.go +++ b/go/cmd/particleprocessorworker/main.go @@ -1,11 +1,13 @@ package main import ( + "bytes" "context" "fmt" "log" "log/slog" "os" + "os/exec" "strings" "time" @@ -107,6 +109,7 @@ func main() { updateParentLastChildCreatedAt(ctx, change.Doc) transcribeMediaParticle(ctx, depotSvc, speechSvc, change.Doc) + transcodeMediaParticle(ctx, depotSvc, change.Doc) recordFreemiumUsage(ctx, billingSvc, change.Doc) 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) } +// 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 { words := make([]particle.FirestoreTranscriptWord, len(result.Words)) for i, w := range result.Words { diff --git a/go/internal/depot/models.go b/go/internal/depot/models.go index 3f11d1e..8042a72 100644 --- a/go/internal/depot/models.go +++ b/go/internal/depot/models.go @@ -29,6 +29,15 @@ type PrepareUploadResult struct { 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 type Config struct { GoogleServiceAccountEmail string diff --git a/go/internal/depot/service.go b/go/internal/depot/service.go index a9c1290..4026552 100644 --- a/go/internal/depot/service.go +++ b/go/internal/depot/service.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "io" "log/slog" "time" @@ -20,6 +21,7 @@ const ( type Service interface { PrepareUpload(ctx context.Context, input PrepareUploadInput) (*PrepareUploadResult, 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) GetDownloadURL(ctx context.Context, objectID string) (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) } +// 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) { obj, err := s.repo.getByID(ctx, objectID) if err != nil { diff --git a/go/internal/particle/firestore_types.go b/go/internal/particle/firestore_types.go index fd93ff1..16085d5 100644 --- a/go/internal/particle/firestore_types.go +++ b/go/internal/particle/firestore_types.go @@ -35,11 +35,13 @@ type FirestoreTranscript struct { } type FirestoreMediaParticleProperties struct { - ObjectId string `firestore:"object_id"` - MimeType string `firestore:"mime_type"` - DurationMs int `firestore:"duration_ms"` - SizeBytes int `firestore:"size_bytes"` - Transcript *FirestoreTranscript `firestore:"transcript,omitempty"` + ObjectId string `firestore:"object_id"` + MimeType string `firestore:"mime_type"` + DurationMs int `firestore:"duration_ms"` + SizeBytes int `firestore:"size_bytes"` + Transcript *FirestoreTranscript `firestore:"transcript,omitempty"` + TranscodedObjectId string `firestore:"transcoded_object_id,omitempty"` + TranscodedMimeType string `firestore:"transcoded_mime_type,omitempty"` } type FirestoreStreamParticle struct {