diff --git a/NEXT_RELEASE_NOTES.md b/NEXT_RELEASE_NOTES.md index 31cb2ca24..9df11430a 100644 --- a/NEXT_RELEASE_NOTES.md +++ b/NEXT_RELEASE_NOTES.md @@ -33,13 +33,13 @@ The release notes should contain at least the following sections: #### Minimal database schema version -| Schema | CockroachDB | Yugabyte | -|---------|--------------|----------| -| RID | vX.Y.Z | vX.Y.Z | -| SCD | vX.Y.Z | vX.Y.Z | -| AUX | vX.Y.Z | vX.Y.Z | +| Schema | CockroachDB | Yugabyte | +| ------ | ----------- | -------- | +| RID | vX.Y.Z | vX.Y.Z | +| SCD | vX.Y.Z | vX.Y.Z | +| AUX | vX.Y.Z | vX.Y.Z | --------------------------------------------------------------------------------------------------------------------- +--- # Release Notes for vX.Y.Z @@ -50,11 +50,12 @@ The release notes should contain at least the following sections: ## Important information * [deploy] AWS NLB TLS listeners (DSS gateway and Prometheus, for both Helm and Tanka) now use the `ELBSecurityPolicy-TLS13-1-2-Res-2021-06` SSL negotiation policy, which enforces a minimum of TLS 1.2 (aligning with the Google deployment) and restricts TLS 1.2 to AEAD (GCM/ChaCha20) cipher suites. Previously the AWS Load Balancer Controller default (`ELBSecurityPolicy-2016-08`) applied, which allowed TLS 1.0 and 1.1. Clients limited to TLS 1.0 or 1.1, or to CBC cipher suites, will no longer be able to connect. Unlike the Google deployment, no resources need to be provisioned by Terraform or manually to set up the policy. +* [evict] RID cleanup (`db-manager evict`, including the RID cleanup default `CronJob`) now evicts the expired ISAs and subscriptions of all the DSS instances of the pool, instead of only the ones written by the DSS instance identified by `--locality`. The `--locality` flag of `db-manager evict` is deprecated and has no effect. It will be removed in a future release. ## Minimal database schema version -| Schema | CockroachDB | Yugabyte | -|---------|-------------|----------| -| AUX | | | -| RID | | | -| SCD | | | +| Schema | CockroachDB | Yugabyte | +| ------ | ----------- | -------- | +| AUX | | | +| RID | | | +| SCD | | | diff --git a/cmds/db-manager/cleanup/evict.go b/cmds/db-manager/cleanup/evict.go index 4bb446454..2c82cc111 100644 --- a/cmds/db-manager/cleanup/evict.go +++ b/cmds/db-manager/cleanup/evict.go @@ -28,13 +28,16 @@ var ( scdTtl = flags.Duration("scd_ttl", time.Hour*24*112, "time-to-live duration used for determining SCD entries expiration, defaults to 2*56 days") ridTtl = flags.Duration("rid_ttl", time.Minute*30, "time-to-live duration used for determining RID entries expiration, defaults to 30 minutes") 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") + _ = 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() { + if err := flags.MarkDeprecated("locality", "it has no effect anymore, expired RID entities of all DSS instances are evicted"); err != nil { + panic(err) + } EvictCmd.Flags().AddFlagSet(flags) } @@ -120,12 +123,12 @@ func evict(cmd *cobra.Command, _ []string) error { if *checkRidISAs { if *deleteExpired { - expiredISAs, err = ridRepo.DeleteExpiredISAs(ctx, *locality, ridThreshold, ridLimit) + expiredISAs, err = ridRepo.DeleteExpiredISAs(ctx, ridThreshold, ridLimit) if err != nil { return fmt.Errorf("failed to delete expired ISAs: %w", err) } } else { - expiredISAs, err = ridRepo.ListExpiredISAs(ctx, *locality, ridThreshold) + expiredISAs, err = ridRepo.ListExpiredISAs(ctx, ridThreshold) if err != nil { return fmt.Errorf("failed to list expired ISAs: %w", err) } @@ -134,12 +137,12 @@ func evict(cmd *cobra.Command, _ []string) error { if *checkRidSubs { if *deleteExpired { - ridExpiredSub, err = ridRepo.DeleteExpiredSubscriptions(ctx, *locality, ridThreshold, ridLimit) + ridExpiredSub, err = ridRepo.DeleteExpiredSubscriptions(ctx, ridThreshold, ridLimit) if err != nil { return fmt.Errorf("failed to delete expired RID subscriptions: %w", err) } } else { - ridExpiredSub, err = ridRepo.ListExpiredSubscriptions(ctx, *locality, ridThreshold) + ridExpiredSub, err = ridRepo.ListExpiredSubscriptions(ctx, ridThreshold) if err != nil { return fmt.Errorf("failed to list RID expired subscriptions: %w", err) } diff --git a/deploy/services/helm-charts/dss/templates/dss-evict.yaml b/deploy/services/helm-charts/dss/templates/dss-evict.yaml index 618c6ab4b..3ad5fe3b7 100644 --- a/deploy/services/helm-charts/dss/templates/dss-evict.yaml +++ b/deploy/services/helm-charts/dss/templates/dss-evict.yaml @@ -45,7 +45,6 @@ spec: {{- if $.Values.dss.conf.evict.scd.limit }} - --scd_limit={{ $.Values.dss.conf.evict.scd.limit }} {{- end }} - - --locality={{ $.Values.dss.conf.locality }} - --delete=true - --datastore_host={{ $datastoreHost }} - --datastore_port={{ $datastorePort }} @@ -107,7 +106,6 @@ spec: {{- if $.Values.dss.conf.evict.rid.limit }} - --rid_limit={{ $.Values.dss.conf.evict.rid.limit }} {{- end }} - - --locality={{ $.Values.dss.conf.locality }} - --delete=true - --datastore_host={{ $datastoreHost }} - --datastore_port={{ $datastorePort }} diff --git a/deploy/services/tanka/evict.libsonnet b/deploy/services/tanka/evict.libsonnet index 5e4f02923..ba6eeceb5 100644 --- a/deploy/services/tanka/evict.libsonnet +++ b/deploy/services/tanka/evict.libsonnet @@ -29,7 +29,6 @@ local datastoreparameters = import 'datastoreparameters.libsonnet'; rid_isa: false, rid_sub: false, scd_ttl: metadata.evict.scd.ttl, - locality: metadata.locality, delete: true, } + (if metadata.evict.scd.timeout != "" then { timeout: metadata.evict.scd.timeout } else {}) + (if metadata.evict.scd.limit > 0 then { scd_limit: metadata.evict.scd.limit } else {}) @@ -67,7 +66,6 @@ local datastoreparameters = import 'datastoreparameters.libsonnet'; rid_isa: metadata.evict.rid.ISAs, rid_sub: metadata.evict.rid.subscriptions, rid_ttl: metadata.evict.rid.ttl, - locality: metadata.locality, delete: true, } + (if metadata.evict.rid.timeout != "" then { timeout: metadata.evict.rid.timeout } else {}) + (if metadata.evict.rid.limit > 0 then { rid_limit: metadata.evict.rid.limit } else {}) diff --git a/docs/operations/cleanup.md b/docs/operations/cleanup.md index bba11f8f4..85a5ec1fb 100644 --- a/docs/operations/cleanup.md +++ b/docs/operations/cleanup.md @@ -47,12 +47,6 @@ To mitigate contention: - Clean up iteratively, starting with a lower TTL and progressively increasing it. - Use `--rid_limit`/`--scd_limit` to bound the size of each deletion if a batch still causes performance issues. -## Changes in locality - -There may be cases where a DSS instance changes its locality. This commonly occurs if locality wasn't previously required, though routine updates can also trigger a change. - -In such cases, ensure that a cleanup is performed on the older locality (if it was not set before, set an empty one), especially since the automatically deployed cron job (see below) will automatically follow the locality settings of the main service. - ## Usage Extracted from `db-manager evict --help`: @@ -66,7 +60,6 @@ Usage: Flags: --delete set this flag to true to delete the expired entities -h, --help help for evict - --locality string self-identification string of this DSS instance --rid_isa set this flag to true to check for expired RID ISAs (default true) --rid_limit int maximum number of RID entities deleted, defaults to unlimited --rid_sub set this flag to true to check for expired RID subscriptions (default true) @@ -93,6 +86,7 @@ Notes: - By default, expired entities are only listed - `--delete` is required to actually remove them. - `--rid_limit`/`--scd_limit` can only be used together with `--delete` and bound the number of entities deleted by a run, per entity type (e.g. `--scd_limit=1000` deletes up to 1000 operational intents and up to 1000 SCD subscriptions). Which expired entities are deleted first is unspecified. Run the command again to delete the remaining ones. A run without `--delete` always lists all expired entities. +- `--locality` is deprecated and has no effect: RID cleanup considers the expired ISAs and subscriptions of all the DSS instances of the pool. - `--rid_ttl` and `--scd_ttl` accept durations formatted as [Go `time.Duration` strings](https://pkg.go.dev/time#ParseDuration), e.g. `24h`. - `--timeout` accepts the same duration format and bounds the total execution time of the command. - The datastore connection flags match those of the `core-service` command. @@ -103,7 +97,7 @@ Beyond running `db-manager evict` manually, the DSS deployment tooling can sched Shared default: RID cleanup is enabled by default (`*/30 * * * *`, `ttl = 30m`); SCD cleanup is disabled by default (suggested schedule `0 2 * * *`, `ttl = 2688h` - i.e. 2 x 56 days). -The current defaults are structured this way because each DSS instance is intended to manage the cleanup of its own RID objects (with deletion restricted to entities created by that specific instance, based on locality). Conversely, SCD cleanup is a global operation; therefore, DSS operators should coordinate to avoid redundant tasks, or potentially designate a single USS to handle the cleanup process. +RID and SCD cleanups are both global operations: they consider the expired entities of all the DSS instances of the pool, regardless of the instance that created them. DSS operators should therefore coordinate to avoid redundant tasks, or potentially designate a single USS to handle the cleanup process. ### Helm @@ -158,8 +152,8 @@ evict+: { When deploying via terrafrom modules, the parameters are configurable with module variables: -| Terraform variable | Default | -|---------------------------------|------------------| +| Terraform variable | Default | +| --------------------------------- | ------------------ | | `evict_enable_scd_cron` | `false` | | `evict_scd_schedule` | `"0 2 * * *"` | | `evict_scd_ttl` | `"2688h"` | diff --git a/pkg/rid/operations/subscription_test.go b/pkg/rid/operations/subscription_test.go index 252f82dc4..b2a9e5c5d 100644 --- a/pkg/rid/operations/subscription_test.go +++ b/pkg/rid/operations/subscription_test.go @@ -178,11 +178,11 @@ func (r *fakeSubscriptionRepo) MaxSubscriptionCountInCellsByOwner(ctx context.Co return maxValue, nil } -func (r *fakeSubscriptionRepo) ListExpiredSubscriptions(_ context.Context, _ string, _ time.Time) ([]dssmodels.ID, error) { +func (r *fakeSubscriptionRepo) ListExpiredSubscriptions(_ context.Context, _ time.Time) ([]dssmodels.ID, error) { return nil, nil } -func (r *fakeSubscriptionRepo) DeleteExpiredSubscriptions(_ context.Context, _ string, _ time.Time, _ int) ([]dssmodels.ID, error) { +func (r *fakeSubscriptionRepo) DeleteExpiredSubscriptions(_ context.Context, _ time.Time, _ int) ([]dssmodels.ID, error) { return nil, nil } @@ -231,11 +231,11 @@ func (r *fakeSubscriptionRepo) SearchISAs(_ context.Context, cells s2.CellUnion, return isas, nil } -func (r *fakeSubscriptionRepo) ListExpiredISAs(_ context.Context, _ string, _ time.Time) ([]dssmodels.ID, error) { +func (r *fakeSubscriptionRepo) ListExpiredISAs(_ context.Context, _ time.Time) ([]dssmodels.ID, error) { panic("not implemented") } -func (r *fakeSubscriptionRepo) DeleteExpiredISAs(_ context.Context, _ string, _ time.Time, _ int) ([]dssmodels.ID, error) { +func (r *fakeSubscriptionRepo) DeleteExpiredISAs(_ context.Context, _ time.Time, _ int) ([]dssmodels.ID, error) { panic("not implemented") } diff --git a/pkg/rid/repos/isa.go b/pkg/rid/repos/isa.go index ee6e959d0..6a0c18e39 100644 --- a/pkg/rid/repos/isa.go +++ b/pkg/rid/repos/isa.go @@ -29,11 +29,11 @@ type ISA interface { // SearchISAs returns all subscriptions ownded by "owner" in "cells". SearchISAs(ctx context.Context, cells s2.CellUnion, earliest *time.Time, latest *time.Time) ([]*ridmodels.IdentificationServiceArea, error) - // ListExpiredISAs lists the IDs of all expired ISAs based on writer - ListExpiredISAs(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) + // ListExpiredISAs lists the IDs of all expired ISAs. + ListExpiredISAs(ctx context.Context, 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) + // DeleteExpiredISAs deletes up to `limit` expired ISAs and returns the IDs of the deleted ISAs. A limit of 0 means unlimited. + DeleteExpiredISAs(ctx context.Context, 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 3e102c28c..b717523b0 100644 --- a/pkg/rid/repos/subscription.go +++ b/pkg/rid/repos/subscription.go @@ -39,12 +39,12 @@ type Subscription interface { // belonging to the given owner, and returns that number. MaxSubscriptionCountInCellsByOwner(ctx context.Context, cells s2.CellUnion, owner dssmodels.Owner) (int, error) - // ListExpiredSubscriptions lists the IDs of all expired Subscriptions based on writer. - ListExpiredSubscriptions(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) + // ListExpiredSubscriptions lists the IDs of all expired Subscriptions. + ListExpiredSubscriptions(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) - // DeleteExpiredSubscriptions deletes up to `limit` expired Subscriptions based on writer and + // DeleteExpiredSubscriptions deletes up to `limit` expired Subscriptions 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) + 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/rid/store/memstore/identification_service_area.go b/pkg/rid/store/memstore/identification_service_area.go index b1f2e775d..8b9b4632e 100644 --- a/pkg/rid/store/memstore/identification_service_area.go +++ b/pkg/rid/store/memstore/identification_service_area.go @@ -130,12 +130,12 @@ func (r *repo) SearchISAs(_ context.Context, cells s2.CellUnion, earliest *time. } // TODO: Implement when raftstore evict is implemented (#1718) -func (r *repo) ListExpiredISAs(_ context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { - return listExpired(r.state.ISAs, writer, threshold, dssmodels.MaxResultLimit), nil +func (r *repo) ListExpiredISAs(_ context.Context, threshold time.Time) ([]dssmodels.ID, error) { + return listExpired(r.state.ISAs, threshold, dssmodels.MaxResultLimit), nil } // TODO: Implement when raftstore evict is implemented (#1718) -func (r *repo) DeleteExpiredISAs(_ context.Context, writer string, threshold time.Time, limit int) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredISAs(_ context.Context, 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/identification_service_area_test.go b/pkg/rid/store/memstore/identification_service_area_test.go index 2cd3c261b..c91f922b7 100644 --- a/pkg/rid/store/memstore/identification_service_area_test.go +++ b/pkg/rid/store/memstore/identification_service_area_test.go @@ -258,7 +258,7 @@ func TestListExpiredISAs(t *testing.T) { require.NoError(t, err) require.NotNil(t, saOut2) - serviceAreas, err := repo.ListExpiredISAs(ctx, writer, fakeClock.Now().Add(-30*time.Minute)) + serviceAreas, err := repo.ListExpiredISAs(ctx, fakeClock.Now().Add(-30*time.Minute)) require.NoError(t, err) require.Len(t, serviceAreas, 1) } @@ -291,7 +291,7 @@ func TestListExpiredISAsWithEmptyWriter(t *testing.T) { require.NoError(t, err) require.NotNil(t, saOut2) - serviceAreas, err := repo.ListExpiredISAs(ctx, "", fakeClock.Now().Add(-30*time.Minute)) + serviceAreas, err := repo.ListExpiredISAs(ctx, fakeClock.Now().Add(-30*time.Minute)) require.NoError(t, err) require.Len(t, serviceAreas, 1) } diff --git a/pkg/rid/store/memstore/store.go b/pkg/rid/store/memstore/store.go index 381588e7e..ea3a80238 100644 --- a/pkg/rid/store/memstore/store.go +++ b/pkg/rid/store/memstore/store.go @@ -141,28 +141,20 @@ func findForWrite[R versionedRecord](store map[dssmodels.ID]R, id dssmodels.ID, type expiringRecord interface { endTime() *time.Time - writerName() string } func (rec *isaRecord) endTime() *time.Time { return rec.EndTime } -func (rec *isaRecord) writerName() string { return rec.Writer } - func (rec *subscriptionRecord) endTime() *time.Time { return rec.EndTime } -func (rec *subscriptionRecord) writerName() string { return rec.Writer } - -// listExpired returns the IDs of the records whose end time is at or before threshold and -// whose writer matches. A limit of 0 means unlimited. -func listExpired[R expiringRecord](store map[dssmodels.ID]R, writer string, threshold time.Time, limit int) []dssmodels.ID { +// listExpired returns the IDs of the records whose end time is at or before threshold. +// A limit of 0 means unlimited. +func listExpired[R expiringRecord](store map[dssmodels.ID]R, threshold time.Time, limit int) []dssmodels.ID { var out []dssmodels.ID for id, rec := range store { if t := rec.endTime(); t == nil || t.After(threshold) { // TODO: Don't allow endtime to be null, see #1492 continue } - if rec.writerName() != writer { - continue - } out = append(out, id) if limit > 0 && len(out) > limit { // This mimics sqlstore behaviour, but it's not very good. diff --git a/pkg/rid/store/memstore/subscriptions.go b/pkg/rid/store/memstore/subscriptions.go index d10fe6e84..f742790af 100644 --- a/pkg/rid/store/memstore/subscriptions.go +++ b/pkg/rid/store/memstore/subscriptions.go @@ -186,13 +186,13 @@ func (r *repo) MaxSubscriptionCountInCellsByOwner(ctx context.Context, cells s2. } // TODO: Implement when raftstore evict is implemented (#1718) -func (r *repo) ListExpiredSubscriptions(_ context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { +func (r *repo) ListExpiredSubscriptions(_ context.Context, threshold time.Time) ([]dssmodels.ID, error) { // TODO: This mimics sqlstore inconsistency of not limiting results there, compared to ISAs. Should it be normalized? - return listExpired(r.state.Subscriptions, writer, threshold, 0), nil + return listExpired(r.state.Subscriptions, threshold, 0), nil } // TODO: Implement when raftstore evict is implemented (#1718) -func (r *repo) DeleteExpiredSubscriptions(_ context.Context, writer string, threshold time.Time, limit int) ([]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/rid/store/memstore/subscriptions_test.go b/pkg/rid/store/memstore/subscriptions_test.go index c1a191a73..2c588dcdc 100644 --- a/pkg/rid/store/memstore/subscriptions_test.go +++ b/pkg/rid/store/memstore/subscriptions_test.go @@ -312,7 +312,7 @@ func TestListExpiredSubscriptions(t *testing.T) { require.NoError(t, err) require.NotNil(t, subOut2) - subscriptions, err := repo.ListExpiredSubscriptions(ctx, writer, fakeClock.Now().Add(-30*time.Minute)) + subscriptions, err := repo.ListExpiredSubscriptions(ctx, fakeClock.Now().Add(-30*time.Minute)) require.NoError(t, err) require.Len(t, subscriptions, 1) } @@ -345,7 +345,7 @@ func TestListExpiredSubscriptionsWithEmptyWriter(t *testing.T) { require.NoError(t, err) require.NotNil(t, subOut2) - subscriptions, err := repo.ListExpiredSubscriptions(ctx, "", fakeClock.Now().Add(-30*time.Minute)) + subscriptions, err := repo.ListExpiredSubscriptions(ctx, fakeClock.Now().Add(-30*time.Minute)) require.NoError(t, err) require.Len(t, subscriptions, 1) } diff --git a/pkg/rid/store/raftstore/identification_service_area.go b/pkg/rid/store/raftstore/identification_service_area.go index 24506da0f..4170b3164 100644 --- a/pkg/rid/store/raftstore/identification_service_area.go +++ b/pkg/rid/store/raftstore/identification_service_area.go @@ -70,8 +70,8 @@ func (r *repo) SearchISAs(ctx context.Context, cells s2.CellUnion, earliest *tim return r.consensus.HandleReadRequest(ctx, searchISAs, buf) } -func (r *repo) ListExpiredISAs(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { - buf, err := json.Marshal(expiredPayload{Writer: writer, Threshold: threshold}) +func (r *repo) ListExpiredISAs(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) { + buf, err := json.Marshal(expiredPayload{Threshold: threshold}) if err != nil { return nil, stacktrace.Propagate(err, "failed to marshal payload") } @@ -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, limit int) ([]dssmodels.ID, error) { +func (r *repo) DeleteExpiredISAs(_ context.Context, 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/payloads.go b/pkg/rid/store/raftstore/payloads.go index 7db05469a..595237f59 100644 --- a/pkg/rid/store/raftstore/payloads.go +++ b/pkg/rid/store/raftstore/payloads.go @@ -15,7 +15,6 @@ type cellsByOwnerPayload struct { // expiredPayload carries the arguments common to ListExpiredISAs/ListExpiredSubscriptions. type expiredPayload struct { - Writer string `json:"writer"` Threshold time.Time `json:"threshold"` } diff --git a/pkg/rid/store/raftstore/store.go b/pkg/rid/store/raftstore/store.go index ebbe4c40d..3689f3169 100644 --- a/pkg/rid/store/raftstore/store.go +++ b/pkg/rid/store/raftstore/store.go @@ -91,7 +91,7 @@ func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, err if err := json.Unmarshal(proposal.Value, &payload); err != nil { return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", listExpiredISAs) } - return r.Store.GetRepo().ListExpiredISAs(ctx, payload.Writer, payload.Threshold) + return r.Store.GetRepo().ListExpiredISAs(ctx, payload.Threshold) case string(countISAs): return r.Store.GetRepo().CountISAs(ctx) @@ -159,7 +159,7 @@ func (r *repo) Apply(ctx context.Context, proposal consensus.Proposal) (any, err if err := json.Unmarshal(proposal.Value, &payload); err != nil { return nil, stacktrace.Propagate(err, "failed to unmarshal %s payload", listExpiredSubscriptions) } - return r.Store.GetRepo().ListExpiredSubscriptions(ctx, payload.Writer, payload.Threshold) + return r.Store.GetRepo().ListExpiredSubscriptions(ctx, payload.Threshold) case string(countSubscriptions): return r.Store.GetRepo().CountSubscriptions(ctx) diff --git a/pkg/rid/store/raftstore/subscriptions.go b/pkg/rid/store/raftstore/subscriptions.go index 2e767c776..102d262f3 100644 --- a/pkg/rid/store/raftstore/subscriptions.go +++ b/pkg/rid/store/raftstore/subscriptions.go @@ -98,8 +98,8 @@ func (r *repo) MaxSubscriptionCountInCellsByOwner(ctx context.Context, cells s2. return r.consensus.HandleReadRequest(ctx, maxSubscriptionCountInCellsByOwner, buf) } -func (r *repo) ListExpiredSubscriptions(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { - buf, err := json.Marshal(expiredPayload{Writer: writer, Threshold: threshold}) +func (r *repo) ListExpiredSubscriptions(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) { + buf, err := json.Marshal(expiredPayload{Threshold: threshold}) if err != nil { return nil, stacktrace.Propagate(err, "failed to marshal payload") } @@ -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, limit int) ([]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/rid/store/sqlstore/identification_service_area.go b/pkg/rid/store/sqlstore/identification_service_area.go index 41e1bf95c..a20a76f21 100644 --- a/pkg/rid/store/sqlstore/identification_service_area.go +++ b/pkg/rid/store/sqlstore/identification_service_area.go @@ -210,58 +210,24 @@ func (r *repo) SearchISAs(ctx context.Context, cells s2.CellUnion, earliest *tim return r.fetchISAs(ctx, isasInCellsQuery, earliest, latest, dssql.CellUnionToCellIds(cells), dssmodels.MaxResultLimit) } -// ListExpiredISAs lists the IDs of all expired ISAs based on writer. -// The function queries both empty writer and null writer when passing empty string as a writer. -func (r *repo) ListExpiredISAs(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { - if len(writer) == 0 { - isasInCellsQuery := ` - SELECT - id - FROM - identification_service_areas - WHERE - ends_at <= $1 - AND - (writer = '' OR writer IS NULL)` - return dssql.FetchIDs(ctx, r.Queryable, isasInCellsQuery, threshold) - } - - isasInCellsQuery := ` +// ListExpiredISAs lists the IDs of all expired ISAs. +func (r *repo) ListExpiredISAs(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) { + return dssql.FetchIDs(ctx, r.Queryable, ` SELECT id FROM identification_service_areas WHERE - ends_at <= $1 - AND - writer = $2` - return dssql.FetchIDs(ctx, r.Queryable, isasInCellsQuery, threshold, writer) + ends_at <= $1`, threshold) } -// DeleteExpiredISAs deletes up to `limit` expired ISAs based on writer and returns the +// DeleteExpiredISAs deletes up to `limit` expired ISAs 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, limit int) ([]dssmodels.ID, error) { - if len(writer) == 0 { - 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 id IN (SELECT id FROM expired) - RETURNING id`, expiredQuery) - return dssql.FetchIDs(ctx, r.Queryable, deleteExpiredQuery, args...) - } - +func (r *repo) DeleteExpiredISAs(ctx context.Context, threshold time.Time, limit int) ([]dssmodels.ID, error) { expiredQuery, args := dssql.AppendLimitClause(` SELECT id FROM identification_service_areas - WHERE ends_at <= $1 - AND writer = $2`, []any{threshold, writer}, limit) + WHERE ends_at <= $1`, []any{threshold}, limit) deleteExpiredQuery := fmt.Sprintf(` WITH expired AS (%s ) diff --git a/pkg/rid/store/sqlstore/identification_service_area_test.go b/pkg/rid/store/sqlstore/identification_service_area_test.go index 8a370f238..928d6b488 100644 --- a/pkg/rid/store/sqlstore/identification_service_area_test.go +++ b/pkg/rid/store/sqlstore/identification_service_area_test.go @@ -283,7 +283,7 @@ func TestListExpiredISAs(t *testing.T) { require.NoError(t, err) require.NotNil(t, saOut2) - serviceAreas, err := repo.ListExpiredISAs(ctx, writer, fakeClock.Now().Add(-30*time.Minute)) + serviceAreas, err := repo.ListExpiredISAs(ctx, fakeClock.Now().Add(-30*time.Minute)) require.NoError(t, err) require.Len(t, serviceAreas, 1) } @@ -321,11 +321,54 @@ func TestListExpiredISAsWithEmptyWriter(t *testing.T) { require.NoError(t, err) require.NotNil(t, saOut2) - serviceAreas, err := repo.ListExpiredISAs(ctx, "", fakeClock.Now().Add(-30*time.Minute)) + serviceAreas, err := repo.ListExpiredISAs(ctx, fakeClock.Now().Add(-30*time.Minute)) require.NoError(t, err) require.Len(t, serviceAreas, 1) } +func TestListExpiredISAsWithAllWriters(t *testing.T) { + ctx := context.Background() + store, tearDownStore := setUpStore(ctx, t) + defer tearDownStore() + + repo, err := store.Interact(ctx) + require.NoError(t, err) + + fakeClock := clockwork.NewFakeClockAt(time.Now()) + + // Insert ISA with endtime 1 day from now + isa1 := *serviceArea + startTime := fakeClock.Now() + isa1.StartTime = &startTime + endTime := fakeClock.Now().Add(24 * time.Hour) + isa1.EndTime = &endTime + saOut1, err := repo.InsertISA(ctx, &isa1) + require.NoError(t, err) + require.NotNil(t, saOut1) + + // Insert ISAs with endtime to 30 minutes ago, one with a writer and one without + isa2 := *serviceArea + startTime = fakeClock.Now().Add(-1 * time.Hour) + isa2.StartTime = &startTime + endTime = fakeClock.Now().Add(-30 * time.Minute) + isa2.EndTime = &endTime + isa2.ID = dssmodels.ID(uuid.New().String()) + saOut2, err := repo.InsertISA(ctx, &isa2) + require.NoError(t, err) + require.NotNil(t, saOut2) + + isa3 := isa2 + isa3.ID = dssmodels.ID(uuid.New().String()) + isa3.Writer = "" + saOut3, err := repo.InsertISA(ctx, &isa3) + require.NoError(t, err) + require.NotNil(t, saOut3) + + serviceAreas, err := repo.ListExpiredISAs(ctx, fakeClock.Now().Add(-30*time.Minute)) + require.NoError(t, err) + require.ElementsMatch(t, serviceAreas, []dssmodels.ID{isa2.ID, isa3.ID}) +} + func TestStoreCountISAs(t *testing.T) { var ( ctx = context.Background() diff --git a/pkg/rid/store/sqlstore/subscriptions.go b/pkg/rid/store/sqlstore/subscriptions.go index a0429c4e3..81719a489 100644 --- a/pkg/rid/store/sqlstore/subscriptions.go +++ b/pkg/rid/store/sqlstore/subscriptions.go @@ -294,58 +294,24 @@ func (r *repo) SearchSubscriptionsByOwner(ctx context.Context, cells s2.CellUnio return r.process(ctx, query, dssql.CellUnionToCellIds(cells), owner, r.clock.Now(), dssmodels.MaxResultLimit) } -// ListExpiredSubscriptions lists the IDs of all expired Subscriptions based on writer. -// The function queries both empty writer and null writer when passing empty string as a writer. -func (r *repo) ListExpiredSubscriptions(ctx context.Context, writer string, threshold time.Time) ([]dssmodels.ID, error) { - if len(writer) == 0 { - query := ` - SELECT - id - FROM - subscriptions - WHERE - ends_at <= $1 - AND - (writer = '' OR writer IS NULL)` - return dssql.FetchIDs(ctx, r.Queryable, query, threshold) - } - - query := ` +// ListExpiredSubscriptions lists the IDs of all expired Subscriptions. +func (r *repo) ListExpiredSubscriptions(ctx context.Context, threshold time.Time) ([]dssmodels.ID, error) { + return dssql.FetchIDs(ctx, r.Queryable, ` SELECT id FROM subscriptions WHERE - ends_at <= $1 - AND - writer = $2` - return dssql.FetchIDs(ctx, r.Queryable, query, threshold, writer) + ends_at <= $1`, threshold) } -// 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, limit int) ([]dssmodels.ID, error) { - if len(writer) == 0 { - 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 id IN (SELECT id FROM expired) - RETURNING id`, expiredQuery) - return dssql.FetchIDs(ctx, r.Queryable, deleteExpiredQuery, args...) - } - +// DeleteExpiredSubscriptions deletes up to `limit` expired Subscriptions and returns the +// IDs of the deleted Subscriptions. A limit of 0 means unlimited. +func (r *repo) DeleteExpiredSubscriptions(ctx context.Context, threshold time.Time, limit int) ([]dssmodels.ID, error) { expiredQuery, args := dssql.AppendLimitClause(` SELECT id FROM subscriptions - WHERE ends_at <= $1 - AND writer = $2`, []any{threshold, writer}, limit) + WHERE ends_at <= $1`, []any{threshold}, limit) deleteExpiredQuery := fmt.Sprintf(` WITH expired AS (%s ) diff --git a/pkg/rid/store/sqlstore/subscriptions_test.go b/pkg/rid/store/sqlstore/subscriptions_test.go index 99bd6a161..743b589ce 100644 --- a/pkg/rid/store/sqlstore/subscriptions_test.go +++ b/pkg/rid/store/sqlstore/subscriptions_test.go @@ -343,7 +343,7 @@ func TestListExpiredSubscriptions(t *testing.T) { require.NoError(t, err) require.NotNil(t, subOut2) - subscriptions, err := repo.ListExpiredSubscriptions(ctx, writer, fakeClock.Now().Add(-30*time.Minute)) + subscriptions, err := repo.ListExpiredSubscriptions(ctx, fakeClock.Now().Add(-30*time.Minute)) require.NoError(t, err) require.Len(t, subscriptions, 1) } @@ -381,7 +381,7 @@ func TestListExpiredSubscriptionsWithEmptyWriter(t *testing.T) { require.NoError(t, err) require.NotNil(t, subOut2) - subscriptions, err := repo.ListExpiredSubscriptions(ctx, "", fakeClock.Now().Add(-30*time.Minute)) + subscriptions, err := repo.ListExpiredSubscriptions(ctx, fakeClock.Now().Add(-30*time.Minute)) require.NoError(t, err) require.Len(t, subscriptions, 1) } diff --git a/test/evict/evict_helper.py b/test/evict/evict_helper.py index 6b2c6ceae..f4f4b8645 100644 --- a/test/evict/evict_helper.py +++ b/test/evict/evict_helper.py @@ -18,7 +18,6 @@ def run_evict( rid_ttl: str | None = None, scd_limit: int | None = None, rid_limit: int | None = None, - locality: str = "local_dev", delete: bool = False, ): db_hostname = os.environ.get("DB_HOSTNAME", "local-dss-crdb") @@ -35,8 +34,6 @@ def run_evict( f"--scd_sub={str(scd_sub).lower()}", f"--rid_isa={str(rid_isa).lower()}", f"--rid_sub={str(rid_sub).lower()}", - "--locality", - locality, "--datastore_host", db_hostname, "--datastore_port", @@ -94,20 +91,14 @@ def evict_rid_ISAs( self, ttl: str, delete: bool, - locality: str = "local_dev", limit: int | None = None, ): - self.run_evict( - rid_isa=True, delete=delete, rid_ttl=ttl, locality=locality, rid_limit=limit - ) + self.run_evict(rid_isa=True, delete=delete, rid_ttl=ttl, rid_limit=limit) def evict_rid_subscriptions( self, ttl: str, delete: bool, - locality: str = "local_dev", limit: int | None = None, ): - self.run_evict( - rid_sub=True, delete=delete, rid_ttl=ttl, locality=locality, rid_limit=limit - ) + self.run_evict(rid_sub=True, delete=delete, rid_ttl=ttl, rid_limit=limit) diff --git a/test/evict/tests/test_rid_ISA.py b/test/evict/tests/test_rid_ISA.py index c4be18ac4..b89debded 100644 --- a/test/evict/tests/test_rid_ISA.py +++ b/test/evict/tests/test_rid_ISA.py @@ -58,14 +58,6 @@ def test_rid_ISA(qh: QueryHelper, eh: EvictHelper): ) sys.exit(1) - logger.debug("Evicting subscriptions older than 1s on another locality") - eh.evict_rid_ISAs("1s", delete=True, locality="somethingelse") - if not qh.get_rid_ISA(ISA_id): - logger.error( - "❌ Test ISA shall still be present since we used another locality" - ) - sys.exit(1) - logger.debug("Evicting subscriptions older than 1s") eh.evict_rid_ISAs("1s", delete=True) diff --git a/test/evict/tests/test_rid_subscription.py b/test/evict/tests/test_rid_subscription.py index a8d51e0b5..541f8d9ba 100644 --- a/test/evict/tests/test_rid_subscription.py +++ b/test/evict/tests/test_rid_subscription.py @@ -58,14 +58,6 @@ def test_rid_subscription(qh: QueryHelper, eh: EvictHelper): logger.error("❌ Test subscription shall still be present since we evicted ISA") sys.exit(1) - logger.debug("Evicting subscriptions older than 1s on another locality") - eh.evict_rid_subscriptions("1s", delete=True, locality="somethingelse") - if not qh.get_rid_subscription(sub_id): - logger.error( - "❌ Test subscription shall still be present since we used another locality" - ) - sys.exit(1) - logger.debug("Evicting subscriptions older than 1s") eh.evict_rid_subscriptions("1s", delete=True)