Documentation
¶
Index ¶
- Constants
- Variables
- func AuthorizeTenant(ctx context.Context, job Job) error
- func CancelAuthorized(ctx context.Context, queue Queue, jobID string) error
- func RequeueAuthorized(ctx context.Context, queue Queue, admin JobAdmin, jobID string) error
- func StampTenant(ctx context.Context, job *Job) error
- func TenantIDFromContext(ctx context.Context) string
- type EventPayload
- type Handler
- type HandlerFunc
- type Job
- type JobAdmin
- type JobFilter
- type JobState
- type Lease
- type LeaseReleaser
- type LeaseRenewer
- type MemoryReconcilePayload
- type PauseResult
- type Queue
- type QueueMetrics
- type ResumeContinuePayload
- type RunPausedError
- type RunPayload
- type Worker
- type WorkerConfig
- type WorkerOption
Constants ¶
const EventJobType = "event"
const MemoryReconcileJobType = "memory.reconcile"
const ResumeContinueJobType = "resume.continue"
const RunJobType = "run"
Variables ¶
var ( ErrJobNotFound = errors.New("async: job not found") ErrStaleLease = errors.New("async: stale job lease") ErrInvalidJob = errors.New("async: invalid job") )
var ( ErrTenantMismatch = errors.New("async: tenant mismatch") ErrTenantRequired = errors.New("async: tenant identity required") )
Functions ¶
func AuthorizeTenant ¶ added in v0.4.5
AuthorizeTenant enforces ownership for tenant-scoped queue access. Worker and maintenance contexts without a principal remain global unless strict mode is explicitly enabled.
func CancelAuthorized ¶ added in v0.4.5
func RequeueAuthorized ¶ added in v0.4.5
func StampTenant ¶ added in v0.4.5
StampTenant assigns the authenticated principal tenant to a newly enqueued job. A caller cannot overwrite an explicitly supplied, different tenant.
func TenantIDFromContext ¶ added in v0.4.5
TenantIDFromContext returns the authenticated principal's tenant, if any.
Types ¶
type EventPayload ¶ added in v0.1.2
type EventPayload struct {
Type string `json:"type"`
RunID string `json:"run_id,omitempty"`
Payload json.RawMessage `json:"payload,omitempty"`
MaxAttempts int `json:"max_attempts,omitempty"`
Metadata map[string]string `json:"metadata,omitempty"`
Principal identity.Principal `json:"principal,omitempty"`
}
func (EventPayload) Event ¶ added in v0.1.2
func (payload EventPayload) Event() eventrouter.Event
type HandlerFunc ¶
type Job ¶
type Job struct {
ID string `json:"id"`
Type string `json:"type"`
RunID string `json:"run_id,omitempty"`
TenantID string `json:"tenant_id,omitempty"`
Payload json.RawMessage `json:"payload,omitempty"`
State JobState `json:"state"`
Attempts int `json:"attempts"`
MaxAttempts int `json:"max_attempts"`
LastError string `json:"last_error,omitempty"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
AvailableAt time.Time `json:"available_at"`
LeaseWorkerID string `json:"lease_worker_id,omitempty"`
LeaseExpiresAt time.Time `json:"lease_expires_at,omitempty"`
}
func LoadAuthorized ¶ added in v0.4.5
type JobAdmin ¶ added in v0.1.4
type JobAdmin interface {
ListJobs(ctx context.Context, filter JobFilter) ([]Job, error)
Requeue(ctx context.Context, jobID string) error
}
JobAdmin extends Queue with dead-letter inspection and requeue support.
type JobFilter ¶ added in v0.1.4
func ScopeJobFilter ¶ added in v0.4.5
ScopeJobFilter binds an admin listing to the authenticated tenant. Internal maintenance callers without a principal may still request an explicit tenant or leave it empty for a global view.
type LeaseReleaser ¶ added in v0.3.1
LeaseReleaser is an optional Queue capability that returns a leased job to the queued state without counting a failure, so it becomes leaseable by any worker again. Workers use it as a best-effort fallback when a leased job can neither be loaded nor failed (e.g. the queue is mid-outage); queues without it simply let the lease expire, after which the job is leaseable again anyway.
type LeaseRenewer ¶
type MemoryReconcilePayload ¶ added in v0.1.9
type MemoryReconcilePayload struct {
MemoryName string `json:"memory_name"`
Agent string `json:"agent"`
RunID string `json:"run_id,omitempty"`
Principal identity.Principal `json:"principal,omitempty"`
}
func (MemoryReconcilePayload) MarshalJSONBytes ¶ added in v0.1.9
func (payload MemoryReconcilePayload) MarshalJSONBytes() (json.RawMessage, error)
type PauseResult ¶ added in v0.3.0
PauseResult carries metadata persisted when a job pauses for HITL.
type Queue ¶
type Queue interface {
Enqueue(ctx context.Context, job Job) (Job, error)
Lease(ctx context.Context, workerID string, ttl time.Duration) (Lease, bool, error)
Load(ctx context.Context, jobID string) (Job, error)
Complete(ctx context.Context, lease Lease) error
Pause(ctx context.Context, lease Lease, result PauseResult) error
Fail(ctx context.Context, lease Lease, cause error) error
Cancel(ctx context.Context, jobID string) error
}
type QueueMetrics ¶ added in v0.1.4
QueueMetrics summarizes queue depth by state.
func CollectQueueMetrics ¶ added in v0.1.4
func CollectQueueMetrics(ctx context.Context, queue Queue) (QueueMetrics, error)
CollectQueueMetrics counts jobs by state when the queue implements JobAdmin.
type ResumeContinuePayload ¶ added in v0.1.2
type ResumeContinuePayload struct {
Token string `json:"token"`
Decision core.Decision `json:"decision"`
Amendment json.RawMessage `json:"amendment,omitempty"`
RunID string `json:"run_id,omitempty"`
MaxAttempts int `json:"max_attempts,omitempty"`
Metadata map[string]string `json:"metadata,omitempty"`
Principal identity.Principal `json:"principal,omitempty"`
}
type RunPausedError ¶ added in v0.3.0
RunPausedError indicates a run job paused awaiting human approval. Workers should mark the job as paused rather than completed or failed.
func (RunPausedError) Error ¶ added in v0.3.0
func (e RunPausedError) Error() string
type RunPayload ¶
type RunPayload struct {
RunID string `json:"run_id,omitempty"`
ScenarioID string `json:"scenario_id,omitempty"`
Scenario json.RawMessage `json:"scenario,omitempty"`
Agent string `json:"agent,omitempty"`
Prompt string `json:"prompt,omitempty"`
Context json.RawMessage `json:"context,omitempty"`
MaxAttempts int `json:"max_attempts,omitempty"`
Metadata map[string]string `json:"metadata,omitempty"`
Principal identity.Principal `json:"principal,omitempty"`
}
type Worker ¶
type Worker struct {
// contains filtered or unexported fields
}
func NewWorker ¶
func NewWorker(queue Queue, handler Handler, config WorkerConfig, opts ...WorkerOption) (*Worker, error)
type WorkerConfig ¶
type WorkerOption ¶ added in v0.3.1
type WorkerOption func(*Worker)
func WithDrainTimeout ¶ added in v0.3.1
func WithDrainTimeout(d time.Duration) WorkerOption
WithDrainTimeout enables graceful shutdown draining: when the worker context is cancelled the worker stops leasing new jobs, then waits up to d for in-flight jobs to finish before falling back to the cancel-and-requeue path. Lease renewal keeps running during the drain window so drained jobs do not lose their leases. The zero value (the default) preserves the previous behaviour of cancelling in-flight jobs immediately.