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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 29 additions & 0 deletions pkg/common/identifiers.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,3 +35,32 @@ func SanitizeReservationID(id, fallbackPrefix string) string {
}
return s
}

// idempotencyIDTokenLen is how many leading hex characters of the
// idempotency token are folded into a derived reservation ID. 40 hex chars =
// 160 bits, collision-free at any realistic purchase volume, and short enough
// to keep the prefixed result under every AWS reserved-instance/node ID length
// limit (RDS being the tightest).
const idempotencyIDTokenLen = 40

// IdempotentReservationID derives a deterministic, AWS-safe reservation ID from
// an idempotency token (issue #641). The same token always yields the same ID,
// so a re-driven purchase reuses the identical customer-supplied reservation ID
// and AWS rejects the duplicate server-side (RDS/ElastiCache/MemoryDB each
// return a *AlreadyExists* fault). Returns "" when token is empty so the caller
// keeps its prior non-idempotent (timestamp-based) ID behaviour for call sites
// that supply no token (e.g. the CLI path).
//
// prefix should be a short, lowercase, hyphen-terminated service tag (e.g.
// "rds-id-") so the reservation is identifiable in the console; the token is
// hex so the result needs no further sanitisation beyond SanitizeReservationID's
// invariants.
func IdempotentReservationID(prefix, token string) string {
if token == "" {
return ""
}
if len(token) > idempotencyIDTokenLen {
token = token[:idempotencyIDTokenLen]
}
return SanitizeReservationID(prefix+token, prefix)
}
33 changes: 33 additions & 0 deletions pkg/common/identifiers_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
package common

import (
"strings"
"testing"

"github.com/stretchr/testify/assert"
)

func TestIdempotentReservationID_DeterministicAndSafe(t *testing.T) {
token := DeriveIdempotencyToken("exec-1", 0)

a := IdempotentReservationID("rds-id-", token)
b := IdempotentReservationID("rds-id-", token)

assert.Equal(t, a, b, "same token must yield the same reservation ID")
assert.True(t, strings.HasPrefix(a, "rds-id-"), "must carry the prefix for console identifiability")
assert.NotContains(t, a, "--", "must not contain consecutive hyphens")
assert.False(t, strings.HasSuffix(a, "-"), "must not end with a hyphen")
// prefix (7) + 40 hex chars = 47, well under RDS's tightest ID length cap.
assert.LessOrEqual(t, len(a), 60, "must stay under the tightest AWS reservation-ID length cap")
}

func TestIdempotentReservationID_DistinctTokensDistinctIDs(t *testing.T) {
id0 := IdempotentReservationID("rds-id-", DeriveIdempotencyToken("exec-1", 0))
id1 := IdempotentReservationID("rds-id-", DeriveIdempotencyToken("exec-1", 1))
assert.NotEqual(t, id0, id1, "different recs in an execution must get different IDs")
}

func TestIdempotentReservationID_EmptyTokenReturnsEmpty(t *testing.T) {
assert.Equal(t, "", IdempotentReservationID("rds-id-", ""),
"empty token must yield empty so the caller keeps its non-idempotent fallback")
}
17 changes: 17 additions & 0 deletions pkg/common/tokens.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,3 +41,20 @@ func DeriveIdempotencyToken(executionID string, recIndex int) string {
sum := sha256.Sum256([]byte(fmt.Sprintf("%s:%d", executionID, recIndex)))
return hex.EncodeToString(sum[:])
}

// MaskToken returns a log-safe representation of an idempotency/approval token:
// the first 8 characters followed by an ellipsis, never the full value. This
// keeps just enough of the prefix to correlate log lines for a single purchase
// while avoiding emitting the whole caller-supplied token into persistent logs
// (a stable per-execution identifier that should not leak verbatim). An empty
// token yields "(none)"; a token of 8 chars or fewer is returned unchanged
// since there is nothing left to redact.
func MaskToken(token string) string {
if token == "" {
return "(none)"
}
if len(token) <= 8 {
return token
}
return token[:8] + "..."
}
21 changes: 21 additions & 0 deletions pkg/common/tokens_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -64,3 +64,24 @@ func TestDeriveIdempotencyToken_FitsClientTokenLimit(t *testing.T) {
require.NoError(t, err)
assert.Len(t, raw, 32)
}

func TestMaskToken_NeverEmitsFullToken(t *testing.T) {
// CodeRabbit PR #652: the "already exists" skip-purchase log lines must not
// emit the raw idempotency token (a stable per-execution identifier). The
// masked form keeps only an 8-char prefix for log correlation.
full := DeriveIdempotencyToken("exec-abc-123", 0) // 64-char hex digest
masked := MaskToken(full)

assert.NotEqual(t, full, masked, "masked form must differ from the raw token")
assert.NotContains(t, masked, full, "masked output must not contain the full token")
assert.Less(t, len(masked), len(full), "masked output must be shorter than the raw token")
assert.Equal(t, full[:8]+"...", masked, "masked form is an 8-char prefix plus ellipsis")
assert.Len(t, masked, 11, "8 prefix chars + 3-char ellipsis")
}

func TestMaskToken_EmptyAndShort(t *testing.T) {
assert.Equal(t, "(none)", MaskToken(""), "empty token must be reported as (none)")
assert.Equal(t, "abc", MaskToken("abc"), "tokens of <=8 chars have nothing to redact")
assert.Equal(t, "12345678", MaskToken("12345678"), "exactly 8 chars is returned unchanged")
assert.Equal(t, "12345678...", MaskToken("123456789"), "9 chars is truncated to 8 + ellipsis")
}
97 changes: 96 additions & 1 deletion providers/aws/services/elasticache/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,9 @@ package elasticache

import (
"context"
"errors"
"fmt"
"log"
"sort"
"time"

Expand Down Expand Up @@ -124,7 +126,26 @@ func (c *Client) PurchaseCommitment(ctx context.Context, rec common.Recommendati
return result, result.Error
}

reservationID := common.SanitizeReservationID(fmt.Sprintf("elasticache-%s-%d", rec.ResourceType, time.Now().Unix()), "elasticache-reserved-")
// When an idempotency token is supplied (issue #641) the reservation ID is
// derived deterministically from it, so a re-drive sends the identical
// ReservedCacheNodeId and ElastiCache rejects the duplicate server-side
// (ReservedCacheNodeAlreadyExistsFault). Otherwise keep the prior
// timestamp-based ID (non-idempotent path).
reservationID := common.IdempotentReservationID("elasticache-id-", opts.IdempotencyToken)
if reservationID == "" {
reservationID = common.SanitizeReservationID(fmt.Sprintf("elasticache-%s-%d", rec.ResourceType, time.Now().Unix()), "elasticache-reserved-")
}

// Idempotency dedupe guard (issue #641): short-circuit if a reservation
// already exists under the derived ID; fail loud on lookup error.
if existingID, shortCircuit, guardErr := c.idempotencyGuard(ctx, opts.IdempotencyToken, reservationID); guardErr != nil {
result.Error = guardErr
return result, result.Error
} else if shortCircuit {
result.Success = true
result.CommitmentID = existingID
return result, nil
}

input := &elasticache.PurchaseReservedCacheNodesOfferingInput{
ReservedCacheNodesOfferingId: aws.String(offeringID),
Expand All @@ -135,6 +156,11 @@ func (c *Client) PurchaseCommitment(ctx context.Context, rec common.Recommendati

response, err := c.client.PurchaseReservedCacheNodesOffering(ctx, input)
if err != nil {
if existingID, recovered := c.recoverAlreadyExists(ctx, opts.IdempotencyToken, reservationID, err); recovered {
result.Success = true
result.CommitmentID = existingID
return result, nil
}
result.Error = fmt.Errorf("failed to purchase Reserved Cache Node: %w", err)
return result, result.Error
}
Expand All @@ -153,6 +179,75 @@ func (c *Client) PurchaseCommitment(ctx context.Context, rec common.Recommendati
return result, nil
}

// findReservationByID looks for an active or payment-pending reserved cache node
// with the given ReservedCacheNodeId (issue #641), so a re-driven purchase can
// short-circuit instead of buying a second node. Retired/expired nodes are
// excluded (same state filter as GetExistingCommitments).
func (c *Client) findReservationByID(ctx context.Context, reservationID string) (string, bool, error) {
response, err := c.client.DescribeReservedCacheNodes(ctx, &elasticache.DescribeReservedCacheNodesInput{
ReservedCacheNodeId: aws.String(reservationID),
})
if err != nil {
// ElastiCache returns ReservedCacheNodeNotFound for an unknown reservation
// ID; treat that as "not found" (a first-time purchase), not a lookup
// failure. Any other error is a genuine failure.
var notFound *types.ReservedCacheNodeNotFoundFault
if errors.As(err, &notFound) {
return "", false, nil
}
return "", false, fmt.Errorf("failed to describe reserved cache nodes for idempotency check: %w", err)
}
for _, node := range response.ReservedCacheNodes {
state := aws.ToString(node.State)
if state != "active" && state != "payment-pending" {
continue
}
if node.ReservedCacheNodeId != nil {
return aws.ToString(node.ReservedCacheNodeId), true, nil
}
}
return "", false, nil
}

// idempotencyGuard short-circuits a re-drive (issue #641): when token is set, it
// reports (existingID, true, nil) if a reservation already exists under
// reservationID, ("", false, nil) for a first-time purchase, or a fail-loud
// error on lookup failure. With an empty token it is a no-op.
func (c *Client) idempotencyGuard(ctx context.Context, token, reservationID string) (string, bool, error) {
if token == "" {
return "", false, nil
}
existingID, found, lookupErr := c.findReservationByID(ctx, reservationID)
if lookupErr != nil {
return "", false, fmt.Errorf("idempotency lookup failed before ElastiCache purchase (refusing to purchase to avoid a possible double-buy): %w", lookupErr)
}
if found {
log.Printf("ElastiCache reservation for idempotency token %s already exists (%s); skipping purchase (issue #641 re-drive)", common.MaskToken(token), existingID)
return existingID, true, nil
}
return "", false, nil
}

// recoverAlreadyExists handles the native server-side dedupe backstop (issue
// #641): if the by-ID guard missed the existing reservation but AWS rejected the
// duplicate ID with ReservedCacheNodeAlreadyExistsFault, it re-Describes by ID
// and returns (existingID, true) so the re-drive recovers it instead of erroring.
func (c *Client) recoverAlreadyExists(ctx context.Context, token, reservationID string, purchaseErr error) (string, bool) {
if token == "" {
return "", false
}
var already *types.ReservedCacheNodeAlreadyExistsFault
if !errors.As(purchaseErr, &already) {
return "", false
}
existingID, found, lookupErr := c.findReservationByID(ctx, reservationID)
if lookupErr == nil && found {
log.Printf("ElastiCache reservation %s already existed at purchase time; treating as idempotent re-drive (issue #641)", existingID)
return existingID, true
}
return "", false
}

// findOfferingID finds the appropriate Reserved Cache Node offering ID
func (c *Client) findOfferingID(ctx context.Context, rec common.Recommendation) (string, error) {
details, ok := rec.Details.(*common.CacheDetails)
Expand Down
106 changes: 106 additions & 0 deletions providers/aws/services/elasticache/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -411,6 +411,112 @@ func TestClient_ConvertPaymentOption(t *testing.T) {
}
}

func idemRec() common.Recommendation {
return common.Recommendation{
Service: common.ServiceCache,
ResourceType: "cache.m6g.large",
Count: 1,
PaymentOption: "all-upfront",
Term: "1yr",
Details: &common.CacheDetails{Engine: "redis", NodeType: "cache.m6g.large"},
}
}

func expectECOffering(m *MockElastiCacheClient) {
m.On("DescribeReservedCacheNodesOfferings", mock.Anything, mock.Anything).
Return(&elasticache.DescribeReservedCacheNodesOfferingsOutput{
ReservedCacheNodesOfferings: []types.ReservedCacheNodesOffering{
{
ReservedCacheNodesOfferingId: aws.String("offering-1"),
CacheNodeType: aws.String("cache.m6g.large"),
ProductDescription: aws.String("redis"),
OfferingType: aws.String("All Upfront"),
Duration: aws.Int32(31536000),
},
},
}, nil)
}

func TestClient_PurchaseCommitment_Idempotent_GuardShortCircuits(t *testing.T) {
mockEC := &MockElastiCacheClient{}
client := &Client{client: mockEC, region: "eu-west-1"}
token := common.DeriveIdempotencyToken("exec-1", 0)
derivedID := common.IdempotentReservationID("elasticache-id-", token)

expectECOffering(mockEC)
mockEC.On("DescribeReservedCacheNodes", mock.Anything, mock.MatchedBy(func(in *elasticache.DescribeReservedCacheNodesInput) bool {
return aws.ToString(in.ReservedCacheNodeId) == derivedID
})).Return(&elasticache.DescribeReservedCacheNodesOutput{
ReservedCacheNodes: []types.ReservedCacheNode{{ReservedCacheNodeId: aws.String(derivedID), State: aws.String("active")}},
}, nil)

result, err := client.PurchaseCommitment(context.Background(), idemRec(), common.PurchaseOptions{IdempotencyToken: token})
assert.NoError(t, err)
assert.True(t, result.Success)
assert.Equal(t, derivedID, result.CommitmentID)
mockEC.AssertNotCalled(t, "PurchaseReservedCacheNodesOffering", mock.Anything, mock.Anything)
}

func TestClient_PurchaseCommitment_Idempotent_NotFoundProceeds(t *testing.T) {
mockEC := &MockElastiCacheClient{}
client := &Client{client: mockEC, region: "eu-west-1"}
token := common.DeriveIdempotencyToken("exec-2", 0)
derivedID := common.IdempotentReservationID("elasticache-id-", token)

expectECOffering(mockEC)
mockEC.On("DescribeReservedCacheNodes", mock.Anything, mock.Anything).
Return((*elasticache.DescribeReservedCacheNodesOutput)(nil), &types.ReservedCacheNodeNotFoundFault{})
mockEC.On("PurchaseReservedCacheNodesOffering", mock.Anything, mock.MatchedBy(func(in *elasticache.PurchaseReservedCacheNodesOfferingInput) bool {
return aws.ToString(in.ReservedCacheNodeId) == derivedID
})).Return(&elasticache.PurchaseReservedCacheNodesOfferingOutput{
ReservedCacheNode: &types.ReservedCacheNode{ReservedCacheNodeId: aws.String(derivedID)},
}, nil)

result, err := client.PurchaseCommitment(context.Background(), idemRec(), common.PurchaseOptions{IdempotencyToken: token})
assert.NoError(t, err)
assert.True(t, result.Success)
assert.Equal(t, derivedID, result.CommitmentID)
mockEC.AssertExpectations(t)
}

func TestClient_PurchaseCommitment_Idempotent_AlreadyExistsRecovers(t *testing.T) {
mockEC := &MockElastiCacheClient{}
client := &Client{client: mockEC, region: "eu-west-1"}
token := common.DeriveIdempotencyToken("exec-3", 0)
derivedID := common.IdempotentReservationID("elasticache-id-", token)

expectECOffering(mockEC)
mockEC.On("DescribeReservedCacheNodes", mock.Anything, mock.Anything).
Return((*elasticache.DescribeReservedCacheNodesOutput)(nil), &types.ReservedCacheNodeNotFoundFault{}).Once()
mockEC.On("PurchaseReservedCacheNodesOffering", mock.Anything, mock.Anything).
Return((*elasticache.PurchaseReservedCacheNodesOfferingOutput)(nil), &types.ReservedCacheNodeAlreadyExistsFault{})
mockEC.On("DescribeReservedCacheNodes", mock.Anything, mock.Anything).
Return(&elasticache.DescribeReservedCacheNodesOutput{
ReservedCacheNodes: []types.ReservedCacheNode{{ReservedCacheNodeId: aws.String(derivedID), State: aws.String("active")}},
}, nil).Once()

result, err := client.PurchaseCommitment(context.Background(), idemRec(), common.PurchaseOptions{IdempotencyToken: token})
assert.NoError(t, err)
assert.True(t, result.Success)
assert.Equal(t, derivedID, result.CommitmentID)
}

func TestClient_PurchaseCommitment_Idempotent_FailLoudOnLookupError(t *testing.T) {
mockEC := &MockElastiCacheClient{}
client := &Client{client: mockEC, region: "eu-west-1"}
token := common.DeriveIdempotencyToken("exec-4", 0)

expectECOffering(mockEC)
mockEC.On("DescribeReservedCacheNodes", mock.Anything, mock.Anything).
Return((*elasticache.DescribeReservedCacheNodesOutput)(nil), fmt.Errorf("access denied"))

result, err := client.PurchaseCommitment(context.Background(), idemRec(), common.PurchaseOptions{IdempotencyToken: token})
assert.Error(t, err)
assert.False(t, result.Success)
assert.Contains(t, err.Error(), "refusing to purchase")
mockEC.AssertNotCalled(t, "PurchaseReservedCacheNodesOffering", mock.Anything, mock.Anything)
}

func TestCreatePurchaseTags_IncludesPurchaseAutomation(t *testing.T) {
c := &Client{}
rec := common.Recommendation{ResourceType: "cache.m5.large", Region: "us-east-1"}
Expand Down
Loading
Loading