diff --git a/cmd/backend/main.go b/cmd/backend/main.go index 2ab1e4ec..52bd0ed0 100644 --- a/cmd/backend/main.go +++ b/cmd/backend/main.go @@ -314,23 +314,34 @@ func main() { log.Fatal(err) } - // Realtime event transport. Prefer Google Pub/Sub when configured (prod); - // otherwise bridge events to the realtime service over Redis (local dev - // and any env without GCP), reusing the Redis client we already built. - // Exactly one transport is active, so events are never delivered twice. - gcpProjectID := os.Getenv("GCP_PROJECT_ID") - if gcpProjectID != "" { + // Realtime event transport, chosen by PUBSUB_ENABLED — the SAME flag the + // Elixir realtime service reads — so the two sides can never split-brain + // (publisher on Pub/Sub while the subscriber listens on Redis = all events + // dropped). PUBSUB_ENABLED=true => Google Pub/Sub (prod); anything else => + // Redis bridge (local dev / non-GCP). Exactly one transport is active, so + // events are never delivered twice. + if os.Getenv("PUBSUB_ENABLED") == "true" { + gcpProjectID := os.Getenv("GCP_PROJECT_ID") + if gcpProjectID == "" { + log.Fatal("PUBSUB_ENABLED=true requires GCP_PROJECT_ID") + } pubsubClient, err := pubsub.NewClient(ctx, gcpProjectID) if err != nil { sentry.CaptureException(err) - log.Printf("Warning: Failed to initialize Pub/Sub client: %v", err) - } else { - streamingPublisher = pubsub.NewStreamingPublisher(pubsubClient) + log.Fatal("Failed to initialize Pub/Sub client: ", err) } + // Create the realtime topics + "-sub" subscriptions if missing, + // so the Elixir Broadway consumers always have a subscription to read. + if err := pubsubClient.EnsureRealtimeTopology(ctx); err != nil { + sentry.CaptureException(err) + log.Fatal("Failed to provision Pub/Sub topics/subscriptions: ", err) + } + streamingPublisher = pubsub.NewStreamingPublisher(pubsubClient) + log.Println("Realtime events published to Google Pub/Sub") } if streamingPublisher == nil { streamingPublisher = pubsub.NewStreamingPublisher(pubsub.NewRedisBus(cache.Client, "")) - log.Println("Realtime events bridged over Redis (Google Pub/Sub not configured)") + log.Println("Realtime events bridged over Redis (Pub/Sub disabled)") } emailCfg, err := cfg.LoadEmailConfig(ctx) @@ -866,6 +877,10 @@ func main() { // Fan reply + bounce events from the advanced-outreach brain out to // customer webhooks AND third-party integration actions (Slack / CRM). advancedService.WireDispatcher(webhookService) + // Let instant action chains (reply/open/click branches) launch a + // "run_automation" node, the same flow the scheduler runs at a step + // boundary. Backend ingests deliverability + can process replies too. + advancedService.WireAutomationRunner(integrationServiceForHandler) // Wire native (Warmbly-internal) automation actions + realtime now that // the advanced/contact/org services exist (the integration service was // constructed earlier). @@ -938,6 +953,12 @@ func main() { // enqueue the first warmup task. go tasksService.StartWarmupReconciler(ctx, 10*time.Minute) + // Campaign reconciler: re-seed active campaigns whose self-perpetuating + // task chain died (a swallowed enqueue, a worker bounce mid-tick, or a + // crash between send and enqueue). Campaigns have no other bootstrap once + // started, so without this a stranded campaign stops sending forever. + go tasksService.StartCampaignReconciler(ctx, 5*time.Minute) + // Danger zone: schedule + execute delayed deletions (orgs, accounts). dangerZoneRepository := repository.NewDangerZoneRepository(primaryDB.Pool) dangerZoneService = dangerzone.NewService( diff --git a/cmd/consumer/main.go b/cmd/consumer/main.go index c7106db0..f7aa8483 100644 --- a/cmd/consumer/main.go +++ b/cmd/consumer/main.go @@ -152,24 +152,33 @@ func main() { log.Fatal(err) } - // Realtime event transport. Prefer Google Pub/Sub when configured (prod); - // otherwise bridge events to the realtime service over Redis (local dev and - // any env without GCP), reusing the Redis client we already built. Exactly - // one transport is active, so events are never delivered twice. - gcpProjectID := os.Getenv("GCP_PROJECT_ID") + // Realtime event transport, chosen by PUBSUB_ENABLED — the SAME flag the + // backend and the Elixir realtime service read, so the three services can + // never split-brain. PUBSUB_ENABLED=true => Google Pub/Sub (prod); anything + // else => Redis bridge (local dev / non-GCP). Exactly one transport is active. var streamingPublisher *pubsub.StreamingPublisher - if gcpProjectID != "" { + if os.Getenv("PUBSUB_ENABLED") == "true" { + gcpProjectID := os.Getenv("GCP_PROJECT_ID") + if gcpProjectID == "" { + log.Fatal("PUBSUB_ENABLED=true requires GCP_PROJECT_ID") + } pubsubClient, err := pubsub.NewClient(ctx, gcpProjectID) if err != nil { sentry.CaptureException(err) log.Fatal(err) } defer pubsubClient.Close() + // Idempotently ensure the realtime topics + subscriptions exist (safe to + // run from both backend and consumer; AlreadyExists is treated as success). + if err := pubsubClient.EnsureRealtimeTopology(ctx); err != nil { + sentry.CaptureException(err) + log.Fatal("Failed to provision Pub/Sub topics/subscriptions: ", err) + } streamingPublisher = pubsub.NewStreamingPublisher(pubsubClient) } if streamingPublisher == nil { streamingPublisher = pubsub.NewStreamingPublisher(pubsub.NewRedisBus(redisCache.Client, "")) - log.Println("Realtime events bridged over Redis (Google Pub/Sub not configured)") + log.Println("Realtime events bridged over Redis (Pub/Sub disabled)") } // Repositories @@ -234,6 +243,11 @@ func main() { warmupService, ) advancedService.WireDispatcher(webhookService) + // Reply/open/click instant action chains run in THIS process (inbox ingest + + // tracking consumer), so a "run_automation" node on an instant branch must be + // able to launch the flow here too. Without this it would be stamped sent and + // never fire. Mirrors the scheduler's automationRunner wiring. + advancedService.WireAutomationRunner(integrationServiceC) // Native (Warmbly-internal) automation actions run wherever the event is // dispatched. Reply/bounce/warmup events dispatch in THIS process, so without // wiring native actions here a reply-triggered automation's add_tag /