Files
chorus/internal/core/queue/controller.go
T

66 lines
2.2 KiB
Go

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) Defer(ctx context.Context, generationID uint64, leaseToken string, availableAt time.Time) (bool, error) {
return c.repository.Defer(ctx, generationID, leaseToken, availableAt)
}
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)
}