-
Notifications
You must be signed in to change notification settings - Fork 6
fix(cli): write audit records and check writability on the CSV purchase path #2101
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -421,17 +421,40 @@ func executePurchasePipeline(ctx context.Context, awsCfg aws.Config, recs []comm | |
| } | ||
| result, status := purchaseSingleRec(ctx, awsCfg, rec, i+1, isDryRun, cfg) | ||
| results = append(results, result) | ||
| auditRec := common.NewAuditRecord(runID, rec, result, status, isDryRun, common.PurchaseSourceCLI) | ||
| if err := common.WriteAuditRecord(auditRec, cfg.AuditLog); err != nil { | ||
| log.Printf("Warning: failed to write audit record: %v", err) | ||
| } | ||
| writePurchaseAuditRecord(runID, rec, result, status, isDryRun, cfg.AuditLog) | ||
| if !isDryRun && i < len(recs)-1 && os.Getenv("DISABLE_PURCHASE_DELAY") != "true" { | ||
| time.Sleep(PurchaseDelaySeconds * time.Second) | ||
| } | ||
| } | ||
| return results | ||
| } | ||
|
|
||
| // writePurchaseAuditRecord writes a single purchase's audit record. Shared by | ||
| // both purchase entry points -- executePurchasePipeline (the main pipeline) | ||
| // and processPurchaseLoop (the --input-csv path) -- so every recommendation | ||
| // that reaches a purchase attempt, dry-run or real, is recorded to | ||
| // cfg.AuditLog regardless of which one produced it. Before #1609, | ||
| // processPurchaseLoop never wrote a record at all, so CSV-mode purchases left | ||
| // no audit trail: on a partial failure there was no durable, per-recommendation | ||
| // record of which rows succeeded, so an operator could only re-run the whole | ||
| // file, which is a double purchase for the rows that already succeeded. | ||
| func writePurchaseAuditRecord(runID string, rec common.Recommendation, result common.PurchaseResult, status string, isDryRun bool, auditLogPath string) { | ||
| auditRec := common.NewAuditRecord(runID, rec, result, status, isDryRun, common.PurchaseSourceCLI) | ||
| if err := common.WriteAuditRecord(auditRec, auditLogPath); err != nil { | ||
| log.Printf("Warning: failed to write audit record: %v", err) | ||
| } | ||
| } | ||
|
|
||
| // purchaseAuditStatus derives the audit status for a completed (non-dry-run) | ||
| // purchase attempt. Callers handle the dry-run ("skipped") case separately, | ||
| // since that never reaches a PurchaseResult from an actual API call. | ||
| func purchaseAuditStatus(result common.PurchaseResult) string { | ||
| if result.Success { | ||
| return "success" | ||
| } | ||
| return "error" | ||
| } | ||
|
|
||
| // purchaseSingleRec executes or dry-runs a single purchase and returns the result + audit status. | ||
| func purchaseSingleRec(ctx context.Context, awsCfg aws.Config, rec common.Recommendation, index int, isDryRun bool, cfg Config) (purchaseResult common.PurchaseResult, auditStatus string) { | ||
| AppLogger.Printf(" [%d] %s %s %s (count=%d)\n", index, rec.Service, rec.Region, rec.ResourceType, rec.Count) | ||
|
|
@@ -450,12 +473,11 @@ func purchaseSingleRec(ctx context.Context, awsCfg aws.Config, rec common.Recomm | |
| } | ||
|
|
||
| result := executePurchase(ctx, rec, rec.Region, index, serviceClient, cfg) | ||
| status := "success" | ||
| if !result.Success { | ||
| status = "error" | ||
| AppLogger.Printf(" ❌ %v\n", result.Error) | ||
| } else { | ||
| status := purchaseAuditStatus(result) | ||
| if result.Success { | ||
| AppLogger.Printf(" ✅ %s\n", result.CommitmentID) | ||
| } else { | ||
| AppLogger.Printf(" ❌ %v\n", result.Error) | ||
| } | ||
| return result, status | ||
| } | ||
|
|
@@ -489,6 +511,47 @@ func runCSVPathOrFatal(ctx context.Context, cfg Config) { | |
| } | ||
| } | ||
|
|
||
| // prepareCSVPurchaseRun validates and loads everything runToolFromCSV needs | ||
| // before the per-service purchase loop: the audit log writability, the CSV | ||
| // file, filtering/sizing, and the AWS config. Extracted to keep | ||
| // runToolFromCSV under the project's gocyclo budget. | ||
| // | ||
| // The audit-log check runs first and before any cloud API call, matching the | ||
| // non-CSV path (CheckAuditLogWritable in runToolMultiService). Before #1609 | ||
| // this check ran only on the non-CSV path, so a CSV-mode purchase run could | ||
| // reach real purchase calls with no way to have written a durable, | ||
| // per-recommendation audit record even in principle. | ||
| // | ||
| // A nil recs with a nil error means "nothing to process after filtering", | ||
| // which the caller treats as success rather than an error. | ||
| func prepareCSVPurchaseRun(ctx context.Context, cfg Config, csvModeCoverage float64) (recs []common.Recommendation, awsCfg aws.Config, runID string, err error) { | ||
| if err = CheckAuditLogWritable(cfg.AuditLog); err != nil { | ||
| return nil, aws.Config{}, "", fmt.Errorf("cannot write audit log: %w", err) | ||
| } | ||
|
|
||
| AppLogger.Printf("📄 Reading recommendations from CSV: %s\n", cfg.CSVInput) | ||
| recs, err = loadRecommendationsFromCSV(cfg.CSVInput) | ||
| if err != nil { | ||
| return nil, aws.Config{}, "", fmt.Errorf("failed to read CSV file: %w", err) | ||
| } | ||
| AppLogger.Printf("✅ Loaded %d recommendations from CSV\n", len(recs)) | ||
|
|
||
| recs, err = filterAndAdjustRecommendations(recs, csvModeCoverage, cfg) | ||
| if err != nil { | ||
| return nil, aws.Config{}, "", err | ||
| } | ||
| if len(recs) == 0 { | ||
| return nil, aws.Config{}, "", nil | ||
| } | ||
|
|
||
| awsCfg, err = loadAWSConfig(ctx, cfg) | ||
| if err != nil { | ||
| return nil, aws.Config{}, "", fmt.Errorf("failed to load AWS config: %w", err) | ||
| } | ||
|
|
||
| return recs, awsCfg, uuid.New().String(), nil | ||
| } | ||
|
|
||
| // runToolFromCSV processes recommendations from a CSV input file. | ||
| // It returns an error instead of exiting so the orchestration glue is | ||
| // unit-testable; the caller (runCSVPathOrFatal) turns errors fatal. | ||
|
|
@@ -498,32 +561,15 @@ func runToolFromCSV(ctx context.Context, cfg Config) error { | |
|
|
||
| csvModeCoverage := determineCSVCoverage(cfg) | ||
|
|
||
| AppLogger.Printf("📄 Reading recommendations from CSV: %s\n", cfg.CSVInput) | ||
|
|
||
| // Read recommendations from CSV | ||
| recs, err := loadRecommendationsFromCSV(cfg.CSVInput) | ||
| if err != nil { | ||
| return fmt.Errorf("failed to read CSV file: %w", err) | ||
| } | ||
|
|
||
| AppLogger.Printf("✅ Loaded %d recommendations from CSV\n", len(recs)) | ||
|
|
||
| // Filter and adjust recommendations | ||
| recs, err = filterAndAdjustRecommendations(recs, csvModeCoverage, cfg) | ||
| recs, awsCfg, runID, err := prepareCSVPurchaseRun(ctx, cfg, csvModeCoverage) | ||
| if err != nil { | ||
| return err | ||
| } | ||
|
|
||
| if len(recs) == 0 { | ||
| AppLogger.Println("⚠️ No recommendations to process after filtering") | ||
| return nil | ||
| } | ||
|
|
||
| awsCfg, err := loadAWSConfig(ctx, cfg) | ||
| if err != nil { | ||
| return fmt.Errorf("failed to load AWS config: %w", err) | ||
| } | ||
|
|
||
| // Create account alias cache for lookup | ||
| accountCache := NewAccountAliasCache(awsCfg) | ||
|
|
||
|
|
@@ -581,7 +627,7 @@ func runToolFromCSV(ctx context.Context, cfg Config) error { | |
| allAdjustedRecs = append(allAdjustedRecs, recs...) | ||
|
|
||
| // Process purchases for this region | ||
| regionResults := processPurchaseLoop(ctx, recs, region, isDryRun, serviceClient, cfg) | ||
| regionResults := processPurchaseLoop(ctx, recs, region, isDryRun, serviceClient, cfg, runID) | ||
| serviceResults = append(serviceResults, regionResults...) | ||
| } | ||
|
|
||
|
|
@@ -721,8 +767,11 @@ func processService(ctx context.Context, awsCfg aws.Config, recClient provider.R | |
| return serviceRecs, serviceResults | ||
| } | ||
|
|
||
| // processPurchaseLoop processes purchases for a single region (used by CSV mode). | ||
| func processPurchaseLoop(ctx context.Context, recs []common.Recommendation, region string, isDryRun bool, serviceClient provider.ServiceClient, cfg Config) []common.PurchaseResult { | ||
| // processPurchaseLoop processes purchases for a single region (used by CSV | ||
| // mode). runID groups every recommendation processed across the whole CSV | ||
| // run into one audit trail, matching how executePurchasePipeline (the main | ||
| // pipeline) generates one runID per invocation. | ||
| func processPurchaseLoop(ctx context.Context, recs []common.Recommendation, region string, isDryRun bool, serviceClient provider.ServiceClient, cfg Config, runID string) []common.PurchaseResult { | ||
| results := make([]common.PurchaseResult, 0, len(recs)) | ||
|
|
||
| for j := range recs { | ||
|
|
@@ -731,8 +780,10 @@ func processPurchaseLoop(ctx context.Context, recs []common.Recommendation, regi | |
| AppLogger.Printf(" 💳 Purchasing %d instances\n", rec.Count) | ||
|
|
||
| var result common.PurchaseResult | ||
| var status string | ||
| if isDryRun { | ||
| result = createDryRunResult(rec, region, j+1, cfg) | ||
| status = "skipped" | ||
| } else { | ||
| // Ask for confirmation before proceeding with purchases (only on first item) | ||
| if j == 0 { | ||
|
|
@@ -744,20 +795,25 @@ func processPurchaseLoop(ctx context.Context, recs []common.Recommendation, regi | |
| } | ||
|
|
||
| if !ConfirmPurchase(totalInstances, totalSavings, cfg.SkipConfirmation) { | ||
| // User canceled - return canceled results for all | ||
| // User canceled - return canceled results for all. No audit | ||
| // record is written for a declined run, matching the | ||
| // non-CSV path: runPurchaseAndReport returns before ever | ||
| // calling executePurchasePipeline when the user declines. | ||
| return createCancelledResults(recs, region, cfg) | ||
| } | ||
| } | ||
|
|
||
| // Execute actual purchase | ||
| result = executePurchase(ctx, rec, region, j+1, serviceClient, cfg) | ||
| status = purchaseAuditStatus(result) | ||
|
|
||
| // Add delay between purchases to avoid rate limiting | ||
| if j < len(recs)-1 && os.Getenv("DISABLE_PURCHASE_DELAY") != "true" { | ||
| time.Sleep(PurchaseDelaySeconds * time.Second) | ||
| } | ||
| } | ||
|
|
||
| writePurchaseAuditRecord(runID, rec, result, status, isDryRun, cfg.AuditLog) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🗄️ Data Integrity & Integration | 🟠 Major | ⚡ Quick win Write each CSV audit record before the purchase delay. For every real purchase except the last in a region, 🤖 Prompt for AI Agents |
||
| results = append(results, result) | ||
|
|
||
| if result.Success { | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy lift
Stop the purchase run when an audit write fails.
If the audit log fills or becomes unwritable after preflight,
writePurchaseAuditRecordlogs a warning and the CSV loop continues purchasing. Those purchases can finish without the records needed to reconcile a partial run. Return the write error to the purchase pipeline and stop further purchases; report the failure to the CLI caller. The pinned audit writer explicitly returns I/O errors, while its writability check only opens and closes the file. (github.com)🤖 Prompt for AI Agents