package alertworker import ( "context" "encoding/json" "errors" "fmt" "strings" "time" "github.com/itworx/pulse/internal/alert" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" ) type PostgresLeaseStore struct { Pool *pgxpool.Pool } func (s PostgresLeaseStore) Acquire(ctx context.Context, jobType, jobKey string, scheduledAt time.Time, owner string, now time.Time, ttl time.Duration) (Lease, bool, error) { if s.Pool == nil { return Lease{}, false, ErrInvalidConfig } tx, err := s.Pool.BeginTx(ctx, pgx.TxOptions{}) if err != nil { return Lease{}, false, fmt.Errorf("begin evaluator lease: %w", err) } defer func() { _ = tx.Rollback(ctx) }() leaseID := alert.NewID() leaseUntil := now.Add(ttl) var insertedID string err = tx.QueryRow(ctx, `INSERT INTO job_runs (id,job_type,job_key,scheduled_at,started_at,status,lease_owner,lease_until) VALUES ($1,$2,$3,$4,$5,'running',$6,$7) ON CONFLICT (job_type,job_key,scheduled_at) DO NOTHING RETURNING id`, leaseID, jobType, jobKey, scheduledAt.UTC(), now.UTC(), owner, leaseUntil.UTC()).Scan(&insertedID) if err == nil { if err := tx.Commit(ctx); err != nil { return Lease{}, false, fmt.Errorf("commit evaluator lease: %w", err) } return Lease{ID: insertedID, JobKey: jobKey, ScheduledAt: scheduledAt.UTC(), Owner: owner}, true, nil } if !errors.Is(err, pgx.ErrNoRows) { return Lease{}, false, fmt.Errorf("insert evaluator lease: %w", err) } var status string var existingUntil *time.Time if err := tx.QueryRow(ctx, `SELECT status,lease_until FROM job_runs WHERE job_type=$1 AND job_key=$2 AND scheduled_at=$3 FOR UPDATE`, jobType, jobKey, scheduledAt.UTC()).Scan(&status, &existingUntil); err != nil { return Lease{}, false, fmt.Errorf("read evaluator lease: %w", err) } if status != "running" || (existingUntil != nil && existingUntil.After(now)) { if err := tx.Commit(ctx); err != nil { return Lease{}, false, err } return Lease{}, false, nil } tag, err := tx.Exec(ctx, `UPDATE job_runs SET status='running',started_at=$1,completed_at=NULL,error_code=NULL,lease_owner=$2,lease_until=$3 WHERE job_type=$4 AND job_key=$5 AND scheduled_at=$6 AND status='running' AND (lease_until IS NULL OR lease_until <= $7)`, now.UTC(), owner, leaseUntil.UTC(), jobType, jobKey, scheduledAt.UTC(), now.UTC()) if err != nil { return Lease{}, false, fmt.Errorf("renew evaluator lease: %w", err) } if tag.RowsAffected() != 1 { if err := tx.Commit(ctx); err != nil { return Lease{}, false, err } return Lease{}, false, nil } if err := tx.Commit(ctx); err != nil { return Lease{}, false, fmt.Errorf("commit evaluator lease renewal: %w", err) } return Lease{ID: leaseID, JobKey: jobKey, ScheduledAt: scheduledAt.UTC(), Owner: owner}, true, nil } func (s PostgresLeaseStore) Complete(ctx context.Context, lease Lease, status, errorCode string) error { if s.Pool == nil { return ErrInvalidConfig } if status != "completed" && status != "failed" && status != "canceled" { return errors.New("invalid evaluator job status") } errorCode = strings.TrimSpace(errorCode) if len(errorCode) > 160 { errorCode = errorCode[:160] } counts, _ := json.Marshal(map[string]string{"status": status}) tag, err := s.Pool.Exec(ctx, `UPDATE job_runs SET status=$1,completed_at=now(),counts=$2::jsonb,error_code=NULLIF($3,'') ,lease_owner=NULL,lease_until=NULL WHERE job_type=$4 AND job_key=$5 AND scheduled_at=$6 AND lease_owner=$7 AND status='running'`, status, counts, errorCode, jobType, lease.JobKey, lease.ScheduledAt.UTC(), lease.Owner) if err != nil { return fmt.Errorf("complete evaluator job: %w", err) } if tag.RowsAffected() != 1 { return ErrLeaseLost } return nil }