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
5 changes: 5 additions & 0 deletions internal/analytics/collector_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"time"

"github.com/LeanerCloud/CUDly/internal/config"
"github.com/LeanerCloud/CUDly/pkg/common"
"github.com/LeanerCloud/CUDly/pkg/ladder"
"github.com/jackc/pgx/v5"
"github.com/stretchr/testify/assert"
Expand Down Expand Up @@ -320,6 +321,10 @@ func (m *mockConfigStore) GetRIExchangeDailySpend(ctx context.Context, date time
func (m *mockConfigStore) CancelAllPendingExchanges(ctx context.Context) (int64, error) {
return 0, nil
}

func (m *mockConfigStore) CancelPendingExchangesByOrigin(_ context.Context, _ common.ExchangeOrigin) (int64, error) {
return 0, nil
}
func (m *mockConfigStore) GetStaleProcessingExchanges(ctx context.Context, olderThan time.Duration) ([]config.RIExchangeRecord, error) {
return nil, nil
}
Expand Down
9 changes: 9 additions & 0 deletions internal/config/interfaces.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (

"github.com/jackc/pgx/v5"

"github.com/LeanerCloud/CUDly/pkg/common"
"github.com/LeanerCloud/CUDly/pkg/ladder"
)

Expand Down Expand Up @@ -209,6 +210,14 @@ type StoreInterface interface {
FailRIExchange(ctx context.Context, id string, errorMsg string) error
GetRIExchangeDailySpend(ctx context.Context, date time.Time) (string, error)
CancelAllPendingExchanges(ctx context.Context) (int64, error)
// CancelPendingExchangesByOrigin cancels only pending records whose origin
// matches:
// - common.ExchangeOriginStandalone: cancels WHERE ladder_run_id IS NULL
// - common.ExchangeOriginLadder: cancels WHERE ladder_run_id IS NOT NULL
// The origin is validated at the boundary and an unknown value fails loud.
// This prevents the standalone ri_exchange_reshape task from wiping out
// ladder-linked pending reshapes and vice versa (gap G10 / issue #1348).
CancelPendingExchangesByOrigin(ctx context.Context, origin common.ExchangeOrigin) (int64, error)
GetStaleProcessingExchanges(ctx context.Context, olderThan time.Duration) ([]RIExchangeRecord, error)

// Cloud accounts
Expand Down
88 changes: 73 additions & 15 deletions internal/config/store_postgres.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import (
"time"

"github.com/LeanerCloud/CUDly/internal/database"
"github.com/LeanerCloud/CUDly/pkg/common"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgconn"
Expand Down Expand Up @@ -2200,8 +2201,8 @@ func (s *PostgresStore) SaveRIExchangeRecord(ctx context.Context, record *RIExch
source_instance_type, source_count, target_offering_id,
target_instance_type, target_count, payment_due,
status, approval_token, error, mode, completed_at, expires_at,
created_at, updated_at, created_by_user_id
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20)
created_at, updated_at, created_by_user_id, ladder_run_id
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21)
`

_, err := s.db.Exec(ctx, query,
Expand All @@ -2225,6 +2226,7 @@ func (s *PostgresStore) SaveRIExchangeRecord(ctx context.Context, record *RIExch
record.CreatedAt,
record.UpdatedAt,
record.CreatedByUserID,
record.LadderRunID,
)

if err != nil {
Expand All @@ -2242,7 +2244,7 @@ func (s *PostgresStore) GetRIExchangeRecord(ctx context.Context, id string) (*RI
target_instance_type, target_count, payment_due::text,
status, approval_token, error, mode,
created_at, updated_at, completed_at, expires_at,
created_by_user_id, approved_by
created_by_user_id, approved_by, ladder_run_id
FROM ri_exchange_history
WHERE id = $1
`
Expand All @@ -2267,7 +2269,7 @@ func (s *PostgresStore) GetRIExchangeRecordByToken(ctx context.Context, token st
target_instance_type, target_count, payment_due::text,
status, approval_token, error, mode,
created_at, updated_at, completed_at, expires_at,
created_by_user_id, approved_by
created_by_user_id, approved_by, ladder_run_id
FROM ri_exchange_history
WHERE approval_token = $1
`
Expand All @@ -2292,7 +2294,7 @@ func (s *PostgresStore) GetRIExchangeHistory(ctx context.Context, since time.Tim
target_instance_type, target_count, payment_due::text,
status, approval_token, error, mode,
created_at, updated_at, completed_at, expires_at,
created_by_user_id, approved_by
created_by_user_id, approved_by, ladder_run_id
FROM ri_exchange_history
WHERE created_at >= $1
ORDER BY created_at DESC
Expand All @@ -2317,7 +2319,7 @@ func (s *PostgresStore) TransitionRIExchangeStatus(ctx context.Context, id strin
target_instance_type, target_count, payment_due::text,
status, approval_token, error, mode,
created_at, updated_at, completed_at, expires_at,
created_by_user_id, approved_by
created_by_user_id, approved_by, ladder_run_id
`

records, err := s.queryRIExchangeRecords(ctx, query, id, fromStatus, toStatus, actor)
Expand Down Expand Up @@ -2432,7 +2434,9 @@ func (s *PostgresStore) GetRIExchangeDailySpend(ctx context.Context, date time.T
return total, nil
}

// CancelAllPendingExchanges cancels all pending RI exchange records.
// CancelAllPendingExchanges cancels all pending RI exchange records regardless
// of origin. Kept for interface compatibility; new callers should prefer
// CancelPendingExchangesByOrigin to avoid cross-origin contamination.
func (s *PostgresStore) CancelAllPendingExchanges(ctx context.Context) (int64, error) {
query := `
UPDATE ri_exchange_history
Expand All @@ -2448,6 +2452,52 @@ func (s *PostgresStore) CancelAllPendingExchanges(ctx context.Context) (int64, e
return result.RowsAffected(), nil
}

// CancelPendingExchangesByOrigin cancels only pending records that match the
// given origin (gap G10 / issue #1348):
// - common.ExchangeOriginStandalone: cancels WHERE ladder_run_id IS NULL
// - common.ExchangeOriginLadder: cancels WHERE ladder_run_id IS NOT NULL
//
// The origin is validated at this boundary; an unknown value fails loud rather
// than silently cancelling the wrong partition on a money path.
//
// DELIBERATE COARSE PARTITION: the ExchangeOriginLadder branch cancels EVERY
// ladder-linked pending record (ladder_run_id IS NOT NULL) across ALL ladder
// runs and configs, not just the current run's. This is acceptable today
// because the ladder never creates pending exchange records: buildRIExchangeConfig
// forces Mode=auto, which completes or fails immediately without leaving a
// pending row. Per-run / per-config scoping needs an additional selection key
// (e.g. the specific ladder_run_id or config_id) and is tracked in TODO(#1367).
func (s *PostgresStore) CancelPendingExchangesByOrigin(ctx context.Context, origin common.ExchangeOrigin) (int64, error) {
if err := origin.Validate(); err != nil {
return 0, fmt.Errorf("CancelPendingExchangesByOrigin: %w", err)
}

var query string
switch origin {
case common.ExchangeOriginLadder:
query = `
UPDATE ri_exchange_history
SET status = 'cancelled', updated_at = NOW()
WHERE status = 'pending'
AND ladder_run_id IS NOT NULL
`
default: // common.ExchangeOriginStandalone (validated non-unknown above)
query = `
UPDATE ri_exchange_history
SET status = 'cancelled', updated_at = NOW()
WHERE status = 'pending'
AND ladder_run_id IS NULL
`
}

result, err := s.db.Exec(ctx, query)
if err != nil {
return 0, fmt.Errorf("failed to cancel pending exchanges by origin %q: %w", origin, err)
}

return result.RowsAffected(), nil
}

// GetStaleProcessingExchanges returns processing exchanges older than the given duration.
func (s *PostgresStore) GetStaleProcessingExchanges(ctx context.Context, olderThan time.Duration) ([]RIExchangeRecord, error) {
query := `
Expand All @@ -2456,7 +2506,7 @@ func (s *PostgresStore) GetStaleProcessingExchanges(ctx context.Context, olderTh
target_instance_type, target_count, payment_due::text,
status, approval_token, error, mode,
created_at, updated_at, completed_at, expires_at,
created_by_user_id, approved_by
created_by_user_id, approved_by, ladder_run_id
FROM ri_exchange_history
WHERE status = 'processing' AND updated_at < NOW() - $1::interval
`
Expand All @@ -2477,7 +2527,7 @@ func (s *PostgresStore) queryRIExchangeRecords(ctx context.Context, query string
var record RIExchangeRecord
var approvalToken, errStr sql.NullString
var completedAt, expiresAt sql.NullTime
var createdByUserID, approvedBy sql.NullString
var createdByUserID, approvedBy, ladderRunID sql.NullString

err := rows.Scan(
&record.ID,
Expand All @@ -2501,6 +2551,7 @@ func (s *PostgresStore) queryRIExchangeRecords(ctx context.Context, query string
&expiresAt,
&createdByUserID,
&approvedBy,
&ladderRunID,
)
if err != nil {
return nil, fmt.Errorf("failed to scan ri exchange record: %w", err)
Expand All @@ -2518,12 +2569,9 @@ func (s *PostgresStore) queryRIExchangeRecords(ctx context.Context, query string
if expiresAt.Valid {
record.ExpiresAt = &expiresAt.Time
}
if createdByUserID.Valid {
record.CreatedByUserID = &createdByUserID.String
}
if approvedBy.Valid {
record.ApprovedBy = &approvedBy.String
}
record.CreatedByUserID = nullPtrFromNullString(createdByUserID)
record.ApprovedBy = nullPtrFromNullString(approvedBy)
record.LadderRunID = nullPtrFromNullString(ladderRunID)

records = append(records, record)
}
Expand Down Expand Up @@ -3243,3 +3291,13 @@ func nullStringFromString(s string) sql.NullString {
}
return sql.NullString{String: s, Valid: true}
}

// nullPtrFromNullString converts a scanned sql.NullString to a *string for
// optional pointer fields. Returns nil when the DB value is NULL.
func nullPtrFromNullString(ns sql.NullString) *string {
if !ns.Valid {
return nil
}
s := ns.String
return &s
}
16 changes: 16 additions & 0 deletions internal/config/store_postgres_nil_db_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@ import (
"time"

"github.com/stretchr/testify/assert"

"github.com/LeanerCloud/CUDly/pkg/common"
)

// TestPostgresStore_SaveRIExchangeRecord_GeneratesID verifies that SaveRIExchangeRecord
Expand Down Expand Up @@ -165,6 +167,20 @@ func TestPostgresStore_CancelAllPendingExchanges_NilDB(t *testing.T) {
assert.True(t, panicked, "expected panic with nil db connection")
}

// TestPostgresStore_CancelPendingExchangesByOrigin_NilDB exercises the method entry.
func TestPostgresStore_CancelPendingExchangesByOrigin_NilDB(t *testing.T) {
store := NewPostgresStore(nil)
ctx := context.Background()

// Pass a VALID origin so the method reaches the nil db (an invalid origin
// would fail loud at Validate before touching the connection).
panicked := callWithRecover(func() {
_, _ = store.CancelPendingExchangesByOrigin(ctx, common.ExchangeOriginStandalone)
})

assert.True(t, panicked, "expected panic with nil db connection")
}

// TestPostgresStore_GetStaleProcessingExchanges_NilDB exercises the method entry.
func TestPostgresStore_GetStaleProcessingExchanges_NilDB(t *testing.T) {
store := NewPostgresStore(nil)
Expand Down
78 changes: 77 additions & 1 deletion internal/config/store_postgres_pgxmock_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@ import (
"github.com/pashagolub/pgxmock/v4"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"

"github.com/LeanerCloud/CUDly/pkg/common"
)

// newMock creates a pgxmock pool with regexp query matching.
Expand Down Expand Up @@ -988,6 +990,7 @@ func riExchangeRow(now time.Time) []interface{} {
"manual",
now, now, sql.NullTime{}, sql.NullTime{},
sql.NullString{}, sql.NullString{}, // created_by_user_id, approved_by
sql.NullString{}, // ladder_run_id (NULL for standalone)
}
}

Expand All @@ -997,7 +1000,7 @@ var riExchangeCols = []string{
"target_instance_type", "target_count", "payment_due",
"status", "approval_token", "error", "mode",
"created_at", "updated_at", "completed_at", "expires_at",
"created_by_user_id", "approved_by",
"created_by_user_id", "approved_by", "ladder_run_id",
}

func TestPGXMock_GetRIExchangeRecord_Success(t *testing.T) {
Expand Down Expand Up @@ -1032,6 +1035,7 @@ func TestPGXMock_GetRIExchangeRecord_WithTimestamps(t *testing.T) {
"auto",
now, now, sql.NullTime{Valid: true, Time: now}, sql.NullTime{Valid: true, Time: now},
sql.NullString{}, sql.NullString{}, // created_by_user_id, approved_by
sql.NullString{}, // ladder_run_id
}
rows := pgxmock.NewRows(riExchangeCols).AddRow(row...)
mock.ExpectQuery("SELECT").WithArgs(pgxmock.AnyArg()).WillReturnRows(rows)
Expand Down Expand Up @@ -1815,6 +1819,78 @@ func TestPGXMock_CancelAllPendingExchanges_Success(t *testing.T) {
assert.Equal(t, int64(3), n)
}

// ─── CancelPendingExchangesByOrigin ──────────────────────────────────────────

// The regex pins each test to its own WHERE clause so a branch mix-up (both
// branches emitting the same SQL) fails the test. "ladder_run_id IS NULL" is
// NOT a substring of "ladder_run_id IS NOT NULL", so the two matchers are
// mutually exclusive.
func TestPGXMock_CancelPendingExchangesByOrigin_Standalone(t *testing.T) {
mock := newMock(t)
store := storeWith(mock)
ctx := context.Background()

// Standalone origin must target WHERE ladder_run_id IS NULL only.
mock.ExpectExec("ladder_run_id IS NULL").WillReturnResult(pgxmock.NewResult("UPDATE", 2))

n, err := store.CancelPendingExchangesByOrigin(ctx, common.ExchangeOriginStandalone)
require.NoError(t, err)
assert.Equal(t, int64(2), n)
assert.NoError(t, mock.ExpectationsWereMet())
}

func TestPGXMock_CancelPendingExchangesByOrigin_Ladder(t *testing.T) {
mock := newMock(t)
store := storeWith(mock)
ctx := context.Background()

// Ladder origin must target WHERE ladder_run_id IS NOT NULL only.
mock.ExpectExec("ladder_run_id IS NOT NULL").WillReturnResult(pgxmock.NewResult("UPDATE", 1))

n, err := store.CancelPendingExchangesByOrigin(ctx, common.ExchangeOriginLadder)
require.NoError(t, err)
assert.Equal(t, int64(1), n)
assert.NoError(t, mock.ExpectationsWereMet())
}

func TestPGXMock_CancelPendingExchangesByOrigin_UnknownOrigin_FailsLoud(t *testing.T) {
mock := newMock(t)
store := storeWith(mock)
ctx := context.Background()

// An unknown origin must fail at the boundary WITHOUT issuing any query
// (no ExpectExec registered → ExpectationsWereMet stays satisfied).
n, err := store.CancelPendingExchangesByOrigin(ctx, common.ExchangeOrigin("bogus"))
require.Error(t, err)
assert.Contains(t, err.Error(), "unknown exchange origin")
assert.Equal(t, int64(0), n)
assert.NoError(t, mock.ExpectationsWereMet())
}

// ─── SaveRIExchangeRecord ─────────────────────────────────────────────────────

func TestPGXMock_SaveRIExchangeRecord_InsertColumnAlignment(t *testing.T) {
mock := newMock(t)
store := storeWith(mock)
ctx := context.Background()

// 21 columns/placeholders: original 20 + ladder_run_id ($21, issue #1348).
// A column/placeholder count drift makes WithArgs(anyArgsCfg(21)) fail.
mock.ExpectExec("INSERT INTO ri_exchange_history").WithArgs(anyArgsCfg(21)...).
WillReturnResult(pgxmock.NewResult("INSERT", 1))

ladderRunID := "run-123"
err := store.SaveRIExchangeRecord(ctx, &RIExchangeRecord{
AccountID: "acc-1",
Region: "us-east-1",
Status: "pending",
Mode: "manual",
LadderRunID: &ladderRunID,
})
require.NoError(t, err)
assert.NoError(t, mock.ExpectationsWereMet())
}

// ─── CleanupOldExecutions ─────────────────────────────────────────────────────

func TestPGXMock_CleanupOldExecutions_Success(t *testing.T) {
Expand Down
15 changes: 14 additions & 1 deletion internal/config/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -804,7 +804,20 @@ type RIExchangeRecord struct {
CreatedByUserID *string `json:"created_by_user_id,omitempty"`
// ApprovedBy carries the email of the session user who approved the exchange
// via the dashboard Approve button (issue #300). Nil for token-authed approvals.
ApprovedBy *string `json:"approved_by,omitempty"`
ApprovedBy *string `json:"approved_by,omitempty"`
// LadderRunID links this exchange record to the ladder run that created it
// (cudly-ladder engine). Nil for standalone ri_exchange_reshape task records.
// The database column ri_exchange_history.ladder_run_id was added in migration
// 000080 and is the authoritative source for origin scoping in
// CancelPendingExchangesByOrigin.
//
// KNOWN/ACCEPTABLE: the FK is ON DELETE SET NULL (migration 000080), so
// deleting a ladder_runs row nulls this column and reclassifies the record
// as standalone. A still-pending reshape then becomes standalone-cancellable
// (the standalone-origin sweep would cancel it). This is acceptable: a
// deleted run has no owner to approve its pendings, so cancelling them on the
// next standalone sweep is the safe outcome, not a leak.
LadderRunID *string `json:"ladder_run_id,omitempty"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
CompletedAt *time.Time `json:"completed_at,omitempty"`
Expand Down
Loading
Loading