Files
probo/pkg/coredata/access_review_campaign_source.go
Bryan Frimin 68a64647bf Batch-load and validate campaign sources before merging
MergeByCampaignID joined access_review_sources directly from coredata
to resolve live sources, so unrecognized or out-of-scope IDs were
silently dropped instead of erroring, and in the worst case (every ID
invalid) the NOT MATCHED BY SOURCE clause deleted every existing
campaign source. The syncCampaignSources ErrResourceNotFound check
was therefore unreachable dead code.

Add AccessReviewSources.LoadByIDs, matching the existing scoped
LoadByIDs pattern (id = ANY(@ids) plus a resolved-count check), and
have CreateCampaign, AddCampaignSource, and syncCampaignSources
resolve and validate sources up front. MergeByCampaignID now takes
the already-loaded sources and builds its desired-state CTE from an
unnest() of their values instead of joining access_review_sources,
keeping the merge inside the campaign-source entity boundary.

Signed-off-by: Bryan Frimin <bryan@probo.com>
2026-07-27 18:56:24 +02:00

383 lines
10 KiB
Go

// Copyright (c) 2026 Probo Inc <hello@probo.com>.
//
// Permission is hereby granted, free of charge, to any person obtaining a copy
// of this software and associated documentation files (the "Software"), to deal
// in the Software without restriction, including without limitation the rights
// to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
// copies of the Software, and to permit persons to whom the Software is
// furnished to do so, subject to the following conditions:
//
// The above copyright notice and this permission notice shall be included in
// all copies or substantial portions of the Software.
//
// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
// IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
// FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
// AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
// LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
// SOFTWARE.
package coredata
import (
"context"
"errors"
"fmt"
"maps"
"time"
"github.com/jackc/pgx/v5"
"go.gearno.de/kit/pg"
"go.probo.inc/probo/pkg/gid"
"go.probo.inc/probo/pkg/iam/policy"
)
type (
// AccessReviewCampaignSource is the per-campaign snapshot of an access
// source. It captures the source identity (name, connector) at
// the time the source was scoped into the campaign so that the review's
// data survives even if the live access source is later deleted. Access
// entries and fetch attempts reference this snapshot, not the live source.
AccessReviewCampaignSource struct {
ID gid.GID `db:"id"`
OrganizationID gid.GID `db:"organization_id"`
TenantID gid.TenantID `db:"tenant_id"`
AccessReviewCampaignID gid.GID `db:"access_review_campaign_id"`
AccessReviewSourceID *gid.GID `db:"access_review_source_id"`
Name string `db:"name"`
ConnectorID *gid.GID `db:"connector_id"`
CreatedAt time.Time `db:"created_at"`
UpdatedAt time.Time `db:"updated_at"`
}
AccessReviewCampaignSources []*AccessReviewCampaignSource
)
func (s *AccessReviewCampaignSource) AuthorizationAttributes(
ctx context.Context,
conn pg.Querier,
resourceIDs []gid.GID,
) (policy.AttributesByID, error) {
q := `SELECT id, organization_id FROM access_review_campaign_sources WHERE id = ANY(@resource_ids::text[])`
args := pgx.StrictNamedArgs{
"resource_ids": resourceIDs,
}
rows, err := conn.Query(ctx, q, args)
if err != nil {
return nil, fmt.Errorf("cannot query authorization attributes: %w", err)
}
defer rows.Close()
attrsByID := make(policy.AttributesByID)
for rows.Next() {
var id, organizationID gid.GID
if err := rows.Scan(&id, &organizationID); err != nil {
return nil, fmt.Errorf("cannot scan authorization attributes: %w", err)
}
attrsByID[id] = policy.Attributes{
"organization_id": organizationID.String(),
}
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("cannot iterate authorization attributes: %w", err)
}
return attrsByID, nil
}
// Upsert inserts the snapshot or refreshes its denormalized identity from the
// live source. The generated ID is preserved across upserts because it is not
// part of the conflict target, so entries that already reference the snapshot
// keep pointing at the same row.
func (s *AccessReviewCampaignSource) Upsert(
ctx context.Context,
conn pg.Tx,
scope Scoper,
) error {
q := `
INSERT INTO access_review_campaign_sources (
id,
organization_id,
tenant_id,
access_review_campaign_id,
access_review_source_id,
name,
connector_id,
created_at,
updated_at
) VALUES (
@id,
@organization_id,
@tenant_id,
@access_review_campaign_id,
@access_review_source_id,
@name,
@connector_id,
@created_at,
@updated_at
)
ON CONFLICT (access_review_campaign_id, access_review_source_id) DO UPDATE SET
name = EXCLUDED.name,
connector_id = EXCLUDED.connector_id,
updated_at = EXCLUDED.updated_at
RETURNING id
`
args := pgx.StrictNamedArgs{
"id": s.ID,
"organization_id": s.OrganizationID,
"tenant_id": scope.GetTenantID(),
"access_review_campaign_id": s.AccessReviewCampaignID,
"access_review_source_id": s.AccessReviewSourceID,
"name": s.Name,
"connector_id": s.ConnectorID,
"created_at": s.CreatedAt,
"updated_at": s.UpdatedAt,
}
if err := conn.QueryRow(ctx, q, args).Scan(&s.ID); err != nil {
return fmt.Errorf("cannot upsert campaign source: %w", err)
}
return nil
}
// MergeByCampaignID syncs scoped campaign source snapshots against the given,
// already-loaded and validated live access sources: upserts a snapshot for
// each source and deletes snapshots no longer in the set. Callers are
// responsible for resolving accessReviewSources (existence, tenant, and
// organization checks) before calling this method.
func (sources *AccessReviewCampaignSources) MergeByCampaignID(
ctx context.Context,
conn pg.Tx,
scope Scoper,
campaignID gid.GID,
accessReviewSources AccessReviewSources,
) error {
ids := make([]string, 0, len(accessReviewSources))
organizationIDs := make([]string, 0, len(accessReviewSources))
names := make([]string, 0, len(accessReviewSources))
connectorIDs := make([]*string, 0, len(accessReviewSources))
seen := make(map[gid.GID]struct{}, len(accessReviewSources))
for _, source := range accessReviewSources {
if _, ok := seen[source.ID]; ok {
continue
}
seen[source.ID] = struct{}{}
var connectorID *string
if source.ConnectorID != nil {
s := source.ConnectorID.String()
connectorID = &s
}
ids = append(ids, source.ID.String())
organizationIDs = append(organizationIDs, source.OrganizationID.String())
names = append(names, source.Name)
connectorIDs = append(connectorIDs, connectorID)
}
now := time.Now()
q := `
WITH desired_sources AS (
SELECT
id AS access_review_source_id,
organization_id,
name,
connector_id
FROM unnest(
@access_review_source_ids::text[],
@organization_ids::text[],
@names::text[],
@connector_ids::text[]
) AS t(id, organization_id, name, connector_id)
)
MERGE INTO access_review_campaign_sources AS target
USING desired_sources AS source
ON
%s
AND target.access_review_campaign_id = @access_review_campaign_id
AND target.access_review_source_id = source.access_review_source_id
WHEN MATCHED THEN
UPDATE SET
name = source.name,
connector_id = source.connector_id,
updated_at = @now
WHEN NOT MATCHED THEN
INSERT (
id,
organization_id,
tenant_id,
access_review_campaign_id,
access_review_source_id,
name,
connector_id,
created_at,
updated_at
)
VALUES (
generate_gid(decode_base64_unpadded(@tenant_id), @access_review_campaign_source_entity_type),
source.organization_id,
@tenant_id,
@access_review_campaign_id,
source.access_review_source_id,
source.name,
source.connector_id,
@now,
@now
)
WHEN NOT MATCHED BY SOURCE
AND %s
AND target.access_review_campaign_id = @access_review_campaign_id THEN
DELETE
`
q = fmt.Sprintf(q, scope.SQLFragment(), scope.SQLFragment())
args := pgx.StrictNamedArgs{
"access_review_campaign_id": campaignID,
"access_review_source_ids": ids,
"organization_ids": organizationIDs,
"names": names,
"connector_ids": connectorIDs,
"access_review_campaign_source_entity_type": AccessReviewCampaignSourceEntityType,
"tenant_id": scope.GetTenantID(),
"now": now,
}
maps.Copy(args, scope.SQLArguments())
if _, err := conn.Exec(ctx, q, args); err != nil {
return fmt.Errorf("cannot merge campaign sources: %w", err)
}
return nil
}
func (s *AccessReviewCampaignSource) LoadByID(
ctx context.Context,
conn pg.Querier,
scope Scoper,
id gid.GID,
) error {
q := `
SELECT
id,
organization_id,
tenant_id,
access_review_campaign_id,
access_review_source_id,
name,
connector_id,
created_at,
updated_at
FROM access_review_campaign_sources
WHERE
%s
AND id = @id
LIMIT 1
`
q = fmt.Sprintf(q, scope.SQLFragment())
args := pgx.StrictNamedArgs{"id": id}
maps.Copy(args, scope.SQLArguments())
rows, err := conn.Query(ctx, q, args)
if err != nil {
return fmt.Errorf("cannot query campaign source: %w", err)
}
result, err := pgx.CollectExactlyOneRow(rows, pgx.RowToStructByName[AccessReviewCampaignSource])
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return ErrResourceNotFound
}
return fmt.Errorf("cannot collect campaign source: %w", err)
}
*s = result
return nil
}
func (s *AccessReviewCampaignSource) DeleteByCampaignIDAndAccessReviewSourceID(
ctx context.Context,
conn pg.Tx,
scope Scoper,
campaignID gid.GID,
accessSourceID gid.GID,
) error {
q := `
DELETE FROM access_review_campaign_sources
WHERE
%s
AND access_review_campaign_id = @access_review_campaign_id
AND access_review_source_id = @access_review_source_id
`
q = fmt.Sprintf(q, scope.SQLFragment())
args := pgx.StrictNamedArgs{
"access_review_campaign_id": campaignID,
"access_review_source_id": accessSourceID,
}
maps.Copy(args, scope.SQLArguments())
if _, err := conn.Exec(ctx, q, args); err != nil {
return fmt.Errorf("cannot delete campaign source: %w", err)
}
return nil
}
func (sources *AccessReviewCampaignSources) LoadByCampaignID(
ctx context.Context,
conn pg.Querier,
scope Scoper,
campaignID gid.GID,
) error {
q := `
SELECT
id,
organization_id,
tenant_id,
access_review_campaign_id,
access_review_source_id,
name,
connector_id,
created_at,
updated_at
FROM access_review_campaign_sources
WHERE
%s
AND access_review_campaign_id = @access_review_campaign_id
ORDER BY name ASC
`
q = fmt.Sprintf(q, scope.SQLFragment())
args := pgx.StrictNamedArgs{"access_review_campaign_id": campaignID}
maps.Copy(args, scope.SQLArguments())
rows, err := conn.Query(ctx, q, args)
if err != nil {
return fmt.Errorf("cannot query campaign sources: %w", err)
}
result, err := pgx.CollectRows(rows, pgx.RowToAddrOfStructByName[AccessReviewCampaignSource])
if err != nil {
return fmt.Errorf("cannot collect campaign sources: %w", err)
}
*sources = result
return nil
}