diff --git a/go/cmd/emailnotifierjob/main.go b/go/cmd/emailnotifierjob/main.go index 14be3e6..fb39163 100644 --- a/go/cmd/emailnotifierjob/main.go +++ b/go/cmd/emailnotifierjob/main.go @@ -3,11 +3,15 @@ package main import ( "context" "fmt" - "log/slog" "os" "time" + "github.com/flowy-live/llink/internal/utils/flog" + "cloud.google.com/go/firestore" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + pbaero "github.com/flowy-live/llink/genproto/aero" pbpusher "github.com/flowy-live/llink/genproto/llink/pusher" "github.com/flowy-live/llink/internal/db" @@ -15,12 +19,10 @@ import ( "github.com/flowy-live/llink/internal/network" "github.com/flowy-live/llink/internal/particle" "github.com/flowy-live/llink/internal/utils" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials/insecure" ) const ( - maxActivityAge = 24 * time.Hour // ignore streams idle longer than this + maxActivityAge = 24 * time.Hour // ignore streams idle longer than this unreadThreshold = 10 * time.Minute // grace window before a message is "unread" emailCooldown = 12 * time.Hour // min gap between emails to the same user ) @@ -34,7 +36,7 @@ func main() { gcpProject := utils.MustGetEnv("GCP_PROJECT") firestoreClient, err := firestore.NewClient(ctx, gcpProject) if err != nil { - slog.Error("failed to create Firestore client", "error", err) + flog.Error("failed to create Firestore client", "error", err) os.Exit(1) } defer firestoreClient.Close() @@ -42,7 +44,7 @@ func main() { aeroAddr := utils.MustGetEnv("AERO_ADDR") aeroConn, err := grpc.NewClient(aeroAddr, grpc.WithTransportCredentials(insecure.NewCredentials())) if err != nil { - slog.Error("failed to connect to aero", "error", err) + flog.Error("failed to connect to aero", "error", err) os.Exit(1) } defer aeroConn.Close() @@ -51,7 +53,7 @@ func main() { 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) + flog.Error("failed to connect to pusher", "error", err) os.Exit(1) } defer pusherConn.Close() @@ -60,12 +62,12 @@ func main() { humanSvc := human.NewService(db.Pool()) networkSvc := network.NewReader(db.Pool()) - slog.Info("starting email notification cycle") + flog.Info("starting email notification cycle") if err := runNotificationCycle(ctx, firestoreClient, aeroSvc, pusherSvc, humanSvc, networkSvc); err != nil { - slog.Error("notification cycle failed", "error", err) + flog.Error("notification cycle failed", "error", err) os.Exit(1) } - slog.Info("email notification cycle complete") + flog.Info("email notification cycle complete") } func runNotificationCycle( @@ -100,7 +102,7 @@ func runNotificationCycle( if err != nil { return fmt.Errorf("failed to get online humans: %w", err) } - slog.Info("gathered online presence", "onlineCount", len(onlineResp.HumanIds), "humanIds", onlineResp.HumanIds) + flog.Info("gathered online presence", "onlineCount", len(onlineResp.HumanIds), "humanIds", onlineResp.HumanIds) for _, id := range onlineResp.HumanIds { allOnline[id] = true } @@ -108,7 +110,7 @@ func runNotificationCycle( for _, net := range networks { streams, err := getOpenStreams(ctx, fsClient, net.ID) if err != nil { - slog.Error("failed to query streams", "networkId", net.ID, "error", err) + flog.Error("failed to query streams", "networkId", net.ID, "error", err) continue } @@ -141,7 +143,7 @@ func runNotificationCycle( sentCount := 0 for humanId, count := range behindCounts { if allOnline[humanId] { - slog.Info("human online...skipping email", "humanId", humanId) + flog.Info("human online...skipping email", "humanId", humanId) continue } @@ -164,18 +166,18 @@ func runNotificationCycle( } if err := sendNotificationEmail(ctx, aeroSvc, h, count); err != nil { - slog.Error("failed to send email", "humanId", humanId, "error", err) + flog.Error("failed to send email", "humanId", humanId, "error", err) continue } if err := humanSvc.UpdateLastEmailNotificationSentAt(ctx, humanId, now); err != nil { - slog.Error("failed to update last_email_notification_sent_at", "humanId", humanId, "error", err) + flog.Error("failed to update last_email_notification_sent_at", "humanId", humanId, "error", err) } sentCount++ } - slog.Info("notification cycle summary", + flog.Info("notification cycle summary", "networks", len(networks), "humansBehind", len(behindCounts), "emailsSent", sentCount, @@ -198,7 +200,7 @@ func getOpenStreams(ctx context.Context, client *firestore.Client, networkId str for _, doc := range docs { var s particle.FirestoreStreamParticle if err := doc.DataTo(&s); err != nil { - slog.Warn("failed to unmarshal stream particle", "docId", doc.Ref.ID, "error", err) + flog.Warn("failed to unmarshal stream particle", "docId", doc.Ref.ID, "error", err) continue } streams = append(streams, s) diff --git a/go/cmd/memberreconciler/main.go b/go/cmd/memberreconciler/main.go index 02b4528..bc172b0 100644 --- a/go/cmd/memberreconciler/main.go +++ b/go/cmd/memberreconciler/main.go @@ -6,17 +6,19 @@ package main import ( "context" - "log/slog" "os" "slices" "time" + "github.com/flowy-live/llink/internal/utils/flog" + "cloud.google.com/go/firestore" + "google.golang.org/api/iterator" + "github.com/flowy-live/llink/internal/db" "github.com/flowy-live/llink/internal/human" "github.com/flowy-live/llink/internal/network" "github.com/flowy-live/llink/internal/utils" - "google.golang.org/api/iterator" ) func main() { @@ -28,7 +30,7 @@ func main() { gcpProject := utils.MustGetEnv("GCP_PROJECT") fs, err := firestore.NewClient(ctx, gcpProject) if err != nil { - slog.Error("failed to create Firestore client", "error", err) + flog.Error("failed to create Firestore client", "error", err) os.Exit(1) } defer fs.Close() @@ -39,10 +41,10 @@ func main() { started := time.Now() written, scanned, err := reconcile(ctx, fs, humanSvc, networkSvc) if err != nil { - slog.Error("reconciliation failed", "error", err, "elapsed", time.Since(started)) + flog.Error("reconciliation failed", "error", err, "elapsed", time.Since(started)) os.Exit(1) } - slog.Info("reconciliation complete", + flog.Info("reconciliation complete", "humans_scanned", scanned, "humans_written", written, "elapsed", time.Since(started), @@ -109,7 +111,7 @@ func snapshotMirror(ctx context.Context, fs *firestore.Client) (map[string][]str Networks []string `firestore:"networks"` } if err := doc.DataTo(&data); err != nil { - slog.Warn("skipping malformed mirror doc", "id", doc.Ref.ID, "error", err) + flog.Warn("skipping malformed mirror doc", "id", doc.Ref.ID, "error", err) continue } out[doc.Ref.ID] = data.Networks diff --git a/go/cmd/orion/main.go b/go/cmd/orion/main.go index b261048..11395bc 100644 --- a/go/cmd/orion/main.go +++ b/go/cmd/orion/main.go @@ -3,14 +3,19 @@ package main import ( "context" "fmt" - "log/slog" "net/http" "os" "strings" + "github.com/flowy-live/llink/internal/utils/flog" + "cloud.google.com/go/firestore" "cloud.google.com/go/storage" firebase "firebase.google.com/go/v4" + "github.com/redis/go-redis/v9" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + pbaero "github.com/flowy-live/llink/genproto/aero" "github.com/flowy-live/llink/internal" "github.com/flowy-live/llink/internal/auth" @@ -27,9 +32,6 @@ import ( "github.com/flowy-live/llink/internal/particle" "github.com/flowy-live/llink/internal/utils" "github.com/flowy-live/llink/internal/waitlist" - "github.com/redis/go-redis/v9" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials/insecure" ) func redisForAuth() *redis.Client { @@ -48,19 +50,19 @@ func main() { ctx := context.Background() storageClient, err := storage.NewClient(ctx) if err != nil { - slog.Error("failed to create GCS client", "error", err) + flog.Error("failed to create GCS client", "error", err) os.Exit(1) } defer storageClient.Close() aeroAddr := utils.MustGetEnv("AERO_ADDR") if aeroAddr == "" { - slog.Error("must provide AERO_ADDR") + flog.Error("must provide AERO_ADDR") os.Exit(1) } aeroServer, err := grpc.NewClient(aeroAddr, grpc.WithTransportCredentials(insecure.NewCredentials())) if err != nil { - slog.Error("connection to aero server invalid", "error", err) + flog.Error("connection to aero server invalid", "error", err) os.Exit(1) } defer aeroServer.Close() @@ -69,12 +71,12 @@ func main() { gcpProject := utils.MustGetEnv("GCP_PROJECT") fbApp, err := firebase.NewApp(ctx, &firebase.Config{ProjectID: gcpProject}) if err != nil { - slog.Error("failed to init Firebase Admin app", "error", err) + flog.Error("failed to init Firebase Admin app", "error", err) os.Exit(1) } fbAuth, err := fbApp.Auth(ctx) if err != nil { - slog.Error("failed to create Firebase auth client", "error", err) + flog.Error("failed to create Firebase auth client", "error", err) os.Exit(1) } @@ -90,13 +92,13 @@ func main() { CancelURL: utils.MustGetEnv("BILLING_CANCEL_URL"), }) if err != nil { - slog.Error("failed to initialize billing service", "error", err) + flog.Error("failed to initialize billing service", "error", err) os.Exit(1) } firestoreClient, err := firestore.NewClient(ctx, gcpProject) if err != nil { - slog.Error("failed to create Firestore client", "error", err) + flog.Error("failed to create Firestore client", "error", err) os.Exit(1) } defer firestoreClient.Close() @@ -202,9 +204,9 @@ func main() { muxWithCors := middleware.CORS(allowedOrigins)(mux) addr := fmt.Sprintf("0.0.0.0:%s", port) - slog.Info("running server", "addr", addr) + flog.Info("running server", "addr", addr) if err := http.ListenAndServe(addr, muxWithCors); err != nil { - slog.Error("server failed", "error", err) + flog.Error("server failed", "error", err) os.Exit(1) } } diff --git a/go/cmd/particleprocessorworker/main.go b/go/cmd/particleprocessorworker/main.go index 51e5da0..9e9ae53 100644 --- a/go/cmd/particleprocessorworker/main.go +++ b/go/cmd/particleprocessorworker/main.go @@ -4,11 +4,15 @@ import ( "context" "fmt" "log" - "log/slog" "os" "strings" "time" + "github.com/flowy-live/llink/internal/utils/flog" + + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + "github.com/flowy-live/llink/internal/billing" "github.com/flowy-live/llink/internal/db" "github.com/flowy-live/llink/internal/depot" @@ -18,8 +22,6 @@ import ( "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" @@ -48,7 +50,7 @@ func main() { storageClient, err := storage.NewClient(ctx) if err != nil { - slog.Error("failed to create GCS client", "error", err) + flog.Error("failed to create GCS client", "error", err) os.Exit(1) } defer storageClient.Close() @@ -83,7 +85,7 @@ func main() { panic(fmt.Errorf("error: %w", err)) } if err != nil { - slog.Error("error in processing snapshot", "error", err) + flog.Error("error in processing snapshot", "error", err) continue } @@ -100,15 +102,15 @@ func main() { processed, err := processingRepo.IsProcessed(ctx, particleID) if err != nil { - slog.Error("failed to check processing status", "particleID", particleID, "error", err) + flog.Error("failed to check processing status", "particleID", particleID, "error", err) continue } if processed { - slog.Debug("skipping already processed particle", "particleID", particleID) + flog.Debug("skipping already processed particle", "particleID", particleID) continue } - slog.Debug("processing particle", "particleID", particleID, "data", change.Doc.Data()) + flog.Debug("processing particle", "particleID", particleID, "data", change.Doc.Data()) // Side effects below are best-effort — failures don't prevent // marking the particle as processed. @@ -120,7 +122,7 @@ func main() { 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) + flog.Error("failed to mark particle as processed", "particleID", particleID, "error", err) } } } @@ -132,30 +134,30 @@ func transcribeMediaParticle(ctx context.Context, depotSvc depot.Service, speech var mediaParticle particle.FirestoreMediaParticle err := doc.DataTo(&mediaParticle) if err != nil { - slog.Error("unable to marshal particle data", "error", err) + flog.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) + flog.Error("invalid particle type", "error", err) return "" } if particleType != particle.TypeMedia { - slog.Info("received a particle of type", "particle type", particleType) + flog.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) + flog.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) + flog.Error("failed to transcribe media", "error", err, "particleID", doc.Ref.ID) return "" } @@ -167,11 +169,11 @@ func transcribeMediaParticle(ctx context.Context, depotSvc depot.Service, speech }, }, firestore.MergeAll) if err != nil { - slog.Error("failed to update transcript in firestore", "error", err, "particleID", doc.Ref.ID) + flog.Error("failed to update transcript in firestore", "error", err, "particleID", doc.Ref.ID) return "" } - slog.Info("transcribed media particle", "particleID", doc.Ref.ID) + flog.Info("transcribed media particle", "particleID", doc.Ref.ID) return transcript.Transcript } @@ -215,17 +217,17 @@ func toFirestoreTranscript(result *speech.TranscriptResult) particle.FirestoreTr 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) + flog.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) + flog.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) + flog.Error("invalid particle type", "error", err, "particleID", doc.Ref.ID) return } // Containers don't count toward the daily message cap. @@ -235,12 +237,12 @@ func recordFreemiumUsage(ctx context.Context, billingSvc billing.Service, doc *f networkID, err := particle.NetworkIDFromParticlePath(doc.Ref.Path) if err != nil { - slog.Error("failed to derive network id", "error", err, "path", doc.Ref.Path) + flog.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) + flog.Error("failed to increment daily usage", "error", err, "networkID", networkID, "particleID", doc.Ref.ID) } } @@ -252,12 +254,12 @@ func loadParentParticle(ctx context.Context, doc *firestore.DocumentSnapshot) *f } parentParticleDocRef := parentChildrenCollectionRef.Parent if parentParticleDocRef == nil { - slog.Error("particle has no parent document", "particleID", doc.Ref.ID) + flog.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) + flog.Error("failed to get parent particle", "error", err, "particleID", doc.Ref.ID) return nil } return parentParticleDoc @@ -272,13 +274,13 @@ func updateParentLastChildCreatedAt(ctx context.Context, doc *firestore.Document var streamParticle particle.FirestoreStreamParticle if err := parent.DataTo(&streamParticle); err != nil { - slog.Error("failed to parse stream particle", "error", err) + flog.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) + flog.Error("invalid particle type", "error", err) return } @@ -288,7 +290,7 @@ func updateParentLastChildCreatedAt(ctx context.Context, doc *firestore.Document childCreatedAt, err := doc.DataAt("created_at") if err != nil { - slog.Error("failed to read child created_at", "error", err, "particleID", doc.Ref.ID) + flog.Error("failed to read child created_at", "error", err, "particleID", doc.Ref.ID) return } @@ -299,7 +301,7 @@ func updateParentLastChildCreatedAt(ctx context.Context, doc *firestore.Document }, }) if err != nil { - slog.Error("unable to update parent particle `last_child_created_at`", "error", err) + flog.Error("unable to update parent particle `last_child_created_at`", "error", err) } } @@ -315,7 +317,7 @@ func notifyForParticle( transcript string, ) { if parent == nil { - slog.Info("notify: skip — no parent", "particleID", doc.Ref.ID) + flog.Info("notify: skip — no parent", "particleID", doc.Ref.ID) return } @@ -323,31 +325,31 @@ func notifyForParticle( typeName, _ := typeStr.(string) pType, err := particle.ParseParticleType(typeName) if err != nil { - slog.Info("notify: skip — unparseable particle type", + flog.Info("notify: skip — unparseable particle type", "particleID", doc.Ref.ID, "type", typeName, "error", err) return } if pType == particle.TypeStream || pType == particle.TypeFolder { - slog.Info("notify: skip — container particle", + flog.Info("notify: skip — container particle", "particleID", doc.Ref.ID, "type", pType) return } var parentStream particle.FirestoreStreamParticle if err := parent.DataTo(&parentStream); err != nil { - slog.Error("notify: failed to parse parent stream", "error", err) + flog.Error("notify: failed to parse parent stream", "error", err) return } parentType, err := particle.ParseParticleType(parentStream.Type) if err != nil || parentType != particle.TypeStream { - slog.Info("notify: skip — parent isn't a stream", + flog.Info("notify: skip — parent isn't a stream", "particleID", doc.Ref.ID, "parentType", parentType, "parseErr", err) 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) + flog.Error("notify: failed to derive network id", "error", err, "path", doc.Ref.Path) return } @@ -370,7 +372,7 @@ func notifyForParticle( 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) + flog.Warn("notify: failed to look up sender", "error", err, "humanID", senderHumanID) } } @@ -385,7 +387,7 @@ func notifyForParticle( StreamVisibleTo: parentStream.VisibleTo, Body: previewForParticle(pType, doc, transcript), }); err != nil { - slog.Error("notify: dispatch failed", "error", err, "particleID", doc.Ref.ID, "networkID", networkID) + flog.Error("notify: dispatch failed", "error", err, "particleID", doc.Ref.ID, "networkID", networkID) } } diff --git a/go/cmd/pusherservice/main.go b/go/cmd/pusherservice/main.go index 2dee5ca..1d9155e 100644 --- a/go/cmd/pusherservice/main.go +++ b/go/cmd/pusherservice/main.go @@ -3,13 +3,14 @@ package main import ( "context" "fmt" - "log/slog" "net" "net/http" "os" "os/signal" "syscall" + "github.com/flowy-live/llink/internal/utils/flog" + "github.com/flowy-live/llink/internal" "github.com/flowy-live/llink/internal/auth" "github.com/flowy-live/llink/internal/db" @@ -17,8 +18,9 @@ import ( "github.com/flowy-live/llink/internal/pusher" "github.com/flowy-live/llink/internal/utils" - pbpusher "github.com/flowy-live/llink/genproto/llink/pusher" "google.golang.org/grpc" + + pbpusher "github.com/flowy-live/llink/genproto/llink/pusher" ) func main() { @@ -56,15 +58,15 @@ func main() { // --- gRPC server (internal presence queries) --- grpcListener, err := net.Listen("tcp", fmt.Sprintf("0.0.0.0:%s", grpcPort)) if err != nil { - slog.Error("failed to listen for gRPC", "port", grpcPort, "error", err) + flog.Error("failed to listen for gRPC", "port", grpcPort, "error", err) os.Exit(1) } grpcServer := grpc.NewServer() pbpusher.RegisterPusherServiceServer(grpcServer, server) go func() { - slog.Info("gRPC server listening", "port", grpcPort) + flog.Info("gRPC server listening", "port", grpcPort) if err := grpcServer.Serve(grpcListener); err != nil { - slog.Error("gRPC server failed", "error", err) + flog.Error("gRPC server failed", "error", err) } }() @@ -79,9 +81,9 @@ func main() { httpServer := &http.Server{Addr: httpAddr, Handler: mux} go func() { - slog.Info("HTTP server listening", "addr", httpAddr) + flog.Info("HTTP server listening", "addr", httpAddr) if err := httpServer.ListenAndServe(); err != nil && err != http.ErrServerClosed { - slog.Error("HTTP server failed", "error", err) + flog.Error("HTTP server failed", "error", err) os.Exit(1) } }() @@ -91,10 +93,10 @@ func main() { signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT) <-sigCh - slog.Info("shutting down...") + flog.Info("shutting down...") cancel() grpcServer.GracefulStop() httpServer.Shutdown(context.Background()) - slog.Info("shutdown complete") + flog.Info("shutdown complete") } diff --git a/go/cmd/transcodebackfill/main.go b/go/cmd/transcodebackfill/main.go index df931df..363bb89 100644 --- a/go/cmd/transcodebackfill/main.go +++ b/go/cmd/transcodebackfill/main.go @@ -5,10 +5,11 @@ package main import ( "context" - "log/slog" "os" "time" + "github.com/flowy-live/llink/internal/utils/flog" + "cloud.google.com/go/firestore" "cloud.google.com/go/storage" "google.golang.org/api/iterator" @@ -27,7 +28,7 @@ func main() { storageClient, err := storage.NewClient(ctx) if err != nil { - slog.Error("failed to create GCS client", "error", err) + flog.Error("failed to create GCS client", "error", err) os.Exit(1) } defer storageClient.Close() @@ -40,7 +41,7 @@ func main() { gcpProject := utils.MustGetEnv("GCP_PROJECT") fs, err := firestore.NewClient(ctx, gcpProject) if err != nil { - slog.Error("failed to create Firestore client", "error", err) + flog.Error("failed to create Firestore client", "error", err) os.Exit(1) } defer fs.Close() @@ -49,7 +50,7 @@ func main() { stats, err := run(ctx, fs, depotSvc) elapsed := time.Since(started) - slog.Info("transcode_backfill_summary", + flog.Info("transcode_backfill_summary", "scanned", stats.scanned, "media", stats.media, "transcoded", stats.transcoded, @@ -61,7 +62,7 @@ func main() { ) if err != nil { - slog.Error("transcode backfill aborted", "error", err) + flog.Error("transcode backfill aborted", "error", err) os.Exit(1) } if stats.failures > 0 { @@ -116,7 +117,7 @@ func run(ctx context.Context, fs *firestore.Client, depotSvc depot.Service) (sta switch { case result.Err != nil: s.failures++ - slog.Error("transcode_backfill_failure", + flog.Error("transcode_backfill_failure", "particleID", doc.Ref.ID, "path", doc.Ref.Path, "error", result.Err, @@ -129,7 +130,7 @@ func run(ctx context.Context, fs *firestore.Client, depotSvc depot.Service) (sta s.skippedNonMedia++ default: s.transcoded++ - slog.Info("transcode_backfill_progress", + flog.Info("transcode_backfill_progress", "particleID", doc.Ref.ID, "transcodedObjectID", result.TranscodedObjectID, "outputMime", result.OutputMimeType, diff --git a/go/go.mod b/go/go.mod index acc81d9..2f54f7f 100644 --- a/go/go.mod +++ b/go/go.mod @@ -13,6 +13,7 @@ require ( github.com/livekit/protocol v1.45.1 github.com/livekit/server-sdk-go/v2 v2.16.1 github.com/redis/go-redis/v9 v9.17.2 + github.com/sirupsen/logrus v1.9.3 github.com/stretchr/testify v1.11.1 github.com/stripe/stripe-go/v85 v85.0.1 github.com/testcontainers/testcontainers-go v0.40.0 @@ -147,7 +148,6 @@ require ( github.com/prometheus/procfs v0.19.2 // indirect github.com/puzpuzpuz/xsync/v3 v3.5.1 // indirect github.com/shirou/gopsutil/v4 v4.25.6 // indirect - github.com/sirupsen/logrus v1.9.3 // indirect github.com/spiffe/go-spiffe/v2 v2.6.0 // indirect github.com/tklauser/go-sysconf v0.3.12 // indirect github.com/tklauser/numcpus v0.6.1 // indirect diff --git a/go/internal/auth/service.go b/go/internal/auth/service.go index 244e35c..acc0e39 100644 --- a/go/internal/auth/service.go +++ b/go/internal/auth/service.go @@ -5,15 +5,17 @@ import ( "encoding/json" "errors" "fmt" - "log/slog" "strings" "time" + "github.com/flowy-live/llink/internal/utils/flog" + firebaseauth "firebase.google.com/go/v4/auth" - pbaero "github.com/flowy-live/llink/genproto/aero" - "github.com/flowy-live/llink/internal/utils" "github.com/redis/go-redis/v9" "go.jetify.com/typeid" + + pbaero "github.com/flowy-live/llink/genproto/aero" + "github.com/flowy-live/llink/internal/utils" ) const ( @@ -86,7 +88,7 @@ func (a *authServiceImpl) MintFirebaseCustomToken(ctx context.Context, humanId s func (a *authServiceImpl) IsSystemAdmin(ctx context.Context, email string) bool { formattedEmail, err := utils.NormalizeEmail(email) if err != nil { - slog.Error("problem validating email", "error", err) + flog.Error("problem validating email", "error", err) return false } @@ -110,7 +112,7 @@ func (a *authServiceImpl) RequestSignInCode(ctx context.Context, email string) e err = a.redisClient.Set(ctx, formattedEmail, code, codeExpiry).Err() if err != nil { - slog.Error("error setting code in redis", "error", err) + flog.Error("error setting code in redis", "error", err) return fmt.Errorf("error storing sign-in code: %w", err) } @@ -129,7 +131,7 @@ func (a *authServiceImpl) RequestSignInCode(ctx context.Context, email string) e return fmt.Errorf("an error occurred while sending the email: %w", err) } - slog.Info("sent sign in code", "email", formattedEmail) + flog.Info("sent sign in code", "email", formattedEmail) return nil } @@ -144,7 +146,7 @@ func (a *authServiceImpl) VerifySignInCode(ctx context.Context, email, code, hum if errors.Is(err, redis.Nil) { return "", ErrInvalidCode } - slog.Error("error getting code from redis", "error", err) + flog.Error("error getting code from redis", "error", err) return "", fmt.Errorf("error verifying code: %w", err) } @@ -153,7 +155,7 @@ func (a *authServiceImpl) VerifySignInCode(ctx context.Context, email, code, hum } if err := a.redisClient.Del(ctx, formattedEmail).Err(); err != nil { - slog.Error("error deleting code from redis", "error", err) + flog.Error("error deleting code from redis", "error", err) } token, err := a.createSession(ctx, formattedEmail, humanId) diff --git a/go/internal/billing/webhooks.go b/go/internal/billing/webhooks.go index 9dab258..fd1452f 100644 --- a/go/internal/billing/webhooks.go +++ b/go/internal/billing/webhooks.go @@ -4,9 +4,10 @@ import ( "context" "encoding/json" "fmt" - "log/slog" "time" + "github.com/flowy-live/llink/internal/utils/flog" + "github.com/stripe/stripe-go/v85" "github.com/stripe/stripe-go/v85/webhook" ) @@ -17,7 +18,7 @@ func (s *serviceImpl) HandleWebhook(ctx context.Context, payload []byte, signatu return fmt.Errorf("verify stripe signature: %w", err) } - slog.Info("stripe webhook", "type", event.Type, "id", event.ID) + flog.Info("stripe webhook", "type", event.Type, "id", event.ID) switch event.Type { case "checkout.session.completed": diff --git a/go/internal/db/postgres.go b/go/internal/db/postgres.go index 672371f..52536d1 100644 --- a/go/internal/db/postgres.go +++ b/go/internal/db/postgres.go @@ -2,9 +2,10 @@ package db import ( "context" - "log/slog" "os" + "github.com/flowy-live/llink/internal/utils/flog" + "github.com/jackc/pgx/v5/pgxpool" ) @@ -21,24 +22,24 @@ func Pool() *pgxpool.Pool { func Init() { connString := os.Getenv("LLINK_POSTGRES_CONNECTION_URL") if connString == "" { - slog.Error("must provide LLINK_POSTGRES_CONNECTION_URL in env") + flog.Error("must provide LLINK_POSTGRES_CONNECTION_URL in env") os.Exit(1) } dbpool, err := pgxpool.New(context.Background(), connString) if err != nil { - slog.Error("unable to create connection pool", "error", err) + flog.Error("unable to create connection pool", "error", err) os.Exit(1) } var greeting string err = dbpool.QueryRow(context.Background(), "select 'Hello, world!'").Scan(&greeting) if err != nil { - slog.Error("queryRow failed", "error", err) + flog.Error("queryRow failed", "error", err) os.Exit(1) } - slog.Info("successfully connected to database", "greeting", greeting) + flog.Info("successfully connected to database", "greeting", greeting) db = dbpool } diff --git a/go/internal/depot/service.go b/go/internal/depot/service.go index c87221d..92657fc 100644 --- a/go/internal/depot/service.go +++ b/go/internal/depot/service.go @@ -5,9 +5,10 @@ import ( "errors" "fmt" "io" - "log/slog" "time" + "github.com/flowy-live/llink/internal/utils/flog" + "cloud.google.com/go/storage" "github.com/google/uuid" "github.com/jackc/pgx/v5/pgxpool" @@ -49,7 +50,7 @@ func NewService(pool *pgxpool.Pool, storageClient *storage.Client, config Config } if config.GoogleServiceAccountEmail == "" { - slog.Error("GoogleServiceAccountEmail is not set in config. Signed URLs may not work if the storage client is not properly authenticated with a service account.") + flog.Error("GoogleServiceAccountEmail is not set in config. Signed URLs may not work if the storage client is not properly authenticated with a service account.") panic("GoogleServiceAccountEmail is required for signed URL generation") } @@ -102,10 +103,10 @@ func (s *serviceImpl) PrepareUpload(ctx context.Context, input PrepareUploadInpu Headers: []string{contentLengthHeader}, }) if err != nil { - slog.Error("failed to generate signed upload URL", "error", err, "bucket", s.bucketName, "object_key", objectKey) + flog.Error("failed to generate signed upload URL", "error", err, "bucket", s.bucketName, "object_key", objectKey) // Roll back the placeholder row. if delErr := s.repo.delete(ctx, created.ID); delErr != nil { - slog.Warn("failed to cleanup db record after signed URL failure", "error", delErr, "object_id", created.ID) + flog.Warn("failed to cleanup db record after signed URL failure", "error", delErr, "object_id", created.ID) } return nil, err } @@ -134,7 +135,7 @@ func (s *serviceImpl) ConfirmUpload(ctx context.Context, objectID string) (*Obje if errors.Is(err, storage.ErrObjectNotExist) { return nil, errors.Join(ErrNotFound, errors.New("object not found in storage")) } - slog.Error("failed to get GCS object attrs", "error", err, "bucket", obj.BucketName, "object_key", obj.ObjectKey) + flog.Error("failed to get GCS object attrs", "error", err, "bucket", obj.BucketName, "object_key", obj.ObjectKey) return nil, err } @@ -170,13 +171,13 @@ func (s *serviceImpl) CreateFromReader(ctx context.Context, input CreateFromRead if _, err := io.Copy(w, body); err != nil { // Always release the writer; surface the copy error, not Close's. if cerr := w.Close(); cerr != nil { - slog.Warn("failed to close GCS writer after copy failure", "error", cerr, "object_key", objectKey) + flog.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) + flog.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) + flog.Error("failed to close GCS writer", "error", err, "bucket", s.bucketName, "object_key", objectKey) return nil, err } @@ -193,7 +194,7 @@ func (s *serviceImpl) CreateFromReader(ctx context.Context, input CreateFromRead if err != nil { // Best-effort: drop the now-untracked GCS object. 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) + flog.Warn("failed to clean up GCS object after db create failure", "error", delErr, "object_key", objectKey) } return nil, err } @@ -227,7 +228,7 @@ func (s *serviceImpl) GetDownloadURL(ctx context.Context, objectID string) (stri Expires: time.Now().Add(s.downloadURLExpiry), }) if err != nil { - slog.Error("failed to generate signed download URL", "error", err, "bucket", obj.BucketName, "object_key", obj.ObjectKey) + flog.Error("failed to generate signed download URL", "error", err, "bucket", obj.BucketName, "object_key", obj.ObjectKey) return "", err } @@ -246,7 +247,7 @@ func (s *serviceImpl) Delete(ctx context.Context, objectID string) error { // GCS first so we don't strand an object after the row vanishes; missing object is fine. gcsErr := s.storageClient.Bucket(obj.BucketName).Object(obj.ObjectKey).Delete(ctx) if gcsErr != nil && !errors.Is(gcsErr, storage.ErrObjectNotExist) { - slog.Error("failed to delete object from GCS", "error", gcsErr, "bucket", obj.BucketName, "object_key", obj.ObjectKey) + flog.Error("failed to delete object from GCS", "error", gcsErr, "bucket", obj.BucketName, "object_key", obj.ObjectKey) return gcsErr } diff --git a/go/internal/handler/billing.go b/go/internal/handler/billing.go index 23423ec..501e468 100644 --- a/go/internal/handler/billing.go +++ b/go/internal/handler/billing.go @@ -5,9 +5,10 @@ import ( "errors" "fmt" "io" - "log/slog" "net/http" + "github.com/flowy-live/llink/internal/utils/flog" + "github.com/flowy-live/llink/internal/billing" "github.com/flowy-live/llink/internal/middleware" "github.com/flowy-live/llink/internal/network" @@ -42,7 +43,7 @@ func (h *Handler) GetNetworkUsage(w http.ResponseWriter, r *http.Request) { isMember, err := h.networkSvc.IsMember(r.Context(), networkID, humanId) if err != nil { - slog.Error("failed to check network membership", "error", err, "network_id", networkID) + flog.Error("failed to check network membership", "error", err, "network_id", networkID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -53,7 +54,7 @@ func (h *Handler) GetNetworkUsage(w http.ResponseWriter, r *http.Request) { usage, err := h.billingSvc.GetUsage(r.Context(), networkID) if err != nil { - slog.Error("failed to get network usage", "error", err, "network_id", networkID) + flog.Error("failed to get network usage", "error", err, "network_id", networkID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -69,7 +70,7 @@ func (h *Handler) GetNetworkBilling(w http.ResponseWriter, r *http.Request) { status, err := h.billingSvc.GetStatus(r.Context(), net.ID) if err != nil { - slog.Error("failed to get billing status", "error", err, "network_id", net.ID) + flog.Error("failed to get billing status", "error", err, "network_id", net.ID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -97,14 +98,14 @@ func (h *Handler) CreateCheckoutSession(w http.ResponseWriter, r *http.Request) adminHuman, err := h.humanSvc.GetByID(r.Context(), adminHumanId) if err != nil { - slog.Error("failed to load admin human", "error", err, "human_id", adminHumanId) + flog.Error("failed to load admin human", "error", err, "human_id", adminHumanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } seats, err := h.networkSvc.CountSeats(r.Context(), net.ID) if err != nil { - slog.Error("failed to count seats", "error", err, "network_id", net.ID) + flog.Error("failed to count seats", "error", err, "network_id", net.ID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -117,7 +118,7 @@ func (h *Handler) CreateCheckoutSession(w http.ResponseWriter, r *http.Request) Seats: seats, }) if err != nil { - slog.Error("failed to create checkout session", "error", err, "network_id", net.ID) + flog.Error("failed to create checkout session", "error", err, "network_id", net.ID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -137,7 +138,7 @@ func (h *Handler) CreatePortalSession(w http.ResponseWriter, r *http.Request) { return } if err != nil { - slog.Error("failed to create portal session", "error", err, "network_id", net.ID) + flog.Error("failed to create portal session", "error", err, "network_id", net.ID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -150,14 +151,14 @@ const maxStripeWebhookBytes = 1 << 20 // 1 MiB func (h *Handler) HandleStripeWebhook(w http.ResponseWriter, r *http.Request) { payload, err := io.ReadAll(io.LimitReader(r.Body, maxStripeWebhookBytes)) if err != nil { - slog.Warn("stripe webhook: failed to read body", "error", err) + flog.Warn("stripe webhook: failed to read body", "error", err) http.Error(w, "bad request", http.StatusBadRequest) return } signature := r.Header.Get("Stripe-Signature") if err := h.billingSvc.HandleWebhook(r.Context(), payload, signature); err != nil { - slog.Error("stripe webhook failed", "error", err) + flog.Error("stripe webhook failed", "error", err) formattedErr := fmt.Errorf("webhook processing failed: %w", err) http.Error(w, formattedErr.Error(), http.StatusBadRequest) return @@ -187,7 +188,7 @@ func (h *Handler) loadNetworkForAdmin(w http.ResponseWriter, r *http.Request) (* return nil, "", false } if err != nil { - slog.Error("failed to load network", "error", err, "network_id", networkID) + flog.Error("failed to load network", "error", err, "network_id", networkID) http.Error(w, "internal server error", http.StatusInternalServerError) return nil, "", false } @@ -203,6 +204,6 @@ func (h *Handler) loadNetworkForAdmin(w http.ResponseWriter, r *http.Request) (* func writeJSON(w http.ResponseWriter, body any) { w.Header().Set("Content-Type", "application/json") if err := json.NewEncoder(w).Encode(body); err != nil { - slog.Error("failed to write json", "error", err) + flog.Error("failed to write json", "error", err) } } diff --git a/go/internal/handler/handler.go b/go/internal/handler/handler.go index d87a62c..0ec211a 100644 --- a/go/internal/handler/handler.go +++ b/go/internal/handler/handler.go @@ -5,12 +5,15 @@ import ( "encoding/json" "errors" "fmt" - "log/slog" "net/http" "strings" "time" + "github.com/flowy-live/llink/internal/utils/flog" + "cloud.google.com/go/firestore" + "github.com/livekit/protocol/webhook" + "github.com/flowy-live/llink/internal/auth" "github.com/flowy-live/llink/internal/billing" "github.com/flowy-live/llink/internal/depot" @@ -22,7 +25,6 @@ import ( "github.com/flowy-live/llink/internal/particle" "github.com/flowy-live/llink/internal/utils" "github.com/flowy-live/llink/internal/waitlist" - "github.com/livekit/protocol/webhook" ) type Handler struct { @@ -187,13 +189,13 @@ func (h *Handler) RequestSignInCode(w http.ResponseWriter, r *http.Request) { _, err := h.humanSvc.GetOrCreateByEmail(r.Context(), req.Email) if err != nil { - slog.Error("failed to get or create human", "error", err, "email", req.Email) + flog.Error("failed to get or create human", "error", err, "email", req.Email) http.Error(w, "internal server error", http.StatusInternalServerError) return } if err := h.authSvc.RequestSignInCode(r.Context(), req.Email); err != nil { - slog.Error("failed to request sign-in code", "error", err, "email", req.Email) + flog.Error("failed to request sign-in code", "error", err, "email", req.Email) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -220,7 +222,7 @@ func (h *Handler) SignIn(w http.ResponseWriter, r *http.Request) { http.Error(w, "human not found", http.StatusNotFound) return } - slog.Error("failed to get human for sign-in", "error", err, "email", req.Email) + flog.Error("failed to get human for sign-in", "error", err, "email", req.Email) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -231,7 +233,7 @@ func (h *Handler) SignIn(w http.ResponseWriter, r *http.Request) { http.Error(w, "invalid code", http.StatusUnauthorized) return } - slog.Error("failed to verify sign-in code", "error", err, "email", req.Email) + flog.Error("failed to verify sign-in code", "error", err, "email", req.Email) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -256,7 +258,7 @@ func (h *Handler) FirebaseToken(w http.ResponseWriter, r *http.Request) { token, err := h.authSvc.MintFirebaseCustomToken(r.Context(), humanId) if err != nil { - slog.Error("failed to mint Firebase custom token", "error", err, "humanId", humanId) + flog.Error("failed to mint Firebase custom token", "error", err, "humanId", humanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -273,7 +275,7 @@ func (h *Handler) SignOut(w http.ResponseWriter, r *http.Request) { } if err := h.authSvc.SignOut(r.Context(), token); err != nil { - slog.Error("failed to sign out", "error", err) + flog.Error("failed to sign out", "error", err) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -294,7 +296,7 @@ func (h *Handler) GetCurrentHuman(w http.ResponseWriter, r *http.Request) { http.Error(w, "human not found", http.StatusNotFound) return } - slog.Error("failed to get current human", "error", err, "email", email) + flog.Error("failed to get current human", "error", err, "email", email) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -323,7 +325,7 @@ func (h *Handler) UpdateSettings(w http.ResponseWriter, r *http.Request) { if req.EmailNotificationsEnabled != nil { if err := h.humanSvc.UpdateEmailNotificationsEnabled(r.Context(), humanId, *req.EmailNotificationsEnabled); err != nil { - slog.Error("failed to update email notifications setting", "error", err, "humanId", humanId) + flog.Error("failed to update email notifications setting", "error", err, "humanId", humanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -355,14 +357,14 @@ func (h *Handler) CreateNetwork(w http.ResponseWriter, r *http.Request) { http.Error(w, "name cannot be empty", http.StatusBadRequest) return } - slog.Error("failed to create network", "error", err, "humanId", humanId, "name", req.Name) + flog.Error("failed to create network", "error", err, "humanId", humanId, "name", req.Name) http.Error(w, "internal server error", http.StatusInternalServerError) return } resp, err := h.networkToDTO(r.Context(), net) if err != nil { - slog.Error("failed to convert network to DTO", "error", err, "network_id", net.ID) + flog.Error("failed to convert network to DTO", "error", err, "network_id", net.ID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -381,7 +383,7 @@ func (h *Handler) ListNetworks(w http.ResponseWriter, r *http.Request) { networks, err := h.networkSvc.ListForHuman(r.Context(), humanId) if err != nil { - slog.Error("failed to list networks", "error", err, "humanId", humanId) + flog.Error("failed to list networks", "error", err, "humanId", humanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -390,7 +392,7 @@ func (h *Handler) ListNetworks(w http.ResponseWriter, r *http.Request) { for _, net := range networks { dto, err := h.networkToDTO(r.Context(), net) if err != nil { - slog.Warn("failed to convert network to DTO in list", "error", err, "network_id", net.ID) + flog.Warn("failed to convert network to DTO in list", "error", err, "network_id", net.ID) continue } resp = append(resp, dto) @@ -415,7 +417,7 @@ func (h *Handler) GetNetwork(w http.ResponseWriter, r *http.Request) { isMember, err := h.networkSvc.IsMember(r.Context(), networkID, humanId) if err != nil { - slog.Error("failed to check network membership", "error", err, "network_id", networkID, "humanId", humanId) + flog.Error("failed to check network membership", "error", err, "network_id", networkID, "humanId", humanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -430,14 +432,14 @@ func (h *Handler) GetNetwork(w http.ResponseWriter, r *http.Request) { http.Error(w, "network not found", http.StatusNotFound) return } - slog.Error("failed to get network", "error", err, "network_id", networkID) + flog.Error("failed to get network", "error", err, "network_id", networkID) http.Error(w, "internal server error", http.StatusInternalServerError) return } resp, err := h.networkToDTO(r.Context(), net) if err != nil { - slog.Error("failed to convert network to DTO", "error", err, "network_id", networkID) + flog.Error("failed to convert network to DTO", "error", err, "network_id", networkID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -463,7 +465,7 @@ func (h *Handler) AddMembersToNetwork(w http.ResponseWriter, r *http.Request) { isMember, err := h.networkSvc.IsMember(r.Context(), networkID, humanId) if err != nil { - slog.Error("failed to check network membership", "error", err, "network_id", networkID, "humanId", humanId) + flog.Error("failed to check network membership", "error", err, "network_id", networkID, "humanId", humanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -497,7 +499,7 @@ func (h *Handler) AddMembersToNetwork(w http.ResponseWriter, r *http.Request) { inviteEmails = append(inviteEmails, normalized) continue } - slog.Error("failed to look up human by email", "error", err, "email", normalized) + flog.Error("failed to look up human by email", "error", err, "email", normalized) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -506,7 +508,7 @@ func (h *Handler) AddMembersToNetwork(w http.ResponseWriter, r *http.Request) { if len(memberHumanIds) > 0 { if err := h.networkSvc.AddMembers(r.Context(), networkID, memberHumanIds); err != nil { - slog.Error("failed to add members to network", "error", err, "network_id", networkID) + flog.Error("failed to add members to network", "error", err, "network_id", networkID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -514,7 +516,7 @@ func (h *Handler) AddMembersToNetwork(w http.ResponseWriter, r *http.Request) { if len(inviteEmails) > 0 { if err := h.networkSvc.InviteByEmail(r.Context(), networkID, inviteEmails); err != nil { - slog.Error("failed to invite members to network", "error", err, "network_id", networkID) + flog.Error("failed to invite members to network", "error", err, "network_id", networkID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -522,14 +524,14 @@ func (h *Handler) AddMembersToNetwork(w http.ResponseWriter, r *http.Request) { net, err := h.networkSvc.GetByID(r.Context(), networkID) if err != nil { - slog.Error("failed to get network after adding members", "error", err, "network_id", networkID) + flog.Error("failed to get network after adding members", "error", err, "network_id", networkID) http.Error(w, "internal server error", http.StatusInternalServerError) return } resp, err := h.networkToDTO(r.Context(), net) if err != nil { - slog.Error("failed to convert network to DTO", "error", err, "network_id", networkID) + flog.Error("failed to convert network to DTO", "error", err, "network_id", networkID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -558,7 +560,7 @@ func (h *Handler) RemoveMemberFromNetwork(w http.ResponseWriter, r *http.Request } if err := h.networkSvc.RemoveMember(r.Context(), net.ID, memberHumanId); err != nil { - slog.Error("failed to remove member from network", "error", err, "network_id", net.ID, "memberHumanId", memberHumanId) + flog.Error("failed to remove member from network", "error", err, "network_id", net.ID, "memberHumanId", memberHumanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -581,7 +583,7 @@ func (h *Handler) ListInvitationsForNetwork(w http.ResponseWriter, r *http.Reque isMember, err := h.networkSvc.IsMember(r.Context(), networkID, humanId) if err != nil { - slog.Error("failed to check network membership", "error", err, "network_id", networkID, "humanId", humanId) + flog.Error("failed to check network membership", "error", err, "network_id", networkID, "humanId", humanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -592,7 +594,7 @@ func (h *Handler) ListInvitationsForNetwork(w http.ResponseWriter, r *http.Reque invitations, err := h.networkSvc.ListInvitationsForNetwork(r.Context(), networkID) if err != nil { - slog.Error("failed to list invitations", "error", err, "network_id", networkID) + flog.Error("failed to list invitations", "error", err, "network_id", networkID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -620,7 +622,7 @@ func (h *Handler) ListMyInvitations(w http.ResponseWriter, r *http.Request) { invitations, err := h.networkSvc.ListInvitationsForEmail(r.Context(), email) if err != nil { - slog.Error("failed to list invitations for email", "error", err, "email", email) + flog.Error("failed to list invitations for email", "error", err, "email", email) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -663,7 +665,7 @@ func (h *Handler) AcceptInvitation(w http.ResponseWriter, r *http.Request) { } if err := h.networkSvc.AcceptInvitation(r.Context(), req.NetworkId, email, humanId); err != nil { - slog.Error("failed to accept invitation", "error", err, "network_id", req.NetworkId, "email", email) + flog.Error("failed to accept invitation", "error", err, "network_id", req.NetworkId, "email", email) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -686,7 +688,7 @@ func (h *Handler) RevokeInvitation(w http.ResponseWriter, r *http.Request) { isMember, err := h.networkSvc.IsMember(r.Context(), networkID, humanId) if err != nil { - slog.Error("failed to check network membership", "error", err, "network_id", networkID, "humanId", humanId) + flog.Error("failed to check network membership", "error", err, "network_id", networkID, "humanId", humanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -707,7 +709,7 @@ func (h *Handler) RevokeInvitation(w http.ResponseWriter, r *http.Request) { } if err := h.networkSvc.RevokeInvitation(r.Context(), networkID, req.Email); err != nil { - slog.Error("failed to revoke invitation", "error", err, "network_id", networkID, "email", req.Email) + flog.Error("failed to revoke invitation", "error", err, "network_id", networkID, "email", req.Email) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -729,7 +731,7 @@ func (h *Handler) DownloadParticleMedia(w http.ResponseWriter, r *http.Request) downloadURL, err := h.depotSvc.GetDownloadURL(r.Context(), objectID) if err != nil { - slog.Error("failed to get download URL", "error", err, "object_id", objectID) + flog.Error("failed to get download URL", "error", err, "object_id", objectID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -763,7 +765,7 @@ func (h *Handler) PrepareUpload(w http.ResponseWriter, r *http.Request) { isMember, err := h.networkSvc.IsMember(r.Context(), req.NetworkId, humanId) if err != nil { - slog.Error("failed to check network membership", "error", err, "network_id", req.NetworkId, "humanId", humanId) + flog.Error("failed to check network membership", "error", err, "network_id", req.NetworkId, "humanId", humanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -785,7 +787,7 @@ func (h *Handler) PrepareUpload(w http.ResponseWriter, r *http.Request) { http.Error(w, err.Error(), http.StatusBadRequest) return } - slog.Error("failed to prepare upload", "error", err, "network_id", req.NetworkId, "name", req.Name) + flog.Error("failed to prepare upload", "error", err, "network_id", req.NetworkId, "name", req.Name) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -823,7 +825,7 @@ func (h *Handler) ConfirmUpload(w http.ResponseWriter, r *http.Request) { http.Error(w, err.Error(), http.StatusBadRequest) return } - slog.Error("failed to confirm upload", "error", err, "object_id", objectID) + flog.Error("failed to confirm upload", "error", err, "object_id", objectID) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -886,7 +888,7 @@ func (h *Handler) AddToWaitlist(w http.ResponseWriter, r *http.Request) { http.Error(w, "already in the waitlist", http.StatusConflict) return } - slog.Error("failed to add to waitlist", "error", err, "email", req.Email) + flog.Error("failed to add to waitlist", "error", err, "email", req.Email) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -912,7 +914,7 @@ func (h *Handler) GetWaitlist(w http.ResponseWriter, r *http.Request) { entries, err := h.waitlistSvc.GetWaitlist(r.Context(), filter) if err != nil { - slog.Error("failed to get waitlist", "error", err) + flog.Error("failed to get waitlist", "error", err) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -945,7 +947,7 @@ func (h *Handler) GetWaitlistEntry(w http.ResponseWriter, r *http.Request) { http.Error(w, "entry not found", http.StatusNotFound) return } - slog.Error("failed to get waitlist entry", "error", err, "email", email) + flog.Error("failed to get waitlist entry", "error", err, "email", email) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -977,7 +979,7 @@ func (h *Handler) InviteWaitlistEntrant(w http.ResponseWriter, r *http.Request) http.Error(w, "entry not found", http.StatusNotFound) return } - slog.Error("failed to invite waitlist entrant", "error", err, "email", req.Email) + flog.Error("failed to invite waitlist entrant", "error", err, "email", req.Email) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -1019,7 +1021,7 @@ func (h *Handler) networkToDTO(ctx context.Context, n *network.Network) (Network for _, memberHumanId := range n.MemberHumanIds { hum, err := h.humanSvc.GetByID(ctx, memberHumanId) if err != nil { - slog.Warn("failed to look up network member", "humanId", memberHumanId, "error", err) + flog.Warn("failed to look up network member", "humanId", memberHumanId, "error", err) continue } humans = append(humans, humanToDTO(hum)) @@ -1070,7 +1072,7 @@ func (h *Handler) GetLivekitToken(w http.ResponseWriter, r *http.Request) { token, err := h.livekitClient.GetJoinToken(roomName, humanId, humanEmail) if err != nil { - slog.Error("failed to generate livekit token", "error", err, "humanId", humanId, "roomName", roomName) + flog.Error("failed to generate livekit token", "error", err, "humanId", humanId, "roomName", roomName) http.Error(w, "failed to generate token", http.StatusInternalServerError) return } @@ -1084,13 +1086,13 @@ func (h *Handler) GetLivekitToken(w http.ResponseWriter, r *http.Request) { func (h *Handler) HandleLivekitWebhook(w http.ResponseWriter, r *http.Request) { event, err := webhook.ReceiveWebhookEvent(r, h.livekitClient.KeyProvider()) if err != nil { - slog.Error("failed to verify livekit webhook", "error", err) + flog.Error("failed to verify livekit webhook", "error", err) http.Error(w, "unauthorized", http.StatusUnauthorized) return } eventType := event.GetEvent() - slog.Info("received livekit webhook", "event", eventType, "room", event.GetRoom().GetName()) + flog.Info("received livekit webhook", "event", eventType, "room", event.GetRoom().GetName()) switch eventType { case "participant_joined", "participant_left", "room_finished": @@ -1103,7 +1105,7 @@ func (h *Handler) HandleLivekitWebhook(w http.ResponseWriter, r *http.Request) { roomName := event.GetRoom().GetName() parts := strings.SplitN(roomName, "/", 2) if len(parts) != 2 { - slog.Error("invalid room name format", "room", roomName) + flog.Error("invalid room name format", "room", roomName) http.Error(w, "invalid room name", http.StatusBadRequest) return } @@ -1120,7 +1122,7 @@ func (h *Handler) HandleLivekitWebhook(w http.ResponseWriter, r *http.Request) { // Authoritative list avoids drift from missed/out-of-order webhooks. participants, err := h.livekitClient.ListParticipants(ctx, roomName) if err != nil { - slog.Error("failed to list participants", "error", err, "room", roomName) + flog.Error("failed to list participants", "error", err, "room", roomName) // 200 to suppress LiveKit retries. w.WriteHeader(http.StatusOK) return @@ -1136,7 +1138,7 @@ func (h *Handler) HandleLivekitWebhook(w http.ResponseWriter, r *http.Request) { {Path: "huddle_active_participants", Value: participantIds}, }) if err != nil { - slog.Error("failed to update huddle participants in firestore", "error", err, "path", docPath) + flog.Error("failed to update huddle participants in firestore", "error", err, "path", docPath) } w.WriteHeader(http.StatusOK) diff --git a/go/internal/handler/metadata.go b/go/internal/handler/metadata.go index cccae0f..61c83f6 100644 --- a/go/internal/handler/metadata.go +++ b/go/internal/handler/metadata.go @@ -4,7 +4,6 @@ import ( "context" "errors" "io" - "log/slog" "net" "net/http" "net/url" @@ -12,6 +11,8 @@ import ( "strings" "sync" "time" + + "github.com/flowy-live/llink/internal/utils/flog" ) // LinkMetadata mirrors the TS shape in `js/desktop/src/lib/link-metadata.ts`. @@ -63,7 +64,7 @@ func (h *Handler) GetLinkMetadata(w http.ResponseWriter, r *http.Request) { meta, err := fetchLinkMetadata(r.Context(), parsed) if err != nil { - slog.Warn("link metadata fetch failed", "url", raw, "error", err) + flog.Warn("link metadata fetch failed", "url", raw, "error", err) writeJSON(w, nil) return } diff --git a/go/internal/handler/push_token.go b/go/internal/handler/push_token.go index edff810..4b9fc4e 100644 --- a/go/internal/handler/push_token.go +++ b/go/internal/handler/push_token.go @@ -3,9 +3,10 @@ package handler import ( "encoding/json" "errors" - "log/slog" "net/http" + "github.com/flowy-live/llink/internal/utils/flog" + "github.com/flowy-live/llink/internal/human/pushnotify" "github.com/flowy-live/llink/internal/middleware" ) @@ -45,7 +46,7 @@ func (h *Handler) RegisterPushToken(w http.ResponseWriter, r *http.Request) { http.Error(w, err.Error(), http.StatusBadRequest) return } - slog.Error("failed to register push token", "error", err, "humanId", humanId) + flog.Error("failed to register push token", "error", err, "humanId", humanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } @@ -74,7 +75,7 @@ func (h *Handler) UnregisterPushToken(w http.ResponseWriter, r *http.Request) { err := h.pushTokenSvc.Unregister(r.Context(), humanId, req.Token) if err != nil && !errors.Is(err, pushnotify.ErrNotFound) { - slog.Error("failed to unregister push token", "error", err, "humanId", humanId) + flog.Error("failed to unregister push token", "error", err, "humanId", humanId) http.Error(w, "internal server error", http.StatusInternalServerError) return } diff --git a/go/internal/human/pushnotify/notifier.go b/go/internal/human/pushnotify/notifier.go index ee5be37..2f9676a 100644 --- a/go/internal/human/pushnotify/notifier.go +++ b/go/internal/human/pushnotify/notifier.go @@ -4,7 +4,8 @@ import ( "context" "errors" "fmt" - "log/slog" + + "github.com/flowy-live/llink/internal/utils/flog" "github.com/flowy-live/llink/internal/network" ) @@ -55,7 +56,7 @@ func NewNotifier(networkR network.Reader, tokens Service, expo *ExpoClient) *Not func (n *Notifier) NotifyParticleCreated(ctx context.Context, in NotifyInput) error { if in.NetworkID == "" || in.ParticleID == "" { - slog.Info("pushnotify: skip — missing ids", + flog.Info("pushnotify: skip — missing ids", "networkID", in.NetworkID, "particleID", in.ParticleID, ) @@ -71,7 +72,7 @@ func (n *Notifier) NotifyParticleCreated(ctx context.Context, in NotifyInput) er recipientsBeforeSenderFilter := len(recipients) recipients = filterOut(recipients, in.SenderHumanID) if len(recipients) == 0 { - slog.Info("pushnotify: skip — no recipients", + flog.Info("pushnotify: skip — no recipients", "networkID", in.NetworkID, "particleID", in.ParticleID, "senderHumanID", in.SenderHumanID, @@ -87,7 +88,7 @@ func (n *Notifier) NotifyParticleCreated(ctx context.Context, in NotifyInput) er return fmt.Errorf("token lookup: %w", err) } if len(tokens) == 0 { - slog.Info("pushnotify: skip — no tokens for recipients", + flog.Info("pushnotify: skip — no tokens for recipients", "networkID", in.NetworkID, "particleID", in.ParticleID, "recipients", len(recipients), @@ -98,7 +99,7 @@ func (n *Notifier) NotifyParticleCreated(ctx context.Context, in NotifyInput) er msgs := buildMessages(tokens, in) tickets, sendErr := n.expo.Send(ctx, msgs) - slog.Info("pushnotify: dispatch", + flog.Info("pushnotify: dispatch", "networkID", in.NetworkID, "particleID", in.ParticleID, "recipients", len(recipients), @@ -127,14 +128,14 @@ func (n *Notifier) cleanupDeadTokens(ctx context.Context, msgs []Message, ticket code, _ := t.Details["error"].(string) if code != ExpoErrorDeviceNotRegistered { if t.Status == "error" { - slog.Warn("pushnotify: ticket error", "code", code, "message", t.Message, "to", msgs[i].To) + flog.Warn("pushnotify: ticket error", "code", code, "message", t.Message, "to", msgs[i].To) } continue } if err := n.tokens.DeleteByToken(ctx, msgs[i].To); err != nil && !errors.Is(err, ErrNotFound) { - slog.Error("pushnotify: failed to delete dead token", "error", err, "token", msgs[i].To) + flog.Error("pushnotify: failed to delete dead token", "error", err, "token", msgs[i].To) } else { - slog.Info("pushnotify: removed unregistered token", "token", msgs[i].To) + flog.Info("pushnotify: removed unregistered token", "token", msgs[i].To) } } } diff --git a/go/internal/media/helpers.go b/go/internal/media/helpers.go index d75c33c..46476a1 100644 --- a/go/internal/media/helpers.go +++ b/go/internal/media/helpers.go @@ -4,10 +4,11 @@ import ( "bytes" "context" "errors" - "log/slog" "os" "os/exec" "strings" + + "github.com/flowy-live/llink/internal/utils/flog" ) func IsIOSPlayableMime(mime string) bool { @@ -56,7 +57,7 @@ func TranscodeToMp4(ctx context.Context, input TranscodeInput) (*TranscodeOutput tmp, err := os.CreateTemp("", "transcode-*"+outputExt) if err != nil { - slog.Error("transcode: failed to create temp file", "error", err) + flog.Error("transcode: failed to create temp file", "error", err) return nil, err } tmpPath := tmp.Name() diff --git a/go/internal/middleware/auth.go b/go/internal/middleware/auth.go index 7dedfec..9217e57 100644 --- a/go/internal/middleware/auth.go +++ b/go/internal/middleware/auth.go @@ -2,10 +2,11 @@ package middleware import ( "context" - "log/slog" "net/http" "strings" + "github.com/flowy-live/llink/internal/utils/flog" + "github.com/flowy-live/llink/internal/auth" ) @@ -62,7 +63,7 @@ func Auth(authSvc auth.AuthService) func(http.Handler) http.Handler { } if err := authSvc.ExtendSession(r.Context(), token); err != nil { - slog.Warn("failed to extend session", "error", err) + flog.Warn("failed to extend session", "error", err) } ctx := WithEmail(r.Context(), session.Email) diff --git a/go/internal/network/service.go b/go/internal/network/service.go index 734a7bd..dd96f95 100644 --- a/go/internal/network/service.go +++ b/go/internal/network/service.go @@ -5,17 +5,19 @@ import ( "errors" "fmt" "html" - "log/slog" "slices" "strings" + "github.com/flowy-live/llink/internal/utils/flog" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + pbaero "github.com/flowy-live/llink/genproto/aero" "github.com/flowy-live/llink/internal/billing" "github.com/flowy-live/llink/internal/constants" "github.com/flowy-live/llink/internal/livestore" "github.com/flowy-live/llink/internal/utils" - "github.com/jackc/pgx/v5" - "github.com/jackc/pgx/v5/pgxpool" ) var ErrNotFound = errors.New("network not found") @@ -91,7 +93,7 @@ func (s *serviceImpl) Create(ctx context.Context, name, adminHumanId string) (*N }, }) if err != nil { - slog.Warn("unable to send admin update email", "error", err) + flog.Warn("unable to send admin update email", "error", err) } return network, nil @@ -149,13 +151,13 @@ func (s *serviceImpl) RemoveMember(ctx context.Context, networkID, humanId strin // Postgres is the source of truth: failures are logged and the reconciler heals drift. func (s *serviceImpl) mirrorAddMembership(ctx context.Context, humanId, networkID string) { if err := s.pub.Add(ctx, humanId, networkID); err != nil { - slog.Error("membership publish add failed", "error", err, "humanId", humanId, "networkID", networkID) + flog.Error("membership publish add failed", "error", err, "humanId", humanId, "networkID", networkID) } } func (s *serviceImpl) mirrorRemoveMembership(ctx context.Context, humanId, networkID string) { if err := s.pub.Remove(ctx, humanId, networkID); err != nil { - slog.Error("membership publish remove failed", "error", err, "humanId", humanId, "networkID", networkID) + flog.Error("membership publish remove failed", "error", err, "humanId", humanId, "networkID", networkID) } } @@ -214,7 +216,7 @@ func (s *serviceImpl) InviteByEmail(ctx context.Context, networkID string, email }, }) if err != nil { - slog.Warn("unable to send email notification", "email", email, "network", network.Name) + flog.Warn("unable to send email notification", "email", email, "network", network.Name) } } return nil diff --git a/go/internal/particle/transcode.go b/go/internal/particle/transcode.go index fd75385..afdba52 100644 --- a/go/internal/particle/transcode.go +++ b/go/internal/particle/transcode.go @@ -3,12 +3,14 @@ package particle import ( "context" "fmt" - "log/slog" "os" "strings" "time" + "github.com/flowy-live/llink/internal/utils/flog" + "cloud.google.com/go/firestore" + "github.com/flowy-live/llink/internal/depot" "github.com/flowy-live/llink/internal/media" ) @@ -44,17 +46,17 @@ type TranscodeResult struct { func Transcode(ctx context.Context, depotSvc depot.Service, doc *firestore.DocumentSnapshot) TranscodeResult { var mp FirestoreMediaParticle if err := doc.DataTo(&mp); err != nil { - slog.Error("transcode: unable to marshal particle data", "error", err) + flog.Error("transcode: unable to marshal particle data", "error", err) return TranscodeResult{Err: fmt.Errorf("unmarshal particle: %w", err)} } particleType, err := ParseParticleType(mp.Type) if err != nil { - slog.Error("transcode: invalid particle type", "error", err) + flog.Error("transcode: invalid particle type", "error", err) return TranscodeResult{Err: fmt.Errorf("parse type: %w", err)} } if particleType != TypeMedia { - slog.Info("transcode: particle is not of type media") + flog.Info("transcode: particle is not of type media") return TranscodeResult{Skipped: true, SkipReason: SkipReasonNotMedia} } @@ -63,19 +65,19 @@ func Transcode(ctx context.Context, depotSvc depot.Service, doc *firestore.Docum } if media.IsIOSPlayableMime(mp.Properties.MimeType) { - slog.Info("transcode: skipping because already playable on ios") + flog.Info("transcode: skipping because already playable on ios") return TranscodeResult{Skipped: true, SkipReason: SkipReasonIOSPlayable} } sourceURL, err := depotSvc.GetDownloadURL(ctx, mp.Properties.ObjectId) if err != nil { - slog.Error("transcode: failed to get download URL", "error", err, "object_id", mp.Properties.ObjectId) + flog.Error("transcode: failed to get download URL", "error", err, "object_id", mp.Properties.ObjectId) return TranscodeResult{Err: fmt.Errorf("download URL: %w", err)} } networkID, err := NetworkIDFromParticlePath(doc.Ref.Path) if err != nil { - slog.Error("transcode: failed to derive network id", "error", err, "path", doc.Ref.Path) + flog.Error("transcode: failed to derive network id", "error", err, "path", doc.Ref.Path) return TranscodeResult{Err: fmt.Errorf("network id: %w", err)} } @@ -86,14 +88,14 @@ func Transcode(ctx context.Context, depotSvc depot.Service, doc *firestore.Docum MimeType: mp.Properties.MimeType, }) if err != nil { - slog.Error("transcode: ffmpeg failed", "error", err) + flog.Error("transcode: ffmpeg failed", "error", err) return TranscodeResult{Err: fmt.Errorf("ffmpeg: %w", err)} } defer os.Remove(transcodeOutput.TempLocalFilePath) f, err := os.Open(transcodeOutput.TempLocalFilePath) if err != nil { - slog.Error("transcode: failed to open transcoded file", "error", err, "path", transcodeOutput.TempLocalFilePath) + flog.Error("transcode: failed to open transcoded file", "error", err, "path", transcodeOutput.TempLocalFilePath) return TranscodeResult{Err: fmt.Errorf("open temp file: %w", err)} } defer f.Close() @@ -104,7 +106,7 @@ func Transcode(ctx context.Context, depotSvc depot.Service, doc *firestore.Docum ContentType: transcodeOutput.OutputMimeType, }, f) if err != nil { - slog.Error("transcode: failed to upload transcoded object", "error", err, "particleID", doc.Ref.ID) + flog.Error("transcode: failed to upload transcoded object", "error", err, "particleID", doc.Ref.ID) return TranscodeResult{Err: fmt.Errorf("upload: %w", err)} } @@ -114,11 +116,11 @@ func Transcode(ctx context.Context, depotSvc depot.Service, doc *firestore.Docum "transcoded_mime_type": transcodeOutput.OutputMimeType, }, }, firestore.MergeAll); err != nil { - slog.Error("transcode: failed to update particle in firestore", "error", err, "particleID", doc.Ref.ID) + flog.Error("transcode: failed to update particle in firestore", "error", err, "particleID", doc.Ref.ID) return TranscodeResult{Err: fmt.Errorf("firestore update: %w", err)} } - slog.Info("transcoded media particle", "particleID", doc.Ref.ID, "transcoded_object_id", newObj.ID, "mime", transcodeOutput.OutputMimeType) + flog.Info("transcoded media particle", "particleID", doc.Ref.ID, "transcoded_object_id", newObj.ID, "mime", transcodeOutput.OutputMimeType) return TranscodeResult{ TranscodedObjectID: newObj.ID, OutputMimeType: transcodeOutput.OutputMimeType, diff --git a/go/internal/pusher/conn.go b/go/internal/pusher/conn.go index 6a73912..cb00c01 100644 --- a/go/internal/pusher/conn.go +++ b/go/internal/pusher/conn.go @@ -3,9 +3,10 @@ package pusher import ( "context" "encoding/json" - "log/slog" "sync" + "github.com/flowy-live/llink/internal/utils/flog" + "nhooyr.io/websocket" ) @@ -37,16 +38,16 @@ func (c *Conn) ReadPump(ctx context.Context, hub *Hub) { _, data, err := c.ws.Read(ctx) if err != nil { if ctx.Err() != nil { - slog.Info("websocket context cancelled", "connId", c.id, "humanId", c.humanID, "error", ctx.Err()) + flog.Info("websocket context cancelled", "connId", c.id, "humanId", c.humanID, "error", ctx.Err()) return } - slog.Warn("websocket read error", "connId", c.id, "humanId", c.humanID, "error", err) + flog.Warn("websocket read error", "connId", c.id, "humanId", c.humanID, "error", err) return } if string(data) == "ping" { if err := c.ws.Write(ctx, websocket.MessageText, []byte("pong")); err != nil { - slog.Warn("websocket pong write error", "connId", c.id, "error", err) + flog.Warn("websocket pong write error", "connId", c.id, "error", err) } continue } @@ -92,7 +93,7 @@ func (c *Conn) WritePump(ctx context.Context) { return } if err := c.ws.Write(ctx, websocket.MessageText, data); err != nil { - slog.Debug("websocket write error", "connId", c.id, "error", err) + flog.Debug("websocket write error", "connId", c.id, "error", err) return } } @@ -103,14 +104,14 @@ func (c *Conn) WritePump(ctx context.Context) { func (c *Conn) Send(msg ServerMessage) { data, err := json.Marshal(msg) if err != nil { - slog.Error("failed to marshal server message", "error", err) + flog.Error("failed to marshal server message", "error", err) return } select { case c.send <- data: default: - slog.Warn("slow client, closing connection", "connId", c.id, "humanId", c.humanID) + flog.Warn("slow client, closing connection", "connId", c.id, "humanId", c.humanID) c.Close() } } diff --git a/go/internal/pusher/hub.go b/go/internal/pusher/hub.go index 3ecc608..cb0e5ec 100644 --- a/go/internal/pusher/hub.go +++ b/go/internal/pusher/hub.go @@ -3,7 +3,8 @@ package pusher import ( "context" "encoding/json" - "log/slog" + + "github.com/flowy-live/llink/internal/utils/flog" ) type subscribeRequest struct { @@ -110,7 +111,7 @@ func (h *Hub) handleSubscribe(ctx context.Context, req *subscribeRequest) { presence, err := h.bridge.Subscribe(ctx, req.channelID, req.conn.id, req.conn.humanID) if err != nil { - slog.Error("redis subscribe failed", "channelId", req.channelID, "error", err) + flog.Error("redis subscribe failed", "channelId", req.channelID, "error", err) // Fall back to local-only presence. presence = ch.localHumanIDs() } @@ -144,7 +145,7 @@ func (h *Hub) handleUnsubscribe(ctx context.Context, req *unsubscribeRequest) { } if err := h.bridge.Unsubscribe(ctx, req.channelID, req.conn.id, req.conn.humanID); err != nil { - slog.Error("redis unsubscribe failed", "channelId", req.channelID, "error", err) + flog.Error("redis unsubscribe failed", "channelId", req.channelID, "error", err) } // Only emit leave once the humanID has no remaining tabs on this pod. @@ -198,7 +199,7 @@ func (h *Hub) handleDisconnect(ctx context.Context, conn *Conn) { ch.removeMember(conn) if err := h.bridge.Unsubscribe(ctx, channelID, conn.id, conn.humanID); err != nil { - slog.Error("redis unsubscribe on disconnect failed", "channelId", channelID, "error", err) + flog.Error("redis unsubscribe on disconnect failed", "channelId", channelID, "error", err) } if !ch.hasHumanID(conn.humanID) { diff --git a/go/internal/pusher/redis_bridge.go b/go/internal/pusher/redis_bridge.go index 403b15f..03e4c66 100644 --- a/go/internal/pusher/redis_bridge.go +++ b/go/internal/pusher/redis_bridge.go @@ -4,10 +4,11 @@ import ( "context" "encoding/json" "fmt" - "log/slog" "strings" "time" + "github.com/flowy-live/llink/internal/utils/flog" + "github.com/redis/go-redis/v9" ) @@ -200,7 +201,7 @@ func (rb *RedisBridge) handlePubSubMessage(msg *redis.Message) { var event redisEvent if err := json.Unmarshal([]byte(msg.Payload), &event); err != nil { - slog.Error("failed to parse pub/sub event", "error", err) + flog.Error("failed to parse pub/sub event", "error", err) return } @@ -257,7 +258,7 @@ func (rb *RedisBridge) cleanupStalePods(ctx context.Context) { for { keys, nextCursor, err := rb.client.Scan(ctx, cursor, channelConnsPrefix+"*"+channelConnsSuffix, 100).Result() if err != nil { - slog.Error("failed to scan channel keys", "error", err) + flog.Error("failed to scan channel keys", "error", err) return } @@ -292,7 +293,7 @@ func (rb *RedisBridge) cleanupStalePods(ctx context.Context) { for podID := range knownPods { if !alivePods[podID] { - slog.Info("cleaning up stale pod", "podId", podID) + flog.Info("cleaning up stale pod", "podId", podID) rb.cleanupPod(ctx, podID) } } @@ -344,11 +345,11 @@ func (rb *RedisBridge) connField(connID string) string { func (rb *RedisBridge) publishEvent(ctx context.Context, channelID string, event redisEvent) { data, err := json.Marshal(event) if err != nil { - slog.Error("failed to marshal event", "error", err) + flog.Error("failed to marshal event", "error", err) return } if err := rb.client.Publish(ctx, pubsubPrefix+channelID, data).Err(); err != nil { - slog.Error("failed to publish event", "channelId", channelID, "error", err) + flog.Error("failed to publish event", "channelId", channelID, "error", err) } } diff --git a/go/internal/pusher/server.go b/go/internal/pusher/server.go index cc32816..6a7c468 100644 --- a/go/internal/pusher/server.go +++ b/go/internal/pusher/server.go @@ -2,13 +2,14 @@ package pusher import ( "context" - "log/slog" "net/http" - "github.com/flowy-live/llink/genproto/llink/pusher" - "github.com/flowy-live/llink/internal/auth" "github.com/google/uuid" "nhooyr.io/websocket" + + pbpusher "github.com/flowy-live/llink/genproto/llink/pusher" + "github.com/flowy-live/llink/internal/auth" + "github.com/flowy-live/llink/internal/utils/flog" ) // Server handles WebSocket upgrades and gRPC presence queries. @@ -49,14 +50,14 @@ func (s *Server) HandleWebSocket(w http.ResponseWriter, r *http.Request) { InsecureSkipVerify: true, }) if err != nil { - slog.Error("websocket accept failed", "error", err) + flog.Error("websocket accept failed", "error", err) return } connID := uuid.New().String() conn := newConn(connID, session.HumanId, ws) - slog.Info("websocket connected", "connId", connID, "humanId", session.HumanId) + flog.Info("websocket connected", "connId", connID, "humanId", session.HumanId) // Use the server context, not r.Context(): after upgrade the HTTP request // context can be cancelled by load balancers and nhooyr/websocket would @@ -70,7 +71,7 @@ func (s *Server) HandleWebSocket(w http.ResponseWriter, r *http.Request) { go conn.WritePump(ctx) conn.ReadPump(ctx, s.hub) - slog.Info("websocket disconnected", "connId", connID, "humanId", session.HumanId) + flog.Info("websocket disconnected", "connId", connID, "humanId", session.HumanId) } func (s *Server) GetOnlineHumanIds(ctx context.Context, _ *pbpusher.GetOnlineHumanIdsRequest) (*pbpusher.GetOnlineHumanIdsResponse, error) { diff --git a/go/internal/redis.go b/go/internal/redis.go index 2fb052f..f67ed15 100644 --- a/go/internal/redis.go +++ b/go/internal/redis.go @@ -3,17 +3,19 @@ package internal import ( "context" "fmt" - "log/slog" "os" - "github.com/flowy-live/llink/internal/utils" + "github.com/flowy-live/llink/internal/utils/flog" + "github.com/redis/go-redis/v9" + + "github.com/flowy-live/llink/internal/utils" ) func ConnectAndTestRedis(db int) *redis.Client { redisHost := utils.MustGetEnv("REDIS_HOST") if redisHost == "" { - slog.Error("must provide REDIS_HOST") + flog.Error("must provide REDIS_HOST") os.Exit(1) } redisAddr := fmt.Sprintf("%s:%s", redisHost, "6379") @@ -32,7 +34,7 @@ func ConnectAndTestRedis(db int) *redis.Client { if err != nil { panic(err) } - slog.Debug("redis test", "key", val) + flog.Debug("redis test", "key", val) if val != "value" { panic("unexpected value") } diff --git a/go/internal/speech/service.go b/go/internal/speech/service.go index 0a88c75..eff0929 100644 --- a/go/internal/speech/service.go +++ b/go/internal/speech/service.go @@ -2,7 +2,8 @@ package speech import ( "context" - "log/slog" + + "github.com/flowy-live/llink/internal/utils/flog" dgapi "github.com/deepgram/deepgram-go-sdk/v3/pkg/api/listen/v1/rest" interfaces "github.com/deepgram/deepgram-go-sdk/v3/pkg/client/interfaces" @@ -56,7 +57,7 @@ func (s *speechServiceImpl) Transcribe(ctx context.Context, mediaUrl string) (*T response, err := s.deepgramClient.FromURL(ctx, mediaUrl, options) if err != nil { - slog.Error("failed to transcribe prerecorded media", "error", err) + flog.Error("failed to transcribe prerecorded media", "error", err) return nil, err } diff --git a/go/internal/testhelper/postgres.go b/go/internal/testhelper/postgres.go index b9c8df3..d59fcef 100644 --- a/go/internal/testhelper/postgres.go +++ b/go/internal/testhelper/postgres.go @@ -3,10 +3,11 @@ package testhelper import ( "context" "fmt" - "log/slog" "os" "time" + "github.com/flowy-live/llink/internal/utils/flog" + "github.com/golang-migrate/migrate/v4" _ "github.com/golang-migrate/migrate/v4/database/postgres" _ "github.com/golang-migrate/migrate/v4/source/file" @@ -39,20 +40,20 @@ func SetupTestDB() *pgxpool.Pool { ), ) if err != nil { - slog.Error("failed to start postgres container", "error", err) + flog.Error("failed to start postgres container", "error", err) os.Exit(1) } // Get connection URL host, err := container.Host(ctx) if err != nil { - slog.Error("failed to get container host", "error", err) + flog.Error("failed to get container host", "error", err) os.Exit(1) } port, err := container.MappedPort(ctx, "5432") if err != nil { - slog.Error("failed to get container port", "error", err) + flog.Error("failed to get container port", "error", err) os.Exit(1) } @@ -62,20 +63,20 @@ func SetupTestDB() *pgxpool.Pool { // Run migrations m, err := migrate.New("file://../../migrations", connectionURL) if err != nil { - slog.Error("failed to create migrate instance", "error", err) + flog.Error("failed to create migrate instance", "error", err) os.Exit(1) } defer m.Close() if err := m.Up(); err != nil && err != migrate.ErrNoChange { - slog.Error("failed to run migrations", "error", err) + flog.Error("failed to run migrations", "error", err) os.Exit(1) } // Create connection pool dbPool, err = pgxpool.New(ctx, connectionURL) if err != nil { - slog.Error("failed to create connection pool", "error", err) + flog.Error("failed to create connection pool", "error", err) os.Exit(1) } @@ -89,7 +90,7 @@ func TeardownTestDB() { } if container != nil { if err := container.Terminate(ctx); err != nil { - slog.Error("failed to terminate container", "error", err) + flog.Error("failed to terminate container", "error", err) } } } diff --git a/go/internal/utils/env.go b/go/internal/utils/env.go index 568d3ef..4877ace 100644 --- a/go/internal/utils/env.go +++ b/go/internal/utils/env.go @@ -2,7 +2,8 @@ package utils import ( "os" - "log/slog" + + "github.com/flowy-live/llink/internal/utils/flog" ) // EnvVar enumerates the env vars referenced via this package. @@ -15,7 +16,7 @@ func MustGetEnv[T string | EnvVar](key T) string { keyString := string(key) value := os.Getenv(keyString) if value == "" { - slog.Errorf("Missing required environment variable %s", key) + flog.Error("missing required environment variable", "key", key) panic("Missing required environment variable") } @@ -26,7 +27,7 @@ func MustGetEnv[T string | EnvVar](key T) string { func GetEnv(key string) string { value := os.Getenv(key) if value == "" { - slog.Warnf("Missing optional environment variable %s", key) + flog.Warn("missing optional environment variable", "key", key) } return value } diff --git a/go/internal/utils/flog/flog.go b/go/internal/utils/flog/flog.go new file mode 100644 index 0000000..59b45ac --- /dev/null +++ b/go/internal/utils/flog/flog.go @@ -0,0 +1,58 @@ +// Package flog is a structured logger formatted for GCP Cloud Logging. +// +// The API mirrors log/slog: a message string followed by alternating key/value +// pairs. Each kv pair becomes a JSON field in the emitted record, which Cloud +// Logging promotes to a queryable jsonPayload field. +// +// flog.Info("particle processed", "particleID", id, "elapsed", dt) +package flog + +import ( + "os" + + "github.com/sirupsen/logrus" +) + +// Prepares global logger for GCP Cloud Logging's quirks +var log = func() *logrus.Logger { + logger := logrus.New() + logger.SetFormatter(&logrus.JSONFormatter{ + FieldMap: logrus.FieldMap{ + logrus.FieldKeyLevel: "severity", + logrus.FieldKeyMsg: "message", + logrus.FieldKeyTime: "timestamp", + }, + }) + logger.SetOutput(os.Stderr) + logger.SetLevel(logrus.DebugLevel) + return logger +}() + +// fields converts alternating key/value args into a logrus.Fields map. A +// trailing odd arg or non-string key is stored under "!BADKEY" to match +// slog's behavior. +func fields(args []any) logrus.Fields { + if len(args) == 0 { + return nil + } + f := make(logrus.Fields, len(args)/2+1) + for i := 0; i < len(args); i += 2 { + if i+1 >= len(args) { + f["!BADKEY"] = args[i] + break + } + key, ok := args[i].(string) + if !ok { + f["!BADKEY"] = args[i] + continue + } + f[key] = args[i+1] + } + return f +} + +func Debug(msg string, args ...any) { log.WithFields(fields(args)).Debug(msg) } +func Info(msg string, args ...any) { log.WithFields(fields(args)).Info(msg) } +func Warn(msg string, args ...any) { log.WithFields(fields(args)).Warn(msg) } +func Error(msg string, args ...any) { log.WithFields(fields(args)).Error(msg) } +func Fatal(msg string, args ...any) { log.WithFields(fields(args)).Fatal(msg) } diff --git a/go/internal/waitlist/service.go b/go/internal/waitlist/service.go index 4711079..c5f753c 100644 --- a/go/internal/waitlist/service.go +++ b/go/internal/waitlist/service.go @@ -4,12 +4,14 @@ import ( "context" "errors" "fmt" - "log/slog" "time" + "github.com/flowy-live/llink/internal/utils/flog" + + "github.com/jackc/pgx/v5/pgxpool" + pbaero "github.com/flowy-live/llink/genproto/aero" "github.com/flowy-live/llink/internal/utils" - "github.com/jackc/pgx/v5/pgxpool" ) type WaitlistEntry struct { @@ -83,7 +85,7 @@ func (s *serviceImpl) AddToWaitlist(ctx context.Context, email string, metadata }, }) if err != nil { - slog.Error("failed to add to waitlist", "error", err, "email", email) + flog.Error("failed to add to waitlist", "error", err, "email", email) } return nil