package agenttask import ( "context" "fmt" ) // leaseClaim is an immutable handle for one durable lease. The manager returns // it from claim operations and tracks it inside the leaseSet so the renewal // supervisor and fence validators can reference a single source of truth for // scope, owner, token, and subject identity. type leaseClaim struct { scope string owner string token string subject string } func (m *Manager) claimDevice(ctx context.Context) (*leaseClaim, error) { now := m.clock.Now() token := fmt.Sprintf("%s/device/%d", m.config.OwnerID, now.UnixNano()) claimed, err := mutateDecision(m, ctx, func(state *ManagerState) (bool, error) { if state.DeviceLease != nil && state.DeviceLease.ExpiresAt.After(now) { return false, nil } state.DeviceLease = &LeaseRecord{ OwnerID: m.config.OwnerID, Token: token, ExpiresAt: now.Add(m.config.LeaseDuration), } return true, nil }) if err != nil { return nil, err } if !claimed { return nil, deviceLeaseError() } return &leaseClaim{ scope: "device", owner: m.config.OwnerID, token: token, subject: "", }, nil } func (m *Manager) renewDevice(ctx context.Context, token string) error { now := m.clock.Now() _, err := mutateDecision(m, ctx, func(state *ManagerState) (struct{}, error) { if state.DeviceLease == nil || state.DeviceLease.Token != token { return struct{}{}, fmt.Errorf("device lease token mismatch: want %s, have %v", token, state.DeviceLease) } state.DeviceLease.ExpiresAt = now.Add(m.config.LeaseDuration) return struct{}{}, nil }) return err } func (m *Manager) validateDevice(ctx context.Context, token string) error { state, err := m.load(ctx) if err != nil { return err } if state.DeviceLease == nil || state.DeviceLease.Token != token { return fmt.Errorf("%w: device lease no longer matches token %s", ErrLeaseLost, token) } return nil } func (m *Manager) releaseDevice(ctx context.Context, token string) { _ = m.mutate(ctx, func(state *ManagerState) error { if state.DeviceLease == nil || state.DeviceLease.Token != token { return nil } state.DeviceLease = nil return nil }) } func (m *Manager) claimWorkspace( ctx context.Context, workspaceID WorkspaceID, ) (*leaseClaim, error) { now := m.clock.Now() token := fmt.Sprintf("%s/workspace/%s/%d", m.config.OwnerID, workspaceID, now.UnixNano()) claimed, err := mutateDecision(m, ctx, func(state *ManagerState) (bool, error) { lease, exists := state.WorkspaceLeases[workspaceID] if exists && lease.ExpiresAt.After(now) { return false, nil } state.WorkspaceLeases[workspaceID] = LeaseRecord{ OwnerID: m.config.OwnerID, Token: token, ExpiresAt: now.Add(m.config.LeaseDuration), } return true, nil }) if err != nil { return nil, err } if !claimed { return nil, nil } return &leaseClaim{ scope: "workspace", owner: m.config.OwnerID, token: token, subject: string(workspaceID), }, nil } func (m *Manager) renewWorkspace(ctx context.Context, token, subject string) error { now := m.clock.Now() _, err := mutateDecision(m, ctx, func(state *ManagerState) (struct{}, error) { lease, ok := state.WorkspaceLeases[WorkspaceID(subject)] if !ok || lease.Token != token { return struct{}{}, fmt.Errorf("workspace lease token mismatch for %s", subject) } lease.ExpiresAt = now.Add(m.config.LeaseDuration) state.WorkspaceLeases[WorkspaceID(subject)] = lease return struct{}{}, nil }) return err } func (m *Manager) validateWorkspace(ctx context.Context, token, subject string) error { state, err := m.load(ctx) if err != nil { return err } lease, ok := state.WorkspaceLeases[WorkspaceID(subject)] if !ok || lease.Token != token { return fmt.Errorf("%w: workspace lease %s no longer matches", ErrLeaseLost, subject) } return nil } func (m *Manager) releaseWorkspace(ctx context.Context, token, subject string) { _ = m.mutate(ctx, func(state *ManagerState) error { lease, ok := state.WorkspaceLeases[WorkspaceID(subject)] if ok && lease.Token == token { delete(state.WorkspaceLeases, WorkspaceID(subject)) } return nil }) } func deviceLeaseError() error { return fmt.Errorf("%w: another live owner retained the durable lease", ErrDeviceLeaseHeld) } func (m *Manager) claimProject(ctx context.Context, projectID ProjectID) (*leaseClaim, error) { now := m.clock.Now() token := fmt.Sprintf("%s/%s/%d", m.config.OwnerID, projectID, now.UnixNano()) claimed, err := mutateDecision(m, ctx, func(state *ManagerState) (bool, error) { project, ok := state.Projects[projectID] if !ok || project.Status != ProjectStatusRunning { return false, nil } if project.Lease != nil && project.Lease.ExpiresAt.After(now) { return false, nil } project.Lease = &LeaseRecord{ OwnerID: m.config.OwnerID, Token: token, ExpiresAt: now.Add(m.config.LeaseDuration), } project.UpdatedAt = now state.Projects[projectID] = project return true, nil }) if err != nil { return nil, err } if !claimed { return nil, nil } return &leaseClaim{ scope: "project", owner: m.config.OwnerID, token: token, subject: string(projectID), }, nil } func (m *Manager) renewProject(ctx context.Context, token, subject string) error { now := m.clock.Now() _, err := mutateDecision(m, ctx, func(state *ManagerState) (struct{}, error) { project, ok := state.Projects[ProjectID(subject)] if !ok || project.Lease == nil || project.Lease.Token != token { return struct{}{}, fmt.Errorf("project lease token mismatch for %s", subject) } project.Lease.ExpiresAt = now.Add(m.config.LeaseDuration) state.Projects[ProjectID(subject)] = project return struct{}{}, nil }) return err } func (m *Manager) validateProject(ctx context.Context, token, subject string) error { state, err := m.load(ctx) if err != nil { return err } project, ok := state.Projects[ProjectID(subject)] if !ok || project.Lease == nil || project.Lease.Token != token { return fmt.Errorf("%w: project lease %s no longer matches", ErrLeaseLost, subject) } return nil } func (m *Manager) releaseProject(ctx context.Context, token, subject string) { _ = m.mutate(ctx, func(state *ManagerState) error { project, ok := state.Projects[ProjectID(subject)] if !ok || project.Lease == nil || project.Lease.Token != token { return nil } project.Lease = nil project.UpdatedAt = m.clock.Now() state.Projects[ProjectID(subject)] = project return nil }) } func (m *Manager) claimIntegration( ctx context.Context, workspaceID WorkspaceID, ) (*leaseClaim, error) { now := m.clock.Now() token := fmt.Sprintf("%s/integration/%s/%d", m.config.OwnerID, workspaceID, now.UnixNano()) claimed, err := mutateDecision(m, ctx, func(state *ManagerState) (bool, error) { lease, exists := state.IntegrationLeases[workspaceID] if exists && lease.ExpiresAt.After(now) { return false, nil } state.IntegrationLeases[workspaceID] = LeaseRecord{ OwnerID: m.config.OwnerID, Token: token, ExpiresAt: now.Add(m.config.LeaseDuration), } return true, nil }) if err != nil { return nil, err } if !claimed { return nil, nil } return &leaseClaim{ scope: "integration", owner: m.config.OwnerID, token: token, subject: string(workspaceID), }, nil } func (m *Manager) renewIntegration(ctx context.Context, token, subject string) error { now := m.clock.Now() _, err := mutateDecision(m, ctx, func(state *ManagerState) (struct{}, error) { lease, ok := state.IntegrationLeases[WorkspaceID(subject)] if !ok || lease.Token != token { return struct{}{}, fmt.Errorf("integration lease token mismatch for %s", subject) } lease.ExpiresAt = now.Add(m.config.LeaseDuration) state.IntegrationLeases[WorkspaceID(subject)] = lease return struct{}{}, nil }) return err } func (m *Manager) validateIntegration(ctx context.Context, token, subject string) error { state, err := m.load(ctx) if err != nil { return err } lease, ok := state.IntegrationLeases[WorkspaceID(subject)] if !ok || lease.Token != token { return fmt.Errorf("%w: integration lease %s no longer matches", ErrLeaseLost, subject) } return nil } func (m *Manager) releaseIntegration(ctx context.Context, token, subject string) { _ = m.mutate(ctx, func(state *ManagerState) error { lease, ok := state.IntegrationLeases[WorkspaceID(subject)] if ok && lease.Token == token { delete(state.IntegrationLeases, WorkspaceID(subject)) } return nil }) } // releaseExact intentionally bypasses guarded-context fencing: cleanup must be // able to remove only the token this manager acquired even after ownership has // been lost. Each scope-specific release independently preserves a successor. func (m *Manager) releaseExact(ctx context.Context, claim *leaseClaim) { if claim == nil { return } ctx = unfencedLeaseContext(ctx) switch claim.scope { case "device": m.releaseDevice(ctx, claim.token) case "project": m.releaseProject(ctx, claim.token, claim.subject) case "workspace": m.releaseWorkspace(ctx, claim.token, claim.subject) case "integration": m.releaseIntegration(ctx, claim.token, claim.subject) } }