Files
ITWorx Pulse release export bd774932d5
Public source validation / validate (push) Failing after 3m8s
Publish ITWorx Pulse source
2026-09-03 02:09:19 +02:00

197 lines
7.6 KiB
Go

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)
}
}