mirror of
https://github.com/XTLS/Xray-core.git
synced 2026-10-04 23:03:33 +00:00
refactor: move execution context and timeout management into pool
This commit is contained in:
+42
-19
@@ -2,7 +2,9 @@ package lua
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
glua "github.com/yuin/gopher-lua"
|
||||
)
|
||||
@@ -13,8 +15,9 @@ const maxIdleStates = 16
|
||||
// keeps up to maxIdleStates idle states until Close. Acquire/Release callers
|
||||
// decide reusability; WithState uses its callback's error.
|
||||
type Pool struct {
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
timeout time.Duration
|
||||
|
||||
factory LStateFactory
|
||||
idle []*glua.LState
|
||||
@@ -26,7 +29,11 @@ type Pool struct {
|
||||
}
|
||||
|
||||
// NewPool tests the factory by creating one state during initialization.
|
||||
func NewPool(ctx context.Context, factory LStateFactory) (*Pool, error) {
|
||||
func NewPool(ctx context.Context, timeout time.Duration, factory LStateFactory) (*Pool, error) {
|
||||
if timeout <= 0 {
|
||||
return nil, errors.New("Lua pool timeout must be positive")
|
||||
}
|
||||
|
||||
poolCtx, cancel := context.WithCancel(ctx)
|
||||
|
||||
state, err := factory(poolCtx)
|
||||
@@ -35,20 +42,26 @@ func NewPool(ctx context.Context, factory LStateFactory) (*Pool, error) {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Pool{ctx: poolCtx, cancel: cancel, factory: factory, idle: []*glua.LState{state}, top: state.GetTop()}, nil
|
||||
}
|
||||
|
||||
// Context is cancelled by Close. Query contexts should derive from it.
|
||||
func (p *Pool) Context() context.Context {
|
||||
return p.ctx
|
||||
return &Pool{ctx: poolCtx, cancel: cancel, timeout: timeout, factory: factory, idle: []*glua.LState{state}, top: state.GetTop()}, nil
|
||||
}
|
||||
|
||||
// Acquire returns an initialized exclusive state, growing the pool if necessary.
|
||||
func (p *Pool) Acquire() (*glua.LState, error) {
|
||||
// ctx is passed to the factory for state creation; nil uses the pool context.
|
||||
func (p *Pool) Acquire(ctx context.Context) (*glua.LState, error) {
|
||||
p.mu.Lock()
|
||||
if p.closed || p.ctx.Err() != nil {
|
||||
if p.closed {
|
||||
p.mu.Unlock()
|
||||
return nil, p.ctx.Err()
|
||||
return nil, errors.New("Lua pool is closed")
|
||||
}
|
||||
if err := p.ctx.Err(); err != nil {
|
||||
p.mu.Unlock()
|
||||
return nil, err
|
||||
}
|
||||
if ctx == nil {
|
||||
ctx = p.ctx
|
||||
} else if err := ctx.Err(); err != nil {
|
||||
p.mu.Unlock()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
p.active.Add(1)
|
||||
@@ -65,7 +78,7 @@ func (p *Pool) Acquire() (*glua.LState, error) {
|
||||
// TODO: Limit the total number of states. When the limit is reached, wait
|
||||
// for a Release instead of creating another state; allow the wait to be
|
||||
// cancelled by the caller or by Close.
|
||||
state, err := p.factory(p.ctx)
|
||||
state, err := p.factory(ctx)
|
||||
if err != nil {
|
||||
p.active.Done()
|
||||
return nil, err
|
||||
@@ -74,15 +87,24 @@ func (p *Pool) Acquire() (*glua.LState, error) {
|
||||
return state, nil
|
||||
}
|
||||
|
||||
// WithState runs work on an exclusive state and releases it afterward. A state
|
||||
// is reusable only when work succeeds; a panic closes it before propagating.
|
||||
func (p *Pool) WithState(work func(*glua.LState) error) error {
|
||||
state, err := p.Acquire()
|
||||
// WithState runs work on an exclusive state and releases it afterward.
|
||||
// Nil ctx and zero timeout use pool defaults. The timeout starts after acquisition.
|
||||
func (p *Pool) WithState(ctx context.Context, timeout time.Duration, work func(*glua.LState) error) error {
|
||||
state, err := p.Acquire(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if ctx == nil {
|
||||
ctx = p.ctx
|
||||
}
|
||||
if timeout == 0 {
|
||||
timeout = p.timeout
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(ctx, timeout)
|
||||
state.SetContext(ctx)
|
||||
reusable := false
|
||||
defer func() {
|
||||
cancel()
|
||||
p.Release(state, reusable)
|
||||
}()
|
||||
err = work(state)
|
||||
@@ -90,9 +112,10 @@ func (p *Pool) WithState(work func(*glua.LState) error) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// Release returns a healthy state to the pool and closes a failed or cancelled one.
|
||||
// Release resets a state for reuse or closes it.
|
||||
func (p *Pool) Release(state *glua.LState, reusable bool) {
|
||||
if reusable {
|
||||
state.RemoveContext()
|
||||
state.SetTop(p.top)
|
||||
p.mu.Lock()
|
||||
if !p.closed && p.ctx.Err() == nil && len(p.idle) < maxIdleStates {
|
||||
@@ -110,7 +133,7 @@ func (p *Pool) Release(state *glua.LState, reusable bool) {
|
||||
p.active.Done()
|
||||
}
|
||||
|
||||
// Close cancels active work, closes idle states, and waits for borrowed states.
|
||||
// Close cancels the pool context, closes idle states, and waits for borrowed states.
|
||||
func (p *Pool) Close() {
|
||||
p.mu.Lock()
|
||||
if !p.closed {
|
||||
|
||||
Reference in New Issue
Block a user