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
4 changes: 4 additions & 0 deletions internal/analytics/collector_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -455,6 +455,10 @@ func (m *mockConfigStore) ClaimMarketplaceListingSlot(_ context.Context, _ strin
return true, nil
}

func (m *mockConfigStore) ClaimRIExchangeIdempotencyKey(_ context.Context, _ string, _ time.Duration) (bool, error) {
return true, nil
}

func (m *mockConfigStore) StampOfferingClass(_ context.Context, _, _ string) error {
return nil
}
Expand Down
75 changes: 69 additions & 6 deletions internal/api/handler_ri_exchange.go
Original file line number Diff line number Diff line change
Expand Up @@ -572,7 +572,10 @@ func validateAzureExchangeTargets(targets []AzureExchangeTargetBody, subscriptio
}
scope := azureBillingScopeID(subscriptionID)
for i, t := range targets {
if t.SKU == "" {
// Trimmed, like location below: the #1642 idempotency fingerprint
// normalizes surrounding whitespace away, so a blank-but-not-empty
// SKU would reach it as "" and collapse onto other blank spellings.
if strings.TrimSpace(t.SKU) == "" {
return NewClientError(400, fmt.Sprintf("targets[%d].sku is required", i))
}
// Blank-but-not-empty is rejected too: targetLocations trims before
Expand Down Expand Up @@ -1203,10 +1206,28 @@ func (h *Handler) executeAzureExchange(ctx context.Context, req *events.LambdaFu
return nil, err
}

// Submit-time idempotency (#1642). The re-quote above mints a FRESH
// session on every request, so Azure's own session-level replay
// protection cannot see two POSTs of one logical exchange as duplicates:
// without this claim a client that times out mid-LRO and retries commits
// the exchange twice, each half individually under the cap. Taken here,
// last, so no gate rejection ever leaves a claim behind.
err = h.claimExchangeSubmit(ctx, azureExchangeIdempotencyKey(body))
if err != nil {
return nil, err
}

result, err := client.ExecuteExchange(ctx, preview.SessionID)
if err != nil {
logging.Errorf("azure exchange execution failed: %v", err)
return nil, mapAzureExchangeError("exchange execution failed", err)
// Non-4xx failures here are ambiguous: BeginPost may already have
// submitted the exchange, and a ctx cancellation mid-poll looks
// identical to one that never reached Azure. Say so rather than
// reporting a flat failure that invites the retry this claim now
// refuses (#1642).
return nil, mapAzureExchangeError(
"exchange execution failed and may already have been submitted to Azure; "+
"verify the reservation state in the Azure portal before retrying", err)
}

logging.Infof("azure ri-exchange executed: subscription=%s session=%s status=%s", body.SubscriptionID, result.SessionID, result.Status)
Expand Down Expand Up @@ -1616,8 +1637,15 @@ func firstNonEmptyCurrency(instances []ec2svc.ConvertibleRI) string {
}

// validateTargets checks each entry in targets for a non-empty, UUID-shaped
// offering_id. Extracted so both getExchangeQuote and validateExecuteExchangeBody
// share the same check without exceeding the gocyclo threshold.
// offering_id and a positive count. Extracted so both getExchangeQuote and
// validateExecuteExchangeBody share the same check without exceeding the
// gocyclo threshold.
//
// The count check mirrors pkg/exchange.validateTargets, which applies the same
// rule but only once ExecuteExchange is already running. Repeating it here
// moves the refusal ahead of the #1642 submit claim -- a request that can never
// commit must not leave a claim behind -- and turns what pkg/exchange would
// surface as an opaque 500 into a 400 naming the offending field.
func validateTargets(targets []ExchangeTargetBody) error {
for i, t := range targets {
if t.OfferingID == "" {
Expand All @@ -1630,6 +1658,25 @@ func validateTargets(targets []ExchangeTargetBody) error {
"did you paste an instance type by mistake?",
i, t.OfferingID))
}
if t.Count < 1 {
return NewClientError(400, fmt.Sprintf("targets[%d].count must be >= 1, got %d", i, t.Count))
}
}
return nil
}

// validateExchangeRIIDs checks the source list and each id in it. Each id
// individually, not just the list length: a blank entry reaches the #1642
// submit fingerprint as an empty component, so two different blank spellings
// of one request would claim the same key.
func validateExchangeRIIDs(riIDs []string) error {
if len(riIDs) == 0 {
return NewClientError(400, "ri_ids is required")
}
for i, id := range riIDs {
if strings.TrimSpace(id) == "" {
return NewClientError(400, fmt.Sprintf("ri_ids[%d] is empty", i))
}
}
return nil
}
Expand Down Expand Up @@ -1684,15 +1731,21 @@ func (h *Handler) getExchangeQuote(ctx context.Context, req *events.LambdaFuncti
// cyclomatic-complexity threshold; every branch here becomes a
// separate test case so the logic stays inspectable.
func validateExecuteExchangeBody(body ExchangeExecuteRequestBody) error {
if len(body.RIIDs) == 0 {
return NewClientError(400, "ri_ids is required")
if err := validateExchangeRIIDs(body.RIIDs); err != nil {
return err
}
if len(body.Targets) == 0 && body.TargetOfferingID == "" {
return NewClientError(400, "either targets[] or target_offering_id is required")
}
if err := validateTargets(body.Targets); err != nil {
return err
}
// The legacy singleton's count, which validateTargets above does not see.
// Same reasoning as the targets[] count: refuse before the submit claim
// rather than inside ExecuteExchange, after it.
if len(body.Targets) == 0 && body.TargetCount < 1 {
return NewClientError(400, fmt.Sprintf("target_count must be >= 1, got %d", body.TargetCount))
}
if body.MaxPaymentDueUSD == "" {
return NewClientError(400, "max_payment_due_usd is required as a safety guardrail")
}
Expand Down Expand Up @@ -1761,6 +1814,16 @@ func (h *Handler) executeExchange(ctx context.Context, req *events.LambdaFunctio
return nil, err
}

// Submit-time idempotency (#1642). AcceptReservedInstancesExchangeQuote
// carries no ClientToken, so AWS will happily accept the same exchange
// twice; a client that times out and retries otherwise double-spends with
// each half individually under the cap. Taken last, after every gate, so
// no rejected request leaves a claim behind.
err = h.claimExchangeSubmit(ctx, awsExchangeIdempotencyKey(cloudAccountID, body))
if err != nil {
return nil, err
}

exchangeID, quote, err := exchange.ExecuteExchange(ctx, exchange.ExchangeExecuteRequest{
Region: region,
ReservedIDs: body.RIIDs,
Expand Down
8 changes: 4 additions & 4 deletions internal/api/handler_ri_exchange_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1522,7 +1522,7 @@ func TestExecuteExchange_EmptyRegionReturns400(t *testing.T) {

_, err := h.executeExchange(context.Background(), &events.LambdaFunctionURLRequest{
Headers: map[string]string{"authorization": "Bearer test-token"},
Body: `{"ri_ids":["ri-123"],"target_offering_id":"off-1","max_payment_due_usd":"10.00"}`,
Body: `{"ri_ids":["ri-123"],"target_offering_id":"off-1","target_count":1,"max_payment_due_usd":"10.00"}`,
})
require.Error(t, err)
ce, ok := IsClientError(err)
Expand Down Expand Up @@ -1571,7 +1571,7 @@ func TestExecuteExchange_PermissionConstraintsDenied(t *testing.T) {
}
_, err := h.executeExchange(ctx, &events.LambdaFunctionURLRequest{
Headers: map[string]string{"authorization": "Bearer exchange-token"},
Body: `{"ri_ids":["ri-123"],"target_offering_id":"off-1","max_payment_due_usd":"250.50","region":"eu-central-1"}`,
Body: `{"ri_ids":["ri-123"],"target_offering_id":"off-1","target_count":1,"max_payment_due_usd":"250.50","region":"eu-central-1"}`,
})
require.Error(t, err)
ce, ok := IsClientError(err)
Expand Down Expand Up @@ -1606,7 +1606,7 @@ func TestExecuteExchange_AccountResolutionErrorFailsClosed(t *testing.T) {
}
_, err := h.executeExchange(ctx, &events.LambdaFunctionURLRequest{
Headers: map[string]string{"authorization": "Bearer exchange-token"},
Body: `{"ri_ids":["ri-123"],"target_offering_id":"off-1","max_payment_due_usd":"250.50","region":"eu-central-1"}`,
Body: `{"ri_ids":["ri-123"],"target_offering_id":"off-1","target_count":1,"max_payment_due_usd":"250.50","region":"eu-central-1"}`,
})
require.Error(t, err)
assert.Contains(t, err.Error(), "resolve cloud account scope")
Expand Down Expand Up @@ -1643,7 +1643,7 @@ func TestExecuteExchange_UnattributedAccountStillConstrained(t *testing.T) {
}
_, err := h.executeExchange(ctx, &events.LambdaFunctionURLRequest{
Headers: map[string]string{"authorization": "Bearer exchange-token"},
Body: `{"ri_ids":["ri-123"],"target_offering_id":"off-1","max_payment_due_usd":"250.50","region":"eu-central-1"}`,
Body: `{"ri_ids":["ri-123"],"target_offering_id":"off-1","target_count":1,"max_payment_due_usd":"250.50","region":"eu-central-1"}`,
})
require.Error(t, err)
ce, ok := IsClientError(err)
Expand Down
33 changes: 33 additions & 0 deletions internal/api/openapi.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -775,6 +775,13 @@ paths:
the fresh quote carries policy errors, omits a net payable amount,
is denominated in a different currency than requested, or exceeds
`max_payment_due`.
Submits are deduplicated for 15 minutes on a fingerprint of
`subscription_id` plus the sources and targets with their quantities,
so a retry after a client timeout returns 409 rather than committing
the exchange a second time. The spend cap and currency are not part of
that fingerprint: raising the cap does not make it a different
purchase. See the 409 response: it does not assert that the earlier
submit succeeded, only that it claimed the window.
parameters:
- $ref: '#/components/parameters/CSRFToken'
requestBody:
Expand Down Expand Up @@ -820,6 +827,8 @@ paths:
$ref: '#/components/responses/Forbidden'
'404':
$ref: '#/components/responses/NotFound'
'409':
$ref: '#/components/responses/Conflict'
'422':
$ref: '#/components/responses/UnprocessableEntity'

Expand Down Expand Up @@ -917,6 +926,13 @@ paths:
exchanges submitted to AWS cannot be rolled back. Non-admin users must
be explicitly granted `execute:ri-exchange` via a custom group; there
is no default user-role grant.
Submits are deduplicated for 15 minutes on a fingerprint of the
deployment's cloud account and region plus the source RIs and targets
with their counts, so a retry after a client timeout returns 409 rather
than committing the exchange a second time. `max_payment_due_usd` is
not part of that fingerprint: raising the cap does not make it a
different purchase. See the 409 response: it does not assert that the
earlier submit succeeded, only that it claimed the window.
parameters:
- $ref: '#/components/parameters/CSRFToken'
requestBody:
Expand Down Expand Up @@ -951,6 +967,8 @@ paths:
$ref: '#/components/responses/Unauthorized'
'403':
$ref: '#/components/responses/Forbidden'
'409':
$ref: '#/components/responses/Conflict'

/api/ri-exchange/config:
get:
Expand Down Expand Up @@ -1881,6 +1899,21 @@ components:
application/json:
schema:
$ref: '#/components/schemas/Error'
Conflict:
description: >
An identical submit already holds the idempotency claim, so THIS
request was not executed. The claim is retained unconditionally once
taken, including when the earlier submit failed, so this response
asserts nothing about that submit's outcome: it may still be running,
it may have committed, its outcome may be unresolved (the provider
call can fail after the operation was already submitted), or it may
have failed without committing anything. Do not read a 409 as
confirmation that the earlier submit succeeded. Verify its outcome
with the provider before resubmitting.
content:
application/json:
schema:
$ref: '#/components/schemas/Error'
Comment thread
coderabbitai[bot] marked this conversation as resolved.
RateLimited:
description: Rate limit exceeded
content:
Expand Down
Loading
Loading