879 lines
24 KiB
Go
879 lines
24 KiB
Go
package agenttask
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"reflect"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
type ManagerConfig struct {
|
|
OwnerID string
|
|
LeaseDuration time.Duration
|
|
MaxReworkAttempts uint32
|
|
MaxFailureAttempts uint32
|
|
StateWriteAttempts int
|
|
}
|
|
|
|
type Manager struct {
|
|
config ManagerConfig
|
|
clock Clock
|
|
store StateStore
|
|
workflow WorkflowAdapter
|
|
selector Selector
|
|
isolation IsolationBackend
|
|
invoker ProviderInvoker
|
|
recovery RecoveryInspector
|
|
evidence WorkflowEvidence
|
|
reviewer Reviewer
|
|
integrator Integrator
|
|
events EventSink
|
|
scheduler *Scheduler
|
|
|
|
// renewalTicks is a package-private test seam. Production leaves it nil
|
|
// and uses a time.Ticker; tests can drive renewal without wall-clock waits.
|
|
renewalTicks func(time.Duration) <-chan time.Time
|
|
|
|
reconcileMu sync.Mutex
|
|
deliveryMu sync.Mutex
|
|
activeMu sync.Mutex
|
|
activeRuns map[ProjectID]context.CancelFunc
|
|
}
|
|
|
|
type deliveryErrorContextKey struct{}
|
|
|
|
type deliveryErrorCollector struct {
|
|
mu sync.Mutex
|
|
errs []error
|
|
}
|
|
|
|
func withDeliveryErrors(ctx context.Context) (context.Context, *deliveryErrorCollector) {
|
|
collector := &deliveryErrorCollector{}
|
|
return context.WithValue(ctx, deliveryErrorContextKey{}, collector), collector
|
|
}
|
|
|
|
func recordDeliveryError(ctx context.Context, err error) {
|
|
if err == nil {
|
|
return
|
|
}
|
|
collector, _ := ctx.Value(deliveryErrorContextKey{}).(*deliveryErrorCollector)
|
|
if collector == nil {
|
|
return
|
|
}
|
|
collector.mu.Lock()
|
|
collector.errs = append(collector.errs, err)
|
|
collector.mu.Unlock()
|
|
}
|
|
|
|
func (c *deliveryErrorCollector) Err() error {
|
|
if c == nil {
|
|
return nil
|
|
}
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
return errors.Join(c.errs...)
|
|
}
|
|
|
|
// leaseSet owns the collection of active durable lease claims for one
|
|
// reconciliation. A background supervisor renews each claim by CAS at a
|
|
// bounded fraction of LeaseDuration; if any renewal cannot prove the token
|
|
// still matches, the guarded context is cancelled so that no late result can
|
|
// commit durable state under a stale owner.
|
|
type leaseSet struct {
|
|
mu sync.Mutex
|
|
claims []*leaseClaim
|
|
cancel context.CancelFunc
|
|
manager *Manager
|
|
wg sync.WaitGroup
|
|
}
|
|
|
|
type leaseSetContextKey struct{}
|
|
type leaseFenceBypassContextKey struct{}
|
|
|
|
// renewAll renews every tracked claim by CAS. It returns the first error so
|
|
// the supervisor can cancel the guarded context and stop the loop.
|
|
func (ls *leaseSet) renewAll(ctx context.Context) error {
|
|
ls.mu.Lock()
|
|
defer ls.mu.Unlock()
|
|
for _, claim := range ls.claims {
|
|
if err := ls.manager.renewOneLease(ctx, claim); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Validate confirms that every tracked claim still matches the current
|
|
// durable state. It returns ErrLeaseLost for the first mismatch so callers
|
|
// can reject an external result that arrived after ownership was lost.
|
|
func (ls *leaseSet) Validate(ctx context.Context) error {
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
default:
|
|
}
|
|
ls.mu.Lock()
|
|
claims := make([]*leaseClaim, len(ls.claims))
|
|
copy(claims, ls.claims)
|
|
ls.mu.Unlock()
|
|
|
|
for _, claim := range claims {
|
|
if err := ls.manager.validateOneLease(ctx, claim); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Add registers a newly acquired claim with the set so the supervisor also
|
|
// renews it and the fence validator also checks it.
|
|
func (ls *leaseSet) Add(claim *leaseClaim) {
|
|
if claim == nil {
|
|
return
|
|
}
|
|
ls.mu.Lock()
|
|
ls.claims = append(ls.claims, claim)
|
|
ls.mu.Unlock()
|
|
}
|
|
|
|
// Remove prevents any later renewal pass from observing claim. It waits for
|
|
// an in-progress renewal pass through the set lock before the caller performs
|
|
// the exact-token release.
|
|
func (ls *leaseSet) Remove(claim *leaseClaim) {
|
|
if claim == nil {
|
|
return
|
|
}
|
|
ls.mu.Lock()
|
|
defer ls.mu.Unlock()
|
|
for index, current := range ls.claims {
|
|
if current == claim {
|
|
ls.claims = append(ls.claims[:index], ls.claims[index+1:]...)
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
func (ls *leaseSet) snapshotClaims() []*leaseClaim {
|
|
ls.mu.Lock()
|
|
defer ls.mu.Unlock()
|
|
claims := make([]*leaseClaim, len(ls.claims))
|
|
copy(claims, ls.claims)
|
|
return claims
|
|
}
|
|
|
|
// Close cancels the supervisor and the guarded context and waits for the
|
|
// renewal goroutine to exit. It does not touch durable state; callers
|
|
// release each lease by exact token after closing.
|
|
func (ls *leaseSet) Close() {
|
|
if ls == nil || ls.cancel == nil {
|
|
return
|
|
}
|
|
ls.cancel()
|
|
ls.wg.Wait()
|
|
}
|
|
|
|
// maintainLeases starts a reconciliation-owned lease supervisor. The returned
|
|
// context is cancelled when any renewal fails, so that every mutation under
|
|
// it is guaranteed to observe a token the supervisor still holds.
|
|
func (m *Manager) maintainLeases(ctx context.Context, initial *leaseClaim) (context.Context, *leaseSet) {
|
|
ctx, cancel := context.WithCancel(ctx)
|
|
ls := &leaseSet{
|
|
claims: []*leaseClaim{initial},
|
|
cancel: cancel,
|
|
manager: m,
|
|
}
|
|
ctx = context.WithValue(ctx, leaseSetContextKey{}, ls)
|
|
ls.wg.Add(1)
|
|
go ls.runRenewalLoop(ctx)
|
|
return ctx, ls
|
|
}
|
|
|
|
// runRenewalLoop ticks at a bounded fraction of LeaseDuration and renews all
|
|
// tracked claims. A single CAS miss cancels the guarded context and stops.
|
|
func (ls *leaseSet) runRenewalLoop(ctx context.Context) {
|
|
defer ls.wg.Done()
|
|
interval := ls.manager.config.LeaseDuration / 3
|
|
if interval < time.Millisecond {
|
|
interval = time.Millisecond
|
|
}
|
|
if ticks := ls.manager.renewalTicks; ticks != nil {
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticks(interval):
|
|
if err := ls.renewAll(unfencedLeaseContext(ctx)); err != nil {
|
|
ls.cancel()
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
ticker := time.NewTicker(interval)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
if err := ls.renewAll(unfencedLeaseContext(ctx)); err != nil {
|
|
ls.cancel()
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// Wait blocks until the renewal supervisor has fully stopped.
|
|
func (ls *leaseSet) Wait() {
|
|
ls.wg.Wait()
|
|
}
|
|
|
|
// renewOneLease renews a single claim by CAS. It returns ErrLeaseLost when
|
|
// the current state no longer holds the tracked token.
|
|
func (m *Manager) renewOneLease(ctx context.Context, claim *leaseClaim) error {
|
|
switch claim.scope {
|
|
case "device":
|
|
return m.renewDevice(ctx, claim.token)
|
|
case "workspace":
|
|
return m.renewWorkspace(ctx, claim.token, claim.subject)
|
|
case "project":
|
|
return m.renewProject(ctx, claim.token, claim.subject)
|
|
case "integration":
|
|
return m.renewIntegration(ctx, claim.token, claim.subject)
|
|
default:
|
|
return fmt.Errorf("agenttask: unknown lease scope %q", claim.scope)
|
|
}
|
|
}
|
|
|
|
// validateOneLease checks that a single claim still matches current state.
|
|
func (m *Manager) validateOneLease(ctx context.Context, claim *leaseClaim) error {
|
|
switch claim.scope {
|
|
case "device":
|
|
return m.validateDevice(ctx, claim.token)
|
|
case "workspace":
|
|
return m.validateWorkspace(ctx, claim.token, claim.subject)
|
|
case "project":
|
|
return m.validateProject(ctx, claim.token, claim.subject)
|
|
case "integration":
|
|
return m.validateIntegration(ctx, claim.token, claim.subject)
|
|
default:
|
|
return fmt.Errorf("agenttask: unknown lease scope %q", claim.scope)
|
|
}
|
|
}
|
|
|
|
func NewManager(
|
|
config ManagerConfig,
|
|
clock Clock,
|
|
store StateStore,
|
|
workflow WorkflowAdapter,
|
|
selector Selector,
|
|
isolation IsolationBackend,
|
|
invoker ProviderInvoker,
|
|
recovery RecoveryInspector,
|
|
evidence WorkflowEvidence,
|
|
reviewer Reviewer,
|
|
integrator Integrator,
|
|
events EventSink,
|
|
) (*Manager, error) {
|
|
if err := validateIdentity("manager_owner", config.OwnerID); err != nil {
|
|
return nil, err
|
|
}
|
|
if store == nil || workflow == nil || selector == nil || isolation == nil ||
|
|
invoker == nil || recovery == nil || evidence == nil || reviewer == nil || integrator == nil {
|
|
return nil, errors.New("agenttask: every execution port is required; unsafe fallback is disabled")
|
|
}
|
|
if clock == nil {
|
|
clock = systemClock{}
|
|
}
|
|
if events == nil {
|
|
events = nopEventSink{}
|
|
}
|
|
if config.LeaseDuration <= 0 {
|
|
config.LeaseDuration = 30 * time.Second
|
|
}
|
|
if config.MaxReworkAttempts == 0 {
|
|
config.MaxReworkAttempts = 3
|
|
}
|
|
if config.MaxFailureAttempts == 0 {
|
|
config.MaxFailureAttempts = 10
|
|
}
|
|
if config.StateWriteAttempts <= 0 {
|
|
config.StateWriteAttempts = 32
|
|
}
|
|
return &Manager{
|
|
config: config,
|
|
clock: clock,
|
|
store: store,
|
|
workflow: workflow,
|
|
selector: selector,
|
|
isolation: isolation,
|
|
invoker: invoker,
|
|
recovery: recovery,
|
|
evidence: evidence,
|
|
reviewer: reviewer,
|
|
integrator: integrator,
|
|
events: events,
|
|
scheduler: NewScheduler(),
|
|
activeRuns: make(map[ProjectID]context.CancelFunc),
|
|
}, nil
|
|
}
|
|
|
|
func (m *Manager) StartProject(ctx context.Context, req StartRequest) error {
|
|
if err := validateStartRequest(req); err != nil {
|
|
return err
|
|
}
|
|
ctx, deliveryErrors := withDeliveryErrors(ctx)
|
|
if err := m.flushPendingEvents(ctx); err != nil {
|
|
recordDeliveryError(ctx, err)
|
|
return deliveryErrors.Err()
|
|
}
|
|
autoResume := true
|
|
if req.AutoResumeInterrupted != nil {
|
|
autoResume = *req.AutoResumeInterrupted
|
|
}
|
|
intent := StartIntent{
|
|
CommandID: req.CommandID,
|
|
ProjectID: req.ProjectID,
|
|
WorkspaceID: req.WorkspaceID,
|
|
MilestoneID: req.MilestoneID,
|
|
WorkflowRevision: req.WorkflowRevision,
|
|
ConfigRevision: req.ConfigRevision,
|
|
GrantRevision: req.GrantRevision,
|
|
AutoResumeInterrupted: autoResume,
|
|
StartedAt: m.clock.Now(),
|
|
}
|
|
err := m.mutate(ctx, func(state *ManagerState) error {
|
|
if previous, ok := state.Commands[req.CommandID]; ok {
|
|
if sameCommandIntent(previous.Intent, intent) {
|
|
return nil
|
|
}
|
|
return fmt.Errorf("agenttask: command %q was already used with different immutable input", req.CommandID)
|
|
}
|
|
project := state.Projects[req.ProjectID]
|
|
if project.ProjectID != "" && project.WorkspaceID != "" &&
|
|
project.WorkspaceID != req.WorkspaceID {
|
|
return fmt.Errorf("agenttask: project %q workspace identity changed", req.ProjectID)
|
|
}
|
|
project.ProjectID = req.ProjectID
|
|
project.WorkspaceID = req.WorkspaceID
|
|
project.Status = ProjectStatusStarted
|
|
project.Intent = &intent
|
|
project.Blocker = nil
|
|
project.UpdatedAt = m.clock.Now()
|
|
if project.Works == nil {
|
|
project.Works = make(map[WorkUnitID]WorkRecord)
|
|
}
|
|
state.Commands[req.CommandID] = CommandRecord{Intent: intent}
|
|
state.Projects[req.ProjectID] = project
|
|
return nil
|
|
})
|
|
if err == nil {
|
|
m.emit(ctx, Event{
|
|
Type: EventManualStart,
|
|
ProjectID: req.ProjectID,
|
|
WorkspaceID: req.WorkspaceID,
|
|
CommandID: req.CommandID,
|
|
WorkflowRevision: req.WorkflowRevision,
|
|
Detail: string(req.MilestoneID),
|
|
})
|
|
}
|
|
return errors.Join(err, deliveryErrors.Err())
|
|
}
|
|
|
|
func (m *Manager) StopProject(ctx context.Context, projectID ProjectID) error {
|
|
if err := validateIdentity("project", string(projectID)); err != nil {
|
|
return err
|
|
}
|
|
ctx, deliveryErrors := withDeliveryErrors(ctx)
|
|
if err := m.flushPendingEvents(ctx); err != nil {
|
|
recordDeliveryError(ctx, err)
|
|
return deliveryErrors.Err()
|
|
}
|
|
m.activeMu.Lock()
|
|
cancel := m.activeRuns[projectID]
|
|
m.activeMu.Unlock()
|
|
hadActiveRun := cancel != nil
|
|
if cancel != nil {
|
|
cancel()
|
|
}
|
|
var commandID CommandID
|
|
var workflowRev WorkflowRevision
|
|
err := m.mutate(ctx, func(state *ManagerState) error {
|
|
project, ok := state.Projects[projectID]
|
|
if !ok {
|
|
return fmt.Errorf("agenttask: project %q is not registered in manager state", projectID)
|
|
}
|
|
if project.Intent != nil {
|
|
commandID = project.Intent.CommandID
|
|
workflowRev = project.Intent.WorkflowRevision
|
|
}
|
|
project.Status = ProjectStatusStopped
|
|
liveLease := project.Lease != nil &&
|
|
project.Lease.ExpiresAt.After(m.clock.Now())
|
|
if !hadActiveRun && !liveLease {
|
|
project.Lease = nil
|
|
if lease, ok := state.WorkspaceLeases[project.WorkspaceID]; ok &&
|
|
lease.OwnerID == m.config.OwnerID {
|
|
delete(state.WorkspaceLeases, project.WorkspaceID)
|
|
}
|
|
if lease, ok := state.IntegrationLeases[project.WorkspaceID]; ok &&
|
|
lease.OwnerID == m.config.OwnerID {
|
|
delete(state.IntegrationLeases, project.WorkspaceID)
|
|
}
|
|
}
|
|
project.UpdatedAt = m.clock.Now()
|
|
for id, work := range project.Works {
|
|
if !work.State.Terminal() {
|
|
preStop := work.State
|
|
if preStop != WorkStateStopped {
|
|
work.ResumeStage = preStop
|
|
}
|
|
if err := transitionWork(&work, WorkStateStopped); err != nil {
|
|
return err
|
|
}
|
|
work.UpdatedAt = m.clock.Now()
|
|
project.Works[id] = work
|
|
}
|
|
}
|
|
state.Projects[projectID] = project
|
|
return nil
|
|
})
|
|
if err == nil {
|
|
m.emit(ctx, Event{
|
|
Type: EventStopped,
|
|
ProjectID: projectID,
|
|
CommandID: commandID,
|
|
WorkflowRevision: workflowRev,
|
|
})
|
|
}
|
|
return errors.Join(err, deliveryErrors.Err())
|
|
}
|
|
|
|
func mutateDecision[T any](
|
|
m *Manager,
|
|
ctx context.Context,
|
|
change func(*ManagerState) (T, error),
|
|
) (T, error) {
|
|
var zero T
|
|
var conflict error
|
|
for range m.config.StateWriteAttempts {
|
|
state, revision, err := m.store.Load(ctx)
|
|
if err != nil {
|
|
return zero, err
|
|
}
|
|
if err := validateManagerState(state); err != nil {
|
|
return zero, err
|
|
}
|
|
if err := validateContextClaims(ctx, state); err != nil {
|
|
return zero, err
|
|
}
|
|
next := cloneState(state)
|
|
if next.SchemaVersion != currentSchemaVersion {
|
|
return zero, fmt.Errorf("agenttask: unsupported state schema %d", next.SchemaVersion)
|
|
}
|
|
result, err := change(&next)
|
|
if err != nil {
|
|
return zero, err
|
|
}
|
|
if err := validateManagerState(next); err != nil {
|
|
return zero, err
|
|
}
|
|
_, err = m.store.CompareAndSwap(ctx, revision, next)
|
|
if errors.Is(err, ErrRevisionConflict) {
|
|
conflict = err
|
|
continue
|
|
}
|
|
if err != nil {
|
|
return zero, err
|
|
}
|
|
return result, nil
|
|
}
|
|
return zero, fmt.Errorf("agenttask: state CAS retry budget exhausted: %w", conflict)
|
|
}
|
|
|
|
func validateContextClaims(ctx context.Context, state ManagerState) error {
|
|
if ctx.Value(leaseFenceBypassContextKey{}) != nil {
|
|
return nil
|
|
}
|
|
leases, _ := ctx.Value(leaseSetContextKey{}).(*leaseSet)
|
|
if leases == nil {
|
|
return nil
|
|
}
|
|
for _, claim := range leases.snapshotClaims() {
|
|
if err := claimMatchesState(claim, state); err != nil {
|
|
leases.cancel()
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func claimMatchesState(claim *leaseClaim, state ManagerState) error {
|
|
if claim == nil {
|
|
return fmt.Errorf("%w: missing lease claim", ErrLeaseLost)
|
|
}
|
|
matched := false
|
|
switch claim.scope {
|
|
case "device":
|
|
matched = state.DeviceLease != nil &&
|
|
state.DeviceLease.OwnerID == claim.owner && state.DeviceLease.Token == claim.token
|
|
case "project":
|
|
project, ok := state.Projects[ProjectID(claim.subject)]
|
|
matched = ok && project.Lease != nil &&
|
|
project.Lease.OwnerID == claim.owner && project.Lease.Token == claim.token
|
|
case "workspace":
|
|
lease, ok := state.WorkspaceLeases[WorkspaceID(claim.subject)]
|
|
matched = ok && lease.OwnerID == claim.owner && lease.Token == claim.token
|
|
case "integration":
|
|
lease, ok := state.IntegrationLeases[WorkspaceID(claim.subject)]
|
|
matched = ok && lease.OwnerID == claim.owner && lease.Token == claim.token
|
|
default:
|
|
return fmt.Errorf("%w: unknown lease scope %q", ErrLeaseLost, claim.scope)
|
|
}
|
|
if !matched {
|
|
return fmt.Errorf("%w: %s lease %q no longer matches its exact token", ErrLeaseLost, claim.scope, claim.subject)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func unfencedLeaseContext(ctx context.Context) context.Context {
|
|
return context.WithValue(context.WithoutCancel(ctx), leaseFenceBypassContextKey{}, true)
|
|
}
|
|
|
|
func (m *Manager) mutate(
|
|
ctx context.Context,
|
|
change func(*ManagerState) error,
|
|
) error {
|
|
_, err := mutateDecision(m, ctx, func(state *ManagerState) (struct{}, error) {
|
|
return struct{}{}, change(state)
|
|
})
|
|
return err
|
|
}
|
|
|
|
func (m *Manager) load(ctx context.Context) (ManagerState, error) {
|
|
state, _, err := m.store.Load(ctx)
|
|
if err != nil {
|
|
return ManagerState{}, err
|
|
}
|
|
if err := validateManagerState(state); err != nil {
|
|
return ManagerState{}, err
|
|
}
|
|
state = cloneState(state)
|
|
if state.SchemaVersion != currentSchemaVersion {
|
|
return ManagerState{}, fmt.Errorf("agenttask: unsupported state schema %d", state.SchemaVersion)
|
|
}
|
|
return state, nil
|
|
}
|
|
|
|
func durableIdentity(domain string, components ...string) string {
|
|
var sb strings.Builder
|
|
sb.WriteString(domain)
|
|
for _, c := range components {
|
|
sb.WriteString("/")
|
|
sb.WriteString(strconv.Itoa(len(c)))
|
|
sb.WriteString(":")
|
|
sb.WriteString(c)
|
|
}
|
|
return sb.String()
|
|
}
|
|
|
|
func normalizeEventIdentity(event Event) Event {
|
|
if event.EventID == "" {
|
|
event.EventID = durableIdentity(
|
|
"event-v1",
|
|
string(event.Type),
|
|
string(event.ProjectID),
|
|
string(event.WorkspaceID),
|
|
string(event.WorkUnitID),
|
|
string(event.CommandID),
|
|
string(event.WorkflowRevision),
|
|
string(event.AttemptID),
|
|
strconv.FormatUint(uint64(event.Ordinal), 10),
|
|
string(event.ChangeSetID),
|
|
event.ChangeSetRevision,
|
|
strconv.FormatUint(uint64(event.IntegrationAttempt), 10),
|
|
string(event.State),
|
|
event.ProviderID,
|
|
event.ProfileID,
|
|
string(event.WriteSetKind),
|
|
string(event.IsolationMode),
|
|
event.Detail,
|
|
)
|
|
}
|
|
return event
|
|
}
|
|
|
|
func sameLogicalEvent(left, right Event) bool {
|
|
left.Timestamp = time.Time{}
|
|
right.Timestamp = time.Time{}
|
|
return reflect.DeepEqual(left, right)
|
|
}
|
|
|
|
func (m *Manager) emit(ctx context.Context, event Event) {
|
|
// Direct construction is retained only for package-level identity tests.
|
|
// NewManager requires a durable store for every production manager.
|
|
if m.store == nil {
|
|
event = normalizeEventIdentity(event)
|
|
if event.Timestamp.IsZero() {
|
|
event.Timestamp = m.clock.Now()
|
|
}
|
|
if m.events != nil {
|
|
recordDeliveryError(ctx, m.events.Emit(ctx, event))
|
|
}
|
|
return
|
|
}
|
|
_, err := m.enqueueEvent(ctx, event)
|
|
if err == nil {
|
|
err = m.flushPendingEvents(ctx)
|
|
}
|
|
recordDeliveryError(ctx, err)
|
|
}
|
|
|
|
func (m *Manager) enqueueEvent(ctx context.Context, event Event) (EventDelivery, error) {
|
|
event = normalizeEventIdentity(event)
|
|
var conflict error
|
|
for range m.config.StateWriteAttempts {
|
|
state, revision, err := m.store.Load(ctx)
|
|
if err != nil {
|
|
return EventDelivery{}, err
|
|
}
|
|
if err := validateManagerState(state); err != nil {
|
|
return EventDelivery{}, err
|
|
}
|
|
if err := validateContextClaims(ctx, state); err != nil {
|
|
return EventDelivery{}, err
|
|
}
|
|
if pending, ok := state.PendingEvents[event.EventID]; ok {
|
|
if !sameLogicalEvent(pending.Event, event) {
|
|
return EventDelivery{}, fmt.Errorf(
|
|
"agenttask: pending event %q was reused with conflicting logical content",
|
|
event.EventID,
|
|
)
|
|
}
|
|
return cloneEventDelivery(pending), nil
|
|
}
|
|
project, ok := state.Projects[event.ProjectID]
|
|
if !ok {
|
|
return EventDelivery{}, fmt.Errorf(
|
|
"agenttask: event %q project %q is absent from committed state",
|
|
event.EventID,
|
|
event.ProjectID,
|
|
)
|
|
}
|
|
if event.Timestamp.IsZero() {
|
|
event.Timestamp = m.clock.Now()
|
|
}
|
|
project = cloneProject(project)
|
|
delivery := EventDelivery{
|
|
Event: event,
|
|
EvidenceRevision: revision,
|
|
Project: &project,
|
|
}
|
|
if event.WorkUnitID != "" {
|
|
work, ok := project.Works[event.WorkUnitID]
|
|
if !ok {
|
|
return EventDelivery{}, fmt.Errorf(
|
|
"agenttask: event %q work %q is absent from committed state",
|
|
event.EventID,
|
|
event.WorkUnitID,
|
|
)
|
|
}
|
|
work = cloneWork(work)
|
|
delivery.Work = &work
|
|
}
|
|
if err := validateEventDelivery(event.EventID, delivery); err != nil {
|
|
return EventDelivery{}, err
|
|
}
|
|
next := cloneState(state)
|
|
next.PendingEvents[event.EventID] = delivery
|
|
if _, err := m.store.CompareAndSwap(ctx, revision, next); err != nil {
|
|
if errors.Is(err, ErrRevisionConflict) {
|
|
conflict = err
|
|
continue
|
|
}
|
|
return EventDelivery{}, err
|
|
}
|
|
return cloneEventDelivery(delivery), nil
|
|
}
|
|
return EventDelivery{}, fmt.Errorf(
|
|
"agenttask: pending event CAS retry budget exhausted: %w",
|
|
conflict,
|
|
)
|
|
}
|
|
|
|
func (m *Manager) flushPendingEvents(ctx context.Context) error {
|
|
if m.store == nil {
|
|
return nil
|
|
}
|
|
m.deliveryMu.Lock()
|
|
defer m.deliveryMu.Unlock()
|
|
|
|
state, _, err := m.store.Load(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := validateManagerState(state); err != nil {
|
|
return err
|
|
}
|
|
eventIDs := make([]string, 0, len(state.PendingEvents))
|
|
for eventID := range state.PendingEvents {
|
|
eventIDs = append(eventIDs, eventID)
|
|
}
|
|
sort.Strings(eventIDs)
|
|
for _, eventID := range eventIDs {
|
|
if err := m.flushEventDeliveryLocked(ctx, eventID); err != nil {
|
|
return fmt.Errorf("agenttask: deliver pending event %q: %w", eventID, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (m *Manager) flushEventDeliveryLocked(ctx context.Context, eventID string) error {
|
|
state, _, err := m.store.Load(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := validateManagerState(state); err != nil {
|
|
return err
|
|
}
|
|
if err := validateContextClaims(ctx, state); err != nil {
|
|
return err
|
|
}
|
|
delivery, ok := state.PendingEvents[eventID]
|
|
if !ok {
|
|
return nil
|
|
}
|
|
if err := m.events.Emit(ctx, delivery.Event); err != nil {
|
|
return err
|
|
}
|
|
var conflict error
|
|
for range m.config.StateWriteAttempts {
|
|
state, revision, err := m.store.Load(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := validateManagerState(state); err != nil {
|
|
return err
|
|
}
|
|
if err := validateContextClaims(ctx, state); err != nil {
|
|
return err
|
|
}
|
|
current, ok := state.PendingEvents[eventID]
|
|
if !ok {
|
|
return nil
|
|
}
|
|
if !reflect.DeepEqual(current, delivery) {
|
|
return fmt.Errorf(
|
|
"agenttask: pending event %q changed before acknowledgement",
|
|
eventID,
|
|
)
|
|
}
|
|
next := cloneState(state)
|
|
delete(next.PendingEvents, eventID)
|
|
if _, err := m.store.CompareAndSwap(ctx, revision, next); err != nil {
|
|
if errors.Is(err, ErrRevisionConflict) {
|
|
conflict = err
|
|
continue
|
|
}
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
return fmt.Errorf(
|
|
"agenttask: pending event acknowledgement CAS retry budget exhausted: %w",
|
|
conflict,
|
|
)
|
|
}
|
|
|
|
func sameCommandIntent(left, right StartIntent) bool {
|
|
left.StartedAt = time.Time{}
|
|
right.StartedAt = time.Time{}
|
|
return reflect.DeepEqual(left, right)
|
|
}
|
|
|
|
func (m *Manager) beginProjectRun(
|
|
parent context.Context,
|
|
projectID ProjectID,
|
|
) (context.Context, func()) {
|
|
ctx, cancel := context.WithCancel(parent)
|
|
monitorDone := make(chan struct{})
|
|
go m.monitorProjectRun(ctx, projectID, cancel, monitorDone)
|
|
m.activeMu.Lock()
|
|
if previous := m.activeRuns[projectID]; previous != nil {
|
|
previous()
|
|
}
|
|
m.activeRuns[projectID] = cancel
|
|
m.activeMu.Unlock()
|
|
return ctx, func() {
|
|
cancel()
|
|
<-monitorDone
|
|
m.activeMu.Lock()
|
|
delete(m.activeRuns, projectID)
|
|
m.activeMu.Unlock()
|
|
}
|
|
}
|
|
|
|
func (m *Manager) monitorProjectRun(
|
|
ctx context.Context,
|
|
projectID ProjectID,
|
|
cancel context.CancelFunc,
|
|
done chan<- struct{},
|
|
) {
|
|
defer close(done)
|
|
interval := m.config.LeaseDuration / 10
|
|
if interval < 25*time.Millisecond {
|
|
interval = 25 * time.Millisecond
|
|
}
|
|
if interval > 250*time.Millisecond {
|
|
interval = 250 * time.Millisecond
|
|
}
|
|
ticker := time.NewTicker(interval)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
state, err := m.load(ctx)
|
|
if err != nil {
|
|
cancel()
|
|
return
|
|
}
|
|
project, ok := state.Projects[projectID]
|
|
if !ok || project.Status == ProjectStatusStopped {
|
|
cancel()
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func attemptID(workID WorkUnitID, attempt uint32) AttemptID {
|
|
return AttemptID(durableIdentity("attempt-v1", string(workID), strconv.FormatUint(uint64(attempt), 10)))
|
|
}
|
|
|
|
func dispatchKey(projectID ProjectID, workID WorkUnitID, attempt uint32) string {
|
|
return durableIdentity("dispatch-v1", string(projectID), string(workID), strconv.FormatUint(uint64(attempt), 10))
|
|
}
|
|
|
|
func reviewKey(projectID ProjectID, workID WorkUnitID, attempt uint32, artifact ArtifactID) string {
|
|
return durableIdentity("review-v1", string(projectID), string(workID), strconv.FormatUint(uint64(attempt), 10), string(artifact))
|
|
}
|
|
|
|
func integrationKey(
|
|
projectID ProjectID,
|
|
workID WorkUnitID,
|
|
changeSet ChangeSetIdentity,
|
|
attempt IntegrationAttempt,
|
|
) string {
|
|
return durableIdentity(
|
|
"integrate-v1",
|
|
string(projectID), string(workID), string(changeSet.ID), string(changeSet.Revision), strconv.FormatUint(uint64(attempt), 10),
|
|
)
|
|
}
|