464 lines
14 KiB
Go
464 lines
14 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log"
|
|
"log/slog"
|
|
"os"
|
|
"strings"
|
|
"time"
|
|
|
|
pbpusher "github.com/flowy-live/llink/genproto/llink/pusher"
|
|
"github.com/flowy-live/llink/internal/billing"
|
|
"github.com/flowy-live/llink/internal/db"
|
|
"github.com/flowy-live/llink/internal/depot"
|
|
"github.com/flowy-live/llink/internal/human"
|
|
"github.com/flowy-live/llink/internal/human/pushnotify"
|
|
"github.com/flowy-live/llink/internal/network"
|
|
"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"
|
|
"google.golang.org/grpc/codes"
|
|
"google.golang.org/grpc/credentials/insecure"
|
|
"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())
|
|
billingSvc := billing.NewServiceForWorker(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)
|
|
|
|
pusherAddr := utils.MustGetEnv("PUSHER_GRPC_ADDR")
|
|
pusherConn, err := grpc.NewClient(pusherAddr, grpc.WithTransportCredentials(insecure.NewCredentials()))
|
|
if err != nil {
|
|
slog.Error("failed to connect to pusher", "error", err)
|
|
os.Exit(1)
|
|
}
|
|
defer pusherConn.Close()
|
|
pusherClient := pbpusher.NewPusherServiceClient(pusherConn)
|
|
|
|
humanSvc := human.NewService(db.Pool())
|
|
networkReader := network.NewReader(db.Pool())
|
|
pushTokenSvc := pushnotify.NewService(db.Pool())
|
|
// EXPO_ACCESS_TOKEN is required: with Enhanced Security enabled on the Expo
|
|
// project, sends without it fail; without it, anyone holding one of our
|
|
// push tokens could spam our users via the public Expo endpoint.
|
|
expoClient := pushnotify.NewExpoClient(utils.MustGetEnv("EXPO_ACCESS_TOKEN"))
|
|
notifier := pushnotify.NewNotifier(networkReader, pushTokenSvc, pusherClient, expoClient)
|
|
|
|
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 ---
|
|
// All of them do not stop us from marking the particle as processed.
|
|
parentDoc := loadParentParticle(ctx, change.Doc)
|
|
updateParentLastChildCreatedAt(ctx, change.Doc, parentDoc)
|
|
transcript := transcribeMediaParticle(ctx, depotSvc, speechSvc, change.Doc)
|
|
particle.Transcode(ctx, depotSvc, change.Doc)
|
|
recordFreemiumUsage(ctx, billingSvc, change.Doc)
|
|
notifyForParticle(ctx, notifier, humanSvc, change.Doc, parentDoc, transcript)
|
|
|
|
if err := processingRepo.MarkProcessed(ctx, particleID); err != nil {
|
|
slog.Error("failed to mark particle as processed", "particleID", particleID, "error", err)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// transcribeMediaParticle transcribes a media particle, writes the structured
|
|
// transcript to Firestore, and returns the raw transcript text. Returns "" for
|
|
// non-media particles or on any error (errors are logged internally).
|
|
func transcribeMediaParticle(ctx context.Context, depotSvc depot.Service, speechSvc speech.SpeechService, doc *firestore.DocumentSnapshot) string {
|
|
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)
|
|
return transcript.Transcript
|
|
}
|
|
|
|
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,
|
|
}
|
|
}
|
|
|
|
// recordFreemiumUsage bumps the network's daily message counter for non-container
|
|
// particles. Idempotent via the surrounding processed_particles guard: the worker
|
|
// only reaches this path on first-seen particles, so a crash/restart won't
|
|
// double-count.
|
|
func recordFreemiumUsage(ctx context.Context, billingSvc billing.Service, doc *firestore.DocumentSnapshot) {
|
|
rawType, err := doc.DataAt("type")
|
|
if err != nil {
|
|
slog.Error("failed to read particle type", "error", err, "particleID", doc.Ref.ID)
|
|
return
|
|
}
|
|
typeStr, ok := rawType.(string)
|
|
if !ok {
|
|
slog.Error("particle type is not a string", "particleID", doc.Ref.ID, "type", rawType)
|
|
return
|
|
}
|
|
particleType, err := particle.ParseParticleType(typeStr)
|
|
if err != nil {
|
|
slog.Error("invalid particle type", "error", err, "particleID", doc.Ref.ID)
|
|
return
|
|
}
|
|
// Containers (stream/folder) don't count as "messages" for the daily cap.
|
|
if particleType == particle.TypeStream || particleType == particle.TypeFolder {
|
|
return
|
|
}
|
|
|
|
networkID, err := particle.NetworkIDFromParticlePath(doc.Ref.Path)
|
|
if err != nil {
|
|
slog.Error("failed to derive network id", "error", err, "path", doc.Ref.Path)
|
|
return
|
|
}
|
|
|
|
if err := billingSvc.IncrementDailyUsage(ctx, networkID, doc.CreateTime); err != nil {
|
|
slog.Error("failed to increment daily usage", "error", err, "networkID", networkID, "particleID", doc.Ref.ID)
|
|
}
|
|
}
|
|
|
|
// loadParentParticle fetches the immediate parent particle doc for `doc`.
|
|
// Returns nil (and logs) if the path doesn't have a parent or the read fails.
|
|
func loadParentParticle(ctx context.Context, doc *firestore.DocumentSnapshot) *firestore.DocumentSnapshot {
|
|
parentChildrenCollectionRef := doc.Ref.Parent
|
|
if parentChildrenCollectionRef == nil {
|
|
return nil
|
|
}
|
|
parentParticleDocRef := parentChildrenCollectionRef.Parent
|
|
if parentParticleDocRef == nil {
|
|
slog.Error("particle has no parent document", "particleID", doc.Ref.ID)
|
|
return nil
|
|
}
|
|
parentParticleDoc, err := parentParticleDocRef.Get(ctx)
|
|
if err != nil {
|
|
slog.Error("failed to get parent particle", "error", err, "particleID", doc.Ref.ID)
|
|
return nil
|
|
}
|
|
return parentParticleDoc
|
|
}
|
|
|
|
// 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, parent *firestore.DocumentSnapshot) {
|
|
if parent == nil {
|
|
return
|
|
}
|
|
|
|
var streamParticle particle.FirestoreStreamParticle
|
|
if err := parent.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
|
|
}
|
|
|
|
_, err = parent.Ref.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`", "error", err)
|
|
}
|
|
}
|
|
|
|
// notifyForParticle dispatches a push notification for a newly-created particle.
|
|
// Skips containers (streams/folders) and particles whose parent isn't a stream
|
|
// (notifications are only sent for stream messages today). The transcript arg
|
|
// is used as the preview body for media particles when available.
|
|
func notifyForParticle(
|
|
ctx context.Context,
|
|
notifier *pushnotify.Notifier,
|
|
humanSvc human.Service,
|
|
doc *firestore.DocumentSnapshot,
|
|
parent *firestore.DocumentSnapshot,
|
|
transcript string,
|
|
) {
|
|
if parent == nil {
|
|
return
|
|
}
|
|
|
|
typeStr, _ := doc.DataAt("type")
|
|
typeName, _ := typeStr.(string)
|
|
pType, err := particle.ParseParticleType(typeName)
|
|
if err != nil {
|
|
return
|
|
}
|
|
if pType == particle.TypeStream || pType == particle.TypeFolder {
|
|
return
|
|
}
|
|
|
|
var parentStream particle.FirestoreStreamParticle
|
|
if err := parent.DataTo(&parentStream); err != nil {
|
|
slog.Error("notify: failed to parse parent stream", "error", err)
|
|
return
|
|
}
|
|
parentType, err := particle.ParseParticleType(parentStream.Type)
|
|
if err != nil || parentType != particle.TypeStream {
|
|
return
|
|
}
|
|
|
|
networkID, err := particle.NetworkIDFromParticlePath(doc.Ref.Path)
|
|
if err != nil {
|
|
slog.Error("notify: failed to derive network id", "error", err, "path", doc.Ref.Path)
|
|
return
|
|
}
|
|
|
|
senderHumanID := parentStream.CreatedByHumanId
|
|
if v, err := doc.DataAt("created_by_human_id"); err == nil {
|
|
if s, ok := v.(string); ok && s != "" {
|
|
senderHumanID = s
|
|
}
|
|
}
|
|
|
|
streamName := ""
|
|
if v, err := parent.DataAt("properties.name"); err == nil {
|
|
if s, ok := v.(string); ok {
|
|
streamName = s
|
|
}
|
|
}
|
|
|
|
senderEmailPrefix := ""
|
|
if senderHumanID != "" {
|
|
if sender, err := humanSvc.GetByID(ctx, senderHumanID); err == nil {
|
|
senderEmailPrefix = sender.EmailPrefix
|
|
} else {
|
|
slog.Warn("notify: failed to look up sender", "error", err, "humanID", senderHumanID)
|
|
}
|
|
}
|
|
|
|
if err := notifier.NotifyParticleCreated(ctx, pushnotify.NotifyInput{
|
|
NetworkID: networkID,
|
|
SenderHumanID: senderHumanID,
|
|
SenderEmailPrefix: senderEmailPrefix,
|
|
ParticleID: doc.Ref.ID,
|
|
ParticleKind: string(pType),
|
|
StreamID: parent.Ref.ID,
|
|
StreamName: streamName,
|
|
StreamVisibleTo: parentStream.VisibleTo,
|
|
Body: previewForParticle(pType, doc, transcript),
|
|
}); err != nil {
|
|
slog.Error("notify: dispatch failed", "error", err, "particleID", doc.Ref.ID, "networkID", networkID)
|
|
}
|
|
}
|
|
|
|
// previewForParticle builds the visible notification body. Kept short — push
|
|
// previews truncate aggressively on lockscreens. For media, prefers the
|
|
// transcript text (already computed by transcribeMediaParticle in the same
|
|
// processing step) and falls back to the generic "Sent a ..." line if speech
|
|
// recognition produced nothing.
|
|
func previewForParticle(pType particle.ParticleType, doc *firestore.DocumentSnapshot, transcript string) string {
|
|
switch pType {
|
|
case particle.TypeText:
|
|
if v, err := doc.DataAt("properties.content"); err == nil {
|
|
if s, ok := v.(string); ok {
|
|
return truncatePreview(s, 140)
|
|
}
|
|
}
|
|
return "Sent a message"
|
|
case particle.TypeMedia:
|
|
if t := strings.TrimSpace(transcript); t != "" {
|
|
return truncatePreview(t, 140)
|
|
}
|
|
mime := ""
|
|
if v, err := doc.DataAt("properties.mime_type"); err == nil {
|
|
if s, ok := v.(string); ok {
|
|
mime = s
|
|
}
|
|
}
|
|
if strings.HasPrefix(mime, "video/") {
|
|
return "Sent a video"
|
|
}
|
|
return "Sent a voice message"
|
|
case particle.TypeFile:
|
|
return "Sent a file"
|
|
case particle.TypeQuest:
|
|
if v, err := doc.DataAt("properties.title"); err == nil {
|
|
if s, ok := v.(string); ok && s != "" {
|
|
return "Quest: " + truncatePreview(s, 120)
|
|
}
|
|
}
|
|
return "Added a quest"
|
|
case particle.TypePaper:
|
|
if v, err := doc.DataAt("properties.title"); err == nil {
|
|
if s, ok := v.(string); ok && s != "" {
|
|
return "Paper: " + truncatePreview(s, 120)
|
|
}
|
|
}
|
|
return "Added a paper"
|
|
default:
|
|
return "New activity"
|
|
}
|
|
}
|
|
|
|
func truncatePreview(s string, n int) string {
|
|
s = strings.TrimSpace(s)
|
|
if len(s) <= n {
|
|
return s
|
|
}
|
|
return s[:n] + "…"
|
|
}
|