package workerruntime import ( "context" "os" "testing" "time" "github.com/itworx/pulse/internal/database" "github.com/itworx/pulse/internal/probe" "github.com/itworx/pulse/internal/systemstatus" ) // TestPostgreSQLLeaseStoreAndJobHealth exercises the real job_runs coordination // and the status projection the API reads back. It is skipped unless // PULSE_TEST_DATABASE_URL points at a disposable PostgreSQL instance. func TestPostgreSQLLeaseStoreAndJobHealth(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) } store := PostgresLeaseStore{Pool: pool} job := Job{Name: "discovery", Component: systemstatus.ComponentWorker, Interval: time.Minute, Timeout: time.Second, Run: func(context.Context) (Outcome, error) { return Outcome{}, nil }} jobKey := "integration-" + newID() now := time.Date(2026, 8, 4, 12, 0, 0, 0, time.UTC) t.Cleanup(func() { cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 20*time.Second) defer cleanupCancel() _, _ = pool.Exec(cleanupCtx, `DELETE FROM job_runs WHERE job_key=$1`, jobKey) }) lease, acquired, err := store.Acquire(ctx, job.JobType(), jobKey, now, "worker-a", now, time.Minute) if err != nil || !acquired { t.Fatalf("first acquire = %v err = %v", acquired, err) } if _, acquired, err := store.Acquire(ctx, job.JobType(), jobKey, now, "worker-b", now, time.Minute); err != nil || acquired { t.Fatalf("a held window was acquired twice: %v err = %v", acquired, err) } if err := store.Complete(ctx, lease, StatusCompleted, "", map[string]int64{"items": 3}); err != nil { t.Fatal(err) } if err := store.Complete(ctx, lease, StatusCompleted, "", nil); err != ErrLeaseLost { t.Fatalf("completing twice = %v, want ErrLeaseLost", err) } if _, acquired, err := store.Acquire(ctx, job.JobType(), jobKey, now, "worker-b", now, time.Minute); err != nil || acquired { t.Fatalf("a completed window was reacquired: %v err = %v", acquired, err) } // An expired lease is reclaimable so a crashed worker cannot block the job. expired, acquired, err := store.Acquire(ctx, job.JobType(), jobKey, now.Add(time.Minute), "worker-a", now, time.Second) if err != nil || !acquired { t.Fatalf("expired-window acquire = %v err = %v", acquired, err) } if _, acquired, err := store.Acquire(ctx, job.JobType(), jobKey, expired.ScheduledAt, "worker-b", now.Add(10*time.Second), time.Minute); err != nil || !acquired { t.Fatalf("expired lease was not reclaimed: %v err = %v", acquired, err) } health, err := ReadJobHealth(ctx, pool, []Job{job}) if err != nil { t.Fatal(err) } if len(health) != 1 || health[0].Component != systemstatus.ComponentWorker || health[0].LastRunAt.IsZero() { t.Fatalf("job health = %#v", health) } // A job type that never ran is absent, so its component stays Unknown. absent, err := ReadJobHealth(ctx, pool, []Job{{Name: "never-scheduled", Component: systemstatus.ComponentProbes, Interval: time.Minute, Timeout: time.Second, Run: job.Run}}) if err != nil { t.Fatal(err) } if len(absent) != 0 { t.Fatalf("a job that never ran reported health: %#v", absent) } } // TestPostgreSQLProbeAndAliasStoresAreIdempotent verifies that replaying a // probe batch and a container alias snapshot creates no duplicate rows. func TestPostgreSQLProbeAndAliasStoresAreIdempotent(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, serviceID, probeID, entityID := newID(), newID(), newID(), newID() now := time.Now().UTC().Truncate(time.Microsecond) seed := []struct { query string args []any }{ {`INSERT INTO data_sources (id,type,name,configuration_ref) VALUES ($1,'agent','worker-test','test')`, []any{sourceID}}, {`INSERT INTO entities (id,entity_type,canonical_name,display_name,first_seen_at) VALUES ($1,'container','worker/test','worker-test',$2)`, []any{entityID, now}}, {`INSERT INTO services (id,name) VALUES ($1,'worker-test-service')`, []any{serviceID}}, {`INSERT INTO probes (id,service_id,name,probe_type,target,interval_seconds,timeout_seconds) VALUES ($1,$2,'worker-test-probe','tcp','{"host":"example.internal","port":443}'::jsonb,60,5)`, []any{probeID, serviceID}}, } for _, statement := range seed { if _, err := pool.Exec(ctx, statement.query, statement.args...); err != nil { t.Fatalf("seed %q: %v", statement.query, err) } } t.Cleanup(func() { cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 20*time.Second) defer cleanupCancel() for _, statement := range []string{ `DELETE FROM probe_results WHERE probe_id=$1`, `DELETE FROM probes WHERE id=$1`, } { _, _ = pool.Exec(cleanupCtx, statement, probeID) } _, _ = pool.Exec(cleanupCtx, `DELETE FROM services WHERE id=$1`, serviceID) _, _ = pool.Exec(cleanupCtx, `DELETE FROM container_aliases WHERE source_id=$1`, sourceID) _, _ = pool.Exec(cleanupCtx, `DELETE FROM entities WHERE id=$1`, entityID) _, _ = pool.Exec(cleanupCtx, `DELETE FROM data_sources WHERE id=$1`, sourceID) }) probeStore := PostgresProbeStore{Pool: pool, SourceID: sourceID} due, err := probeStore.ListDue(ctx, now, MaxProbeBatch) if err != nil { t.Fatal(err) } found := false for _, definition := range due { if definition.ID == probeID { found = true if definition.Interval != time.Minute || definition.Timeout != 5*time.Second || definition.Target.Host != "example.internal" { t.Fatalf("decoded probe = %#v", definition) } } } if !found { t.Fatal("a probe without results is not due") } batch := []probe.Result{{ProbeID: probeID, ObservedAt: now, CompletedAt: now, State: "up", Attempts: 1}} saved, err := probeStore.SaveResults(ctx, batch) if err != nil { t.Fatal(err) } if saved != 1 { t.Fatalf("saved = %d, want 1", saved) } saved, err = probeStore.SaveResults(ctx, batch) if err != nil { t.Fatal(err) } if saved != 0 { t.Fatalf("replayed save = %d, want 0", saved) } if due, err := probeStore.ListDue(ctx, now, MaxProbeBatch); err == nil { for _, definition := range due { if definition.ID == probeID { t.Fatal("a probe with a fresh result is still due") } } } else { t.Fatal(err) } aliasStore := PostgresContainerAliasStore{Pool: pool} record := ContainerAliasRecord{State: "running", Health: "healthy"} record.EntityID, record.SourceID, record.RuntimeID, record.Name = entityID, sourceID, "runtime-1", "worker-test" record.FirstSeenAt, record.LastSeenAt, record.Active = now, now, true for range 2 { if err := aliasStore.Save(ctx, sourceID, []ContainerAliasRecord{record}); err != nil { t.Fatal(err) } } stored, err := aliasStore.List(ctx, sourceID) if err != nil { t.Fatal(err) } if len(stored) != 1 || stored[0].RuntimeID != "runtime-1" || stored[0].State != "running" || !stored[0].Active { t.Fatalf("stored aliases = %#v", stored) } if err := aliasStore.Save(ctx, sourceID, nil); err != nil { t.Fatal(err) } if remaining, err := aliasStore.List(ctx, sourceID); err != nil || len(remaining) != 0 { t.Fatalf("aliases after prune = %#v err = %v", remaining, err) } }