66 lines
2.2 KiB
Go
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)
|
|
}
|