package queue import ( "context" "errors" "sync" "time" "git.ilapage.cn/OPC/chorus/internal/core/model" ) var ErrClaimsStopped = errors.New("queue claims are stopped") type Controller struct { repository Repository mutex sync.RWMutex stopped bool } func NewController(repository Repository) *Controller { return &Controller{repository: repository} } func (c *Controller) StopClaims() { c.mutex.Lock() defer c.mutex.Unlock() c.stopped = true } func (c *Controller) ClaimNext(ctx context.Context, owner string, leaseDuration time.Duration) (*Claim, error) { c.mutex.RLock() defer c.mutex.RUnlock() if c.stopped { return nil, ErrClaimsStopped } return c.repository.ClaimNext(ctx, owner, leaseDuration) } func (c *Controller) AssignProvider(ctx context.Context, generationID uint64, leaseToken string, providerModelID uint64) (bool, error) { return c.repository.AssignProvider(ctx, generationID, leaseToken, providerModelID) } func (c *Controller) BeginProviderAttempt(ctx context.Context, generationID uint64, leaseToken string, routeMemberID, providerModelID uint64) (model.Attempt, bool, error) { return c.repository.BeginProviderAttempt(ctx, generationID, leaseToken, routeMemberID, providerModelID) } func (c *Controller) FinishProviderAttempt(ctx context.Context, generationID uint64, leaseToken string, attempt model.Attempt) (bool, error) { return c.repository.FinishProviderAttempt(ctx, generationID, leaseToken, attempt) } func (c *Controller) Succeed(ctx context.Context, generationID uint64, leaseToken string, outputs []model.GenerationOutput, attempt model.Attempt) (bool, error) { return c.repository.Succeed(ctx, generationID, leaseToken, outputs, attempt) } func (c *Controller) Fail(ctx context.Context, generationID uint64, leaseToken, code, message string, attempt model.Attempt) (bool, error) { return c.repository.Fail(ctx, generationID, leaseToken, code, message, attempt) } func (c *Controller) Inputs(ctx context.Context, generationID uint64) ([]model.GenerationInput, error) { return c.repository.Inputs(ctx, generationID) }