diff --git a/go/Dockerfile.emailnotifierjob b/go/Dockerfile.emailnotifierjob
new file mode 100644
index 0000000..5bd7815
--- /dev/null
+++ b/go/Dockerfile.emailnotifierjob
@@ -0,0 +1,20 @@
+# golang two stage build
+FROM golang:1.25 AS first-stage
+
+WORKDIR /app
+
+COPY go.mod go.sum ./
+RUN go mod download && go mod verify
+
+COPY . .
+
+WORKDIR /app/cmd/emailnotifierjob
+RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -o main
+RUN ls
+
+FROM alpine:latest AS second-stage
+
+WORKDIR /app
+COPY --from=first-stage /app/cmd/emailnotifierjob .
+RUN echo "copied over binary to production stage"
+CMD ["./main"]
diff --git a/go/Makefile b/go/Makefile
index d9ab0d9..df45e5b 100644
--- a/go/Makefile
+++ b/go/Makefile
@@ -46,7 +46,7 @@ migrate-prod:
SKAFFOLD_DEFAULT_REPO=$(PROD_REPO) skaffold run -p migrations --kube-context prod --tail
# ---- Deploy ----
-# Use MODULE=orion or MODULE=worker or MODULE=pusher to deploy a single service, e.g.:
+# Use MODULE=orion or MODULE=worker or MODULE=pusher or MODULE=emailnotifierjob to deploy a single service, e.g.:
# make deploy-dev MODULE=orion
.PHONY: deploy-dev
diff --git a/go/README.md b/go/README.md
index fdd6085..2100f92 100644
--- a/go/README.md
+++ b/go/README.md
@@ -13,6 +13,7 @@ make deploy-prod
make deploy-dev MODULE=orion
make deploy-dev MODULE=worker
make deploy-dev MODULE=pusher
+make deploy-dev MODULE=emailnotifierjob
```
## Migrations
diff --git a/go/cmd/emailnotifierjob/main.go b/go/cmd/emailnotifierjob/main.go
new file mode 100644
index 0000000..e600089
--- /dev/null
+++ b/go/cmd/emailnotifierjob/main.go
@@ -0,0 +1,293 @@
+package main
+
+import (
+ "context"
+ "fmt"
+ "log/slog"
+ "os"
+ "strings"
+ "time"
+
+ "cloud.google.com/go/firestore"
+ pbaero "github.com/flowy-live/llink/genproto/aero"
+ pbpusher "github.com/flowy-live/llink/genproto/llink/pusher"
+ "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/particle"
+ "github.com/flowy-live/llink/internal/utils"
+ "google.golang.org/grpc"
+ "google.golang.org/grpc/credentials/insecure"
+)
+
+const (
+ // Only notify about streams with activity in the last 24 hours
+ maxActivityAge = 24 * time.Hour
+ // Minimum time a message must be unread before we consider notifying
+ unreadThreshold = 10 * time.Minute
+ // Minimum time between emails to the same user
+ emailCooldown = 1 * time.Hour
+)
+
+func main() {
+ ctx := context.Background()
+
+ // Initialize Postgres
+ db.Init()
+ defer db.Cleanup()
+
+ // Initialize Firestore
+ gcpProject := utils.MustGetEnv("GCP_PROJECT")
+ firestoreClient, err := firestore.NewClient(ctx, gcpProject)
+ if err != nil {
+ slog.Error("failed to create Firestore client", "error", err)
+ os.Exit(1)
+ }
+ defer firestoreClient.Close()
+
+ // Initialize aero (email) gRPC client
+ 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)
+ os.Exit(1)
+ }
+ defer aeroConn.Close()
+ aeroSvc := pbaero.NewPrimaryClient(aeroConn)
+
+ // Initialize pusher gRPC client
+ 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()
+ pusherSvc := pbpusher.NewPusherServiceClient(pusherConn)
+
+ // Initialize services
+ humanSvc := human.NewService(db.Pool())
+ networkSvc := network.NewService(db.Pool())
+
+ slog.Info("starting email notification cycle")
+ if err := runNotificationCycle(ctx, firestoreClient, aeroSvc, pusherSvc, humanSvc, networkSvc); err != nil {
+ slog.Error("notification cycle failed", "error", err)
+ os.Exit(1)
+ }
+ slog.Info("email notification cycle complete")
+}
+
+func runNotificationCycle(
+ ctx context.Context,
+ fsClient *firestore.Client,
+ aeroSvc pbaero.PrimaryClient,
+ pusherSvc pbpusher.PusherServiceClient,
+ humanSvc human.Service,
+ networkSvc network.Service,
+) error {
+ now := time.Now()
+
+ // Load all networks
+ networks, err := networkSvc.ListAll(ctx)
+ if err != nil {
+ return fmt.Errorf("listing networks: %w", err)
+ }
+
+ // Load all humans into a lookup map
+ allHumans, err := humanSvc.ListAll(ctx)
+ if err != nil {
+ return fmt.Errorf("listing humans: %w", err)
+ }
+ humansById := make(map[string]*human.Human, len(allHumans))
+ for _, h := range allHumans {
+ humansById[h.ID] = h
+ }
+
+ // Track which streams each human is behind on, and the latest activity across those streams
+ behindCounts := map[string]int{}
+ latestActivity := map[string]time.Time{}
+
+ // Get all currently connected humans (empty channel_ids = all connections)
+ allOnline := map[string]bool{}
+ presenceResp, err := pusherSvc.BulkGetPresence(ctx, &pbpusher.BulkGetPresenceRequest{})
+ if err != nil {
+ slog.Error("failed to get connected humans", "error", err)
+ // Continue without presence data — worst case we send an email to someone who's online
+ } else if ch, ok := presenceResp.Presences["_all"]; ok {
+ for _, id := range ch.HumanIds {
+ allOnline[id] = true
+ }
+ }
+
+ for _, net := range networks {
+ // Build set of all member humanIds for this network (members + admin)
+ networkMembers := make(map[string]bool, len(net.MemberHumanIds)+1)
+ for _, id := range net.MemberHumanIds {
+ networkMembers[id] = true
+ }
+ networkMembers[net.AdminHumanId] = true
+
+ // Query Firestore for open streams in this network
+ streams, err := getOpenStreams(ctx, fsClient, net.ID)
+ if err != nil {
+ slog.Error("failed to query streams", "networkId", net.ID, "error", err)
+ continue
+ }
+
+ for _, stream := range streams {
+ if stream.LastChildCreatedAt == nil {
+ continue
+ }
+ // Skip streams with no recent activity
+ if now.Sub(*stream.LastChildCreatedAt) > maxActivityAge {
+ continue
+ }
+ // Skip if the latest message is too fresh (within threshold)
+ if now.Sub(*stream.LastChildCreatedAt) < unreadThreshold {
+ continue
+ }
+
+ // Resolve members from visible_to
+ members := resolveMembers(stream.VisibleTo, networkMembers)
+
+ for humanId := range members {
+ marker, hasMarker := stream.PlaybackMarkers[humanId]
+ if hasMarker && !marker.Before(*stream.LastChildCreatedAt) {
+ continue // up to date
+ }
+ // No marker or marker is behind → this human is behind on this stream
+ behindCounts[humanId]++
+ if stream.LastChildCreatedAt.After(latestActivity[humanId]) {
+ latestActivity[humanId] = *stream.LastChildCreatedAt
+ }
+ }
+ }
+
+ }
+
+ // Send notifications
+ sentCount := 0
+ for humanId, count := range behindCounts {
+ // Skip online users
+ if allOnline[humanId] {
+ continue
+ }
+
+ h, ok := humansById[humanId]
+ if !ok {
+ continue
+ }
+
+ // Skip if notifications disabled
+ if !h.EmailNotificationsEnabled {
+ continue
+ }
+
+ // Skip if no new activity since last notification
+ if h.LastEmailNotificationSentAt != nil && !latestActivity[humanId].After(*h.LastEmailNotificationSentAt) {
+ continue
+ }
+
+ // Enforce cooldown between emails to the same user
+ if h.LastEmailNotificationSentAt != nil && now.Sub(*h.LastEmailNotificationSentAt) < emailCooldown {
+ continue
+ }
+
+ // Send email
+ if err := sendNotificationEmail(ctx, aeroSvc, h, count); err != nil {
+ slog.Error("failed to send email", "humanId", humanId, "error", err)
+ continue
+ }
+
+ // Update last sent timestamp
+ if err := humanSvc.UpdateLastEmailNotificationSentAt(ctx, humanId, now); err != nil {
+ slog.Error("failed to update last_email_notification_sent_at", "humanId", humanId, "error", err)
+ }
+
+ sentCount++
+ }
+
+ slog.Info("notification cycle summary",
+ "networks", len(networks),
+ "humansBehind", len(behindCounts),
+ "emailsSent", sentCount,
+ )
+ return nil
+}
+
+// getOpenStreams queries Firestore for all open stream particles in a network.
+func getOpenStreams(ctx context.Context, client *firestore.Client, networkId string) ([]particle.FirestoreStreamParticle, error) {
+ collPath := fmt.Sprintf("networks/%s/children", networkId)
+ docs, err := client.Collection(collPath).
+ Where("type", "==", "stream").
+ Where("status", "==", "open").
+ Documents(ctx).
+ GetAll()
+ if err != nil {
+ return nil, err
+ }
+
+ streams := make([]particle.FirestoreStreamParticle, 0, len(docs))
+ 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)
+ continue
+ }
+ streams = append(streams, s)
+ }
+ return streams, nil
+}
+
+// resolveMembers expands visible_to entries into a set of humanIds.
+// "human:{id}" adds that id directly. "network:{id}" expands to all network members.
+func resolveMembers(visibleTo []string, networkMembers map[string]bool) map[string]bool {
+ members := map[string]bool{}
+ for _, entry := range visibleTo {
+ if strings.HasPrefix(entry, "human:") {
+ humanId := strings.TrimPrefix(entry, "human:")
+ members[humanId] = true
+ } else if strings.HasPrefix(entry, "network:") {
+ // Expand to all network members
+ for id := range networkMembers {
+ members[id] = true
+ }
+ }
+ }
+ return members
+}
+
+func sendNotificationEmail(ctx context.Context, aeroSvc pbaero.PrimaryClient, h *human.Human, streamCount int) error {
+ streamsWord := "stream"
+ if streamCount != 1 {
+ streamsWord = "streams"
+ }
+
+ subject := fmt.Sprintf("You have unseen messages in %d %s", streamCount, streamsWord)
+ html := buildEmailHTML(h.EmailPrefix, streamCount, streamsWord)
+
+ _, err := aeroSvc.ShootEmail(ctx, &pbaero.ShootEmailRequest{
+ ToEmails: []string{h.Email},
+ Subject: subject,
+ TemplateData: &pbaero.ShootEmailRequest_SimpleHtmlData{
+ SimpleHtmlData: &pbaero.SimpleHtmlData{
+ Html: html,
+ },
+ },
+ })
+ return err
+}
+
+func buildEmailHTML(name string, count int, streamsWord string) string {
+ return fmt.Sprintf(`
+
+
+ Hi %s,
+ You have unread messages in %d %s on Flowy.llink.
+ Open the app to catch up with your team.
+
+ Best,
Flowy Team
+
+
+`, name, count, streamsWord)
+}
diff --git a/go/cmd/orion/main.go b/go/cmd/orion/main.go
index 25fd578..6cd2183 100644
--- a/go/cmd/orion/main.go
+++ b/go/cmd/orion/main.go
@@ -122,6 +122,9 @@ func main() {
mux.Handle("POST /auth/sign-out", withAuth(h.SignOut))
mux.Handle("GET /auth/me", withAuth(h.GetCurrentHuman))
+ // Settings
+ mux.Handle("PATCH /humans/me/settings", withAuth(h.UpdateSettings))
+
// Networks
mux.Handle("POST /networks", withAuth(h.CreateNetwork))
mux.Handle("GET /networks", withAuth(h.ListNetworks))
diff --git a/go/genproto/llink/pusher/pusher.pb.go b/go/genproto/llink/pusher/pusher.pb.go
index beb85d5..727ec72 100644
--- a/go/genproto/llink/pusher/pusher.pb.go
+++ b/go/genproto/llink/pusher/pusher.pb.go
@@ -25,6 +25,7 @@ type BulkGetPresenceRequest struct {
sizeCache protoimpl.SizeCache
unknownFields protoimpl.UnknownFields
+ // Channel IDs to query. Empty means "all connected humans".
ChannelIds []string `protobuf:"bytes,1,rep,name=channel_ids,json=channelIds,proto3" json:"channel_ids,omitempty"`
}
diff --git a/go/genproto/llink/pusher/pusher_grpc.pb.go b/go/genproto/llink/pusher/pusher_grpc.pb.go
index 26f6600..0cf7816 100644
--- a/go/genproto/llink/pusher/pusher_grpc.pb.go
+++ b/go/genproto/llink/pusher/pusher_grpc.pb.go
@@ -26,6 +26,9 @@ const (
//
// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
type PusherServiceClient interface {
+ // BulkGetPresence returns presence for the given channels.
+ // When channel_ids is empty, returns all connected humanIDs across all channels
+ // under the key "_all" — useful for checking overall online status.
BulkGetPresence(ctx context.Context, in *BulkGetPresenceRequest, opts ...grpc.CallOption) (*BulkGetPresenceResponse, error)
}
@@ -51,6 +54,9 @@ func (c *pusherServiceClient) BulkGetPresence(ctx context.Context, in *BulkGetPr
// All implementations must embed UnimplementedPusherServiceServer
// for forward compatibility.
type PusherServiceServer interface {
+ // BulkGetPresence returns presence for the given channels.
+ // When channel_ids is empty, returns all connected humanIDs across all channels
+ // under the key "_all" — useful for checking overall online status.
BulkGetPresence(context.Context, *BulkGetPresenceRequest) (*BulkGetPresenceResponse, error)
mustEmbedUnimplementedPusherServiceServer()
}
diff --git a/go/internal/handler/handler.go b/go/internal/handler/handler.go
index 45f9077..2925942 100644
--- a/go/internal/handler/handler.go
+++ b/go/internal/handler/handler.go
@@ -50,10 +50,11 @@ func NewHandler(authSvc auth.AuthService, humanSvc human.Service, networkSvc net
// Response DTOs
type Human struct {
- Id string `json:"id"`
- Email string `json:"email"`
- EmailPrefix string `json:"email_prefix"`
- CreatedAt time.Time `json:"created_at"`
+ Id string `json:"id"`
+ Email string `json:"email"`
+ EmailPrefix string `json:"email_prefix"`
+ EmailNotificationsEnabled bool `json:"email_notifications_enabled"`
+ CreatedAt time.Time `json:"created_at"`
}
type Network struct {
@@ -267,6 +268,35 @@ func (h *Handler) GetCurrentHuman(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(dto)
}
+type UpdateSettingsRequest struct {
+ EmailNotificationsEnabled *bool `json:"email_notifications_enabled"`
+}
+
+// UpdateSettings updates the authenticated human's settings
+func (h *Handler) UpdateSettings(w http.ResponseWriter, r *http.Request) {
+ humanId, ok := middleware.HumanIdFromContext(r.Context())
+ if !ok {
+ http.Error(w, "unauthorized", http.StatusUnauthorized)
+ return
+ }
+
+ var req UpdateSettingsRequest
+ if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
+ http.Error(w, "invalid request body", http.StatusBadRequest)
+ return
+ }
+
+ 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)
+ http.Error(w, "internal server error", http.StatusInternalServerError)
+ return
+ }
+ }
+
+ w.WriteHeader(http.StatusNoContent)
+}
+
// ============================================================================
// Network Handlers
// ============================================================================
@@ -953,10 +983,11 @@ func waitlistEntryToDTO(e *waitlist.WaitlistEntry) WaitlistEntryResponse {
func humanToDTO(h *human.Human) Human {
return Human{
- Id: h.ID,
- Email: h.Email,
- EmailPrefix: h.EmailPrefix,
- CreatedAt: h.CreatedAt,
+ Id: h.ID,
+ Email: h.Email,
+ EmailPrefix: h.EmailPrefix,
+ EmailNotificationsEnabled: h.EmailNotificationsEnabled,
+ CreatedAt: h.CreatedAt,
}
}
diff --git a/go/internal/human/models.go b/go/internal/human/models.go
index 66b3d08..a17f592 100644
--- a/go/internal/human/models.go
+++ b/go/internal/human/models.go
@@ -3,8 +3,10 @@ package human
import "time"
type Human struct {
- ID string
- Email string
- EmailPrefix string
- CreatedAt time.Time
+ ID string
+ Email string
+ EmailPrefix string
+ EmailNotificationsEnabled bool
+ LastEmailNotificationSentAt *time.Time
+ CreatedAt time.Time
}
diff --git a/go/internal/human/repository.go b/go/internal/human/repository.go
index 9314ab4..8045cfc 100644
--- a/go/internal/human/repository.go
+++ b/go/internal/human/repository.go
@@ -4,6 +4,7 @@ import (
"context"
"errors"
"strings"
+ "time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
@@ -33,6 +34,9 @@ type repository interface {
getByID(ctx context.Context, id string) (*Human, error)
create(ctx context.Context, email string) (*Human, error)
exists(ctx context.Context, email string) (bool, error)
+ listAll(ctx context.Context) ([]*Human, error)
+ updateEmailNotificationsEnabled(ctx context.Context, id string, enabled bool) error
+ updateLastEmailNotificationSentAt(ctx context.Context, id string, t time.Time) error
}
type repositoryImpl struct {
@@ -46,9 +50,9 @@ func newRepository(pool *pgxpool.Pool) repository {
func (r *repositoryImpl) getByEmail(ctx context.Context, email string) (*Human, error) {
var h Human
err := r.pool.QueryRow(ctx,
- `SELECT id, email, created_at FROM humans WHERE email = $1`,
+ `SELECT id, email, email_notifications_enabled, last_email_notification_sent_at, created_at FROM humans WHERE email = $1`,
email,
- ).Scan(&h.ID, &h.Email, &h.CreatedAt)
+ ).Scan(&h.ID, &h.Email, &h.EmailNotificationsEnabled, &h.LastEmailNotificationSentAt, &h.CreatedAt)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, errNotFound
@@ -62,9 +66,9 @@ func (r *repositoryImpl) getByEmail(ctx context.Context, email string) (*Human,
func (r *repositoryImpl) getByID(ctx context.Context, id string) (*Human, error) {
var h Human
err := r.pool.QueryRow(ctx,
- `SELECT id, email, created_at FROM humans WHERE id = $1`,
+ `SELECT id, email, email_notifications_enabled, last_email_notification_sent_at, created_at FROM humans WHERE id = $1`,
id,
- ).Scan(&h.ID, &h.Email, &h.CreatedAt)
+ ).Scan(&h.ID, &h.Email, &h.EmailNotificationsEnabled, &h.LastEmailNotificationSentAt, &h.CreatedAt)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, errNotFound
@@ -84,9 +88,9 @@ func (r *repositoryImpl) create(ctx context.Context, email string) (*Human, erro
var h Human
err = r.pool.QueryRow(ctx,
`INSERT INTO humans (id, email) VALUES ($1, $2)
- RETURNING id, email, created_at`,
+ RETURNING id, email, email_notifications_enabled, last_email_notification_sent_at, created_at`,
id.String(), email,
- ).Scan(&h.ID, &h.Email, &h.CreatedAt)
+ ).Scan(&h.ID, &h.Email, &h.EmailNotificationsEnabled, &h.LastEmailNotificationSentAt, &h.CreatedAt)
if err != nil {
return nil, err
}
@@ -106,3 +110,52 @@ func (r *repositoryImpl) exists(ctx context.Context, email string) (bool, error)
}
return exists, nil
}
+
+func (r *repositoryImpl) listAll(ctx context.Context) ([]*Human, error) {
+ rows, err := r.pool.Query(ctx,
+ `SELECT id, email, email_notifications_enabled, last_email_notification_sent_at, created_at FROM humans`,
+ )
+ if err != nil {
+ return nil, err
+ }
+ defer rows.Close()
+
+ var humans []*Human
+ for rows.Next() {
+ var h Human
+ if err := rows.Scan(&h.ID, &h.Email, &h.EmailNotificationsEnabled, &h.LastEmailNotificationSentAt, &h.CreatedAt); err != nil {
+ return nil, err
+ }
+ h.EmailPrefix = emailPrefix(h.Email)
+ humans = append(humans, &h)
+ }
+ return humans, rows.Err()
+}
+
+func (r *repositoryImpl) updateEmailNotificationsEnabled(ctx context.Context, id string, enabled bool) error {
+ result, err := r.pool.Exec(ctx,
+ `UPDATE humans SET email_notifications_enabled = $2 WHERE id = $1`,
+ id, enabled,
+ )
+ if err != nil {
+ return err
+ }
+ if result.RowsAffected() == 0 {
+ return errNotFound
+ }
+ return nil
+}
+
+func (r *repositoryImpl) updateLastEmailNotificationSentAt(ctx context.Context, id string, t time.Time) error {
+ result, err := r.pool.Exec(ctx,
+ `UPDATE humans SET last_email_notification_sent_at = $2 WHERE id = $1`,
+ id, t,
+ )
+ if err != nil {
+ return err
+ }
+ if result.RowsAffected() == 0 {
+ return errNotFound
+ }
+ return nil
+}
diff --git a/go/internal/human/service.go b/go/internal/human/service.go
index 6ff94e1..211805c 100644
--- a/go/internal/human/service.go
+++ b/go/internal/human/service.go
@@ -3,6 +3,7 @@ package human
import (
"context"
"errors"
+ "time"
"github.com/flowy-live/llink/internal/utils"
"github.com/jackc/pgx/v5/pgxpool"
@@ -16,6 +17,12 @@ type Service interface {
GetByEmail(ctx context.Context, email string) (*Human, error)
// GetByID returns ErrNotFound if no human found
GetByID(ctx context.Context, id string) (*Human, error)
+ // ListAll returns all humans
+ ListAll(ctx context.Context) ([]*Human, error)
+ // UpdateEmailNotificationsEnabled toggles email notification preference
+ UpdateEmailNotificationsEnabled(ctx context.Context, id string, enabled bool) error
+ // UpdateLastEmailNotificationSentAt records when the last notification email was sent
+ UpdateLastEmailNotificationSentAt(ctx context.Context, id string, t time.Time) error
}
type serviceImpl struct {
@@ -62,3 +69,23 @@ func (s *serviceImpl) GetByID(ctx context.Context, id string) (*Human, error) {
}
return h, err
}
+
+func (s *serviceImpl) ListAll(ctx context.Context) ([]*Human, error) {
+ return s.repo.listAll(ctx)
+}
+
+func (s *serviceImpl) UpdateEmailNotificationsEnabled(ctx context.Context, id string, enabled bool) error {
+ err := s.repo.updateEmailNotificationsEnabled(ctx, id, enabled)
+ if errors.Is(err, errNotFound) {
+ return ErrNotFound
+ }
+ return err
+}
+
+func (s *serviceImpl) UpdateLastEmailNotificationSentAt(ctx context.Context, id string, t time.Time) error {
+ err := s.repo.updateLastEmailNotificationSentAt(ctx, id, t)
+ if errors.Is(err, errNotFound) {
+ return ErrNotFound
+ }
+ return err
+}
diff --git a/go/internal/network/repository.go b/go/internal/network/repository.go
index b39bae3..b7db2f7 100644
--- a/go/internal/network/repository.go
+++ b/go/internal/network/repository.go
@@ -35,6 +35,7 @@ type repository interface {
getMemberHumanIds(ctx context.Context, networkID string) ([]string, error)
getNetworksForHuman(ctx context.Context, humanId string) ([]*Network, error)
isMember(ctx context.Context, networkID, humanId string) (bool, error)
+ listAll(ctx context.Context) ([]*Network, error)
// Invitations
createInvitation(ctx context.Context, networkID, email string) error
@@ -202,6 +203,37 @@ func (r *repositoryImpl) isMember(ctx context.Context, networkID, humanId string
return isMember, err
}
+func (r *repositoryImpl) listAll(ctx context.Context) ([]*Network, error) {
+ rows, err := r.pool.Query(ctx,
+ `SELECT id, name, admin_human_id, created_at FROM networks`,
+ )
+ if err != nil {
+ return nil, err
+ }
+ defer rows.Close()
+
+ var networks []*Network
+ for rows.Next() {
+ var n Network
+ if err := rows.Scan(&n.ID, &n.Name, &n.AdminHumanId, &n.CreatedAt); err != nil {
+ return nil, err
+ }
+ networks = append(networks, &n)
+ }
+ if err := rows.Err(); err != nil {
+ return nil, err
+ }
+
+ for _, n := range networks {
+ n.MemberHumanIds, err = r.getMemberHumanIds(ctx, n.ID)
+ if err != nil {
+ return nil, err
+ }
+ }
+
+ return networks, nil
+}
+
// Invitation methods
func (r *repositoryImpl) createInvitation(ctx context.Context, networkID, email string) error {
diff --git a/go/internal/network/service.go b/go/internal/network/service.go
index 8a4ecbb..b3d841e 100644
--- a/go/internal/network/service.go
+++ b/go/internal/network/service.go
@@ -26,6 +26,8 @@ type Service interface {
RemoveMember(ctx context.Context, networkID, humanId string) error
ListForHuman(ctx context.Context, humanId string) ([]*Network, error)
IsMember(ctx context.Context, networkID, humanId string) (bool, error)
+ // ListAll returns all networks with their members
+ ListAll(ctx context.Context) ([]*Network, error)
// Invitations (email-based, for users who haven't registered yet)
InviteByEmail(ctx context.Context, networkID string, emails []string) error
@@ -116,6 +118,10 @@ func (s *serviceImpl) IsMember(ctx context.Context, networkID, humanId string) (
return s.repo.isMember(ctx, networkID, humanId)
}
+func (s *serviceImpl) ListAll(ctx context.Context) ([]*Network, error) {
+ return s.repo.listAll(ctx)
+}
+
// Invitation methods
func (s *serviceImpl) InviteByEmail(ctx context.Context, networkID string, emails []string) error {
diff --git a/go/internal/particle/firestore_types.go b/go/internal/particle/firestore_types.go
index 82d5869..70e21e2 100644
--- a/go/internal/particle/firestore_types.go
+++ b/go/internal/particle/firestore_types.go
@@ -43,12 +43,13 @@ type FirestoreMediaParticleProperties struct {
}
type FirestoreStreamParticle struct {
- CreatedByHumanId string `firestore:"created_by_human_id"`
- Type string `firestore:"type"`
- // Properties FirestoreStreamParticleProperties `firestore:"properties"`
- CreatedAt time.Time `firestore:"created_at,serverTimestamp"`
- LastChildCreatedAt *time.Time `firestore:"last_child_created_at,omitempty"`
- VisibleTo []string `firestore:"visible_to"`
- UpdatedAt *time.Time `firestore:"updated_at,omitempty"`
- HuddleActiveParticipants []string `firestore:"huddle_active_participants,omitempty"`
+ CreatedByHumanId string `firestore:"created_by_human_id"`
+ Type string `firestore:"type"`
+ Status string `firestore:"status"`
+ CreatedAt time.Time `firestore:"created_at,serverTimestamp"`
+ LastChildCreatedAt *time.Time `firestore:"last_child_created_at,omitempty"`
+ VisibleTo []string `firestore:"visible_to"`
+ PlaybackMarkers map[string]time.Time `firestore:"playback_markers,omitempty"`
+ UpdatedAt *time.Time `firestore:"updated_at,omitempty"`
+ HuddleActiveParticipants []string `firestore:"huddle_active_participants,omitempty"`
}
diff --git a/go/internal/pusher/redis_bridge.go b/go/internal/pusher/redis_bridge.go
index 908829b..2212f94 100644
--- a/go/internal/pusher/redis_bridge.go
+++ b/go/internal/pusher/redis_bridge.go
@@ -143,6 +143,41 @@ func (rb *RedisBridge) GetPresence(ctx context.Context, channelIDs []string) (ma
return result, nil
}
+// GetAllConnectedHumanIDs scans all channel connection hashes in Redis and returns
+// the deduplicated set of all humanIDs that have at least one active connection.
+func (rb *RedisBridge) GetAllConnectedHumanIDs(ctx context.Context) ([]string, error) {
+ allHumanIDs := make(map[string]bool)
+ var cursor uint64
+
+ for {
+ keys, nextCursor, err := rb.client.Scan(ctx, cursor, channelConnsPrefix+"*"+channelConnsSuffix, 100).Result()
+ if err != nil {
+ return nil, fmt.Errorf("failed to scan channel keys: %w", err)
+ }
+
+ for _, key := range keys {
+ members, err := rb.client.HVals(ctx, key).Result()
+ if err != nil && err != redis.Nil {
+ continue
+ }
+ for _, humanID := range members {
+ allHumanIDs[humanID] = true
+ }
+ }
+
+ cursor = nextCursor
+ if cursor == 0 {
+ break
+ }
+ }
+
+ result := make([]string, 0, len(allHumanIDs))
+ for id := range allHumanIDs {
+ result = append(result, id)
+ }
+ return result, nil
+}
+
// --- Pub/Sub listener (runs in its own goroutine) ---
// Listen subscribes to Redis Pub/Sub and forwards events to the local hub.
diff --git a/go/internal/pusher/server.go b/go/internal/pusher/server.go
index 80ae174..a4dfb4c 100644
--- a/go/internal/pusher/server.go
+++ b/go/internal/pusher/server.go
@@ -72,7 +72,22 @@ func (s *Server) HandleWebSocket(w http.ResponseWriter, r *http.Request) {
}
// BulkGetPresence implements the gRPC PusherService.
+// When channel_ids is empty, returns all connected humanIDs across all channels
+// under the key "_all" — useful for checking overall online status.
func (s *Server) BulkGetPresence(ctx context.Context, req *pbpusher.BulkGetPresenceRequest) (*pbpusher.BulkGetPresenceResponse, error) {
+ // Empty channel_ids → return all connected humans
+ if len(req.ChannelIds) == 0 {
+ humanIDs, err := s.bridge.GetAllConnectedHumanIDs(ctx)
+ if err != nil {
+ return nil, err
+ }
+ return &pbpusher.BulkGetPresenceResponse{
+ Presences: map[string]*pbpusher.ChannelPresence{
+ "_all": {HumanIds: humanIDs},
+ },
+ }, nil
+ }
+
presence, err := s.bridge.GetPresence(ctx, req.ChannelIds)
if err != nil {
return nil, err
diff --git a/go/k8s/dev/emailnotifierjob.yaml b/go/k8s/dev/emailnotifierjob.yaml
new file mode 100644
index 0000000..4ff5fb9
--- /dev/null
+++ b/go/k8s/dev/emailnotifierjob.yaml
@@ -0,0 +1,42 @@
+apiVersion: batch/v1
+kind: CronJob
+metadata:
+ name: emailnotifierjob
+spec:
+ schedule: "*/5 * * * *"
+ concurrencyPolicy: Forbid
+ successfulJobsHistoryLimit: 3
+ failedJobsHistoryLimit: 3
+ jobTemplate:
+ spec:
+ template:
+ metadata:
+ labels:
+ app: emailnotifierjob
+ spec:
+ serviceAccountName: default-service-account
+ nodeSelector:
+ cloud.google.com/gke-spot: "true"
+ restartPolicy: Never
+ containers:
+ - name: emailnotifierjob
+ image: "emailnotifierjob"
+ resources:
+ requests:
+ memory: "52Mi"
+ cpu: 50m
+ limits:
+ memory: "52Mi"
+ cpu: 50m
+ env:
+ - name: "GCP_PROJECT"
+ value: "flowy-dev-440017"
+ - name: "AERO_ADDR"
+ value: "aero:50051"
+ - name: "PUSHER_GRPC_ADDR"
+ value: "pusher:50051"
+ - name: "LLINK_POSTGRES_CONNECTION_URL"
+ valueFrom:
+ secretKeyRef:
+ name: shared-secrets
+ key: LLINK_POSTGRES_CONNECTION_URL
diff --git a/go/k8s/prod/emailnotifierjob.yaml b/go/k8s/prod/emailnotifierjob.yaml
new file mode 100644
index 0000000..8022db6
--- /dev/null
+++ b/go/k8s/prod/emailnotifierjob.yaml
@@ -0,0 +1,40 @@
+apiVersion: batch/v1
+kind: CronJob
+metadata:
+ name: emailnotifierjob
+spec:
+ schedule: "*/5 * * * *"
+ concurrencyPolicy: Forbid
+ successfulJobsHistoryLimit: 3
+ failedJobsHistoryLimit: 3
+ jobTemplate:
+ spec:
+ template:
+ metadata:
+ labels:
+ app: emailnotifierjob
+ spec:
+ serviceAccountName: default-service-account
+ restartPolicy: Never
+ containers:
+ - name: emailnotifierjob
+ image: "emailnotifierjob"
+ resources:
+ requests:
+ memory: "52Mi"
+ cpu: 50m
+ limits:
+ memory: "52Mi"
+ cpu: 50m
+ env:
+ - name: "GCP_PROJECT"
+ value: "flowy-prod-440017"
+ - name: "AERO_ADDR"
+ value: "aero:50051"
+ - name: "PUSHER_GRPC_ADDR"
+ value: "pusher:50051"
+ - name: "LLINK_POSTGRES_CONNECTION_URL"
+ valueFrom:
+ secretKeyRef:
+ name: shared-secrets
+ key: LLINK_POSTGRES_CONNECTION_URL
diff --git a/go/migrations/000012_human_email_notifications.down.sql b/go/migrations/000012_human_email_notifications.down.sql
new file mode 100644
index 0000000..de59914
--- /dev/null
+++ b/go/migrations/000012_human_email_notifications.down.sql
@@ -0,0 +1,3 @@
+ALTER TABLE humans
+ DROP COLUMN email_notifications_enabled,
+ DROP COLUMN last_email_notification_sent_at;
diff --git a/go/migrations/000012_human_email_notifications.up.sql b/go/migrations/000012_human_email_notifications.up.sql
new file mode 100644
index 0000000..d091e18
--- /dev/null
+++ b/go/migrations/000012_human_email_notifications.up.sql
@@ -0,0 +1,3 @@
+ALTER TABLE humans
+ ADD COLUMN email_notifications_enabled BOOLEAN NOT NULL DEFAULT TRUE,
+ ADD COLUMN last_email_notification_sent_at TIMESTAMPTZ;
diff --git a/go/protocol b/go/protocol
index 9e3e1e0..a8683c0 160000
--- a/go/protocol
+++ b/go/protocol
@@ -1 +1 @@
-Subproject commit 9e3e1e0e6a0cc942be32d6f7e86c49fec596ff87
+Subproject commit a8683c0397d1144506078f6a8bb6991e2d6c4140
diff --git a/go/skaffold.yaml b/go/skaffold.yaml
index 49dd7c0..34bbe9c 100644
--- a/go/skaffold.yaml
+++ b/go/skaffold.yaml
@@ -93,6 +93,41 @@ profiles:
---
apiVersion: skaffold/v4beta11
kind: Config
+metadata:
+ name: emailnotifierjob
+build:
+ local: {}
+ tagPolicy:
+ gitCommit:
+ variant: AbbrevCommitSha
+profiles:
+ - name: dev
+ build:
+ artifacts:
+ - image: emailnotifierjob
+ context: .
+ docker:
+ dockerfile: Dockerfile.emailnotifierjob
+ manifests:
+ rawYaml:
+ - k8s/dev/emailnotifierjob.yaml
+ deploy:
+ kubectl: {}
+ - name: prod
+ build:
+ artifacts:
+ - image: emailnotifierjob
+ context: .
+ docker:
+ dockerfile: Dockerfile.emailnotifierjob
+ manifests:
+ rawYaml:
+ - k8s/prod/emailnotifierjob.yaml
+ deploy:
+ kubectl: {}
+---
+apiVersion: skaffold/v4beta11
+kind: Config
metadata:
name: pusher
build:
diff --git a/js/src/api/client.ts b/js/src/api/client.ts
index 378dfbf..e1f2e98 100644
--- a/js/src/api/client.ts
+++ b/js/src/api/client.ts
@@ -125,6 +125,12 @@ class ApiClient {
return data.url;
}
+ // --- Settings ---
+
+ async updateSettings(data: { email_notifications_enabled?: boolean }): Promise {
+ await this.requestVoid("PATCH", "/humans/me/settings", data);
+ }
+
// --- Depot ---
async prepareUpload(data: PrepareUploadRequest) {
diff --git a/js/src/api/types.ts b/js/src/api/types.ts
index 3bf50ec..5f75f72 100644
--- a/js/src/api/types.ts
+++ b/js/src/api/types.ts
@@ -5,6 +5,7 @@ export const HumanSchema = z.object({
created_at: z.coerce.date(),
email: z.string().email(),
email_prefix: z.string(),
+ email_notifications_enabled: z.boolean(),
});
export type Human = z.infer;
diff --git a/js/src/features/settings-page.tsx b/js/src/features/settings-page.tsx
index a6d76bb..152181a 100644
--- a/js/src/features/settings-page.tsx
+++ b/js/src/features/settings-page.tsx
@@ -1,12 +1,15 @@
+import { useState } from "react";
import { useNavigate } from "react-router-dom";
-import { ChevronRight, LogOut, User, Info, Shield } from "lucide-react";
+import { ChevronRight, LogOut, User, Info, Shield, Mail } from "lucide-react";
import { Avatar, AvatarFallback } from "@/components/ui/avatar";
import { Separator } from "@/components/ui/separator";
+import { Switch } from "@/components/ui/switch";
import { WindowControls } from "@/components/window-controls";
import { Button } from "@/components/ui/button";
import { Muted } from "@/components/ui/typography";
import { ScrollArea } from "@/components/ui/scroll-area";
import { useAuthStore } from "@/stores/auth-store";
+import { apiClient } from "@/api/client";
import { ArrowLeft } from "lucide-react";
interface SettingsRowProps {
@@ -63,6 +66,25 @@ export default function SettingsPage() {
const navigate = useNavigate();
const user = useAuthStore((s) => s.user);
const signOut = useAuthStore((s) => s.signOut);
+ const [emailNotifications, setEmailNotifications] = useState(
+ user?.email_notifications_enabled ?? true,
+ );
+
+ const handleToggleEmailNotifications = async (checked: boolean) => {
+ setEmailNotifications(checked);
+ useAuthStore.setState((state) => ({
+ user: state.user ? { ...state.user, email_notifications_enabled: checked } : null,
+ }));
+ try {
+ await apiClient.updateSettings({ email_notifications_enabled: checked });
+ } catch {
+ // Revert on failure
+ setEmailNotifications(!checked);
+ useAuthStore.setState((state) => ({
+ user: state.user ? { ...state.user, email_notifications_enabled: !checked } : null,
+ }));
+ }
+ };
const initials = user?.email_prefix?.slice(0, 2).toUpperCase() ?? "?";
@@ -115,6 +137,24 @@ export default function SettingsPage() {
+
+
+
+
+
+
+ Email notifications
+
+
+
+
+
+
+
}