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) +}