Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -274,6 +274,23 @@ The corresponding exported constants are `DefaultWorkerQueueCapacity`,
`DefaultWorkerMaxPayloadBytes`, `DefaultWorkerMaxQueueBytes`, `MaxWorkerQueueCapacity`,
`MaxWorkerPayloadBytes`, and `MaxWorkerQueueBytes`.

`WorkerOptions` bounds one worker; `WorkerLimits` bounds the **whole service**, so a
guest cannot exhaust the host by spawning workers without end. It always applies — even
when the host grants `instance.manage` with no `maxInstances` budget — and complements
that core budget rather than replacing it. Pass it at construction:

```go
workerPlugin := workers.New(workers.WithLimits(workers.WorkerLimits{
MaxLiveWorkers: 128, // default 64
MaxQueueBytes: 128 << 20, // default 64 MiB, summed across live workers
}))
```

| Field | Default | Meaning |
| --- | --- | --- |
| `MaxLiveWorkers` | `64` | Max workers live at once. `Spawn` returns `ErrWorkerQuotaExceeded` past it. |
| `MaxQueueBytes` | `64 MiB` | Max total per-worker `MaxQueueBytes` reservation summed across live workers. |

### Errors

All errors are comparable sentinels - match them with `errors.Is`.
Expand All @@ -294,6 +311,7 @@ All errors are comparable sentinels - match them with `errors.Is`.
| `ErrWorkerKilled` | Exit cause: `Kill` was requested. |
| `ErrWorkerParentClosed` | Exit cause: the linked creator instance closed. |
| `ErrWorkerRuntimeClosed` | Exit cause: the runtime shut down. |
| `ErrWorkerQuotaExceeded` | Spawn would exceed the service-wide `WorkerLimits`. |

## Examples

Expand Down
103 changes: 99 additions & 4 deletions workers.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,16 @@ const (
MaxWorkerQueueCapacity uint32 = 1 << 16
MaxWorkerPayloadBytes uint32 = 16 << 20
MaxWorkerQueueBytes uint32 = 64 << 20

// DefaultMaxLiveWorkers bounds the number of simultaneously live workers for
// one service when WorkerLimits.MaxLiveWorkers is left zero. Each worker owns a
// managed instance, a goroutine, and a foreign stack, so an unbounded count is a
// denial-of-service vector; a default keeps Spawn bounded even when the host
// grants instance.manage without a maxInstances budget.
DefaultMaxLiveWorkers uint32 = 64
// DefaultMaxServiceQueueBytes bounds the total queued-payload reservation across
// all live workers when WorkerLimits.MaxQueueBytes is left zero.
DefaultMaxServiceQueueBytes uint64 = 64 << 20
)

type workerError string
Expand All @@ -43,8 +53,33 @@ const (
ErrWorkerKilled workerError = "worker killed"
ErrWorkerParentClosed workerError = "worker parent closed"
ErrWorkerRuntimeClosed workerError = "worker runtime closed"
ErrWorkerQuotaExceeded workerError = "worker service resource quota exceeded"
)

// WorkerLimits bounds the aggregate resources one worker service may hold at
// once, independent of the per-worker WorkerOptions. Zero fields take the
// package defaults. It complements — and does not replace — the core
// instance.manage maxInstances budget: this cap always applies, so workers stay
// bounded even when the host grants instance.manage without a budget.
type WorkerLimits struct {
// MaxLiveWorkers is the maximum number of simultaneously live workers.
MaxLiveWorkers uint32
// MaxQueueBytes is the maximum total per-worker queue-byte reservation summed
// across all live workers (each worker reserves its MaxQueueBytes for its
// lifetime).
MaxQueueBytes uint64
}

func normalizeLimits(l WorkerLimits) WorkerLimits {
if l.MaxLiveWorkers == 0 {
l.MaxLiveWorkers = DefaultMaxLiveWorkers
}
if l.MaxQueueBytes == 0 {
l.MaxQueueBytes = DefaultMaxServiceQueueBytes
}
return l
}

type WorkerOptions struct {
QueueCapacity uint32
MaxPayloadBytes uint32
Expand Down Expand Up @@ -76,9 +111,25 @@ var ServiceKey = plugin.NewServiceKey[*Workers]("wago.workers/v1")

type Plugin struct {
service *Workers
limits WorkerLimits
}

func New() *Plugin { return &Plugin{} }
// Option configures the workers Plugin at construction.
type Option func(*Plugin)

// WithLimits sets the aggregate resource limits for the worker service. Zero
// fields fall back to the package defaults (see WorkerLimits).
func WithLimits(l WorkerLimits) Option { return func(p *Plugin) { p.limits = l } }

// New creates the workers plugin. Pass WithLimits to override the default
// aggregate resource caps.
func New(opts ...Option) *Plugin {
p := &Plugin{}
for _, opt := range opts {
opt(p)
}
return p
}

func (*Plugin) Info() wago.ExtensionInfo {
return wago.ExtensionInfo{
Expand All @@ -100,7 +151,7 @@ func (p *Plugin) Register(reg *wago.Registry) error {
if err != nil {
return err
}
p.service = newWorkers(manager)
p.service = newWorkers(manager, p.limits)
lifecycle.BeforeClose(func(ctx *wago.InstanceContext) { p.service.parentClosing(ctx.Instance) })
return plugin.Provide(reg, ServiceKey, p.service)
}
Expand All @@ -124,6 +175,9 @@ func init() { wago.RegisterExtension(PluginName, func() wago.Extension { return
type Workers struct {
mu sync.Mutex
manager *wago.InstanceManager
limits WorkerLimits
live uint32 // number of workers currently holding a quota reservation
queueBytes uint64 // total per-worker queue-byte reservation currently held
next WorkerID
workers map[WorkerID]*worker
byInstance map[*wago.Instance]*worker
Expand All @@ -133,8 +187,34 @@ type Workers struct {
exitPanics []error
}

func newWorkers(manager *wago.InstanceManager) *Workers {
return &Workers{manager: manager, next: 1, workers: map[WorkerID]*worker{}, byInstance: map[*wago.Instance]*worker{}}
func newWorkers(manager *wago.InstanceManager, limits WorkerLimits) *Workers {
return &Workers{manager: manager, limits: normalizeLimits(limits), next: 1,
workers: map[WorkerID]*worker{}, byInstance: map[*wago.Instance]*worker{}}
}

// reserve claims one live-worker slot and queueBytes of the aggregate queue-byte
// budget, or reports ErrWorkerQuotaExceeded. The reservation is held until the
// worker's goroutine finishes (see release), so a concurrent Spawn cannot exceed
// the ceiling while a worker is still finalizing.
func (w *Workers) reserve(queueBytes uint32) error {
w.mu.Lock()
defer w.mu.Unlock()
if w.closed {
return ErrWorkerRuntimeClosed
}
if w.live >= w.limits.MaxLiveWorkers || uint64(queueBytes) > w.limits.MaxQueueBytes-w.queueBytes {
return ErrWorkerQuotaExceeded
}
w.live++
w.queueBytes += uint64(queueBytes)
return nil
}

func (w *Workers) release(queueBytes uint32) {
w.mu.Lock()
w.live--
w.queueBytes -= uint64(queueBytes)
w.mu.Unlock()
}

func (w *Workers) OnMessage(fns ...func(*MessageContext) error) {
Expand Down Expand Up @@ -180,26 +260,37 @@ func (w *Workers) Spawn(caller wago.HostModule, tableIndex uint32, opts WorkerOp
if err != nil {
return 0, err
}
// Claim aggregate quota before forking so an over-limit Spawn never allocates a
// managed instance. The reservation is released by the worker's goroutine when
// it exits, or here on any failure before the goroutine starts.
if err := w.reserve(opts.MaxQueueBytes); err != nil {
return 0, err
}
child, err := w.manager.Fork(context.Background(), caller)
if errors.Is(err, wago.ErrManagedImportLifetime) {
w.release(opts.MaxQueueBytes)
return 0, ErrWorkerImportLifetime
}
if err != nil {
w.release(opts.MaxQueueBytes)
return 0, err
}
if err := child.ValidateVoidTableEntry(tableIndex); err != nil {
_ = child.Close()
w.release(opts.MaxQueueBytes)
return 0, err
}
w.mu.Lock()
if w.closed {
w.mu.Unlock()
_ = child.Close()
w.release(opts.MaxQueueBytes)
return 0, ErrWorkerRuntimeClosed
}
if w.next == 0 {
w.mu.Unlock()
_ = child.Close()
w.release(opts.MaxQueueBytes)
return 0, ErrWorkerIDExhausted
}
id := w.next
Expand Down Expand Up @@ -443,6 +534,10 @@ func (wr *worker) run() {
fn(ctx)
}()
}
// Release aggregate quota only after the instance is closed and every exit
// observer has run, so a concurrent Spawn cannot exceed the ceiling while this
// worker is still finalizing.
wr.owner.release(wr.maxQueueBytes)
close(wr.done)
}

Expand Down
44 changes: 44 additions & 0 deletions workers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,50 @@ func integrationModule() []byte {
)
}

func TestWorkerLimitsQuota(t *testing.T) {
// Defaults apply when a field is zero.
if got := normalizeLimits(WorkerLimits{}); got.MaxLiveWorkers != DefaultMaxLiveWorkers || got.MaxQueueBytes != DefaultMaxServiceQueueBytes {
t.Fatalf("normalizeLimits(zero) = %+v", got)
}
if got := normalizeLimits(WorkerLimits{MaxLiveWorkers: 3}); got.MaxQueueBytes != DefaultMaxServiceQueueBytes || got.MaxLiveWorkers != 3 {
t.Fatalf("normalizeLimits partial = %+v", got)
}

// MaxLiveWorkers ceiling: the third reservation is rejected, and a release frees a slot.
w := newWorkers(nil, WorkerLimits{MaxLiveWorkers: 2, MaxQueueBytes: 1 << 20})
if err := w.reserve(100); err != nil {
t.Fatalf("reserve 1: %v", err)
}
if err := w.reserve(100); err != nil {
t.Fatalf("reserve 2: %v", err)
}
if err := w.reserve(100); !errors.Is(err, ErrWorkerQuotaExceeded) {
t.Fatalf("reserve 3 = %v, want ErrWorkerQuotaExceeded", err)
}
w.release(100)
if err := w.reserve(100); err != nil {
t.Fatalf("reserve after release: %v", err)
}

// Aggregate queue-byte ceiling is enforced independently and cannot overflow.
w2 := newWorkers(nil, WorkerLimits{MaxLiveWorkers: 100, MaxQueueBytes: 1000})
if err := w2.reserve(600); err != nil {
t.Fatalf("reserve 600: %v", err)
}
if err := w2.reserve(600); !errors.Is(err, ErrWorkerQuotaExceeded) {
t.Fatalf("reserve 600 over budget = %v, want ErrWorkerQuotaExceeded", err)
}
if err := w2.reserve(400); err != nil {
t.Fatalf("reserve exact remaining 400: %v", err)
}

// A closed service rejects reservations.
w2.closed = true
if err := w2.reserve(1); !errors.Is(err, ErrWorkerRuntimeClosed) {
t.Fatalf("reserve on closed = %v, want ErrWorkerRuntimeClosed", err)
}
}

func TestPluginSpawnsCopiesMessageAndStops(t *testing.T) {
p := &integrationPlugin{exits: make(chan WorkerExitContext, 1), messages: make(chan MessageContext, 1)}
rt := wago.NewRuntime()
Expand Down
Loading