Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
65 changes: 57 additions & 8 deletions repository/cloud_identity_repository.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,13 @@ type CloudIdentityRepository interface {
// UpsertIdentity records one identity, keyed on (workspace_id, native_id).
// A repeat scan updates the same row and stamps the generation; it never
// creates a duplicate. Reports whether the row was newly created.
UpsertIdentity(i *models.CloudIdentity) (stored *models.CloudIdentity, created bool, err error)
//
// partialRead says the caller could not read this principal in full — some
// of what it is passing in Attrs is UNKNOWN rather than absent. The row is
// still stamped (it was seen, so it must not be reconciled away), but the
// attrs blob is left as an earlier complete read established it. See the
// implementation for why that distinction cannot be made inside the repo.
UpsertIdentity(i *models.CloudIdentity, partialRead bool) (stored *models.CloudIdentity, created bool, err error)

// UpsertSecret records one secret, keyed on (workspace_id, native_id).
UpsertSecret(s *models.CloudSecret) (stored *models.CloudSecret, created bool, err error)
Expand Down Expand Up @@ -70,7 +76,7 @@ func NewCloudIdentityRepository(db *gorm.DB) CloudIdentityRepository {
return &cloudIdentityRepository{db: db}
}

func (r *cloudIdentityRepository) UpsertIdentity(i *models.CloudIdentity) (*models.CloudIdentity, bool, error) {
func (r *cloudIdentityRepository) UpsertIdentity(i *models.CloudIdentity, partialRead bool) (*models.CloudIdentity, bool, error) {
if i.WorkspaceID == uuid.Nil || i.ConnectorID == uuid.Nil {
return nil, false, errors.New("workspace_id and connector_id are required")
}
Expand Down Expand Up @@ -112,6 +118,39 @@ func (r *cloudIdentityRepository) UpsertIdentity(i *models.CloudIdentity) (*mode
OR excluded.last_used_at > cloud_identity.last_used_at
THEN excluded.last_used_at ELSE cloud_identity.last_used_at END`)
}
// attrs gets the same treatment as last_used_at, for the same reason, and it
// is the caller who must say so. A partial read carries a THIN blob, not a
// corrective one: a throttled iam:GetRole yields Tags=nil, boundary="" and
// Description="", all `omitempty`, so they vanish from the serialised JSON
// and a blind overwrite erases what a complete read established. IAM tags are
// the only ownership signal we collect, so that loss is permanent and silent.
//
// This cannot be decided here. `detail_incomplete` lives inside the attrs
// JSON, not in a column, and it is `omitempty` -- so false is ABSENT from the
// blob and "key missing" would have to be read as complete. Parsing provider
// JSON in a provider-neutral repository to find out is the wrong place for
// that knowledge; the scanner already knows, so it tells us.
//
// Keeping the old blob is only half of it. The blob an earlier COMPLETE read
// stored carries no detail_incomplete key (it is `omitempty`, and it was
// false), so preserving it verbatim would leave the row asserting it is
// fully read when this scan could not confirm that. For a role whose last
// complete read found no boundary, constraintState would then still answer
// `unconstrained` -- a firmer claim than the evidence now supports.
//
// So: merge, do not replace. `||` keeps every key the good read established
// and overlays the one fact this read actually learned -- that it fell
// short. Old tags and boundary survive AND the row stays honest about its
// freshness, which is what makes constraintState degrade to unknown.
//
// This one key is the only provider-specific knowledge in this method, and
// it buys the whole guarantee. Only an UPDATE is affected; a first insert
// still writes whatever the partial read had, because a thin row beats no
// row and there is nothing yet to preserve.
if partialRead {
assignments["attrs"] = gorm.Expr(
`cloud_identity.attrs || jsonb_build_object('detail_incomplete', true)`)
}

err := r.db.Clauses(
clause.OnConflict{
Expand Down Expand Up @@ -141,12 +180,22 @@ func (r *cloudIdentityRepository) UpsertSecret(s *models.CloudSecret) (*models.C
now := time.Now()

assignments := map[string]interface{}{
"connector_id": s.ConnectorID,
"identity_id": s.IdentityID,
"kind": s.Kind,
"created_at": s.ProviderCreatedAt,
"expires_at": s.ExpiresAt,
"status": s.Status,
"connector_id": s.ConnectorID,
"identity_id": s.IdentityID,
"kind": s.Kind,
"created_at": s.ProviderCreatedAt,
"expires_at": s.ExpiresAt,
"status": s.Status,
// Unconditional, unlike UpsertIdentity's attrs, and safe only because no
// caller writes a secret's attrs today: upsertAccessKey builds the row
// from ListAccessKeys alone and never calls SetAWSAttrs, so this blob is
// always "{}" and there is nothing an overwrite could destroy.
//
// If a collector ever populates it -- a key's rotation date from the
// credential report is the obvious candidate -- this acquires the exact
// bug UpsertIdentity's partialRead parameter exists to prevent, and it
// needs the same treatment. It is written here rather than fixed because
// a flag no caller can set is speculative API surface.
"attrs": s.Attrs,
"last_seen_generation": s.LastSeenGeneration,
"last_seen_at": now,
Expand Down
17 changes: 16 additions & 1 deletion repository/cloud_observation_repository.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,22 @@ func (r *cloudObservationRepository) ListObservations(
return nil, 0, err
}
var out []models.CloudObservation
if err := q.Order("observed_at DESC").Limit(clampLimit(f.Limit)).Offset(f.Offset).
// Ends in `id`, per the rule stated on ListIdentities: offset paging is only
// coherent over a TOTAL order, and observed_at alone is not one.
//
// The ties are real, not theoretical. Most Record calls pass their own
// time.Now(), which rarely collides -- but CloudTrail events are stamped
// with e.EventTime, the event's own timestamp, which AWS reports to the
// second. A busy account produces many events per second, and CloudTrail is
// both the highest-volume source here and the one the Events tab pages
// through. Without a tiebreaker, page two can repeat a row page one already
// showed and omit one it never did, and the console's "newest observation
// per key" dedupe relies on this order being stable.
//
// `id DESC` rather than ASC to match the DESC direction of the four
// (workspace_id, <subject>, observed_at DESC) indexes, so the requested
// order stays on the same physical direction those indexes already provide.
if err := q.Order("observed_at DESC, id DESC").Limit(clampLimit(f.Limit)).Offset(f.Offset).
Find(&out).Error; err != nil {
return nil, 0, err
}
Expand Down
23 changes: 23 additions & 0 deletions repository/cloud_permission_repository.go
Original file line number Diff line number Diff line change
Expand Up @@ -131,10 +131,33 @@ func (r *cloudPermissionRepository) UpsertResource(res *models.CloudResource) (*
// bucket's ARN never changes, but the plan may later type a
// service where it can) should not keep showing a stale label.
"name": res.Name,
// The three ARN-derived columns ARE refreshed. 021 added them with
// NOT NULL defaults and a comment saying the next scan would fill
// them in; it did not, because only the Create path wrote them and
// every pre-021 row takes this branch forever. `is_external=false`
// is not "unknown" -- it is a positive claim of locality about what
// may be a cross-account ARN, and idx_cloud_resource_external is
// built to query exactly that column.
//
// Safe to refresh: all three are functions of the ARN, which is
// native_id, i.e. the conflict key -- so a repeat scan derives the
// same values. They move together because 021's CHECKs couple them
// to each other and to `kind` (already refreshed above):
// (NOT is_external OR resource_account <> '') and
// (object_key = '' OR kind = 's3_object'). Updating a subset can
// violate one of those and fail the upsert.
"resource_account": res.ResourceAccount,
"is_external": res.IsExternal,
"object_key": res.ObjectKey,
// sensitivity is NOT refreshed here deliberately -- see
// UpsertPermission's identical note. A later ticket may raise it
// from tags or activity, and this scan must not stamp it back
// down to the rule-based default on every repeat run.
//
// sensitivity_source and sensitivity_reason stay out for the same
// reason: they explain `sensitivity`, so refreshing them while the
// verdict stays pinned would leave a row whose stated reason
// contradicts its own value.
"last_seen_generation": res.LastSeenGeneration,
"last_seen_at": now,
"row_updated_at": now,
Expand Down
51 changes: 46 additions & 5 deletions services/cloud_aws_iam_scan.go
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,10 @@ type AWSIAMScanner struct {
// a caller with no durable run to anchor evidence to (a test, a legacy
// path) degrades to no evidence rather than failing.
evidence *ObservationWriter

// generation, when non-zero, is the number this scan stamps its rows with,
// supplied by the caller instead of derived here. See WithGeneration.
generation int
}

// WithEvidence attaches an observation writer for this run.
Expand All @@ -81,6 +85,26 @@ func (s *AWSIAMScanner) WithEvidence(w *ObservationWriter) *AWSIAMScanner {
return s
}

// WithGeneration fixes the generation this scan stamps its rows with, making
// the caller the single authority for it.
//
// It exists because there were two. cloud_scan_run.generation is assigned once
// at Claim and preserved across re-claims (`CASE WHEN generation > 0`), while
// Scan independently recomputed connector.ScanGeneration+1 on every attempt.
// Those agree until a run is re-claimed after commitScan already advanced the
// connector -- the crashed-worker path -- and then the entity rows land one
// generation ahead of the cloud_observation rows written for the very same
// pass, because the worker built the ObservationWriter from run.Generation.
// Evidence then no longer joins to the inventory it explains, which is the one
// thing cloud_observation.generation exists to guarantee.
//
// Unset (zero) keeps the derived behaviour, which is what callers with no
// cloud_scan_run row -- the integration tests -- rely on.
func (s *AWSIAMScanner) WithGeneration(generation int) *AWSIAMScanner {
s.generation = generation
return s
}

// NewAWSIAMScanner constructs the scanner.
func NewAWSIAMScanner(db *gorm.DB, onboarding *AWSOnboardingService) *AWSIAMScanner {
return &AWSIAMScanner{
Expand Down Expand Up @@ -179,7 +203,17 @@ func (s *AWSIAMScanner) Scan(ctx context.Context, workspaceID, connectorID uuid.
// scan that died left it untouched. Checkpoints at this generation
// therefore mean "the previous attempt was interrupted part-way", and this
// run continues it rather than repeating its work.
generation := connector.ScanGeneration + 1
//
// ...except when the attempt died AFTER commitScan and before Publish, which
// is exactly the crashed-worker case Claim is built to recover. Then the
// connector HAS advanced, this would compute one too many, and the evidence
// the worker is writing against run.Generation would be stamped a generation
// behind the rows it explains. So the run's number wins when there is one;
// see WithGeneration.
generation := s.generation
if generation == 0 {
generation = connector.ScanGeneration + 1
}
resuming, err := s.checkpoints.HasAny(workspaceID, connectorID, generation)
if err != nil {
return nil, err
Expand Down Expand Up @@ -469,7 +503,10 @@ func (s *AWSIAMScanner) upsertRole(
// boundary.
counters["roles_detail_incomplete"]++
}
return s.recordIdentity(identity, counters)
// !DetailComplete is exactly the "partial read" the repository asks about:
// the attrs we just built carry "" / nil for description, boundary and tags
// because GetRole failed, not because the role lacks them.
return s.recordIdentity(identity, !role.DetailComplete, counters)
}

func (s *AWSIAMScanner) upsertUser(
Expand Down Expand Up @@ -497,14 +534,18 @@ func (s *AWSIAMScanner) upsertUser(
}); err != nil {
return nil, err
}
if err := s.recordIdentity(identity, counters); err != nil {
// ListUsers is the whole read for a user today -- there is no GetUser call to
// half-fail -- so a user's attrs are always as complete as this collector
// gets. Not partial. (When GetUser and the user's permissions boundary are
// added, this becomes conditional the way upsertRole's is.)
if err := s.recordIdentity(identity, false, counters); err != nil {
return nil, err
}
return identity, nil
}

func (s *AWSIAMScanner) recordIdentity(identity *models.CloudIdentity, counters map[string]int) error {
stored, created, err := s.identities.UpsertIdentity(identity)
func (s *AWSIAMScanner) recordIdentity(identity *models.CloudIdentity, partialRead bool, counters map[string]int) error {
stored, created, err := s.identities.UpsertIdentity(identity, partialRead)
if err != nil {
return fmt.Errorf("record identity %s: %w", identity.NativeID, err)
}
Expand Down
7 changes: 6 additions & 1 deletion services/cloud_aws_scan_worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,12 @@ func (w *AWSScanWorker) execute(ctx context.Context, run *models.CloudScanRun) e
evidence := NewObservationWriter(
w.db, run.WorkspaceID, run.ConnectorID, run.ID, run.Generation)

scanner := NewAWSIAMScanner(w.db, w.svc).WithEvidence(evidence)
// The same run.Generation the evidence writer above was built with. One
// number, read once, so a row and the observation explaining it can never
// disagree about which pass produced them.
scanner := NewAWSIAMScanner(w.db, w.svc).
WithEvidence(evidence).
WithGeneration(run.Generation)
permissionScanner := NewAWSPermissionScanner(w.db, w.svc).WithEvidence(evidence)
workloadScanner := NewAWSWorkloadScanner(w.db, w.svc).WithEvidence(evidence)

Expand Down
54 changes: 49 additions & 5 deletions services/cloud_aws_workload_scan.go
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,14 @@ type WorkloadSnapshot struct {
Complete bool
// Errors carries what went wrong per surface, for the coverage report.
Errors map[string]string
// ActivityTruncated says the activity pass hit activityIdentityCap and so
// read only a prefix of the account's identities. It is tracked apart from
// Errors because the two mean different things to a reader: nothing was
// denied, the read simply stopped short. It blocks reconciliation exactly
// as an error does, but reports the surface as `partial`, not `denied`.
ActivityTruncated bool
// ActivityTruncatedNote is the operator-facing explanation for the above.
ActivityTruncatedNote string
// Surfaces is the same information as Errors, but in the connector-level
// coverage shape (state + count, not just an error string), so the caller
// can fold it into the overall report instead of dropping it. Keyed
Expand Down Expand Up @@ -171,17 +179,32 @@ func (s *AWSWorkloadScanner) ScanFromSnapshot(

// ---- activity, which is global because IAM is -------------------------
s.scanActivity(ctx, workspaceID, snapshot, out)
if activityErr, failed := out.Errors["activity"]; failed {
// Three outcomes, most severe first. A denial outranks a short read: if some
// identity's report was refused we say so, even if the pass also stopped at
// the cap. Truncation alone is `partial` -- reached, but not exhaustively --
// because nothing was denied, the read simply ran out of budget.
switch {
case out.Errors["activity"] != "":
out.Surfaces["activity"] = models.SurfaceCoverage{
State: models.CloudCoverageDenied, Count: out.UsageWritten, Error: activityErr,
State: models.CloudCoverageDenied, Count: out.UsageWritten,
Error: out.Errors["activity"],
}
} else {
case out.ActivityTruncated:
out.Surfaces["activity"] = models.SurfaceCoverage{
State: models.CloudCoveragePartial, Count: out.UsageWritten,
Error: out.ActivityTruncatedNote,
}
default:
out.Surfaces["activity"] = models.SurfaceCoverage{
State: models.CloudCoverageReached, Count: out.UsageWritten,
}
}

out.Complete = len(out.Errors) == 0 && snapshot.Coverage.Complete()
// A truncated activity read blocks reconciliation for the same reason a
// denied one does: rows for the identities past the cap were never looked
// at, and absence from a read that stopped early is not evidence of
// deletion. This is the only place a bounded read could license a delete.
out.Complete = len(out.Errors) == 0 && !out.ActivityTruncated && snapshot.Coverage.Complete()

if out.Complete {
workloadsRemoved, usageRemoved, err := s.workloads.ReconcileGeneration(
Expand Down Expand Up @@ -671,14 +694,35 @@ func (s *AWSWorkloadScanner) scanActivity(
reader = reader.WithSleep(s.activitySleep)
}

identities, _, err := s.identities.ListIdentities(workspaceID, repositories.CloudIdentityFilter{
identities, total, err := s.identities.ListIdentities(workspaceID, repositories.CloudIdentityFilter{
ConnectorID: &snapshot.ConnectorID,
Limit: activityIdentityCap,
})
if err != nil {
out.Errors["activity"] = err.Error()
return
}
// The cap is a real ceiling, and above it this pass reads a PREFIX of the
// account -- ListIdentities orders by (kind, name, id), so it is the same
// prefix every scan and the same identities are starved every time.
//
// Recording that is not cosmetic. Without it the surface reports `reached`,
// Complete() stays true, and ReconcileGeneration deletes the cloud_usage
// rows carried from the previous generation for every identity past the cap
// -- destroying real "granted but never used" history on the authority of a
// read that never looked at those identities.
//
// This is deliberately NOT an Errors entry, though it gates reconciliation
// just as one does. A later per-identity failure writes that same key and
// would overwrite the message, and a denial should outrank a short read when
// deciding what to show the operator. ScanFromSnapshot consults both.
if total > int64(len(identities)) {
out.ActivityTruncated = true
out.ActivityTruncatedNote = fmt.Sprintf(
"read service activity for %d of %d identities: the per-scan cap is %d, "+
"so the rest were not attempted and nothing may be reconciled away",
len(identities), total, activityIdentityCap)
}

for _, identity := range identities {
activity, err := reader.ServiceActivityFor(ctx, identity.NativeID)
Expand Down
7 changes: 6 additions & 1 deletion services/cloud_gcp_scan_identities.go
Original file line number Diff line number Diff line change
Expand Up @@ -332,7 +332,12 @@ func (s *GCPScanner) upsertServiceAccount(
return gcpWrittenIdentity{}, err
}

stored, _, err := s.identities.UpsertIdentity(identity)
// Not a partial read: GCP lists service accounts with their full shape in one
// call, and a project it could not read fails the scan outright rather than
// yielding a half-populated account. So this always carries a complete blob,
// and attrs is refreshed as before -- Disabled in particular is a state
// transition that must land.
stored, _, err := s.identities.UpsertIdentity(identity, false)
if err != nil {
return gcpWrittenIdentity{}, err
}
Expand Down
Loading
Loading