From c30263f1f95fb3a23e16e59bd0f9c82083ff267b Mon Sep 17 00:00:00 2001 From: akash Date: Thu, 24 Sep 2026 01:07:57 +0530 Subject: [PATCH] Stop the AWS scan writing facts it did not observe Five write-path defects in AWS discovery. Four of them contradict a guarantee the codebase states in its own comments, and two destroy data that a rescan cannot recover. Activity cap no longer licenses a deletion. scanActivity discarded the row count from ListIdentities, so above activityIdentityCap it read a prefix of the account while the surface reported `reached`. Complete() stayed true and ReconcileGeneration deleted the cloud_usage history of every identity past the cap -- deterministically the same identities each scan, since the list is ordered (kind, name, id). It now reports `partial` and blocks reconciliation. Partial, not denied: nothing was refused, the read stopped short, and the operator should be told which. UpsertResource refreshes migration 021's ARN-derived columns. They were written on INSERT only, so every pre-021 resource took the conflict path forever and kept is_external = false -- a positive claim of locality about what may be a cross-account ARN, with a partial index built to query exactly that column. resource_account, is_external and object_key move together because 021's CHECKs couple them to each other and to kind. sensitivity_source and sensitivity_reason stay out: they explain a verdict this scan deliberately does not refresh. UpsertIdentity stops overwriting attrs on a partial read. iam.go says a writer "must not treat a partial role as authoritative and overwrite metadata an earlier complete read established"; it did. A throttled GetRole yields Tags=nil and boundary="", both omitempty, so a blind overwrite erased the tags that are the only ownership signal collected. The caller now says the read fell short -- detail_incomplete lives inside the attrs JSON and is omitempty, so the repository cannot infer it without parsing provider JSON. Merged rather than skipped, so the row keeps the old keys AND records that this read could not confirm them; skipping alone would leave it asserting it was fully read. One generation authority. cloud_scan_run.generation is assigned at Claim and preserved across re-claims, while Scan recomputed connector.ScanGeneration+1 per attempt. They diverge when a run is re-claimed after commitScan advanced the connector -- the crashed-worker path -- stamping entity rows a generation ahead of the observations written for the same pass, so evidence no longer joins to the inventory it explains. The worker now passes run.Generation. Unset keeps the derived behaviour for callers with no run row. ListObservations pages over a total order. observed_at alone is not one, and CloudTrail rows carry e.EventTime, which AWS reports to the second -- 12 rows share a timestamp on a real account here. id DESC matches the direction of the (workspace_id, subject, observed_at DESC) indexes. Tests: each fix has one that fails when the fix is removed, verified by removing it. Three more pin the ways a fix could be worse than the bug: a partially read role is still stamped as seen, a first sighting still lands, and an under-cap read still reconciles. Verified against a live AWS account: 13 resources degraded to pre-021 defaults were repaired by a real scan (the 2 left blank are S3 ARNs, which carry no account segment), and rows, permissions, usage and evidence all published on one generation. --- repository/cloud_identity_repository.go | 65 ++- repository/cloud_observation_repository.go | 17 +- repository/cloud_permission_repository.go | 23 + services/cloud_aws_iam_scan.go | 51 +- services/cloud_aws_scan_worker.go | 7 +- services/cloud_aws_workload_scan.go | 54 +- services/cloud_gcp_scan_identities.go | 7 +- .../integration/cloud_aws_write_path_test.go | 551 ++++++++++++++++++ 8 files changed, 754 insertions(+), 21 deletions(-) create mode 100644 tests/integration/cloud_aws_write_path_test.go diff --git a/repository/cloud_identity_repository.go b/repository/cloud_identity_repository.go index aaa3ac3b..445fe75a 100644 --- a/repository/cloud_identity_repository.go +++ b/repository/cloud_identity_repository.go @@ -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) @@ -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") } @@ -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{ @@ -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, diff --git a/repository/cloud_observation_repository.go b/repository/cloud_observation_repository.go index 025bccc4..7276ca91 100644 --- a/repository/cloud_observation_repository.go +++ b/repository/cloud_observation_repository.go @@ -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, , 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 } diff --git a/repository/cloud_permission_repository.go b/repository/cloud_permission_repository.go index 339d96b3..f6329dea 100644 --- a/repository/cloud_permission_repository.go +++ b/repository/cloud_permission_repository.go @@ -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, diff --git a/services/cloud_aws_iam_scan.go b/services/cloud_aws_iam_scan.go index dc6ff7e5..fe4b242c 100644 --- a/services/cloud_aws_iam_scan.go +++ b/services/cloud_aws_iam_scan.go @@ -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. @@ -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{ @@ -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 @@ -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( @@ -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) } diff --git a/services/cloud_aws_scan_worker.go b/services/cloud_aws_scan_worker.go index 1b83394b..b4f2f6c0 100644 --- a/services/cloud_aws_scan_worker.go +++ b/services/cloud_aws_scan_worker.go @@ -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) diff --git a/services/cloud_aws_workload_scan.go b/services/cloud_aws_workload_scan.go index 4c7c7eed..6304a238 100644 --- a/services/cloud_aws_workload_scan.go +++ b/services/cloud_aws_workload_scan.go @@ -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 @@ -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( @@ -671,7 +694,7 @@ 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, }) @@ -679,6 +702,27 @@ func (s *AWSWorkloadScanner) scanActivity( 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) diff --git a/services/cloud_gcp_scan_identities.go b/services/cloud_gcp_scan_identities.go index 53861f54..02e45457 100644 --- a/services/cloud_gcp_scan_identities.go +++ b/services/cloud_gcp_scan_identities.go @@ -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 } diff --git a/tests/integration/cloud_aws_write_path_test.go b/tests/integration/cloud_aws_write_path_test.go new file mode 100644 index 00000000..958d0ed2 --- /dev/null +++ b/tests/integration/cloud_aws_write_path_test.go @@ -0,0 +1,551 @@ +package integration + +// Regressions for the write-path defects found in the AWS discovery review. +// +// Every test here follows the same shape: establish a good row, then do the +// thing that used to corrupt or delete it, then assert it survived. They are +// written so that reverting the corresponding fix makes them fail -- a test +// that passes with its fix removed proves nothing, and this suite already has +// one documented instance of exactly that. + +import ( + "context" + "encoding/json" + "fmt" + "testing" + "time" + + "github.com/aws/aws-sdk-go-v2/aws" + iamtypes "github.com/aws/aws-sdk-go-v2/service/iam/types" + "github.com/google/uuid" + "gorm.io/gorm" + + "github.com/authsec-ai/authsec/models" + repositories "github.com/authsec-ai/authsec/repository" + "github.com/authsec-ai/authsec/services" +) + +// A throttled iam:GetRole must not erase what a complete read established. +// +// internal/awsdiscovery/iam.go states the guarantee outright: DetailComplete +// false means tags and the permissions boundary are UNKNOWN, not absent, and "a +// writer must not treat a partial role as authoritative". The writer did +// exactly that -- upsertRole builds attrs before it checks DetailComplete, and +// Tags/PermissionsBoundaryARN/Description are all `omitempty`, so they vanish +// from the JSON and the unconditional blob assignment overwrote the good one. +// +// IAM tags are the only ownership signal this system collects. +func TestAThrottledRoleReadDoesNotEraseTagsOrBoundary(t *testing.T) { + db := igaDB(t) + ws := newWorkspace(t, db, "ws-partial-role") + defer cleanIdentities(t, db, ws) + + const roleARN = "arn:aws:iam::429418377036:role/tagged-role" + const boundaryARN = "arn:aws:iam::429418377036:policy/the-ceiling" + + fake := newFakeIAM() + fake.roles = append(fake.roles, iamtypes.Role{ + Arn: aws.String(roleARN), + RoleName: aws.String("tagged-role"), + RoleId: aws.String("AROATAGGED0000000001"), + CreateDate: ago(72 * time.Hour), + Tags: []iamtypes.Tag{ + {Key: aws.String("owner"), Value: aws.String("platform-team")}, + }, + PermissionsBoundary: &iamtypes.AttachedPermissionsBoundary{ + PermissionsBoundaryArn: aws.String(boundaryARN), + }, + }) + + scanner, connectorID := scanFixture(t, db, ws, fake) + + // Scan one: GetRole succeeds, so the row is complete. + if _, err := scanner.Scan(context.Background(), ws, connectorID); err != nil { + t.Fatalf("baseline scan: %v", err) + } + before := awsAttrsOf(t, db, ws, roleARN) + if before.Tags["owner"] != "platform-team" { + t.Fatalf("test setup: baseline should have stored the owner tag, got %#v", before.Tags) + } + if before.PermissionsBoundaryARN != boundaryARN { + t.Fatalf("test setup: baseline should have stored the boundary, got %q", before.PermissionsBoundaryARN) + } + + // Scan two: GetRole is throttled. ListRoles still names the role, so it is + // seen and must not be reconciled away -- but nothing new is known about it. + fake.fail["GetRole"] = throttled("iam:GetRole") + if _, err := scanner.Scan(context.Background(), ws, connectorID); err != nil { + t.Fatalf("partial scan must not fail the run: %v", err) + } + + after := awsAttrsOf(t, db, ws, roleARN) + if after.Tags["owner"] != "platform-team" { + t.Fatalf("a throttled GetRole erased the owner tag: %#v", after.Tags) + } + if after.PermissionsBoundaryARN != boundaryARN { + t.Fatalf("a throttled GetRole erased the boundary ARN: %q", after.PermissionsBoundaryARN) + } + // The other half: preserving the blob must not leave the row claiming it + // was fully read. Without this the earlier complete read's absent + // detail_incomplete key survives, the row looks fresh, and constraintState + // answers from stale evidence as though it were current. + if !after.DetailIncomplete { + t.Fatal("a preserved blob must still record that THIS read fell short") + } + t.Log("PASS: tags and boundary survived, and the row is marked incomplete") +} + +// A later complete read must clear the flag again, or a single throttle would +// pin the role as incomplete forever. +func TestACompleteReReadClearsTheIncompleteFlag(t *testing.T) { + db := igaDB(t) + ws := newWorkspace(t, db, "ws-partial-recover") + defer cleanIdentities(t, db, ws) + + const roleARN = "arn:aws:iam::429418377036:role/recovering-role" + fake := newFakeIAM() + fake.roles = append(fake.roles, iamtypes.Role{ + Arn: aws.String(roleARN), + RoleName: aws.String("recovering-role"), + RoleId: aws.String("AROARECOVER000000001"), + CreateDate: ago(48 * time.Hour), + Tags: []iamtypes.Tag{{Key: aws.String("owner"), Value: aws.String("team-b")}}, + }) + scanner, connectorID := scanFixture(t, db, ws, fake) + + if _, err := scanner.Scan(context.Background(), ws, connectorID); err != nil { + t.Fatalf("baseline scan: %v", err) + } + fake.fail["GetRole"] = throttled("iam:GetRole") + if _, err := scanner.Scan(context.Background(), ws, connectorID); err != nil { + t.Fatalf("throttled scan: %v", err) + } + if !awsAttrsOf(t, db, ws, roleARN).DetailIncomplete { + t.Fatal("test setup: the throttled scan should have marked the role incomplete") + } + + // AWS recovers. + delete(fake.fail, "GetRole") + if _, err := scanner.Scan(context.Background(), ws, connectorID); err != nil { + t.Fatalf("recovery scan: %v", err) + } + + after := awsAttrsOf(t, db, ws, roleARN) + if after.DetailIncomplete { + t.Fatal("a complete read must clear detail_incomplete, not leave it stuck on") + } + if after.Tags["owner"] != "team-b" { + t.Fatalf("the recovery read lost the owner tag: %#v", after.Tags) + } + t.Log("PASS: the flag clears on recovery and the tags are still there") +} + +// The row must still be stamped as seen, or the fix would trade data loss for +// deletion -- a far worse bug. Guards against "fix" by skipping the upsert. +func TestAPartiallyReadRoleIsStillMarkedSeen(t *testing.T) { + db := igaDB(t) + ws := newWorkspace(t, db, "ws-partial-seen") + defer cleanIdentities(t, db, ws) + + const roleARN = "arn:aws:iam::429418377036:role/seen-role" + fake := newFakeIAM() + fake.roles = append(fake.roles, iamtypes.Role{ + Arn: aws.String(roleARN), + RoleName: aws.String("seen-role"), + RoleId: aws.String("AROASEEN000000000001"), + CreateDate: ago(24 * time.Hour), + Tags: []iamtypes.Tag{{Key: aws.String("owner"), Value: aws.String("team-a")}}, + }) + scanner, connectorID := scanFixture(t, db, ws, fake) + + if _, err := scanner.Scan(context.Background(), ws, connectorID); err != nil { + t.Fatalf("baseline scan: %v", err) + } + var firstGen int + db.Raw(`SELECT last_seen_generation FROM cloud_identity WHERE workspace_id = ? AND native_id = ?`, + ws, roleARN).Scan(&firstGen) + + fake.fail["GetRole"] = throttled("iam:GetRole") + if _, err := scanner.Scan(context.Background(), ws, connectorID); err != nil { + t.Fatalf("partial scan: %v", err) + } + + var secondGen int + db.Raw(`SELECT last_seen_generation FROM cloud_identity WHERE workspace_id = ? AND native_id = ?`, + ws, roleARN).Scan(&secondGen) + if secondGen <= firstGen { + t.Fatalf("a partially read role must still be stamped as seen: generation %d -> %d", + firstGen, secondGen) + } + t.Logf("PASS: generation advanced %d -> %d while attrs were preserved", firstGen, secondGen) +} + +// A brand-new role that could only be read partially must still land. There is +// no earlier blob to protect, and a thin row beats no row at all. +func TestAFirstSightingOfAPartialRoleStillWritesItsAttrs(t *testing.T) { + db := igaDB(t) + ws := newWorkspace(t, db, "ws-partial-first") + defer cleanIdentities(t, db, ws) + + const roleARN = "arn:aws:iam::429418377036:role/never-detailed" + fake := newFakeIAM() + fake.roles = append(fake.roles, iamtypes.Role{ + Arn: aws.String(roleARN), + RoleName: aws.String("never-detailed"), + RoleId: aws.String("AROANEVER00000000001"), + CreateDate: ago(time.Hour), + }) + fake.fail["GetRole"] = throttled("iam:GetRole") + + scanner, connectorID := scanFixture(t, db, ws, fake) + if _, err := scanner.Scan(context.Background(), ws, connectorID); err != nil { + t.Fatalf("scan: %v", err) + } + + attrs := awsAttrsOf(t, db, ws, roleARN) + if !attrs.DetailIncomplete { + t.Fatal("a first sighting read only from ListRoles must record detail_incomplete") + } + t.Log("PASS: a first-seen partial role is written and flagged incomplete") +} + +// Migration 021's three ARN-derived columns must be refreshed on conflict. +// +// They were written on INSERT only, so every resource that existed before 021 +// took the conflict path forever and kept is_external = false -- a positive +// claim of locality about what may be a cross-account ARN, with a partial index +// built to query exactly that column. +// +// Driven through the repository because that is where the defect lives. +func TestResourceARNDerivedColumnsAreRefreshedOnConflict(t *testing.T) { + db := igaDB(t) + ws := newWorkspace(t, db, "ws-resource-refresh") + defer func() { + db.Exec(`DELETE FROM cloud_resource WHERE workspace_id = ?`, ws) + db.Exec(`DELETE FROM cloud_connector WHERE workspace_id = ?`, ws) + }() + + svc, _ := newOnboarding(db, okVerifier()) + conn, _, err := svc.Onboard(context.Background(), ws, validInput(mustMint(t, ws)), "admin") + if err != nil { + t.Fatalf("onboard: %v", err) + } + connectorID := conn.ID + repo := repositories.NewCloudPermissionRepository(db) + const arn = "arn:aws:s3:::somebody-elses-bucket" + + // Stand in for a pre-021 row: discovered, but with the columns at their + // NOT NULL defaults because nothing populated them at the time. + stored, created, err2 := repo.UpsertResource(&models.CloudResource{ + WorkspaceID: ws, ConnectorID: connectorID, + Kind: "s3_bucket", NativeID: arn, Name: "somebody-elses-bucket", + ResourceAccount: "", IsExternal: false, ObjectKey: "", + Sensitivity: models.SensitivityLow, SensitivitySource: models.SensitivityFromHeuristic, + SensitivityReason: "baseline", LastSeenGeneration: 1, + }) + if err2 != nil || !created { + t.Fatalf("test setup: seed insert failed (created=%v): %v", created, err2) + } + + // Pretend a later ticket raised the verdict from a customer classification. + // This must survive: sensitivity and its two explaining columns are + // deliberately NOT refreshed. ('customer_classification' because + // cloud_resource_sensitivity_source_chk allows exactly four values.) + if err := db.Exec( + `UPDATE cloud_resource SET sensitivity = ?, sensitivity_source = ?, sensitivity_reason = ? + WHERE id = ?`, + models.SensitivityHigh, "customer_classification", "classified by the customer", stored.ID, + ).Error; err != nil { + t.Fatalf("test setup: raising the verdict failed: %v", err) + } + + // The next scan sees the same ARN and now derives its account correctly. + if _, _, err := repo.UpsertResource(&models.CloudResource{ + WorkspaceID: ws, ConnectorID: connectorID, + Kind: "s3_bucket", NativeID: arn, Name: "somebody-elses-bucket", + ResourceAccount: "999988887777", IsExternal: true, ObjectKey: "", + Sensitivity: models.SensitivityLow, SensitivitySource: models.SensitivityFromHeuristic, + SensitivityReason: "rule default", LastSeenGeneration: 2, + }); err != nil { + t.Fatalf("rescan upsert: %v", err) + } + + var got struct { + ResourceAccount string + IsExternal bool + Sensitivity string + SensitivitySource string + SensitivityReason string + } + db.Raw(`SELECT resource_account, is_external, sensitivity, sensitivity_source, sensitivity_reason + FROM cloud_resource WHERE id = ?`, stored.ID).Scan(&got) + + if got.ResourceAccount != "999988887777" || !got.IsExternal { + t.Fatalf("ARN-derived columns were not refreshed: account=%q is_external=%v", + got.ResourceAccount, got.IsExternal) + } + if got.Sensitivity != models.SensitivityHigh { + t.Fatalf("sensitivity must not be stamped back down: got %q", got.Sensitivity) + } + if got.SensitivitySource != "customer_classification" || + got.SensitivityReason != "classified by the customer" { + t.Fatalf("sensitivity_source/reason must travel with the verdict, not the scan: %q / %q", + got.SensitivitySource, got.SensitivityReason) + } + t.Log("PASS: the three ARN-derived columns refresh; the sensitivity trio does not") +} + +// Above activityIdentityCap the activity pass reads a PREFIX of the account and +// must say so, or reconciliation deletes the cloud_usage rows of every identity +// it never looked at -- on the authority of a read that reported "reached". +// +// Identities are inserted directly: the point is the repository count crossing +// the cap, not how they got there. +func TestATruncatedActivityReadBlocksReconciliation(t *testing.T) { + db := igaDB(t) + ws := newWorkspace(t, db, "ws-activity-cap") + defer cleanWorkloadTables(t, db, ws) + + scanner, snap := activityFixture(t, db, ws) + + // Baseline: a normal pass, under the cap, writes usage and reconciles. + if _, err := scanner.ScanFromSnapshot(context.Background(), ws, snap); err != nil { + t.Fatalf("baseline scan: %v", err) + } + + // A usage row from the previous generation. This is the row the bug deleted. + var idStr string + db.Raw(`SELECT id::text FROM cloud_identity WHERE workspace_id = ? LIMIT 1`, ws).Scan(&idStr) + if idStr == "" { + t.Fatal("test setup: expected the identity scan to have written identities") + } + anyIdentity := uuid.MustParse(idStr) + db.Exec(`INSERT INTO cloud_usage + (id, workspace_id, connector_id, identity_id, service, source, last_used_at, + last_seen_generation, first_seen_at, last_seen_at, row_updated_at) + VALUES (?, ?, ?, ?, 'legacy-service', 'service_last_accessed', NULL, ?, now(), now(), now())`, + uuid.New(), ws, snap.ConnectorID, anyIdentity, snap.Generation) + + // Push the account over the cap. Names are prefixed so they sort AFTER the + // fixture's roles under ORDER BY (kind, name, id) -- i.e. into the starved + // tail, which is the deterministic part of the defect. + for i := 0; i < 520; i++ { + db.Exec(`INSERT INTO cloud_identity + (id, workspace_id, connector_id, kind, native_id, name, enabled, attrs, + last_seen_generation, first_seen_at, last_seen_at, row_updated_at) + VALUES (?, ?, ?, 'iam_role', ?, ?, true, '{}', ?, now(), now(), now())`, + uuid.New(), ws, snap.ConnectorID, + fmt.Sprintf("arn:aws:iam::429418377036:role/zz-bulk-%04d", i), + fmt.Sprintf("zz-bulk-%04d", i), snap.Generation) + } + + var usageBefore int64 + db.Raw(`SELECT count(*) FROM cloud_usage WHERE workspace_id = ?`, ws).Scan(&usageBefore) + + snap.Generation++ + out, err := scanner.ScanFromSnapshot(context.Background(), ws, snap) + if err != nil { + t.Fatalf("a truncated activity read must not fail the scan: %v", err) + } + + if !out.ActivityTruncated { + t.Fatal("reading 500 of 520+ identities must be reported as truncated") + } + if out.Complete { + t.Fatal("a truncated activity read must not license reconciliation") + } + if got := out.Surfaces["activity"].State; got != models.CloudCoveragePartial { + t.Fatalf("truncation is partial, not %q: the read was not denied, it stopped short", got) + } + + var usageAfter int64 + db.Raw(`SELECT count(*) FROM cloud_usage WHERE workspace_id = ?`, ws).Scan(&usageAfter) + if usageAfter < usageBefore { + t.Fatalf("a truncated read deleted usage history: %d -> %d", usageBefore, usageAfter) + } + t.Logf("PASS: truncation reported partial, %d usage row(s) preserved", usageAfter) +} + +// The counterpart: under the cap nothing changes. Without this, "always report +// truncated" would pass the test above and permanently disable reconciliation. +func TestAnUntruncatedActivityReadStillReconciles(t *testing.T) { + db := igaDB(t) + ws := newWorkspace(t, db, "ws-activity-under-cap") + defer cleanWorkloadTables(t, db, ws) + + scanner, snap := activityFixture(t, db, ws) + + out, err := scanner.ScanFromSnapshot(context.Background(), ws, snap) + if err != nil { + t.Fatalf("scan: %v", err) + } + if out.ActivityTruncated { + t.Fatal("a handful of identities is not a truncated read") + } + if got := out.Surfaces["activity"].State; got != models.CloudCoverageReached { + t.Fatalf("an activity read under the cap is reached, got %q", got) + } + if !out.Complete { + t.Fatalf("a clean scan must still reconcile: %v", out.Errors) + } + t.Log("PASS: an under-cap activity read still reports reached and complete") +} + +// activityFixture wires a workload scanner whose ONLY live surface is activity. +// The compute APIs are left unset so they report not-configured rather than +// denied, which keeps these tests about the activity gate and nothing else. +func activityFixture(t *testing.T, db *gorm.DB, ws uuid.UUID) (*services.AWSWorkloadScanner, *services.IAMSnapshot) { + t.Helper() + svc, _ := newOnboarding(db, okVerifier()) + conn, _, err := svc.Onboard(context.Background(), ws, singleRegionInput(mustMint(t, ws)), "admin") + if err != nil { + t.Fatalf("onboard: %v", err) + } + snap, err := services.NewAWSIAMScanner(db, svc).WithIAMAPI(populatedIAM()). + Scan(context.Background(), ws, conn.ID) + if err != nil { + t.Fatalf("identity scan: %v", err) + } + used := time.Now().Add(-90 * 24 * time.Hour) + activity := &fakeActivity{services: []iamtypes.ServiceLastAccessed{{ + ServiceNamespace: aws.String("s3"), + ServiceName: aws.String("Amazon S3"), + LastAuthenticated: &used, + TotalAuthenticatedEntities: aws.Int32(1), + }}} + noSleep := func(context.Context, time.Duration) error { return nil } + l, e, c, p := populatedWorkloads() + return services.NewAWSWorkloadScanner(db, svc). + WithWorkloadAPIs(l, e, c, p). + WithActivityAPI(activity, noSleep), snap +} + +// awsAttrsOf reads one identity's attrs blob back as AWS attributes, straight +// from the column rather than through the scanner, so the assertion is about +// what is actually stored. +func awsAttrsOf(t *testing.T, db *gorm.DB, ws uuid.UUID, nativeID string) models.AWSIdentityAttrs { + t.Helper() + var raw string + if err := db.Raw( + `SELECT attrs::text FROM cloud_identity WHERE workspace_id = ? AND native_id = ?`, + ws, nativeID, + ).Scan(&raw).Error; err != nil { + t.Fatalf("read attrs for %s: %v", nativeID, err) + } + if raw == "" { + t.Fatalf("no identity row for %s", nativeID) + } + var attrs models.AWSIdentityAttrs + if err := json.Unmarshal([]byte(raw), &attrs); err != nil { + t.Fatalf("decode attrs for %s: %v", nativeID, err) + } + return attrs +} + +// The crash-recovery path B5 exists for, end to end. +// +// A worker dies after commitScan advanced the connector but before Publish. +// The run is re-claimed and keeps its generation; the connector has already +// moved. The old code recomputed connector.ScanGeneration+1 and so stamped +// entity rows one generation AHEAD of the cloud_observation rows the worker was +// writing against run.Generation -- evidence that no longer joins to the +// inventory it explains. +// +// This also pins the side effect of the fix: holding the generation steady +// makes the previous attempt's checkpoints visible, so the resume path that was +// previously unreachable now engages. Assert it engages SAFELY -- nothing the +// first attempt wrote may be lost. +func TestAReclaimedRunKeepsOneGenerationForRowsAndEvidence(t *testing.T) { + db := igaDB(t) + ws := newWorkspace(t, db, "ws-reclaim-generation") + defer cleanWorkloadTables(t, db, ws) + + svc, _ := newOnboarding(db, okVerifier()) + conn, _, err := svc.Onboard(context.Background(), ws, singleRegionInput(mustMint(t, ws)), "admin") + if err != nil { + t.Fatalf("onboard: %v", err) + } + runs := repositories.NewCloudScanRunRepository(db) + if _, err := runs.Enqueue(ws, conn.ID, "manual"); err != nil { + t.Fatalf("enqueue: %v", err) + } + + now := time.Now() + first, err := runs.Claim("worker-a", time.Minute, now) + if err != nil || first == nil { + t.Fatalf("first claim: %v %v", first, err) + } + + fake := populatedIAM() + // Attempt one: IAM phase completes (so commitScan advances the connector), + // then the permission phase writes checkpoints, then the worker dies. + snap1, err := services.NewAWSIAMScanner(db, svc).WithIAMAPI(fake). + WithGeneration(first.Generation).Scan(context.Background(), ws, conn.ID) + if err != nil { + t.Fatalf("attempt one iam scan: %v", err) + } + if _, err := services.NewAWSPermissionScanner(db, svc).WithIAMAPI(fake). + ScanFromSnapshot(context.Background(), ws, snap1); err != nil { + t.Fatalf("attempt one permission scan: %v", err) + } + + var connGen int + db.Raw(`SELECT scan_generation FROM cloud_connector WHERE id = ?`, conn.ID).Scan(&connGen) + if connGen != first.Generation { + t.Fatalf("test setup: commitScan should have advanced the connector to %d, got %d", + first.Generation, connGen) + } + var permsBefore int64 + db.Raw(`SELECT count(*) FROM cloud_permission WHERE workspace_id = ?`, ws).Scan(&permsBefore) + if permsBefore == 0 { + t.Fatal("test setup: attempt one should have written permissions") + } + + // worker-a never published. Its lease lapses and worker-b takes the run. + second, err := runs.Claim("worker-b", time.Minute, now.Add(2*time.Minute)) + if err != nil || second == nil { + t.Fatalf("reclaim: %v %v", second, err) + } + if second.Generation != first.Generation { + t.Fatalf("test setup: a reclaimed run keeps its generation, %d -> %d", + first.Generation, second.Generation) + } + + // Attempt two, the way the worker runs it. + snap2, err := services.NewAWSIAMScanner(db, svc).WithIAMAPI(fake). + WithGeneration(second.Generation).Scan(context.Background(), ws, conn.ID) + if err != nil { + t.Fatalf("attempt two iam scan: %v", err) + } + + // THE defect: the rows and the evidence must agree on which pass wrote them. + if snap2.Generation != second.Generation { + t.Fatalf("rows stamped %d while the run's evidence is stamped %d -- "+ + "evidence no longer joins to the inventory it explains", + snap2.Generation, second.Generation) + } + + if _, err := services.NewAWSPermissionScanner(db, svc).WithIAMAPI(fake). + ScanFromSnapshot(context.Background(), ws, snap2); err != nil { + t.Fatalf("attempt two permission scan: %v", err) + } + + // The side effect, checked rather than assumed: resume may skip work, but + // reconciliation at the same generation must not delete what attempt one + // already wrote. + var permsAfter int64 + db.Raw(`SELECT count(*) FROM cloud_permission WHERE workspace_id = ?`, ws).Scan(&permsAfter) + if permsAfter < permsBefore { + t.Fatalf("the resumed attempt destroyed permissions the first one wrote: %d -> %d", + permsBefore, permsAfter) + } + var orphaned int64 + db.Raw(`SELECT count(*) FROM cloud_permission + WHERE workspace_id = ? AND last_seen_generation <> ?`, ws, second.Generation).Scan(&orphaned) + if orphaned != 0 { + t.Fatalf("%d permission row(s) carry a generation other than %d", + orphaned, second.Generation) + } + t.Logf("PASS: one generation (%d) across rows and evidence; %d permission(s) preserved", + second.Generation, permsAfter) +}