mirror of
https://github.com/MHSanaei/3x-ui.git
synced 2026-10-07 22:51:21 +07:00
Fix: traffic writer restart freeze (#4265)
* feat(traffic_writer): enhance traffic writer with concurrency safety and state management
* Revert "feat(traffic_writer): enhance traffic writer with concurrency safety and state management"
This reverts commit e6760ae396.
* feat(traffic_writer): enhance traffic writer with concurrency safety and state management
* feat(web): implement panel-only start/stop methods for in-process restarts
This commit is contained in:
@@ -23,6 +23,7 @@ type trafficWriteRequest struct {
|
||||
var (
|
||||
twMu sync.Mutex
|
||||
twQueue chan *trafficWriteRequest
|
||||
twCtx context.Context
|
||||
twCancel context.CancelFunc
|
||||
twDone chan struct{}
|
||||
)
|
||||
@@ -37,16 +38,26 @@ var (
|
||||
func StartTrafficWriter() {
|
||||
twMu.Lock()
|
||||
defer twMu.Unlock()
|
||||
if twQueue != nil {
|
||||
return
|
||||
|
||||
if twCancel != nil && twDone != nil {
|
||||
select {
|
||||
case <-twDone:
|
||||
clearTrafficWriterState()
|
||||
default:
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
queue := make(chan *trafficWriteRequest, trafficWriterQueueSize)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
done := make(chan struct{})
|
||||
|
||||
twQueue = queue
|
||||
twCtx = ctx
|
||||
twCancel = cancel
|
||||
twDone = done
|
||||
go runTrafficWriter(queue, ctx, done)
|
||||
|
||||
go runTrafficWriter(ctx, queue, done)
|
||||
}
|
||||
|
||||
// StopTrafficWriter cancels the writer context and waits for the goroutine to
|
||||
@@ -56,20 +67,30 @@ func StopTrafficWriter() {
|
||||
twMu.Lock()
|
||||
cancel := twCancel
|
||||
done := twDone
|
||||
twQueue = nil
|
||||
twCancel = nil
|
||||
twDone = nil
|
||||
if cancel == nil || done == nil {
|
||||
twMu.Unlock()
|
||||
return
|
||||
}
|
||||
cancel()
|
||||
twMu.Unlock()
|
||||
|
||||
if cancel != nil {
|
||||
cancel()
|
||||
}
|
||||
if done != nil {
|
||||
<-done
|
||||
<-done
|
||||
|
||||
twMu.Lock()
|
||||
if twDone == done {
|
||||
clearTrafficWriterState()
|
||||
}
|
||||
twMu.Unlock()
|
||||
}
|
||||
|
||||
func runTrafficWriter(queue chan *trafficWriteRequest, ctx context.Context, done chan struct{}) {
|
||||
func clearTrafficWriterState() {
|
||||
twQueue = nil
|
||||
twCtx = nil
|
||||
twCancel = nil
|
||||
twDone = nil
|
||||
}
|
||||
|
||||
func runTrafficWriter(ctx context.Context, queue chan *trafficWriteRequest, done chan struct{}) {
|
||||
defer close(done)
|
||||
for {
|
||||
select {
|
||||
@@ -99,18 +120,43 @@ func safeApply(fn func() error) (err error) {
|
||||
}
|
||||
|
||||
func submitTrafficWrite(fn func() error) error {
|
||||
req := &trafficWriteRequest{apply: fn, done: make(chan error, 1)}
|
||||
|
||||
twMu.Lock()
|
||||
queue := twQueue
|
||||
twMu.Unlock()
|
||||
|
||||
if queue == nil {
|
||||
ctx := twCtx
|
||||
done := twDone
|
||||
if queue == nil || ctx == nil || done == nil {
|
||||
twMu.Unlock()
|
||||
return safeApply(fn)
|
||||
}
|
||||
req := &trafficWriteRequest{apply: fn, done: make(chan error, 1)}
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
twMu.Unlock()
|
||||
return safeApply(fn)
|
||||
default:
|
||||
}
|
||||
|
||||
timer := time.NewTimer(trafficWriterSubmitTimeout)
|
||||
defer timer.Stop()
|
||||
select {
|
||||
case queue <- req:
|
||||
case <-time.After(trafficWriterSubmitTimeout):
|
||||
twMu.Unlock()
|
||||
case <-timer.C:
|
||||
twMu.Unlock()
|
||||
return errors.New("traffic writer queue full")
|
||||
}
|
||||
return <-req.done
|
||||
|
||||
select {
|
||||
case err := <-req.done:
|
||||
return err
|
||||
case <-done:
|
||||
select {
|
||||
case err := <-req.done:
|
||||
return err
|
||||
default:
|
||||
return errors.New("traffic writer stopped before write completed")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user