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 README.md
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,10 @@ poolmgrctl lease claim --pool web --namespace default --addr 127.0.0.1:9090 --in
poolmgrctl events tail --pool web --namespace default --addr 127.0.0.1:9090 --insecure
```

A claim can time out while the manager still commits the lease. To retry a claim safely, pass
the same `--request-id` (e.g. a UUID) on every attempt: while the lease exists, a repeat claim
with that ID returns the same lease instead of claiming another VM.

### Build and test

```sh
Expand Down
165 changes: 91 additions & 74 deletions api/proto/poolmgr/v1alpha1/lease.pb.go

Large diffs are not rendered by default.

8 changes: 8 additions & 0 deletions api/proto/poolmgr/v1alpha1/lease.proto
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,14 @@ service Lease {

message ClaimVMRequest {
PoolRef pool = 1;
// RequestId lets a client retry ClaimVM safely. The manager stores it on
// the lease it creates. A later ClaimVM with the same request_id returns
// that lease again, for as long as it exists, rather than a new one.
//
// Clients that retry ClaimVM must set it. A UUID is a good choice. Empty
// means the request is not retriable, which is the behaviour before this
// field existed. At most 255 bytes.
string request_id = 2;
}

message ClaimVMResponse {
Expand Down
208 changes: 110 additions & 98 deletions api/proto/poolmgr/v1alpha1/types.pb.go

Large diffs are not rendered by default.

3 changes: 3 additions & 0 deletions api/proto/poolmgr/v1alpha1/types.proto
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,9 @@ message LeaseRecord {
// PoolNamespace is the namespace of the pool identified by pool_name (together they form
// the pool's PoolRef identity).
string pool_namespace = 7;
// RequestId is the ClaimVMRequest.request_id that created this lease.
// Empty if the client did not set one.
string request_id = 8;
}

// EventType names the kind of change an Event describes.
Expand Down
2 changes: 1 addition & 1 deletion docs/adr/2026-09-24-idempotent-claims.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

Related: [issue #98](https://github.com/liquidmetal-dev/battery/issues/98).

- Status: Proposed
- Status: Accepted
- Date: 2026-09-24

## Context
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,10 @@ service Events {
`MicroVMStatus.network_interfaces`) plus the `lease_id`. `ClaimVM` fails with a distinct status
(e.g. `RESOURCE_EXHAUSTED`) when no VM is `AVAILABLE`.

`ClaimVMRequest.request_id` makes `ClaimVM` safe to retry: while the lease it created exists, a
repeat request with the same `request_id` returns that lease rather than claiming another VM.
See the [idempotent claims ADR](../adr/2026-09-24-idempotent-claims.md).

## Lease Lifecycle

1. Consumer calls `ClaimVM(pool_name)`. Manager atomically picks an `AVAILABLE` VM, runs the
Expand Down Expand Up @@ -236,7 +240,7 @@ pool's `hook_failure_policy`.
Per-pool gauges/counters (labeled by `pool_name`):
- `poolmgr_pool_size` (target), `poolmgr_pool_available`, `poolmgr_pool_leased`,
`poolmgr_pool_provisioning`, `poolmgr_pool_quarantined`
- `poolmgr_vm_claims_total`, `poolmgr_vm_releases_total{reason="api|expiry"}`
- `poolmgr_vm_claims_total{replayed="true|false"}`, `poolmgr_vm_releases_total{reason="api|expiry"}`
- `poolmgr_vm_provision_duration_seconds` (histogram), `poolmgr_hook_duration_seconds{hook="create|pre_lease"}`
- `poolmgr_hook_failures_total{hook,pool_name}`
- `poolmgr_lease_duration_seconds` (histogram)
Expand Down
110 changes: 103 additions & 7 deletions internal/api/lease.go
Original file line number Diff line number Diff line change
Expand Up @@ -90,11 +90,23 @@ func NewLeaseServer(st store.Store, flint *flintlockclient.Pool, cfg HookExecCon
return &LeaseServer{store: st, flint: flint, cfg: cfg, notifier: notifier, metrics: m}
}

// maxRequestIDLength is the longest ClaimVMRequest.request_id ClaimVM
// accepts, in bytes.
const maxRequestIDLength = 255

// ClaimVM atomically claims an AVAILABLE VM from the named pool, runs the
// pool's pre_lease_commands, and creates a Lease. It fails with
// RESOURCE_EXHAUSTED when no VM is AVAILABLE.
//
// If req.request_id is set and a lease created with it still exists,
// ClaimVM returns that lease again instead of claiming another VM (see
// replayClaim), so a client can retry a claim whose response it lost.
func (s *LeaseServer) ClaimVM(ctx context.Context, req *poolmgrv1alpha1.ClaimVMRequest) (*poolmgrv1alpha1.ClaimVMResponse, error) {
poolName, poolNS := req.GetPool().GetName(), req.GetPool().GetNamespace()
requestID := req.GetRequestId()
if len(requestID) > maxRequestIDLength {
return nil, status.Errorf(codes.InvalidArgument, "request_id is %d bytes, at most %d allowed", len(requestID), maxRequestIDLength)
}

pool, err := s.store.GetPool(ctx, poolName, poolNS)
if errors.Is(err, store.ErrNotFound) {
Expand All @@ -104,6 +116,16 @@ func (s *LeaseServer) ClaimVM(ctx context.Context, req *poolmgrv1alpha1.ClaimVMR
return nil, status.Errorf(codes.Internal, "get pool: %v", err)
}

if requestID != "" {
lease, err := s.store.GetLeaseByRequestID(ctx, requestID)
if err == nil {
return s.replayClaim(ctx, lease, poolName, poolNS)
}
if !errors.Is(err, store.ErrNotFound) {
return nil, status.Errorf(codes.Internal, "get lease by request id: %v", err)
}
}

vm, err := s.store.ClaimAvailableVM(ctx, poolName, poolNS)
if errors.Is(err, store.ErrNoAvailableVM) {
return nil, status.Errorf(codes.ResourceExhausted, "no available vm in pool %s/%s", poolNS, poolName)
Expand Down Expand Up @@ -140,27 +162,101 @@ func (s *LeaseServer) ClaimVM(ctx context.Context, req *poolmgrv1alpha1.ClaimVMR
ClaimedAt: timestamppb.New(now),
LastHeartbeatAt: timestamppb.New(now),
ExpiresAt: timestamppb.New(expiresAt),
RequestId: requestID,
}
if err := s.store.CreateLease(ctx, lease); err != nil {
if errors.Is(err, store.ErrDuplicateRequestID) {
return s.yieldClaim(ctx, pool, vm, requestID)
}
s.applyHookFailurePolicy(ctx, pool, vm)
return nil, status.Errorf(codes.Internal, "create lease: %v", err)
}

reconciler.EmitEvent(ctx, s.store, pool, vm.GetUid(), poolmgrv1alpha1.EventType_VM_CLAIMED)
s.metrics.RecordVMClaim(poolName, poolNS)
s.metrics.RecordVMClaim(poolName, poolNS, false)
s.notifier.NotifyVMClaimed(poolName, poolNS)

// Best-effort: the lease is already committed at this point, so a
// failure to fetch network interfaces just means an empty map in the
// response - the caller can still Heartbeat/ReleaseVM successfully.
return s.claimResponse(ctx, leaseID, vm), nil
}

// replayClaim answers a ClaimVM whose request_id matches an existing lease:
// it returns that lease's response again, without claiming a VM or touching
// the lease's expiry (a replay is not a heartbeat). A request_id reused for
// a different pool is a client bug and fails with INVALID_ARGUMENT.
func (s *LeaseServer) replayClaim(ctx context.Context, lease *poolmgrv1alpha1.LeaseRecord, poolName, poolNS string) (*poolmgrv1alpha1.ClaimVMResponse, error) {
if lease.GetPoolName() != poolName || lease.GetPoolNamespace() != poolNS {
return nil, status.Errorf(codes.InvalidArgument, "request_id %q was used to claim from pool %s/%s, not %s/%s",
lease.GetRequestId(), lease.GetPoolNamespace(), lease.GetPoolName(), poolNS, poolName)
}

vm, err := s.store.GetVM(ctx, lease.GetVmUid())
if errors.Is(err, store.ErrNotFound) {
// Releasing a lease deletes the VM row before the lease row, so this
// lease is on its way out. Once it's gone the request_id is free
// again and a retry makes a fresh claim.
return nil, status.Errorf(codes.Aborted, "lease %s for request_id %q is being released, retry", lease.GetLeaseId(), lease.GetRequestId())
}
if err != nil {
return nil, status.Errorf(codes.Internal, "get vm: %v", err)
}
if vm.GetPhase() == poolmgrv1alpha1.VMPhase_DELETING {
// EnsureVMDeleted marks the VM DELETING before it calls flintlock,
// and a failed delete leaves both rows in place. Treat it like the
// missing-row case above: the lease is ending, so don't replay it.
return nil, status.Errorf(codes.Aborted, "lease %s for request_id %q is being released, retry", lease.GetLeaseId(), lease.GetRequestId())
}

s.metrics.RecordVMClaim(poolName, poolNS, true)
Comment thread
richardcase marked this conversation as resolved.
return s.claimResponse(ctx, lease.GetLeaseId(), vm), nil
}

// yieldClaim handles losing a race with a concurrent ClaimVM for the same
// request_id: both passed the lookup and claimed a VM, and the other one
// created its lease first. It puts vm back to AVAILABLE and returns the
// winner's lease via replayClaim. No VM_CLAIMED event or claim notification
// went out for vm, so there is none to undo.
//
// As in applyHookFailurePolicy, the VM is returned on a context detached
// from ctx so a cancelled RPC can't strand it LEASED with no lease.
func (s *LeaseServer) yieldClaim(ctx context.Context, pool *poolmgrv1alpha1.PoolSpec, vm *poolmgrv1alpha1.VMRecord, requestID string) (*poolmgrv1alpha1.ClaimVMResponse, error) {
cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), s.cfg.CleanupTimeout)
defer cancel()

vm.Phase = poolmgrv1alpha1.VMPhase_AVAILABLE
vm.LeaseId = nil
vm.UpdatedAt = timestamppb.Now()
if err := s.store.UpdateVM(cleanupCtx, vm); err != nil {
s.applyHookFailurePolicy(ctx, pool, vm)
return nil, status.Errorf(codes.Internal, "return vm to available: %v", err)
}

lease, err := s.store.GetLeaseByRequestID(ctx, requestID)
if errors.Is(err, store.ErrNotFound) {
// The winner's lease already ended.
return nil, status.Errorf(codes.Aborted, "concurrent claim for request_id %q already ended, retry", requestID)
}
if err != nil {
return nil, status.Errorf(codes.Internal, "get lease by request id: %v", err)
}
return s.replayClaim(ctx, lease, pool.GetName(), pool.GetNamespace())
}

// claimResponse builds the ClaimVMResponse for leaseID on vm, for both a
// fresh claim and a replay.
//
// Fetching network interfaces is best-effort: the lease is already
// committed at this point, so a failure just means an empty map in the
// response - the caller can still Heartbeat/ReleaseVM successfully.
func (s *LeaseServer) claimResponse(ctx context.Context, leaseID string, vm *poolmgrv1alpha1.VMRecord) *poolmgrv1alpha1.ClaimVMResponse {
hostName := vm.GetFlintlockHost()

var netIfaces map[string]*flintlocktypes.NetworkInterfaceStatus
if client, cerr := s.flint.Client(vm.GetFlintlockHost()); cerr == nil {
if client, cerr := s.flint.Client(hostName); cerr == nil {
if resp, gerr := client.GetMicroVM(ctx, &microvmv1alpha1.GetMicroVMRequest{Uid: vm.GetUid()}); gerr == nil {
netIfaces = resp.GetMicrovm().GetStatus().GetNetworkInterfaces()
}
}

hostName := vm.GetFlintlockHost()
addr, _ := s.flint.Address(hostName)

return &poolmgrv1alpha1.ClaimVMResponse{
Expand All @@ -171,7 +267,7 @@ func (s *LeaseServer) ClaimVM(ctx context.Context, req *poolmgrv1alpha1.ClaimVMR
Name: hostName,
Address: addr,
},
}, nil
}
}

// runPreLeaseHooks transitions vm to PRE_LEASE_HOOK_RUNNING and executes
Expand Down
Loading
Loading