From 27b7c300e2f73a96628d84915bd2a104ab24851e Mon Sep 17 00:00:00 2001 From: Mariem Baccari Date: Mon, 5 Oct 2026 12:38:14 +0200 Subject: [PATCH] [evict] Add limit flags --- cmds/db-manager/cleanup/evict.go | 23 +++++++++-- pkg/rid/operations/subscription_test.go | 4 +- pkg/rid/repos/isa.go | 4 +- pkg/rid/repos/subscription.go | 5 ++- .../memstore/identification_service_area.go | 2 +- pkg/rid/store/memstore/subscriptions.go | 2 +- .../raftstore/identification_service_area.go | 2 +- pkg/rid/store/raftstore/subscriptions.go | 2 +- .../sqlstore/identification_service_area.go | 37 ++++++++++++------ pkg/rid/store/sqlstore/subscriptions.go | 38 ++++++++++++------- pkg/scd/repos/repos.go | 8 ++-- pkg/scd/store/memstore/operational_intents.go | 2 +- pkg/scd/store/memstore/subscriptions.go | 2 +- .../store/raftstore/operational_intents.go | 2 +- pkg/scd/store/raftstore/subscriptions.go | 2 +- pkg/scd/store/sqlstore/operational_intents.go | 25 +++++++----- pkg/scd/store/sqlstore/subscriptions.go | 25 +++++++----- pkg/sql/utils.go | 11 ++++++ 18 files changed, 131 insertions(+), 65 deletions(-) diff --git a/cmds/db-manager/cleanup/evict.go b/cmds/db-manager/cleanup/evict.go index 87b837454..1e7737500 100644 --- a/cmds/db-manager/cleanup/evict.go +++ b/cmds/db-manager/cleanup/evict.go @@ -30,6 +30,8 @@ var ( deleteExpired = flags.Bool("delete", false, "set this flag to true to delete the expired entities") locality = flags.String("locality", "", "self-identification string of this DSS instance") timeout = flags.Duration("timeout", 5*time.Minute, "Timeout for the command") + scdLimit = flags.Int("scd_limit", 0, "maximum number of SCD entities deleted, defaults to unlimited") + ridLimit = flags.Int("rid_limit", 0, "maximum number of RID entities deleted, defaults to unlimited") ) func init() { @@ -41,7 +43,20 @@ func evict(cmd *cobra.Command, _ []string) error { ctx = cmd.Context() scdThreshold = time.Now().Add(-*scdTtl) ridThreshold = time.Now().Add(-*ridTtl) + scdLimit = *scdLimit + ridLimit = *ridLimit ) + if scdLimit < 0 { + return fmt.Errorf("scd_limit must be equal to or greater than 0, got %d", scdLimit) + } + if ridLimit < 0 { + return fmt.Errorf("rid_limit must be equal to or greater than 0, got %d", ridLimit) + } + for _, limitFlag := range []string{"scd_limit", "rid_limit"} { + if cmd.Flags().Changed(limitFlag) && !*deleteExpired { + return fmt.Errorf("%s can only be set together with --delete", limitFlag) + } + } log.Printf("WARNING: The usage of this tool may have an impact on performance when deleting entities. Read more in the README.") ctx, cancel := context.WithTimeout(ctx, *timeout) @@ -72,7 +87,7 @@ func evict(cmd *cobra.Command, _ []string) error { if *checkScdOirs { if *deleteExpired { - expiredOpIntents, err = scdRepo.DeleteExpiredOperationalIntents(ctx, scdThreshold) + expiredOpIntents, err = scdRepo.DeleteExpiredOperationalIntents(ctx, scdThreshold, scdLimit) if err != nil { return fmt.Errorf("failed to delete expired operational intents: %w", err) } @@ -86,7 +101,7 @@ func evict(cmd *cobra.Command, _ []string) error { if *checkScdSubs { if *deleteExpired { - scdExpiredSub, err = scdRepo.DeleteExpiredSubscriptions(ctx, scdThreshold) + scdExpiredSub, err = scdRepo.DeleteExpiredSubscriptions(ctx, scdThreshold, scdLimit) if err != nil { return fmt.Errorf("failed to delete expired SCD subscriptions: %w", err) } @@ -105,7 +120,7 @@ func evict(cmd *cobra.Command, _ []string) error { if *checkRidISAs { if *deleteExpired { - expiredISAs, err = ridRepo.DeleteExpiredISAs(ctx, *locality, ridThreshold) + expiredISAs, err = ridRepo.DeleteExpiredISAs(ctx, *locality, ridThreshold, ridLimit) if err != nil { return fmt.Errorf("failed to delete expired ISAs: %w", err) } @@ -119,7 +134,7 @@ func evict(cmd *cobra.Command, _ []string) error { if *checkRidSubs { if *deleteExpired { - ridExpiredSub, err = ridRepo.DeleteExpiredSubscriptions(ctx, *locality, ridThreshold) + ridExpiredSub, err = ridRepo.DeleteExpiredSubscriptions(ctx, *locality, ridThreshold, ridLimit) if err != nil { return fmt.Errorf("failed to delete expired RID subscriptions: %w", err) } diff --git a/pkg/rid/operations/subscription_test.go b/pkg/rid/operations/subscription_test.go index 6fea929c2..563b751e3 100644 --- a/pkg/rid/operations/subscription_test.go +++ b/pkg/rid/operations/subscription_test.go @@ -184,7 +184,7 @@ func (r *fakeSubscriptionRepo) ListExpiredSubscriptions(_ context.Context, _ str return nil, nil } -func (r *fakeSubscriptionRepo) DeleteExpiredSubscriptions(_ context.Context, _ string, _ time.Time) ([]dssmodels.ID, error) { +func (r *fakeSubscriptionRepo) DeleteExpiredSubscriptions(_ context.Context, _ string, _ time.Time, _ int) ([]dssmodels.ID, error) { return nil, nil } @@ -237,7 +237,7 @@ func (r *fakeSubscriptionRepo) ListExpiredISAs(_ context.Context, _ string, _ ti panic("not implemented") } -func (r *fakeSubscriptionRepo) DeleteExpiredISAs(_ context.Context, _ string, _ time.Time) ([]dssmodels.ID, error) { +func (r *fakeSubscriptionRepo) DeleteExpiredISAs(_ context.Context, _ string, _ time.Time, _ int) ([]dssmodels.ID, error) { panic("not implemented") } diff --git a/pkg/rid/repos/isa.go b/pkg/rid/repos/isa.go index 3f2e1d18e..ee6e959d0 100644 --- a/pkg/rid/repos/isa.go +++ b/pkg/rid/repos/isa.go @@ -32,8 +32,8 @@ type ISA interface { // ListExpiredISAs lists the IDs of all expired ISAs based on writer ListExpiredISAs(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) - // DeleteExpiredISAs deletes all expired ISAs based on writer and returns the IDs of the deleted ISAs. - DeleteExpiredISAs(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) + // DeleteExpiredISAs deletes up to `limit` expired ISAs based on writer and returns the IDs of the deleted ISAs. A limit of 0 means unlimited. + DeleteExpiredISAs(ctx context.Context, writer string, threshold time.Time, limit int) ([]dssmodels.ID, error) // Count the number of existing ISA CountISAs(ctx context.Context) (int64, error) diff --git a/pkg/rid/repos/subscription.go b/pkg/rid/repos/subscription.go index 8a6e241f8..3e102c28c 100644 --- a/pkg/rid/repos/subscription.go +++ b/pkg/rid/repos/subscription.go @@ -42,8 +42,9 @@ type Subscription interface { // ListExpiredSubscriptions lists the IDs of all expired Subscriptions based on writer. ListExpiredSubscriptions(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) - // DeleteExpiredSubscriptions deletes all expired Subscriptions based on writer and returns the IDs of the deleted Subscriptions. - DeleteExpiredSubscriptions(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) + // DeleteExpiredSubscriptions deletes up to `limit` expired Subscriptions based on writer and + // returns the IDs of the deleted Subscriptions. A limit of 0 means unlimited. + DeleteExpiredSubscriptions(ctx context.Context, writer string, threshold time.Time, limit int) ([]dssmodels.ID, error) // Count the number of existing subscriptions CountSubscriptions(ctx context.Context) (int64, error) diff --git a/pkg/rid/store/memstore/identification_service_area.go b/pkg/rid/store/memstore/identification_service_area.go index 0f59f9887..b1f2e775d 100644 --- a/pkg/rid/store/memstore/identification_service_area.go +++ b/pkg/rid/store/memstore/identification_service_area.go @@ -135,7 +135,7 @@ func (r *repo) ListExpiredISAs(_ context.Context, writer string, threshold time. } // TODO: Implement when raftstore evict is implemented (#1718) -func (r *repo) DeleteExpiredISAs(_ context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredISAs(_ context.Context, writer string, threshold time.Time, limit int) ([]dssmodels.ID, error) { return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "DeleteExpiredISAs not implemented for memstore") } diff --git a/pkg/rid/store/memstore/subscriptions.go b/pkg/rid/store/memstore/subscriptions.go index 8f2c1b497..d10fe6e84 100644 --- a/pkg/rid/store/memstore/subscriptions.go +++ b/pkg/rid/store/memstore/subscriptions.go @@ -192,7 +192,7 @@ func (r *repo) ListExpiredSubscriptions(_ context.Context, writer string, thresh } // TODO: Implement when raftstore evict is implemented (#1718) -func (r *repo) DeleteExpiredSubscriptions(_ context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredSubscriptions(_ context.Context, writer string, threshold time.Time, limit int) ([]dssmodels.ID, error) { return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "DeleteExpiredSubscriptions not implemented for memstore") } diff --git a/pkg/rid/store/raftstore/identification_service_area.go b/pkg/rid/store/raftstore/identification_service_area.go index db96d35fe..24506da0f 100644 --- a/pkg/rid/store/raftstore/identification_service_area.go +++ b/pkg/rid/store/raftstore/identification_service_area.go @@ -79,7 +79,7 @@ func (r *repo) ListExpiredISAs(ctx context.Context, writer string, threshold tim return r.consensus.HandleReadRequest(ctx, listExpiredISAs, buf) } -func (r *repo) DeleteExpiredISAs(_ context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredISAs(_ context.Context, writer string, threshold time.Time, limit int) ([]dssmodels.ID, error) { return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "DeleteExpiredISAs not implemented for raftstore") } diff --git a/pkg/rid/store/raftstore/subscriptions.go b/pkg/rid/store/raftstore/subscriptions.go index 5bdaf6669..2e767c776 100644 --- a/pkg/rid/store/raftstore/subscriptions.go +++ b/pkg/rid/store/raftstore/subscriptions.go @@ -107,7 +107,7 @@ func (r *repo) ListExpiredSubscriptions(ctx context.Context, writer string, thre return r.consensus.HandleReadRequest(ctx, listExpiredSubscriptions, buf) } -func (r *repo) DeleteExpiredSubscriptions(_ context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredSubscriptions(_ context.Context, writer string, threshold time.Time, limit int) ([]dssmodels.ID, error) { return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "DeleteExpiredSubscriptions not implemented for raftstore") } diff --git a/pkg/rid/store/sqlstore/identification_service_area.go b/pkg/rid/store/sqlstore/identification_service_area.go index 7fda71ad4..41e1bf95c 100644 --- a/pkg/rid/store/sqlstore/identification_service_area.go +++ b/pkg/rid/store/sqlstore/identification_service_area.go @@ -238,24 +238,37 @@ func (r *repo) ListExpiredISAs(ctx context.Context, writer string, threshold tim return dssql.FetchIDs(ctx, r.Queryable, isasInCellsQuery, threshold, writer) } -// DeleteExpiredISAs deletes all expired ISAs based on writer and returns the IDs of the deleted ISAs. +// DeleteExpiredISAs deletes up to `limit` expired ISAs based on writer and returns the +// IDs of the deleted ISAs. A limit of 0 means unlimited. // The function deletes both empty writer and null writer when passing empty string as a writer. -func (r *repo) DeleteExpiredISAs(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredISAs(ctx context.Context, writer string, threshold time.Time, limit int) ([]dssmodels.ID, error) { if len(writer) == 0 { - deleteExpiredQuery := ` + expiredQuery, args := dssql.AppendLimitClause(` + SELECT id + FROM identification_service_areas + WHERE ends_at <= $1 + AND (writer = '' OR writer IS NULL)`, []any{threshold}, limit) + deleteExpiredQuery := fmt.Sprintf(` + WITH expired AS (%s + ) DELETE FROM identification_service_areas - WHERE ends_at <= $1 - AND (writer = '' OR writer IS NULL) - RETURNING id` - return dssql.FetchIDs(ctx, r.Queryable, deleteExpiredQuery, threshold) + WHERE id IN (SELECT id FROM expired) + RETURNING id`, expiredQuery) + return dssql.FetchIDs(ctx, r.Queryable, deleteExpiredQuery, args...) } - deleteExpiredQuery := ` + expiredQuery, args := dssql.AppendLimitClause(` + SELECT id + FROM identification_service_areas + WHERE ends_at <= $1 + AND writer = $2`, []any{threshold, writer}, limit) + deleteExpiredQuery := fmt.Sprintf(` + WITH expired AS (%s + ) DELETE FROM identification_service_areas - WHERE ends_at <= $1 - AND writer = $2 - RETURNING id` - return dssql.FetchIDs(ctx, r.Queryable, deleteExpiredQuery, threshold, writer) + WHERE id IN (SELECT id FROM expired) + RETURNING id`, expiredQuery) + return dssql.FetchIDs(ctx, r.Queryable, deleteExpiredQuery, args...) } func (r *repo) CountISAs(ctx context.Context) (int64, error) { diff --git a/pkg/rid/store/sqlstore/subscriptions.go b/pkg/rid/store/sqlstore/subscriptions.go index 8e3ab3dfb..a0429c4e3 100644 --- a/pkg/rid/store/sqlstore/subscriptions.go +++ b/pkg/rid/store/sqlstore/subscriptions.go @@ -322,25 +322,37 @@ func (r *repo) ListExpiredSubscriptions(ctx context.Context, writer string, thre return dssql.FetchIDs(ctx, r.Queryable, query, threshold, writer) } -// DeleteExpiredSubscriptions deletes all expired Subscriptions based on writer and returns the -// IDs of the deleted Subscriptions. +// DeleteExpiredSubscriptions deletes up to `limit` expired Subscriptions based on writer and returns +// the IDs of the deleted Subscriptions. A limit of 0 means unlimited. // The function deletes both empty writer and null writer when passing empty string as a writer. -func (r *repo) DeleteExpiredSubscriptions(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredSubscriptions(ctx context.Context, writer string, threshold time.Time, limit int) ([]dssmodels.ID, error) { if len(writer) == 0 { - deleteExpiredQuery := ` + expiredQuery, args := dssql.AppendLimitClause(` + SELECT id + FROM subscriptions + WHERE ends_at <= $1 + AND (writer = '' OR writer IS NULL)`, []any{threshold}, limit) + deleteExpiredQuery := fmt.Sprintf(` + WITH expired AS (%s + ) DELETE FROM subscriptions - WHERE ends_at <= $1 - AND (writer = '' OR writer IS NULL) - RETURNING id` - return dssql.FetchIDs(ctx, r.Queryable, deleteExpiredQuery, threshold) + WHERE id IN (SELECT id FROM expired) + RETURNING id`, expiredQuery) + return dssql.FetchIDs(ctx, r.Queryable, deleteExpiredQuery, args...) } - deleteExpiredQuery := ` + expiredQuery, args := dssql.AppendLimitClause(` + SELECT id + FROM subscriptions + WHERE ends_at <= $1 + AND writer = $2`, []any{threshold, writer}, limit) + deleteExpiredQuery := fmt.Sprintf(` + WITH expired AS (%s + ) DELETE FROM subscriptions - WHERE ends_at <= $1 - AND writer = $2 - RETURNING id` - return dssql.FetchIDs(ctx, r.Queryable, deleteExpiredQuery, threshold, writer) + WHERE id IN (SELECT id FROM expired) + RETURNING id`, expiredQuery) + return dssql.FetchIDs(ctx, r.Queryable, deleteExpiredQuery, args...) } func (r *repo) CountSubscriptions(ctx context.Context) (int64, error) { diff --git a/pkg/scd/repos/repos.go b/pkg/scd/repos/repos.go index 06383d133..e02e6f165 100644 --- a/pkg/scd/repos/repos.go +++ b/pkg/scd/repos/repos.go @@ -34,9 +34,9 @@ type OperationalIntent interface { // Their age is determined by their end time, or by their update time if they do not have an end time. ListExpiredOperationalIntents(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) - // DeleteExpiredOperationalIntents deletes all expired operational intents and returns the IDs of the deleted operational intents. + // DeleteExpiredOperationalIntents deletes up to `limit` expired operational intents and returns the IDs of the deleted operational intents. A limit of 0 means unlimited. // Age is determined by their end time, or by their update time if they do not have an end time. - DeleteExpiredOperationalIntents(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) + DeleteExpiredOperationalIntents(ctx context.Context, threshold time.Time, limit int) ([]dssmodels.ID, error) // Count the number of existing operational intent CountOperationalIntents(ctx context.Context) (int64, error) @@ -77,9 +77,9 @@ type Subscription interface { // Their age is determined by their end time, or by their update time if they do not have an end time. ListExpiredSubscriptions(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) - // DeleteExpiredSubscriptions deletes all expired subscriptions and returns the IDs of the deleted subscriptions. + // DeleteExpiredSubscriptions deletes up to `limit` expired subscriptions and returns the IDs of the deleted subscriptions. A limit of 0 means unlimited. // Age is determined by their end time, or by their update time if they do not have an end time. - DeleteExpiredSubscriptions(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) + DeleteExpiredSubscriptions(ctx context.Context, threshold time.Time, limit int) ([]dssmodels.ID, error) // Count the number of existing subscriptions CountSubscriptions(ctx context.Context) (int64, error) diff --git a/pkg/scd/store/memstore/operational_intents.go b/pkg/scd/store/memstore/operational_intents.go index 1ef606f4e..b27c0c216 100644 --- a/pkg/scd/store/memstore/operational_intents.go +++ b/pkg/scd/store/memstore/operational_intents.go @@ -168,7 +168,7 @@ func (r *repo) ListExpiredOperationalIntents(ctx context.Context, threshold time } // TODO: Implement when raftstore evict is implemented (#1718) -func (r *repo) DeleteExpiredOperationalIntents(_ context.Context, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredOperationalIntents(_ context.Context, threshold time.Time, limit int) ([]dssmodels.ID, error) { return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "DeleteExpiredOperationalIntents not implemented for memstore") } diff --git a/pkg/scd/store/memstore/subscriptions.go b/pkg/scd/store/memstore/subscriptions.go index 7d20c5348..7d0bc0517 100644 --- a/pkg/scd/store/memstore/subscriptions.go +++ b/pkg/scd/store/memstore/subscriptions.go @@ -141,7 +141,7 @@ func (r *repo) ListExpiredSubscriptions(_ context.Context, threshold time.Time) } // TODO: Implement when raftstore evict is implemented (#1718) -func (r *repo) DeleteExpiredSubscriptions(_ context.Context, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredSubscriptions(_ context.Context, threshold time.Time, limit int) ([]dssmodels.ID, error) { return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "DeleteExpiredSubscriptions not implemented for memstore") } diff --git a/pkg/scd/store/raftstore/operational_intents.go b/pkg/scd/store/raftstore/operational_intents.go index 140d2ca9e..97462276b 100644 --- a/pkg/scd/store/raftstore/operational_intents.go +++ b/pkg/scd/store/raftstore/operational_intents.go @@ -77,7 +77,7 @@ func (r *repo) ListExpiredOperationalIntents(ctx context.Context, threshold time return r.consensus.HandleReadRequest(ctx, listExpiredOperationalIntents, buf) } -func (r *repo) DeleteExpiredOperationalIntents(_ context.Context, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredOperationalIntents(_ context.Context, threshold time.Time, limit int) ([]dssmodels.ID, error) { return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "DeleteExpiredOperationalIntents not implemented for raftstore") } diff --git a/pkg/scd/store/raftstore/subscriptions.go b/pkg/scd/store/raftstore/subscriptions.go index c8dcc4f6e..9b5ce2706 100644 --- a/pkg/scd/store/raftstore/subscriptions.go +++ b/pkg/scd/store/raftstore/subscriptions.go @@ -93,7 +93,7 @@ func (r *repo) ListExpiredSubscriptions(ctx context.Context, threshold time.Time return r.consensus.HandleReadRequest(ctx, listExpiredSubscriptions, buf) } -func (r *repo) DeleteExpiredSubscriptions(_ context.Context, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredSubscriptions(_ context.Context, threshold time.Time, limit int) ([]dssmodels.ID, error) { return nil, stacktrace.NewErrorWithCode(dsserr.NotImplemented, "DeleteExpiredSubscriptions not implemented for raftstore") } diff --git a/pkg/scd/store/sqlstore/operational_intents.go b/pkg/scd/store/sqlstore/operational_intents.go index a1f5d73b6..db71806cf 100644 --- a/pkg/scd/store/sqlstore/operational_intents.go +++ b/pkg/scd/store/sqlstore/operational_intents.go @@ -371,18 +371,25 @@ func (s *repo) ListExpiredOperationalIntents(ctx context.Context, threshold time return ids, nil } -// DeleteExpiredOperationalIntents deletes all expired operational intents and returns the IDs of the deleted operational intents. +// DeleteExpiredOperationalIntents deletes up to `limit` expired operational intents and returns the IDs of the deleted operational intents. A limit of 0 means unlimited. // Age is determined by their end time, or by their update time if they do not have an end time. -func (s *repo) DeleteExpiredOperationalIntents(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) { - deleteExpiredQuery := ` +func (s *repo) DeleteExpiredOperationalIntents(ctx context.Context, threshold time.Time, limit int) ([]dssmodels.ID, error) { + expiredQuery, args := dsssql.AppendLimitClause(` + SELECT id + FROM scd_operations + WHERE + (ends_at IS NOT NULL AND ends_at <= $1) + OR + (ends_at IS NULL AND updated_at <= $1)`, []any{threshold}, limit) + + deleteExpiredQuery := fmt.Sprintf(` + WITH expired AS (%s + ) DELETE FROM scd_operations - WHERE - (ends_at IS NOT NULL AND ends_at <= $1) - OR - (ends_at IS NULL AND updated_at <= $1) - RETURNING id` + WHERE id IN (SELECT id FROM expired) + RETURNING id`, expiredQuery) - ids, err := dsssql.FetchIDs(ctx, s.q, deleteExpiredQuery, threshold) + ids, err := dsssql.FetchIDs(ctx, s.q, deleteExpiredQuery, args...) if err != nil { return nil, stacktrace.Propagate(err, "Error deleting expired Operations") } diff --git a/pkg/scd/store/sqlstore/subscriptions.go b/pkg/scd/store/sqlstore/subscriptions.go index 9611dc3b9..924f06163 100644 --- a/pkg/scd/store/sqlstore/subscriptions.go +++ b/pkg/scd/store/sqlstore/subscriptions.go @@ -576,18 +576,25 @@ func (c *repo) ListExpiredSubscriptions(ctx context.Context, threshold time.Time return ids, nil } -// DeleteExpiredSubscriptions deletes all expired subscriptions and returns the IDs of the deleted subscriptions. +// DeleteExpiredSubscriptions deletes up to `limit` expired subscriptions and returns the IDs of the deleted subscriptions. A limit of 0 means unlimited. // Age is determined by their end time, or by their update time if they do not have an end time. -func (c *repo) DeleteExpiredSubscriptions(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) { - deleteExpiredQuery := ` +func (c *repo) DeleteExpiredSubscriptions(ctx context.Context, threshold time.Time, limit int) ([]dssmodels.ID, error) { + expiredQuery, args := dsssql.AppendLimitClause(` + SELECT id + FROM scd_subscriptions + WHERE + (ends_at IS NOT NULL AND ends_at <= $1) + OR + (ends_at IS NULL AND updated_at <= $1)`, []any{threshold}, limit) + + deleteExpiredQuery := fmt.Sprintf(` + WITH expired AS (%s + ) DELETE FROM scd_subscriptions - WHERE - (ends_at IS NOT NULL AND ends_at <= $1) - OR - (ends_at IS NULL AND updated_at <= $1) - RETURNING id` + WHERE id IN (SELECT id FROM expired) + RETURNING id`, expiredQuery) - ids, err := dsssql.FetchIDs(ctx, c.q, deleteExpiredQuery, threshold) + ids, err := dsssql.FetchIDs(ctx, c.q, deleteExpiredQuery, args...) if err != nil { return nil, stacktrace.Propagate(err, "Unable to delete expired Subscriptions") } diff --git a/pkg/sql/utils.go b/pkg/sql/utils.go index 6ce998e0e..b1d330a3e 100644 --- a/pkg/sql/utils.go +++ b/pkg/sql/utils.go @@ -1,6 +1,7 @@ package sql import ( + "fmt" "slices" "time" @@ -41,3 +42,13 @@ func MillisSinceMidnight() int { midnight := time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, time.UTC) return int(now.Sub(midnight).Milliseconds()) } + +// AppendLimitClause appends a "LIMIT $n" clause to query and the limit to args when limit is +// positive, n being the next parameter index after args. query and args are returned unchanged +// when limit <= 0 (unlimited). +func AppendLimitClause(query string, args []any, limit int) (string, []any) { + if limit <= 0 { + return query, args + } + return fmt.Sprintf("%s\nLIMIT $%d", query, len(args)+1), append(args, limit) +}