From 141a5ae8b3bbdb4869305fbeed6aa1fd5cb3556d Mon Sep 17 00:00:00 2001 From: Abhishek Sah Date: Fri, 18 Sep 2026 13:44:55 +0530 Subject: [PATCH 1/2] feat(store): policy, relation, and resource reads skip soft-deleted rows Every read in the three repositories goes through fromLive, or adds the live filter directly where the query aliases the table. The policy repository's small lookups of organization, project, and group titles for audit records use fromLive too. The two updates keyed by id add the filter as well. The ON CONFLICT upserts are unchanged. They move with the unique constraints in a later change. --- internal/store/postgres/policy_repository.go | 20 ++++++------ .../store/postgres/policy_repository_test.go | 31 +++++++++++++++++++ .../store/postgres/relation_repository.go | 8 ++--- .../postgres/relation_repository_test.go | 23 ++++++++++++++ .../store/postgres/resource_repository.go | 8 ++--- .../postgres/resource_repository_test.go | 25 +++++++++++++++ 6 files changed, 97 insertions(+), 18 deletions(-) diff --git a/internal/store/postgres/policy_repository.go b/internal/store/postgres/policy_repository.go index 0003146e18..55a558db03 100644 --- a/internal/store/postgres/policy_repository.go +++ b/internal/store/postgres/policy_repository.go @@ -39,7 +39,7 @@ func (r PolicyRepository) buildListQuery() *goqu.SelectDataset { "p.principal_type", "p.role_id", "p.grant_relation", - ).From(goqu.T(TABLE_POLICIES).As("p")) + ).From(goqu.T(TABLE_POLICIES).As("p")).Where(live("p")) } func (r PolicyRepository) Get(ctx context.Context, id string) (policy.Policy, error) { @@ -181,7 +181,7 @@ func (r PolicyRepository) List(ctx context.Context, flt policy.Filter) ([]policy func (r PolicyRepository) Count(ctx context.Context, flt policy.Filter) (int64, error) { var count int64 - stmt := dialect.Select(goqu.COUNT(goqu.Star()).As("count")).From(goqu.T(TABLE_POLICIES).As("p")) + stmt := dialect.Select(goqu.COUNT(goqu.Star()).As("count")).From(goqu.T(TABLE_POLICIES).As("p")).Where(live("p")) stmt = applyListFilter(stmt, flt) query, params, err := stmt.ToSQL() @@ -292,7 +292,7 @@ func (r PolicyRepository) Update(ctx context.Context, toUpdate policy.Policy) (s "updated_at": goqu.L("now()"), }).Where(goqu.Ex{ "id": toUpdate.ID, - }).Returning("id", "updated_at").ToSQL() + }, live(TABLE_POLICIES)).Returning("id", "updated_at").ToSQL() if err != nil { return "", fmt.Errorf("%w: %s", errQuery, err) } @@ -461,7 +461,7 @@ func (r PolicyRepository) GroupMemberCount(ctx context.Context, groupIDs []strin if len(groupIDs) == 0 { return nil, policy.ErrInvalidID } - stmt := dialect.From("policies"). + stmt := fromLive(TABLE_POLICIES). Select(goqu.I("resource_id").As("id"), goqu.COUNT(goqu.DISTINCT(goqu.I("principal_id"))).As("count")). Where(goqu.Ex{ "resource_type": schema.GroupNamespace, @@ -498,7 +498,7 @@ func (r PolicyRepository) ProjectMemberCount(ctx context.Context, projectIDs []s if len(projectIDs) == 0 { return nil, policy.ErrInvalidID } - stmt := dialect.From("policies"). + stmt := fromLive(TABLE_POLICIES). Select(goqu.I("resource_id").As("id"), goqu.COUNT(goqu.DISTINCT(goqu.I("principal_id"))).As("count")). Where(goqu.Ex{ "resource_type": schema.ProjectNamespace, @@ -534,7 +534,7 @@ func (r PolicyRepository) OrgMemberCount(ctx context.Context, id string) (policy if len(id) == 0 { return policy.MemberCount{}, policy.ErrInvalidID } - stmt := dialect.From("policies"). + stmt := fromLive(TABLE_POLICIES). Select(goqu.I("resource_id").As("id"), goqu.COUNT(goqu.DISTINCT(goqu.I("principal_id"))).As("count")). Where(goqu.Ex{ "resource_type": schema.OrganizationNamespace, @@ -600,7 +600,7 @@ func (r PolicyRepository) buildPolicyAuditRecord(ctx context.Context, tx *sqlx.T // getPolicyByConstraint fetches a policy by unique constraint fields // Returns the policy and true if found, empty policy and false if not found func (r PolicyRepository) getPolicyByConstraint(ctx context.Context, pol policy.Policy) (Policy, bool) { - query, params, _ := dialect.From(TABLE_POLICIES). + query, params, _ := fromLive(TABLE_POLICIES). Select("id", "resource_type", "resource_id", "principal_id", "principal_type", "role_id"). Where(goqu.Ex{ "role_id": pol.RoleID, @@ -625,19 +625,19 @@ func (r PolicyRepository) getResourceInfo(ctx context.Context, tx *sqlx.Tx, reso switch resourceType { case schema.OrganizationNamespace: orgID = resourceID - orgQuery, orgParams, _ := dialect.From(TABLE_ORGANIZATIONS). + orgQuery, orgParams, _ := fromLive(TABLE_ORGANIZATIONS). Select("title"). Where(goqu.Ex{"id": resourceID}). ToSQL() _ = tx.QueryRowContext(ctx, orgQuery, orgParams...).Scan(&resourceName) case schema.ProjectNamespace: - projQuery, projParams, _ := dialect.From(TABLE_PROJECTS). + projQuery, projParams, _ := fromLive(TABLE_PROJECTS). Select("org_id", "title"). Where(goqu.Ex{"id": resourceID}). ToSQL() _ = tx.QueryRowContext(ctx, projQuery, projParams...).Scan(&orgID, &resourceName) case schema.GroupNamespace: - grpQuery, grpParams, _ := dialect.From(TABLE_GROUPS). + grpQuery, grpParams, _ := fromLive(TABLE_GROUPS). Select("org_id", "title"). Where(goqu.Ex{"id": resourceID}). ToSQL() diff --git a/internal/store/postgres/policy_repository_test.go b/internal/store/postgres/policy_repository_test.go index fb5f002aa2..e4e45854e5 100644 --- a/internal/store/postgres/policy_repository_test.go +++ b/internal/store/postgres/policy_repository_test.go @@ -535,3 +535,34 @@ func (s *PolicyRepositoryTestSuite) TestOrgMemberCount() { }) } } + +func (s *PolicyRepositoryTestSuite) TestSkipsSoftDeletedPolicies() { + deleted := s.policies[0] + flt := policy.Filter{PrincipalID: s.userID} + before, err := s.repository.List(s.ctx, flt) + s.Require().NoError(err) + countBefore, err := s.repository.Count(s.ctx, flt) + s.Require().NoError(err) + + _, err = s.client.ExecContext(s.ctx, "UPDATE policies SET deleted_at = now() WHERE id = $1", deleted.ID) + if err != nil { + s.T().Fatal(err) + } + + _, err = s.repository.Get(s.ctx, deleted.ID) + s.Assert().ErrorIs(err, policy.ErrNotExist) + + got, err := s.repository.List(s.ctx, flt) + s.Assert().NoError(err) + s.Assert().Len(got, len(before)-1) + for _, p := range got { + s.Assert().NotEqual(deleted.ID, p.ID) + } + + count, err := s.repository.Count(s.ctx, flt) + s.Assert().NoError(err) + s.Assert().Equal(countBefore-1, count) + + _, err = s.repository.Update(s.ctx, deleted) + s.Assert().ErrorIs(err, policy.ErrNotExist) +} diff --git a/internal/store/postgres/relation_repository.go b/internal/store/postgres/relation_repository.go index 1d79cd68e1..9d93b9007d 100644 --- a/internal/store/postgres/relation_repository.go +++ b/internal/store/postgres/relation_repository.go @@ -59,7 +59,7 @@ func (r RelationRepository) Upsert(ctx context.Context, relationToCreate relatio } func (r RelationRepository) List(ctx context.Context, flt relation.Filter) ([]relation.Relation, error) { - stmt := dialect.Select(&relationCols{}).From(TABLE_RELATIONS) + stmt := fromLive(TABLE_RELATIONS).Select(&relationCols{}) if flt.Subject.ID != "" { stmt = stmt.Where(goqu.Ex{ "subject_id": flt.Subject.ID, @@ -101,7 +101,7 @@ func (r RelationRepository) Get(ctx context.Context, id string) (relation.Relati return relation.Relation{}, relation.ErrInvalidID } - query, params, err := dialect.Select(&relationCols{}).From(TABLE_RELATIONS). + query, params, err := fromLive(TABLE_RELATIONS).Select(&relationCols{}). Where(goqu.Ex{ "id": id, }).ToSQL() @@ -166,7 +166,7 @@ func (r RelationRepository) DeleteByID(ctx context.Context, id string) error { func (r RelationRepository) GetByFields(ctx context.Context, rel relation.Relation) ([]relation.Relation, error) { var fetchedRelations []Relation - stmt := dialect.Select(&relationCols{}).From(TABLE_RELATIONS) + stmt := fromLive(TABLE_RELATIONS).Select(&relationCols{}) if rel.Object.ID != "" { stmt = stmt.Where(goqu.Ex{ "object_id": rel.Object.ID, @@ -230,7 +230,7 @@ func (r RelationRepository) ListByFields(ctx context.Context, rel relation.Relat if len(rel.Object.ID) != 0 { exprs = append(exprs, goqu.Ex{"object_id": rel.Object.ID}) } - query, params, err := dialect.Select(&relationCols{}).From(TABLE_RELATIONS).Where(exprs...).ToSQL() + query, params, err := fromLive(TABLE_RELATIONS).Select(&relationCols{}).Where(exprs...).ToSQL() if err != nil { return nil, fmt.Errorf("%w: %s", errQuery, err) } diff --git a/internal/store/postgres/relation_repository_test.go b/internal/store/postgres/relation_repository_test.go index 632764022d..571229fc2e 100644 --- a/internal/store/postgres/relation_repository_test.go +++ b/internal/store/postgres/relation_repository_test.go @@ -323,3 +323,26 @@ func (s *RelationRepositoryTestSuite) TestDeleteByID() { func TestRelationRepository(t *testing.T) { suite.Run(t, new(RelationRepositoryTestSuite)) } + +func (s *RelationRepositoryTestSuite) TestSkipsSoftDeletedRelations() { + deleted := s.relations[0] + _, err := s.client.ExecContext(s.ctx, "UPDATE relations SET deleted_at = now() WHERE id = $1", deleted.ID) + if err != nil { + s.T().Fatal(err) + } + + _, err = s.repository.Get(s.ctx, deleted.ID) + s.Assert().ErrorIs(err, relation.ErrNotExist) + + got, err := s.repository.List(s.ctx, relation.Filter{Subject: deleted.Subject, Object: deleted.Object}) + s.Assert().NoError(err) + for _, r := range got { + s.Assert().NotEqual(deleted.ID, r.ID) + } + + byFields, err := s.repository.GetByFields(s.ctx, deleted) + s.Assert().NoError(err) + for _, r := range byFields { + s.Assert().NotEqual(deleted.ID, r.ID) + } +} diff --git a/internal/store/postgres/resource_repository.go b/internal/store/postgres/resource_repository.go index bfb324ee54..41e60b1402 100644 --- a/internal/store/postgres/resource_repository.go +++ b/internal/store/postgres/resource_repository.go @@ -101,7 +101,7 @@ func (r ResourceRepository) Create(ctx context.Context, res resource.Resource) ( func (r ResourceRepository) List(ctx context.Context, flt resource.Filter) ([]resource.Resource, error) { var fetchedResources []Resource - sqlStatement := dialect.From(TABLE_RESOURCES) + sqlStatement := fromLive(TABLE_RESOURCES) if flt.ProjectID != "" { sqlStatement = sqlStatement.Where(goqu.Ex{"project_id": flt.ProjectID}) } @@ -149,7 +149,7 @@ func (r ResourceRepository) GetByID(ctx context.Context, id string) (resource.Re return resource.Resource{}, resource.ErrInvalidID } - query, params, err := dialect.From(TABLE_RESOURCES).Where(goqu.Ex{ + query, params, err := fromLive(TABLE_RESOURCES).Where(goqu.Ex{ "id": id, }).ToSQL() if err != nil { @@ -189,7 +189,7 @@ func (r ResourceRepository) Update(ctx context.Context, res resource.Resource) ( "metadata": marshaledMetadata, "updated_at": goqu.L("now()"), }, - ).Where(goqu.Ex{"id": res.ID}).Returning(&ResourceCols{}).ToSQL() + ).Where(goqu.Ex{"id": res.ID}, live(TABLE_RESOURCES)).Returning(&ResourceCols{}).ToSQL() if err != nil { return resource.Resource{}, fmt.Errorf("%w: %s", errQuery, err) } @@ -221,7 +221,7 @@ func (r ResourceRepository) GetByURN(ctx context.Context, urn string) (resource. return resource.Resource{}, resource.ErrInvalidURN } - query, params, err := dialect.Select(&ResourceCols{}).From(TABLE_RESOURCES).Where( + query, params, err := fromLive(TABLE_RESOURCES).Select(&ResourceCols{}).Where( goqu.Ex{ "urn": urn, }).ToSQL() diff --git a/internal/store/postgres/resource_repository_test.go b/internal/store/postgres/resource_repository_test.go index 1bb565848d..87fddc111f 100644 --- a/internal/store/postgres/resource_repository_test.go +++ b/internal/store/postgres/resource_repository_test.go @@ -411,3 +411,28 @@ func (s *ResourceRepositoryTestSuite) TestUpdate() { func TestResourceRepository(t *testing.T) { suite.Run(t, new(ResourceRepositoryTestSuite)) } + +func (s *ResourceRepositoryTestSuite) TestSkipsSoftDeletedResources() { + deleted := s.resources[0] + _, err := s.client.ExecContext(s.ctx, "UPDATE resources SET deleted_at = now() WHERE id = $1", deleted.ID) + if err != nil { + s.T().Fatal(err) + } + + _, err = s.repository.GetByID(s.ctx, deleted.ID) + s.Assert().ErrorIs(err, resource.ErrNotExist) + + _, err = s.repository.GetByURN(s.ctx, deleted.URN) + s.Assert().ErrorIs(err, resource.ErrNotExist) + + got, err := s.repository.List(s.ctx, resource.Filter{ProjectID: deleted.ProjectID}) + s.Assert().NoError(err) + for _, r := range got { + s.Assert().NotEqual(deleted.ID, r.ID) + } + + updated := deleted + updated.Title = "changed" + _, err = s.repository.Update(s.ctx, updated) + s.Assert().ErrorIs(err, resource.ErrNotExist) +} From 6d0ae3496eb3670570d446f82722c04bd24a5273 Mon Sep 17 00:00:00 2001 From: Abhishek Sah Date: Mon, 21 Sep 2026 09:46:28 +0530 Subject: [PATCH 2/2] fix(store): a resource create that reuses a taken id returns conflict The create upsert targets urn only. An insert that reuses an existing id with a different urn fails on the primary key and fell through to a raw duplicate key error. It now maps to resource.ErrConflict, which the handler already turns into a conflict response. With live-only reads the id of a soft-deleted row is invisible to the pre-check, so this path is easier to reach. --- internal/store/postgres/resource_repository.go | 2 ++ internal/store/postgres/resource_repository_test.go | 12 ++++++++++++ 2 files changed, 14 insertions(+) diff --git a/internal/store/postgres/resource_repository.go b/internal/store/postgres/resource_repository.go index 41e60b1402..4fec0982d7 100644 --- a/internal/store/postgres/resource_repository.go +++ b/internal/store/postgres/resource_repository.go @@ -90,6 +90,8 @@ func (r ResourceRepository) Create(ctx context.Context, res resource.Resource) ( return resource.Resource{}, fmt.Errorf("%w: %w", err, resource.ErrInvalidDetail) case errors.Is(err, ErrInvalidTextRepresentation): return resource.Resource{}, fmt.Errorf("%w: %w", err, resource.ErrInvalidUUID) + case errors.Is(err, ErrDuplicateKey): + return resource.Resource{}, resource.ErrConflict default: return resource.Resource{}, err } diff --git a/internal/store/postgres/resource_repository_test.go b/internal/store/postgres/resource_repository_test.go index 87fddc111f..304bab8b91 100644 --- a/internal/store/postgres/resource_repository_test.go +++ b/internal/store/postgres/resource_repository_test.go @@ -435,4 +435,16 @@ func (s *ResourceRepositoryTestSuite) TestSkipsSoftDeletedResources() { updated.Title = "changed" _, err = s.repository.Update(s.ctx, updated) s.Assert().ErrorIs(err, resource.ErrNotExist) + + // the id of a deleted row is still taken; a create that reuses it with a new URN must conflict + _, err = s.repository.Create(s.ctx, resource.Resource{ + ID: deleted.ID, + URN: "urn-reusing-a-deleted-id", + Name: "reused-id", + ProjectID: deleted.ProjectID, + NamespaceID: deleted.NamespaceID, + PrincipalID: deleted.PrincipalID, + PrincipalType: deleted.PrincipalType, + }) + s.Assert().ErrorIs(err, resource.ErrConflict) }