package main import ( "context" "fmt" "log" "log/slog" "os" "time" "github.com/flowy-live/llink/internal/db" "github.com/flowy-live/llink/internal/depot" "github.com/flowy-live/llink/internal/particle" "github.com/flowy-live/llink/internal/speech" "github.com/flowy-live/llink/internal/utils" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" "cloud.google.com/go/firestore" "cloud.google.com/go/storage" ) func createClient(ctx context.Context) *firestore.Client { projectId := utils.MustGetEnv("GCP_PROJECT") client, err := firestore.NewClient(ctx, projectId) if err != nil { log.Fatalf("Failed to create client: %v", err) } return client } // The purpose of the particle processor worker is to listen for new particles // across all streams and perform side effects such as // - generate transcript if the particle is of type media // - send mobile notifications if a client is offline // - update the parent stream's `last_child_created_at` // - generate vector embedding // - synthesize and decide whether ai should generate a particle as a response func main() { ctx := context.Background() db.Init() defer db.Cleanup() processingRepo := particle.NewProcessingRepository(db.Pool()) storageClient, err := storage.NewClient(ctx) if err != nil { slog.Error("failed to create GCS client", "error", err) os.Exit(1) } defer storageClient.Close() gcsBucket := utils.MustGetEnv("GCS_BUCKET") depotSvc := depot.NewService(db.Pool(), storageClient, depot.Config{ GoogleServiceAccountEmail: utils.MustGetEnv("GOOGLE_SERVICE_ACCOUNT_EMAIL"), BucketName: gcsBucket, }) speechSvc := speech.NewSpeechService(ctx) client := createClient(ctx) defer client.Close() cutoff := time.Now().Add(-5 * time.Minute) it := client.CollectionGroup("children"). Where("created_at", ">", cutoff). Snapshots(ctx) for { snap, err := it.Next() if e := status.Code(err); e == codes.DeadlineExceeded || e == codes.Canceled { panic(fmt.Errorf("error: %w", err)) } if err != nil { slog.Error("error in processing snapshot", "error", err) continue } if snap == nil { continue } for _, change := range snap.Changes { if change.Kind != firestore.DocumentAdded { continue } particleID := change.Doc.Ref.ID processed, err := processingRepo.IsProcessed(ctx, particleID) if err != nil { slog.Error("failed to check processing status", "particleID", particleID, "error", err) continue } if processed { slog.Debug("skipping already processed particle", "particleID", particleID) continue } slog.Debug("processing particle", "particleID", particleID, "data", change.Doc.Data()) // --- Perform side effects --- updateParentLastChildCreatedAt(ctx, change.Doc) transcribeMediaParticle(ctx, depotSvc, speechSvc, change.Doc) if err := processingRepo.MarkProcessed(ctx, particleID); err != nil { slog.Error("failed to mark particle as processed", "particleID", particleID, "error", err) } } } } func transcribeMediaParticle(ctx context.Context, depotSvc depot.Service, speechSvc speech.SpeechService, doc *firestore.DocumentSnapshot) { var mediaParticle particle.FirestoreMediaParticle err := doc.DataTo(&mediaParticle) if err != nil { slog.Error("unable to marshal particle data", "error", err) return } particleType, err := particle.ParseParticleType(mediaParticle.Type) if err != nil { slog.Error("invalid particle type", "error", err) return } if particleType != particle.TypeMedia { slog.Info("received a particle of type", "particle type", particleType) return } downloadURL, err := depotSvc.GetDownloadURL(ctx, mediaParticle.Properties.ObjectId) if err != nil { slog.Error("failed to get download URL", "error", err, "object_id", mediaParticle.Properties.ObjectId) return } result, err := speechSvc.Transcribe(ctx, downloadURL) if err != nil { slog.Error("failed to transcribe media", "error", err, "particleID", doc.Ref.ID) return } transcript := toFirestoreTranscript(result) _, err = doc.Ref.Set(ctx, map[string]interface{}{ "properties": map[string]interface{}{ "transcript": transcript, }, }, firestore.MergeAll) if err != nil { slog.Error("failed to update transcript in firestore", "error", err, "particleID", doc.Ref.ID) return } slog.Info("transcribed media particle", "particleID", doc.Ref.ID) } func toFirestoreTranscript(result *speech.TranscriptResult) particle.FirestoreTranscript { words := make([]particle.FirestoreTranscriptWord, len(result.Words)) for i, w := range result.Words { words[i] = particle.FirestoreTranscriptWord{ Word: w.Word, Start: w.Start, End: w.End, } } paragraphs := make([]particle.FirestoreTranscriptParagraph, len(result.Paragraphs)) for i, p := range result.Paragraphs { sentences := make([]particle.FirestoreTranscriptSentence, len(p.Sentences)) for j, s := range p.Sentences { sentences[j] = particle.FirestoreTranscriptSentence{ Text: s.Text, Start: s.Start, End: s.End, } } paragraphs[i] = particle.FirestoreTranscriptParagraph{ Sentences: sentences, Start: p.Start, End: p.End, } } return particle.FirestoreTranscript{ Transcript: result.Transcript, Words: words, Paragraphs: paragraphs, } } // updateParentLastChildCreatedAt updates the parent stream's last_child_created_at // to the child's actual created_at timestamp, so it stays directly comparable with // playback markers (which also store child created_at values). func updateParentLastChildCreatedAt(ctx context.Context, doc *firestore.DocumentSnapshot) { parentChildrenCollectionRef := doc.Ref.Parent if parentChildrenCollectionRef == nil { return } parentParticleDocRef := parentChildrenCollectionRef.Parent if parentParticleDocRef == nil { slog.Error("particle has no parent document", "particleID", doc.Ref.ID) return } parentParticleDoc, err := parentParticleDocRef.Get(ctx) if err != nil { slog.Error("failed to get parent particle", "error", err) return } slog.Info("parent particle is", "parent particle id", parentParticleDoc.Ref.ID) var streamParticle particle.FirestoreStreamParticle if err := parentParticleDoc.DataTo(&streamParticle); err != nil { slog.Error("failed to parse stream particle", "error", err) return } particleType, err := particle.ParseParticleType(streamParticle.Type) if err != nil { slog.Error("invalid particle type", "error", err) return } if particleType != particle.TypeStream { return } // Read the child's created_at — this is the same value that playback markers store childCreatedAt, err := doc.DataAt("created_at") if err != nil { slog.Error("failed to read child created_at", "error", err, "particleID", doc.Ref.ID) return } slog.Info("going to update the last_child_created_at for parent particle") _, err = parentParticleDocRef.Update(ctx, []firestore.Update{ { Path: "last_child_created_at", Value: childCreatedAt, }, }) if err != nil { slog.Error("unable to update parent particle `last_child_created_at`") } }