package discovery import ( "context" "os" "testing" "time" "github.com/itworx/pulse/internal/database" ) // TestPostgreSQLDiscoveryStoreCoordinatesAndDeduplicates exercises the real // job_runs lease and events deduplication semantics. It is skipped unless // PULSE_TEST_DATABASE_URL points at a disposable PostgreSQL instance. func TestPostgreSQLDiscoveryStoreCoordinatesAndDeduplicates(t *testing.T) { dsn := os.Getenv("PULSE_TEST_DATABASE_URL") if dsn == "" { t.Skip("PULSE_TEST_DATABASE_URL is not set") } ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) defer cancel() pool, err := database.NewPool(ctx, database.Config{URL: dsn, MaxConns: 4, MinConns: 1}) if err != nil { t.Fatal(err) } defer pool.Close() if err := database.Migrate(ctx, pool); err != nil { t.Fatal(err) } sourceID := newID() if _, err := pool.Exec(ctx, `INSERT INTO data_sources (id,type,name,configuration_ref) VALUES ($1,'agent','discovery-test','test')`, sourceID); err != nil { t.Fatal(err) } t.Cleanup(func() { cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 20*time.Second) defer cleanupCancel() _, _ = pool.Exec(cleanupCtx, `DELETE FROM events WHERE source_id=$1`, sourceID) _, _ = pool.Exec(cleanupCtx, `DELETE FROM data_sources WHERE id=$1`, sourceID) }) now := time.Date(2026, 8, 4, 12, 0, 0, 0, time.UTC) jobKey := "container:" + sourceID first, err := NewPostgresStore(pool, "worker-a") if err != nil { t.Fatal(err) } first.Now = func() time.Time { return now } second, err := NewPostgresStore(pool, "worker-b") if err != nil { t.Fatal(err) } second.Now = func() time.Time { return now } claimed, err := first.Claim(ctx, jobKey, now) if err != nil || !claimed { t.Fatalf("first claim = %v err = %v", claimed, err) } if claimed, err := second.Claim(ctx, jobKey, now); err != nil || claimed { t.Fatalf("second worker claimed a held window: %v err = %v", claimed, err) } if err := first.Finish(ctx, JobRun{Key: jobKey, Status: "succeeded", Attempts: 1, StartedAt: now, CompletedAt: now}); err != nil { t.Fatal(err) } if claimed, err := second.Claim(ctx, jobKey, now); err != nil || claimed { t.Fatalf("a completed window was reclaimed: %v err = %v", claimed, err) } // A worker that dies mid-run must not block the job forever. expiredKey := jobKey + ":expired" first.LeaseTTL = time.Second if claimed, err := first.Claim(ctx, expiredKey, now); err != nil || !claimed { t.Fatalf("expired-window claim = %v err = %v", claimed, err) } second.Now = func() time.Time { return now.Add(5 * time.Second) } if claimed, err := second.Claim(ctx, expiredKey, now); err != nil || !claimed { t.Fatalf("expired lease was not reclaimed: %v err = %v", claimed, err) } event := Event{SourceID: sourceID, DedupKey: "entity:1:state_changed", Type: "container.state_changed", Summary: "Container state changed.", Severity: "attention", OccurredAt: now} inserted, err := first.Emit(ctx, event) if err != nil || !inserted { t.Fatalf("first emit inserted = %v err = %v", inserted, err) } inserted, err = first.Emit(ctx, event) if err != nil || inserted { t.Fatalf("repeated emit inserted = %v err = %v", inserted, err) } var count int if err := pool.QueryRow(ctx, `SELECT count(*) FROM events WHERE source_id=$1`, sourceID).Scan(&count); err != nil { t.Fatal(err) } if count != 1 { t.Fatalf("event rows = %d, want 1", count) } if _, err := first.Emit(ctx, Event{SourceID: "not-a-uuid", DedupKey: "k", Type: "t", Summary: "s", OccurredAt: now}); err == nil { t.Fatal("an unregistered source must be rejected at the persistence boundary") } }