diff --git a/.markdownlint.json b/.markdownlint.json new file mode 100644 index 0000000..835449c --- /dev/null +++ b/.markdownlint.json @@ -0,0 +1,5 @@ +{ + "MD013": false, + "MD060": false, + "MD029": false +} diff --git a/README.md b/README.md index 3b51007..3b1b850 100644 --- a/README.md +++ b/README.md @@ -16,6 +16,7 @@ A Go library for delayed task processing and state reconciliation. - [Options](#options) - [Delay Strategies](#delay-strategies) - [Groups](#groups) + - [Finalizers](#finalizers) - [Configuration](#configuration) - [Storage](#storage) - [In-Memory](#in-memory) @@ -32,6 +33,7 @@ A Go library for delayed task processing and state reconciliation. - **Retry Strategies**: Fixed delays or exponential backoff - **Tubes**: Route tickets to separate queues - **Groups**: Track batches of related tickets and query their collective progress +- **Finalizers**: Run a ticket automatically after all members of a group have finished - **Automatic Expiration**: Cleanup of completed/expired tickets - **Concurrent Processing**: Configurable worker pools - **Prometheus Metrics**: Per-type and per-tube counters, processing duration and queue wait histograms @@ -160,6 +162,7 @@ kh.Cancel(ctx, id, lymbo.WithKeep(), lymbo.WithErrorReason("cancelled by user")) | `WithGroup(id)` | Assign or transfer ticket to a group | | `WithUpdate(fn)` | Custom ticket modification (executed last) | | `WithResetAttempts()` | Reset attempt counter | +| `AfterGroup(id)` | Register ticket as finalizer for a group (blocked until all current members finish) | ### Delay Strategies @@ -222,6 +225,34 @@ kh.Retry(ctx, t.ID, Tickets submitted without `WithGroup` are ungrouped and never appear in any group query. +### Finalizers + +A finalizer is a ticket blocked until every pending member of a group has finished. Submit group +members first, then submit the finalizer with `AfterGroup`: + +```go +groupID := "order-42-notifications" + +for _, userID := range recipients { + ticket, _ := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "send-notification") + kh.Put(ctx, *ticket, lymbo.WithGroup(groupID), lymbo.WithPayload(userID)) +} + +finalizer, _ := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "notifications-complete") +kh.Put(ctx, *finalizer, lymbo.AfterGroup(groupID)) +``` + +At submission time, the store snapshots all currently pending group members as dependencies. +The finalizer is invisible to workers until every captured dependency reaches a terminal state. +Tickets added to the group after the finalizer is submitted are not captured. + +Notes: + +- Empty group (no pending members) → finalizer is immediately eligible. +- `WithDelay` on the finalizer is evaluated independently of dependencies. +- `AfterGroup("g")` + `WithGroup("g")` on the same `Put` returns `ErrFinalizerInGroup`. +- Re-submitting a finalizer with the same ID is a no-op. + ## Configuration ```go @@ -292,6 +323,7 @@ type Store interface { PollPending(context.Context, PollRequest) (PollResult, error) ExpireTickets(ctx context.Context, limit int, now time.Time) ([]TransitionInfo, error) CountPendingInGroup(ctx context.Context, groupID string) (int, error) + PutAfterGroup(ctx context.Context, ticket Ticket, groupID string) error } ``` diff --git a/errors.go b/errors.go index d0baa74..60ca8bf 100644 --- a/errors.go +++ b/errors.go @@ -10,4 +10,5 @@ var ( ErrTicketIDInvalid = errors.New("ticket ID is invalid") ErrTicketNotFound = errors.New("ticket not found") ErrInvalidStatusTransition = errors.New("invalid status transition") + ErrFinalizerInGroup = errors.New("finalizer must not be a member of the group it finalizes") ) diff --git a/kharon.go b/kharon.go index a0a8375..f74db5e 100644 --- a/kharon.go +++ b/kharon.go @@ -257,9 +257,15 @@ func (k *Kharon) Put(ctx context.Context, t Ticket, opts ...Option) error { if err != nil { return err } + if o.afterGroup != nil && o.groupID != nil && *o.afterGroup == *o.groupID { + return ErrFinalizerInGroup + } if err := beforeUpdate(ctx, &t, o); err != nil { return err } + if o.afterGroup != nil { + return k.store.PutAfterGroup(ctx, t, *o.afterGroup) + } return k.store.Put(ctx, t) } diff --git a/observable.go b/observable.go index c6d0f9e..9bf3def 100644 --- a/observable.go +++ b/observable.go @@ -26,6 +26,14 @@ func (o *observableStore) Put(ctx context.Context, t Ticket) error { return nil } +func (o *observableStore) PutAfterGroup(ctx context.Context, t Ticket, groupID string) error { + if err := o.Store.PutAfterGroup(ctx, t, groupID); err != nil { + return err + } + o.stats.ByKey(t.Type, t.Tube.String()).Inc(stats.Add) + return nil +} + func (o *observableStore) ExpireTickets(ctx context.Context, limit int, now time.Time) ([]TransitionInfo, error) { infos, err := o.Store.ExpireTickets(ctx, limit, now) if err != nil { diff --git a/options.go b/options.go index 5ff9743..34f0228 100644 --- a/options.go +++ b/options.go @@ -59,6 +59,9 @@ type Opts struct { // groupID associates the ticket with a group (effective only on Put). groupID *string + // afterGroup wires the ticket as a finalizer for the given group (effective only on Put). + afterGroup *string + // update allows custom modification of the ticket. update func(ctx context.Context, t *Ticket) error } @@ -156,6 +159,16 @@ func WithGroup(groupID string) Option { } } +// AfterGroup wires the ticket as a finalizer for groupID. +// The finalizer is blocked until all currently pending group members reach a terminal state. +// Only effective when passed to Put. +func AfterGroup(groupID string) Option { + return func(o *Opts) error { + o.afterGroup = &groupID + return nil + } +} + // WithTube sets the tube to transfer the ticket to. func WithTube(tube Tube) Option { return func(o *Opts) error { diff --git a/store.go b/store.go index 05ae512..bacf7a2 100644 --- a/store.go +++ b/store.go @@ -78,6 +78,12 @@ type Store interface { // CountPendingInGroup returns the number of pending tickets with the given group ID. CountPendingInGroup(ctx context.Context, groupID string) (int, error) + + // PutAfterGroup atomically inserts the ticket and wires it as a finalizer for groupID. + // It scans all currently pending group members and creates dep rows blocking the finalizer. + // Returns nil if the ticket ID already exists (idempotent re-submission). + // Returns ErrFinalizerInGroup if the ticket belongs to the same group it finalizes. + PutAfterGroup(ctx context.Context, ticket Ticket, groupID string) error } // TransitionInfo describes a ticket that was affected by a store operation. diff --git a/store/memory/memory.go b/store/memory/memory.go index 36fbcdb..4a16003 100644 --- a/store/memory/memory.go +++ b/store/memory/memory.go @@ -16,6 +16,7 @@ import ( type Store struct { mu sync.RWMutex data map[lymbo.TicketId]lymbo.Ticket + deps map[lymbo.TicketId]map[lymbo.TicketId]struct{} // deps[ticket_id] = set of blocked_by IDs } // Ensure Store implements lymbo.Store interface. @@ -25,6 +26,7 @@ var _ lymbo.Store = (*Store)(nil) func NewStore() *Store { return &Store{ data: make(map[lymbo.TicketId]lymbo.Ticket), + deps: make(map[lymbo.TicketId]map[lymbo.TicketId]struct{}), } } @@ -61,6 +63,49 @@ func (m *Store) Delete(_ context.Context, id lymbo.TicketId) error { defer m.mu.Unlock() delete(m.data, id) + m.removeDep(id) + return nil +} + +// removeDep removes id from all dep sets and deletes its own dep set. +// Must be called with m.mu held. +func (m *Store) removeDep(id lymbo.TicketId) { + for _, depSet := range m.deps { + delete(depSet, id) + } + delete(m.deps, id) +} + +// PutAfterGroup atomically inserts the ticket and wires it as a finalizer for groupID. +func (m *Store) PutAfterGroup(_ context.Context, t lymbo.Ticket, groupID string) error { + if t.ID == "" { + return lymbo.ErrTicketIDEmpty + } + + m.mu.Lock() + defer m.mu.Unlock() + + // REQ 12: no-op if ticket already exists + if _, exists := m.data[t.ID]; exists { + return nil + } + + t.Status = status.Pending + m.data[t.ID] = t + + // REQ 6: scan pending group members, insert dep rows + for id, member := range m.data { + if id == t.ID { + continue + } + if member.GroupId != nil && *member.GroupId == groupID && member.Status == status.Pending { + if m.deps[t.ID] == nil { + m.deps[t.ID] = make(map[lymbo.TicketId]struct{}) + } + m.deps[t.ID][id] = struct{}{} + } + } + return nil } @@ -73,6 +118,7 @@ func (m *Store) DeleteBatch(_ context.Context, ids []lymbo.TicketId) ([]lymbo.Tr if t, ok := m.data[id]; ok { infos = append(infos, lymbo.TransitionInfo{Id: t.ID, Type: t.Type, Tube: t.Tube, Status: t.Status}) delete(m.data, id) + m.removeDep(id) } } return infos, nil @@ -121,6 +167,13 @@ func (m *Store) UpdateBatch(ctx context.Context, updates []lymbo.UpdateSet) ([]l updateOne(&t, us) m.data[t.ID] = t infos = append(infos, lymbo.TransitionInfo{Id: t.ID, Type: t.Type, Tube: t.Tube, Status: t.Status}) + + // REQ 3: remove from all dep sets when ticket reaches terminal state + if t.Status == status.Done || t.Status == status.Failed || t.Status == status.Cancelled { + for _, depSet := range m.deps { + delete(depSet, t.ID) + } + } } return infos, nil @@ -180,6 +233,11 @@ func (m *Store) PollPending(_ context.Context, req lymbo.PollRequest) (lymbo.Pol continue } + // REQ 1: blocked tickets must not be returned + if len(m.deps[t.ID]) > 0 { + continue + } + ready = append(ready, t) } diff --git a/store/postgres/postgres.go b/store/postgres/postgres.go index aaa0340..dd2f2a8 100644 --- a/store/postgres/postgres.go +++ b/store/postgres/postgres.go @@ -206,23 +206,17 @@ func (r *Tickets) DeleteBatch(ctx context.Context, ids []lymbo.TicketId) ([]lymb ticketUUIDs = append(ticketUUIDs, ticketUUID) } - batch := &pgx.Batch{} - for _, ticketUUID := range ticketUUIDs { - batch.Queue(r.queries.delete, ticketUUID) + rows, err := r.db.Query(ctx, r.queries.deleteBatch, ticketUUIDs) + if err != nil { + return nil, err } - - br := r.db.SendBatch(ctx, batch) - defer br.Close() //nolint:errcheck + defer rows.Close() var infos []lymbo.TransitionInfo - for range ticketUUIDs { + for rows.Next() { var id uuid.UUID var ticketType, tube string - err := br.QueryRow().Scan(&id, &ticketType, &tube) - if err != nil { - if errors.Is(err, pgx.ErrNoRows) { - continue - } + if err := rows.Scan(&id, &ticketType, &tube); err != nil { return infos, err } infos = append(infos, lymbo.TransitionInfo{ @@ -231,7 +225,7 @@ func (r *Tickets) DeleteBatch(ctx context.Context, ids []lymbo.TicketId) ([]lymb Tube: lymbo.Tube(tube), }) } - return infos, nil + return infos, rows.Err() } func (r *Tickets) Update(ctx context.Context, id lymbo.TicketId, fn lymbo.UpdateFunc) error { @@ -337,6 +331,9 @@ func (r *Tickets) UpdateBatch(ctx context.Context, updates []lymbo.UpdateSet) ([ if len(updates) == 0 { return nil, nil } + + // Dep rows for terminal transitions are deleted atomically inside the update/backoff + // query templates via a CTE (REQ 3). No explicit transaction needed here. batch := &pgx.Batch{} for _, us := range updates { @@ -348,7 +345,7 @@ func (r *Tickets) UpdateBatch(ctx context.Context, updates []lymbo.UpdateSet) ([ var q string req := []any{ticketUUID} - // status as $2 + // status as $2 switch { case us.Status != nil: req = append(req, sql.NullString{String: us.Status.String(), Valid: true}) @@ -438,6 +435,9 @@ func (r *Tickets) UpdateBatch(ctx context.Context, updates []lymbo.UpdateSet) ([ Status: s, }) } + if err := br.Close(); err != nil { + return infos, err + } return infos, nil } @@ -578,6 +578,39 @@ func (r *Tickets) PollPending(ctx context.Context, req lymbo.PollRequest) (lymbo }, nil } +func (r *Tickets) PutAfterGroup(ctx context.Context, ticket lymbo.Ticket, groupID string) error { + ticketUUID, err := uuid.Parse(ticket.ID.String()) + if err != nil { + return lymbo.ErrTicketIDInvalid + } + + var mtime pgtype.Timestamptz + if ticket.Mtime != nil { + mtime = pgtype.Timestamptz{Time: *ticket.Mtime, Valid: true} + } + tube := ticket.Tube.String() + if tube == "" { + tube = "default" + } + + _, err = r.db.Exec(ctx, r.queries.putAfterGroup, + ticketUUID, // $1 id + ticket.Status.String(), // $2 status + pgtype.Timestamptz{Time: ticket.Runat, Valid: true}, // $3 runat + int16(ticket.Nice), // $4 nice + ticket.Type, // $5 type + pgtype.Timestamptz{Time: ticket.Ctime, Valid: true}, // $6 ctime + mtime, // $7 mtime + int32(ticket.Attempts), // $8 attempts + ticket.Payload, // $9 payload + ticket.ErrorReason, // $10 error_reason + tube, // $11 tube + ticket.GroupId, // $12 group_id (finalizer's own group, may be nil) + groupID, // $13 group to wait for + ) + return err +} + func (r *Tickets) CountPendingInGroup(ctx context.Context, groupID string) (int, error) { var count int err := r.db.QueryRow(ctx, r.queries.countPendingInGroup, groupID).Scan(&count) diff --git a/store/postgres/template.go b/store/postgres/template.go index 318d16f..82b2f8c 100644 --- a/store/postgres/template.go +++ b/store/postgres/template.go @@ -39,6 +39,13 @@ ALTER TABLE {{.TableName}} ADD COLUMN IF NOT EXISTS group_id TEXT DEFAULT NULL; CREATE INDEX IF NOT EXISTS idx_{{.TableName}}_group_pending ON {{.TableName}} (group_id) WHERE group_id IS NOT NULL AND status = 'pending'; +CREATE TABLE IF NOT EXISTS {{.TableName}}_deps ( + ticket_id UUID NOT NULL REFERENCES {{.TableName}}(id), + blocked_by UUID NOT NULL REFERENCES {{.TableName}}(id), + PRIMARY KEY (ticket_id, blocked_by) +); +CREATE INDEX IF NOT EXISTS idx_{{.TableName}}_deps_blocked_by ON {{.TableName}}_deps (blocked_by); + -- Create trigger function CREATE OR REPLACE FUNCTION {{.TableName}}_update_mtime() RETURNS trigger AS $$ @@ -85,7 +92,12 @@ ON CONFLICT (id) DO UPDATE SET var delete = template.Must(template.New("delete").Parse(`-- name: DeleteTicket: DELETE FROM {{.TableName}} WHERE id = $1 RETURNING id, type, tube`)) -var update = template.Must(template.New("update").Parse(`UPDATE {{.TableName}} +var update = template.Must(template.New("update").Parse(`WITH del_deps AS ( + DELETE FROM {{.TableName}}_deps + WHERE blocked_by = $1 + AND $2 IN ('done'::ticket_status, 'failed'::ticket_status, 'cancelled'::ticket_status) +) +UPDATE {{.TableName}} SET status = COALESCE($2, status), nice = COALESCE($3, nice), @@ -99,6 +111,11 @@ RETURNING id, type, tube, status`)) // runat = now() + {jitter} + min(pow({base}, attempt), {max}) var backoff = template.Must(template.New("backoff").Parse(`-- name: BackoffTicket: +WITH del_deps AS ( + DELETE FROM {{.TableName}}_deps + WHERE blocked_by = $1 + AND $2 IN ('done'::ticket_status, 'failed'::ticket_status, 'cancelled'::ticket_status) +) UPDATE {{.TableName}} SET status = COALESCE($2, status), @@ -115,26 +132,49 @@ RETURNING id, type, tube, status`)) var poll = template.Must(template.New("poll").Parse(`-- name: PollTickets: WITH candidates AS ( SELECT t.id, t.runat AS ready_at - FROM {{.TableName}} as t - WHERE t.status = 'pending' AND t.runat <= $1::Timestamptz AND t.tube = ANY($7::text[]) + FROM {{.TableName}} AS t + WHERE t.status = 'pending' + AND t.runat <= $1::Timestamptz + AND t.tube = ANY($7::text[]) + AND NOT EXISTS ( + SELECT 1 FROM {{.TableName}}_deps d WHERE d.ticket_id = t.id LIMIT 1 + ) LIMIT $6 FOR UPDATE SKIP LOCKED -), -rescheduled_tickets AS ( - UPDATE {{.TableName}} as t - SET - attempts = attempts + 1, - runat = $1::Timestamptz + - (GREATEST($2, 0) + LEAST($3, POWER($4, LEAST(t.attempts, $5)))) - * INTERVAL '1 second' - WHERE id IN (SELECT id FROM candidates) - RETURNING id, status, tube, runat, nice, type, ctime, mtime, attempts, payload, error_reason ) -SELECT - r.id, r.status, r.tube, r.runat, r.nice, r.type, r.ctime, r.mtime, r.attempts, r.payload, r.error_reason, - c.ready_at -FROM rescheduled_tickets r -JOIN candidates c ON r.id = c.id;`)) +UPDATE {{.TableName}} AS t +SET + attempts = attempts + 1, + runat = $1::Timestamptz + + (GREATEST($2, 0) + LEAST($3, POWER($4, LEAST(t.attempts, $5)))) + * INTERVAL '1 second' +FROM candidates c +WHERE t.id = c.id +RETURNING + t.id, t.status, t.tube, t.runat, t.nice, t.type, + t.ctime, t.mtime, t.attempts, t.payload, t.error_reason, + c.ready_at;`)) + +var putAfterGroup = template.Must(template.New("putAfterGroup").Parse(`-- name: PutAfterGroup: +WITH inserted AS ( + INSERT INTO {{.TableName}} (id, status, runat, nice, type, ctime, mtime, attempts, payload, error_reason, tube, group_id) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12) + ON CONFLICT (id) DO NOTHING + RETURNING id +) +INSERT INTO {{.TableName}}_deps (ticket_id, blocked_by) +SELECT $1, t.id +FROM {{.TableName}} t +WHERE t.group_id = $13 AND t.status = 'pending' AND t.id != $1 + AND EXISTS (SELECT 1 FROM inserted) +ON CONFLICT DO NOTHING;`)) + +var deleteBatch = template.Must(template.New("deleteBatch").Parse(`-- name: DeleteBatch: +WITH del_deps AS ( + DELETE FROM {{.TableName}}_deps + WHERE ticket_id = ANY($1::uuid[]) OR blocked_by = ANY($1::uuid[]) +) +DELETE FROM {{.TableName}} WHERE id = ANY($1::uuid[]) RETURNING id, type, tube`)) var countPendingInGroup = template.Must(template.New("countPendingInGroup").Parse(`-- name: CountPendingInGroup: SELECT COUNT(*) FROM {{.TableName}} @@ -150,6 +190,8 @@ type Queries struct { get string put string delete string + deleteBatch string + putAfterGroup string update string backoff string poll string @@ -184,6 +226,12 @@ func newQueries(tableName string) (*Queries, error) { if qt.delete, err = exec(delete); err != nil { return nil, fmt.Errorf("failed to execute template `delete`: %w", err) } + if qt.deleteBatch, err = exec(deleteBatch); err != nil { + return nil, fmt.Errorf("failed to execute template `deleteBatch`: %w", err) + } + if qt.putAfterGroup, err = exec(putAfterGroup); err != nil { + return nil, fmt.Errorf("failed to execute template `putAfterGroup`: %w", err) + } if qt.update, err = exec(update); err != nil { return nil, fmt.Errorf("failed to execute template `update`: %w", err) } diff --git a/test/bench_test.go b/test/bench_test.go new file mode 100644 index 0000000..714039c --- /dev/null +++ b/test/bench_test.go @@ -0,0 +1,108 @@ +package lymbo_test + +import ( + "context" + "sync/atomic" + "testing" + "time" + + "github.com/google/uuid" + "github.com/ochaton/lymbo" + "github.com/ochaton/lymbo/store/memory" + "golang.org/x/time/rate" +) + +func BenchmarkMemoryThroughput(b *testing.B) { + ctx := b.Context() + + store := memory.NewStore() + settings := lymbo.DefaultSettings(). + WithMinReactionDelay(100 * time.Microsecond). + WithMaxReactionDelay(5 * time.Millisecond) + kh := lymbo.NewKharon(store, settings, nil) + + var processed atomic.Int64 + + router := lymbo.NewRouter() + router.HandleFunc("bench", func(ctx context.Context, t *lymbo.Ticket) error { + time.Sleep(time.Nanosecond) + processed.Add(1) + return kh.Ack(ctx, t.ID) + }) + + go func() { _ = kh.Run(ctx, router) }() + + b.ResetTimer() + + for i := 0; i < b.N; i++ { + ticket, _ := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "bench") + if err := kh.Put(ctx, *ticket); err != nil { + b.Fatal(err) + } + } + + // wait for all tickets to be processed + deadline := time.Now().Add(30 * time.Second) + for processed.Load() < int64(b.N) { + if time.Now().After(deadline) { + b.Fatalf("timeout: processed %d / %d", processed.Load(), b.N) + } + time.Sleep(time.Millisecond) + } + + b.ReportMetric(float64(b.N)/b.Elapsed().Seconds(), "tickets/sec") +} + +func BenchmarkMemoryQueueWait(b *testing.B) { + const putRatePerSec = 20_000 + + ctx := b.Context() + + store := memory.NewStore() + settings := lymbo.DefaultSettings(). + WithMinReactionDelay(100 * time.Microsecond). + WithMaxReactionDelay(5 * time.Millisecond) + kh := lymbo.NewKharon(store, settings, nil) + + var ( + processed atomic.Int64 + totalWait atomic.Int64 // nanoseconds + ) + + router := lymbo.NewRouter() + router.HandleFunc("bench-wait", func(ctx context.Context, t *lymbo.Ticket) error { + time.Sleep(time.Nanosecond) + totalWait.Add(time.Since(t.ReadyAt).Nanoseconds()) + processed.Add(1) + return kh.Ack(ctx, t.ID) + }) + + go func() { _ = kh.Run(ctx, router) }() + + rl := rate.NewLimiter(putRatePerSec, putRatePerSec) + + b.ResetTimer() + + for i := 0; i < b.N; i++ { + if err := rl.Wait(ctx); err != nil { + b.Fatal(err) + } + ticket, _ := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "bench-wait") + if err := kh.Put(ctx, *ticket); err != nil { + b.Fatal(err) + } + } + + deadline := time.Now().Add(30 * time.Second) + for processed.Load() < int64(b.N) { + if time.Now().After(deadline) { + b.Fatalf("timeout: processed %d / %d", processed.Load(), b.N) + } + time.Sleep(time.Millisecond) + } + b.StopTimer() + + avgWaitNs := totalWait.Load() / int64(b.N) + b.ReportMetric(float64(avgWaitNs), "ns/wait") + b.ReportMetric(float64(b.N)/b.Elapsed().Seconds(), "tickets/sec") +} diff --git a/test/finalizer_testsuite_test.go b/test/finalizer_testsuite_test.go new file mode 100644 index 0000000..26a06b6 --- /dev/null +++ b/test/finalizer_testsuite_test.go @@ -0,0 +1,406 @@ +package lymbo_test + +import ( + "context" + "testing" + "time" + + "github.com/google/uuid" + "github.com/ochaton/lymbo" + "github.com/ochaton/lymbo/status" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// FinalizerTestSuite tests the AfterGroup / finalizer feature. +type FinalizerTestSuite struct { + factory StoreFactory +} + +func NewFinalizerTestSuite(factory StoreFactory) *FinalizerTestSuite { + return &FinalizerTestSuite{factory: factory} +} + +func (s *FinalizerTestSuite) RunAll(t *testing.T) { + t.Run("BlockedInvisibleToPoll", s.TestBlockedInvisibleToPoll) + t.Run("BlockingIgnoresTubeAndRunat", s.TestBlockingIgnoresTubeAndRunat) + t.Run("DoneUnblocks", s.TestDoneUnblocks) + t.Run("FailedUnblocks", s.TestFailedUnblocks) + t.Run("CancelledUnblocks", s.TestCancelledUnblocks) + t.Run("DeleteUnblocks", s.TestDeleteUnblocks) + t.Run("MultiBlockerPartialUnblock", s.TestMultiBlockerPartialUnblock) + t.Run("ZeroMembersImmediatelyEligible", s.TestZeroMembersImmediatelyEligible) + t.Run("LateMemberNotCaptured", s.TestLateMemberNotCaptured) + t.Run("FinalizerInGroupError", s.TestFinalizerInGroupError) + t.Run("IdempotentResubmission", s.TestIdempotentResubmission) + t.Run("WithDelayIndependentGate", s.TestWithDelayIndependentGate) + t.Run("FinalizerIsRegularTicket", s.TestFinalizerIsRegularTicket) +} + +// pollStore calls PollPending directly on the store with sensible defaults. +// TTR is intentionally long so polled tickets don't resurface during the test. +func pollStore(ctx context.Context, t *testing.T, store lymbo.Store) []lymbo.Ticket { + t.Helper() + result, err := store.PollPending(ctx, lymbo.PollRequest{ + Limit: 20, + Now: time.Now(), + TTR: 5 * time.Minute, + BackoffBase: 2.0, + MaxBackoffDelay: 10 * time.Minute, + RequestTubes: []lymbo.Tube{lymbo.DefaultTube}, + }) + require.NoError(t, err) + return result.Tickets +} + +func hasID(tickets []lymbo.Ticket, id lymbo.TicketId) bool { + for _, t := range tickets { + if t.ID == id { + return true + } + } + return false +} + +func terminalUpdate(id lymbo.TicketId, s status.Status) lymbo.UpdateSet { + return lymbo.UpdateSet{Id: id, Status: &s} +} + +// TestBlockedInvisibleToPoll verifies REQ 1–2: blocked ticket never surfaces via PollPending. +func (s *FinalizerTestSuite) TestBlockedInvisibleToPoll(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + member, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *member, lymbo.WithGroup("g"))) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("g"))) + + // Finalizer blocked — only member is polled. + tickets := pollStore(ctx, t, store) + assert.True(t, hasID(tickets, member.ID), "member must be returned") + assert.False(t, hasID(tickets, finalizer.ID), "finalizer must not be returned while blocked") + + // Complete member → dep deleted → finalizer unblocked. + _, err = store.UpdateBatch(ctx, []lymbo.UpdateSet{terminalUpdate(member.ID, status.Done)}) + require.NoError(t, err) + + // Member was rescheduled (TTR=5m) so it won't reappear. Only finalizer should surface. + tickets = pollStore(ctx, t, store) + assert.True(t, hasID(tickets, finalizer.ID), "finalizer must be eligible after member done") + assert.False(t, hasID(tickets, member.ID), "rescheduled member must not reappear") +} + +// TestBlockingIgnoresTubeAndRunat verifies REQ 2: blocking is row-existence only. +// Even with Runat in the past, a blocked ticket must not surface. +func (s *FinalizerTestSuite) TestBlockingIgnoresTubeAndRunat(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + member, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *member, lymbo.WithGroup("g"))) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + finalizer.Runat = time.Now().Add(-10 * time.Second) // clearly in the past + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("g"))) + + tickets := pollStore(ctx, t, store) + assert.False(t, hasID(tickets, finalizer.ID), "finalizer blocked even with past Runat") + assert.True(t, hasID(tickets, member.ID)) +} + +// TestDoneUnblocks verifies REQ 3: Done terminal transition deletes dep rows. +func (s *FinalizerTestSuite) TestDoneUnblocks(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + member, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *member, lymbo.WithGroup("g"))) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("g"))) + + _, err = store.UpdateBatch(ctx, []lymbo.UpdateSet{terminalUpdate(member.ID, status.Done)}) + require.NoError(t, err) + + tickets := pollStore(ctx, t, store) + assert.True(t, hasID(tickets, finalizer.ID), "finalizer must be eligible after member Done") +} + +// TestFailedUnblocks verifies REQ 3: Failed terminal transition deletes dep rows. +func (s *FinalizerTestSuite) TestFailedUnblocks(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + member, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *member, lymbo.WithGroup("g"))) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("g"))) + + _, err = store.UpdateBatch(ctx, []lymbo.UpdateSet{terminalUpdate(member.ID, status.Failed)}) + require.NoError(t, err) + + tickets := pollStore(ctx, t, store) + assert.True(t, hasID(tickets, finalizer.ID), "finalizer must be eligible after member Failed") +} + +// TestCancelledUnblocks verifies REQ 3: Cancelled (kept) transition deletes dep rows. +func (s *FinalizerTestSuite) TestCancelledUnblocks(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + member, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *member, lymbo.WithGroup("g"))) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("g"))) + + _, err = store.UpdateBatch(ctx, []lymbo.UpdateSet{terminalUpdate(member.ID, status.Cancelled)}) + require.NoError(t, err) + + tickets := pollStore(ctx, t, store) + assert.True(t, hasID(tickets, finalizer.ID), "finalizer must be eligible after member Cancelled") +} + +// TestDeleteUnblocks verifies REQ 3+5: DeleteBatch deletes dep rows before the ticket row. +func (s *FinalizerTestSuite) TestDeleteUnblocks(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + member, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *member, lymbo.WithGroup("g"))) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("g"))) + + // Simulate Ack (no-keep path): DeleteBatch must clean deps before deleting ticket. + _, err = store.DeleteBatch(ctx, []lymbo.TicketId{member.ID}) + require.NoError(t, err) + + tickets := pollStore(ctx, t, store) + assert.True(t, hasID(tickets, finalizer.ID), "finalizer must be eligible after member deleted") +} + +// TestMultiBlockerPartialUnblock verifies REQ 3: finalizer stays blocked until all members complete. +func (s *FinalizerTestSuite) TestMultiBlockerPartialUnblock(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + m1, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *m1, lymbo.WithGroup("g"))) + + m2, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *m2, lymbo.WithGroup("g"))) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("g"))) + + // Only m1 done — finalizer still blocked by m2. + _, err = store.UpdateBatch(ctx, []lymbo.UpdateSet{terminalUpdate(m1.ID, status.Done)}) + require.NoError(t, err) + + tickets := pollStore(ctx, t, store) + assert.True(t, hasID(tickets, m2.ID), "m2 must still be polled") + assert.False(t, hasID(tickets, finalizer.ID), "finalizer still blocked by m2") + + // m2 was rescheduled by previous poll. Complete it directly. + _, err = store.UpdateBatch(ctx, []lymbo.UpdateSet{terminalUpdate(m2.ID, status.Done)}) + require.NoError(t, err) + + tickets = pollStore(ctx, t, store) + assert.True(t, hasID(tickets, finalizer.ID), "finalizer eligible after all members done") +} + +// TestZeroMembersImmediatelyEligible verifies REQ 9: no pending members → finalizer immediately eligible. +func (s *FinalizerTestSuite) TestZeroMembersImmediatelyEligible(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + // Group "empty-g" has no pending members. + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("empty-g"))) + + tickets := pollStore(ctx, t, store) + assert.True(t, hasID(tickets, finalizer.ID), "finalizer must be immediately eligible for empty group") +} + +// TestLateMemberNotCaptured verifies REQ 8: members added after PutAfterGroup are not captured. +func (s *FinalizerTestSuite) TestLateMemberNotCaptured(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + m1, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *m1, lymbo.WithGroup("g"))) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("g"))) + + // Late member — added after PutAfterGroup; must not be captured. + late, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *late, lymbo.WithGroup("g"))) + + // Complete m1 (the only captured blocker). + _, err = store.UpdateBatch(ctx, []lymbo.UpdateSet{terminalUpdate(m1.ID, status.Done)}) + require.NoError(t, err) + + // Finalizer is unblocked (no dep on late member). Both finalizer and late are eligible. + tickets := pollStore(ctx, t, store) + assert.True(t, hasID(tickets, finalizer.ID), "finalizer eligible after captured member done") + assert.True(t, hasID(tickets, late.ID), "late member polled independently") +} + +// TestFinalizerInGroupError verifies REQ 7: AfterGroup + WithGroup same ID → ErrFinalizerInGroup. +func (s *FinalizerTestSuite) TestFinalizerInGroupError(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + + err = kh.Put(ctx, *finalizer, lymbo.WithGroup("g"), lymbo.AfterGroup("g")) + require.ErrorIs(t, err, lymbo.ErrFinalizerInGroup) +} + +// TestIdempotentResubmission verifies REQ 12: same finalizer ticket re-submitted → no-op. +func (s *FinalizerTestSuite) TestIdempotentResubmission(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + member, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *member, lymbo.WithGroup("g"))) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("g"))) + // Second submission of the same ticket — must be a no-op. + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("g"))) + + // Finalizer still blocked (dep set unchanged, not duplicated). + tickets := pollStore(ctx, t, store) + require.Len(t, tickets, 1, "only member must be returned") + assert.Equal(t, member.ID, tickets[0].ID) + + // Complete member → finalizer unblocks exactly once. + _, err = store.UpdateBatch(ctx, []lymbo.UpdateSet{terminalUpdate(member.ID, status.Done)}) + require.NoError(t, err) + + tickets = pollStore(ctx, t, store) + require.Len(t, tickets, 1, "only finalizer must surface") + assert.Equal(t, finalizer.ID, tickets[0].ID) +} + +// TestWithDelayIndependentGate verifies REQ 11: WithDelay and AfterGroup are independent gates. +// Finalizer is not eligible until BOTH: all members done AND delay expired. +func (s *FinalizerTestSuite) TestWithDelayIndependentGate(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + member, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "worker") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *member, lymbo.WithGroup("g"))) + + delay := 400 * time.Millisecond + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "finalizer") + require.NoError(t, err) + require.NoError(t, kh.Put(ctx, *finalizer, + lymbo.AfterGroup("g"), + lymbo.WithDelay(lymbo.FixedDelay(delay)), + )) + + // Complete member — dep deleted, but delay not yet expired. + _, err = store.UpdateBatch(ctx, []lymbo.UpdateSet{terminalUpdate(member.ID, status.Done)}) + require.NoError(t, err) + + // Finalizer not eligible yet: delay gate still active. + tickets := pollStore(ctx, t, store) + assert.False(t, hasID(tickets, finalizer.ID), "finalizer must not be eligible before delay expires") + + // Wait for delay to expire. + time.Sleep(delay + 100*time.Millisecond) + + tickets = pollStore(ctx, t, store) + assert.True(t, hasID(tickets, finalizer.ID), "finalizer must be eligible after delay expires") +} + +// TestFinalizerIsRegularTicket verifies REQ 10: finalizer has no special context — normal ticket shape. +func (s *FinalizerTestSuite) TestFinalizerIsRegularTicket(t *testing.T) { + ctx := context.Background() + store, cleanup := s.factory(t) + defer cleanup() + + kh := lymbo.NewKharon(store, lymbo.DefaultSettings(), nil) + + finalizer, err := lymbo.NewTicket(lymbo.TicketId(uuid.NewString()), "my-finalizer-type") + require.NoError(t, err) + finalizer = finalizer.WithPayload(map[string]string{"result": "ok"}) + // Empty group — no members, immediately eligible. + require.NoError(t, kh.Put(ctx, *finalizer, lymbo.AfterGroup("empty-g"))) + + tickets := pollStore(ctx, t, store) + require.Len(t, tickets, 1) + + got := tickets[0] + assert.Equal(t, finalizer.ID, got.ID) + assert.Equal(t, "my-finalizer-type", got.Type) + assert.NotNil(t, got.Payload) + assert.Equal(t, lymbo.DefaultTube, string(got.Tube)) +} diff --git a/test/go.mod b/test/go.mod index 45c00ba..6a7cfe8 100644 --- a/test/go.mod +++ b/test/go.mod @@ -11,6 +11,7 @@ require ( github.com/stretchr/testify v1.11.1 github.com/testcontainers/testcontainers-go v0.41.0 github.com/testcontainers/testcontainers-go/modules/postgres v0.41.0 + golang.org/x/time v0.15.0 ) require ( @@ -68,7 +69,6 @@ require ( golang.org/x/sync v0.20.0 // indirect golang.org/x/sys v0.41.0 // indirect golang.org/x/text v0.35.0 // indirect - golang.org/x/time v0.15.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260406210006-6f92a3bedf2d // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260401024825-9d38bb4040a9 // indirect google.golang.org/grpc v1.79.3 // indirect diff --git a/test/kharon_memory_test.go b/test/kharon_memory_test.go index 5f6b566..b0059c3 100644 --- a/test/kharon_memory_test.go +++ b/test/kharon_memory_test.go @@ -21,3 +21,9 @@ func TestMemoryStore(t *testing.T) { suite := NewStoreTestSuite(memoryStoreFactory) suite.RunAll(t) } + +// TestMemoryFinalizer runs the finalizer test suite against the memory store +func TestMemoryFinalizer(t *testing.T) { + suite := NewFinalizerTestSuite(memoryStoreFactory) + suite.RunAll(t) +} diff --git a/test/kharon_postgres_test.go b/test/kharon_postgres_test.go index 3f2910c..e34c82a 100644 --- a/test/kharon_postgres_test.go +++ b/test/kharon_postgres_test.go @@ -55,11 +55,14 @@ func TestPostgresStore(t *testing.T) { // Factory reuses the shared pool; truncates the table for test isolation factory := func(t *testing.T) (lymbo.Store, func()) { t.Helper() - _, err := pool.Exec(context.Background(), "TRUNCATE tickets") + _, err := pool.Exec(context.Background(), "TRUNCATE tickets CASCADE") require.NoError(t, err) return store, func() {} } suite := NewStoreTestSuite(factory) suite.RunAll(t) + + finalizerSuite := NewFinalizerTestSuite(factory) + finalizerSuite.RunAll(t) }