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
19 changes: 19 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/ladder"
"github.com/jackc/pgx/v5"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
Expand Down Expand Up @@ -747,6 +748,24 @@ func (m *mockConfigStore) GetLadderConfig(_ context.Context, _, _ string) (*conf
func (m *mockConfigStore) UpsertLadderConfig(_ context.Context, cfg *config.LadderConfigDB) (*config.LadderConfigDB, error) {
return cfg, nil
}
func (m *mockConfigStore) SaveLadderRun(_ context.Context, run *config.LadderRunDB) (*config.LadderRunDB, error) {
return run, nil
}
func (m *mockConfigStore) SaveLadderRunWithTranches(_ context.Context, run *config.LadderRunDB, _ []config.LadderTrancheDB) (*config.LadderRunDB, error) {
return run, nil
}
func (m *mockConfigStore) GetLadderRun(_ context.Context, _ string) (*config.LadderRunDB, error) {
return nil, nil
}
func (m *mockConfigStore) SaveLadderTranches(_ context.Context, _ []config.LadderTrancheDB) error {
return nil
}
func (m *mockConfigStore) LatestLadderRunStartedAt(_ context.Context, _ string) (*time.Time, error) {
return nil, nil
}
func (m *mockConfigStore) TransitionLadderRunStatus(_ context.Context, _ string, _ []ladder.RunStatus, _ ladder.RunStatus) (*config.LadderRunDB, error) {
return nil, nil
}
func (m *mockConfigStore) UpdateGlobalConfigAtomic(_ context.Context, apply func(*config.GlobalConfig) error) (*config.GlobalConfig, error) {
cfg := &config.GlobalConfig{}
if err := apply(cfg); err != nil {
Expand Down
36 changes: 36 additions & 0 deletions internal/config/interfaces.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@ import (
"time"

"github.com/jackc/pgx/v5"

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

// StoreInterface defines the methods required for configuration storage.
Expand Down Expand Up @@ -329,4 +331,38 @@ type StoreInterface interface {
GetLadderConfigs(ctx context.Context) ([]LadderConfigDB, error)
GetLadderConfig(ctx context.Context, cloudAccountID, provider string) (*LadderConfigDB, error)
UpsertLadderConfig(ctx context.Context, cfg *LadderConfigDB) (*LadderConfigDB, error)

// Ladder run/tranche persistence (migration 000080/000081, PR-2).
//
// SaveLadderRun inserts a new ladder_runs row, returning the persisted row
// with all DB-stamped fields (id, created_at, updated_at) populated.
// If run.ID is empty, a new UUID is generated before the insert.
//
// GetLadderRun returns the row for the given id, or (nil, nil) when no
// row exists (mirrors GetLadderConfig semantics).
//
// SaveLadderTranches inserts a batch of ladder_tranches rows inside a
// single transaction. Each tranche must carry a non-empty ID; duplicate
// IDs within the batch are rejected at the DB UNIQUE constraint level.
//
// LatestLadderRunStartedAt returns the maximum started_at for the given
// config_id, or nil when no run has been recorded yet. Powers the per-cadence
// self-gate in the scheduler (Q6).
//
// TransitionLadderRunStatus atomically updates the status of a ladder_runs
// row from one of the fromStatuses to toStatus, returning the updated row.
// Returns (nil, nil) when zero rows are affected (CAS race lost or wrong
// current status), so callers can distinguish a race from a hard error.
// Statuses are typed ladder.RunStatus so callers cannot pass an arbitrary
// string that would never match a stored status.
// SaveLadderRunWithTranches inserts the run row and its tranches in ONE
// transaction: a tranche-insert failure rolls back the run row too, so a
// status=planned run never persists without its tranches (which would let
// the cadence gate suppress the retry for a full window).
SaveLadderRun(ctx context.Context, run *LadderRunDB) (*LadderRunDB, error)
SaveLadderRunWithTranches(ctx context.Context, run *LadderRunDB, tranches []LadderTrancheDB) (*LadderRunDB, error)
GetLadderRun(ctx context.Context, id string) (*LadderRunDB, error)
SaveLadderTranches(ctx context.Context, tranches []LadderTrancheDB) error
LatestLadderRunStartedAt(ctx context.Context, configID string) (*time.Time, error)
TransitionLadderRunStatus(ctx context.Context, id string, fromStatuses []ladder.RunStatus, toStatus ladder.RunStatus) (*LadderRunDB, error)
}
273 changes: 273 additions & 0 deletions internal/config/store_postgres_ladder.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@ import (

"github.com/google/uuid"
"github.com/jackc/pgx/v5"

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

// ==========================================
Expand Down Expand Up @@ -177,6 +179,277 @@ func scanLadderConfig(row scannable) (LadderConfigDB, error) {
return cfg, nil
}

// ==========================================
// LADDER RUN / TRANCHE STORE METHODS
// ==========================================

// ladderRunInsertQuery INSERTs one ladder_runs row and RETURNs the full row for
// scanLadderRun. Shared by SaveLadderRun and SaveLadderRunWithTranches so the
// column list has a single definition.
const ladderRunInsertQuery = `
INSERT INTO ladder_runs (
id, config_id, started_at, completed_at, status, mode, cadence,
baseline_usd_hr, target_usd_hr, existing_usd_hr, gap_usd_hr,
plan, total_hourly_commit, total_upfront_cost, estimated_savings,
approval_token_hash, approval_token_expires_at,
approved_by, cancelled_by, fire_at,
created_at, updated_at
) VALUES (
$1, $2, $3, $4, $5, $6, $7,
$8, $9, $10, $11,
$12, $13, $14, $15,
$16, $17, $18, $19, $20,
NOW(), NOW()
)
RETURNING id, config_id, started_at, completed_at, status, mode, cadence,
baseline_usd_hr, target_usd_hr, existing_usd_hr, gap_usd_hr,
plan, total_hourly_commit, total_upfront_cost, estimated_savings,
approval_token_hash, approval_token_expires_at,
approved_by, cancelled_by, fire_at,
created_at, updated_at
`

// ladderRunPlanJSON returns the plan blob to persist, defaulting an empty plan
// to the JSONB '{}' the schema expects (NOT NULL DEFAULT '{}').
func ladderRunPlanJSON(run *LadderRunDB) json.RawMessage {
if len(run.Plan) == 0 {
return json.RawMessage(`{}`)
}
return run.Plan
}

// ladderRunInsertArgs returns the positional args for ladderRunInsertQuery in
// column order. Nullable monetary fields stay as *float64 so NULL != $0.
func ladderRunInsertArgs(run *LadderRunDB, planJSON json.RawMessage) []any {
return []any{
run.ID,
run.ConfigID,
run.StartedAt,
run.CompletedAt,
run.Status,
run.Mode,
run.Cadence,
run.BaselineUSDHr,
run.TargetUSDHr,
run.ExistingUSDHr,
run.GapUSDHr,
planJSON,
run.TotalHourlyCommit,
run.TotalUpfrontCost,
run.EstimatedSavings,
run.ApprovalTokenHash,
run.ApprovalTokenExpiresAt,
run.ApprovedBy,
run.CancelledBy,
run.FireAt,
}
}

// SaveLadderRun inserts a new ladder_runs row. If run.ID is empty a fresh UUID
// is generated. Returns the persisted row with all DB-stamped fields populated.
func (s *PostgresStore) SaveLadderRun(ctx context.Context, run *LadderRunDB) (*LadderRunDB, error) {
if run.ID == "" {
run.ID = uuid.New().String()
}
row := s.db.QueryRow(ctx, ladderRunInsertQuery, ladderRunInsertArgs(run, ladderRunPlanJSON(run))...)
result, err := scanLadderRun(row)
if err != nil {
return nil, fmt.Errorf("failed to insert ladder_run id=%s: %w", run.ID, err)
}
return &result, nil
}

// GetLadderRun returns the ladder_runs row for the given id, or (nil, nil)
// when no row exists.
func (s *PostgresStore) GetLadderRun(ctx context.Context, id string) (*LadderRunDB, error) {
query := `
SELECT id, config_id, started_at, completed_at, status, mode, cadence,
baseline_usd_hr, target_usd_hr, existing_usd_hr, gap_usd_hr,
plan, total_hourly_commit, total_upfront_cost, estimated_savings,
approval_token_hash, approval_token_expires_at,
approved_by, cancelled_by, fire_at,
created_at, updated_at
FROM ladder_runs
WHERE id = $1
`
row := s.db.QueryRow(ctx, query, id)
result, err := scanLadderRun(row)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil
}
return nil, fmt.Errorf("failed to get ladder_run id=%s: %w", id, err)
}
return &result, nil
}

// insertLadderTranchesTx inserts every tranche within the given transaction.
// Each tranche must carry a non-empty ID; duplicate IDs are rejected at the DB
// PRIMARY KEY constraint. Shared by SaveLadderTranches and
// SaveLadderRunWithTranches so both go through identical insert logic.
func insertLadderTranchesTx(ctx context.Context, tx pgx.Tx, tranches []LadderTrancheDB) error {
for i := range tranches {
tr := &tranches[i]
if tr.ID == "" {
return fmt.Errorf("ladder_tranche at index %d has an empty ID; every tranche must carry a unique non-empty ID", i)
}
_, err := tx.Exec(ctx, `
INSERT INTO ladder_tranches (
id, config_id, run_id, layer_type, amount_usd_hr,
term, payment_option, scheduled_date, status, execution_id,
created_at
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, NOW())
`,
tr.ID,
tr.ConfigID,
tr.RunID,
string(tr.LayerType),
tr.AmountUSDHr,
string(tr.Term),
string(tr.PaymentOption),
tr.ScheduledDate,
string(tr.Status),
tr.ExecutionID,
)
if err != nil {
return fmt.Errorf("failed to insert ladder_tranche id=%s run_id=%v: %w", tr.ID, tr.RunID, err)
}
}
return nil
}

// SaveLadderTranches inserts a batch of ladder_tranches rows within a
// single transaction. Each tranche must carry a non-empty ID; duplicate
// IDs are rejected at the DB UNIQUE PRIMARY KEY constraint. An empty
// slice is a no-op.
func (s *PostgresStore) SaveLadderTranches(ctx context.Context, tranches []LadderTrancheDB) error {
if len(tranches) == 0 {
return nil
}
return s.WithTx(ctx, func(tx pgx.Tx) error {
return insertLadderTranchesTx(ctx, tx, tranches)
})
}

// SaveLadderRunWithTranches inserts the ladder_runs row AND its ladder_tranches
// rows in ONE transaction. If any tranche insert fails, the whole transaction
// (including the run row) is rolled back, so a run is never persisted without
// its tranches. This prevents a status=planned run with zero tranches, which
// the cadence self-gate (keyed on any run's started_at) would otherwise use to
// suppress the retry for the full cadence window. If run.ID is empty a fresh
// UUID is generated. Returns the persisted run row.
func (s *PostgresStore) SaveLadderRunWithTranches(ctx context.Context, run *LadderRunDB, tranches []LadderTrancheDB) (*LadderRunDB, error) {
if run.ID == "" {
run.ID = uuid.New().String()
}
var result LadderRunDB
err := s.WithTx(ctx, func(tx pgx.Tx) error {
row := tx.QueryRow(ctx, ladderRunInsertQuery, ladderRunInsertArgs(run, ladderRunPlanJSON(run))...)
scanned, err := scanLadderRun(row)
if err != nil {
return fmt.Errorf("failed to insert ladder_run id=%s: %w", run.ID, err)
}
if err := insertLadderTranchesTx(ctx, tx, tranches); err != nil {
return err
}
result = scanned
return nil
})
if err != nil {
return nil, err
}
return &result, nil
}

// LatestLadderRunStartedAt returns the maximum started_at for the given
// config_id, or nil when no run has been recorded yet. Drives the per-cadence
// self-gate in handleLadderRun (Q6).
func (s *PostgresStore) LatestLadderRunStartedAt(ctx context.Context, configID string) (*time.Time, error) {
query := `SELECT MAX(started_at) FROM ladder_runs WHERE config_id = $1`
var ts *time.Time
err := s.db.QueryRow(ctx, query, configID).Scan(&ts)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil
}
return nil, fmt.Errorf("failed to get latest ladder_run started_at for config_id=%s: %w", configID, err)
}
return ts, nil
}

// TransitionLadderRunStatus atomically transitions a ladder_runs row from
// one of fromStatuses to toStatus via a CAS UPDATE. Returns the updated row
// on success, or (nil, nil) when zero rows are affected (race lost or
// unexpected current status). A hard DB error is returned as a non-nil error.
func (s *PostgresStore) TransitionLadderRunStatus(ctx context.Context, id string, fromStatuses []ladder.RunStatus, toStatus ladder.RunStatus) (*LadderRunDB, error) {
// Convert the typed enum slice to the plain []string that pgx encodes for
// the ANY($3) text[] comparison. Keeping the public signature typed while
// converting here confines the stringly-typed shape to the DB boundary.
from := make([]string, len(fromStatuses))
for i, st := range fromStatuses {
from[i] = string(st)
}
query := `
UPDATE ladder_runs
SET status = $2, updated_at = NOW()
WHERE id = $1 AND status = ANY($3)
RETURNING id, config_id, started_at, completed_at, status, mode, cadence,
baseline_usd_hr, target_usd_hr, existing_usd_hr, gap_usd_hr,
plan, total_hourly_commit, total_upfront_cost, estimated_savings,
approval_token_hash, approval_token_expires_at,
approved_by, cancelled_by, fire_at,
created_at, updated_at
`
row := s.db.QueryRow(ctx, query, id, string(toStatus), from)
result, err := scanLadderRun(row)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil // CAS race lost
}
return nil, fmt.Errorf("failed to transition ladder_run id=%s to status=%s: %w", id, toStatus, err)
}
return &result, nil
}

// scanLadderRun scans a single ladder_runs row from a pgx.Row or pgx.Rows.
// Both implement the scannable interface (store_postgres_registrations.go).
func scanLadderRun(row scannable) (LadderRunDB, error) {
var r LadderRunDB
var planJSON []byte

err := row.Scan(
&r.ID,
&r.ConfigID,
&r.StartedAt,
&r.CompletedAt,
&r.Status,
&r.Mode,
&r.Cadence,
&r.BaselineUSDHr,
&r.TargetUSDHr,
&r.ExistingUSDHr,
&r.GapUSDHr,
&planJSON,
&r.TotalHourlyCommit,
&r.TotalUpfrontCost,
&r.EstimatedSavings,
&r.ApprovalTokenHash,
&r.ApprovalTokenExpiresAt,
&r.ApprovedBy,
&r.CancelledBy,
&r.FireAt,
&r.CreatedAt,
&r.UpdatedAt,
)
if err != nil {
return LadderRunDB{}, err
}
if planJSON != nil {
r.Plan = json.RawMessage(planJSON)
}
return r, nil
}

// marshalRampSchedule ensures the ramp_schedule value stored in the DB is
// valid JSON. Nil/empty input is rejected; already-valid JSON is passed
// through unchanged. Non-JSON input is rejected with a descriptive error.
Expand Down
Loading
Loading