mirror of
https://github.com/MHSanaei/3x-ui.git
synced 2026-07-26 09:52:14 +03:00
Frontend deps: @hookform/resolvers 5.4.0 -> 5.4.3 and react-hook-form 7.82.0 -> 7.83.0. The @typeschema/valibot override is what makes this installable at all. Resolvers 5.4.3 re-declares 25 optional peers for its validator matrix, and npm resolves them into the ideal tree even though none are used here; two of them contradict, since resolvers wants valibot ^1 while @typeschema/main -> @typeschema/valibot pins valibot ^0.39. Both target the same node_modules/valibot, so a plain npm update dies with ERESOLVE. The override settles that one edge and nothing extra lands in node_modules. Backend deps: telego 1.10.0 -> 1.11.1 (Telegram Bot API v10.2, additive only), klauspost/compress 1.19.1, plus the indirect bumps that came with them. VS Code tasks: the golangci-lint and modernize tasks assumed Windows PATH semantics, where PATH is a persistent user variable that every process inherits, so ~/go/bin was always visible. On Linux that directory is exported from ~/.bashrc, which the non-interactive `bash -c` behind a task never sources, and both tasks failed with exit 127. Adds linux/osx option blocks that prepend the Go bin directories and leaves the Windows path untouched, plus tasks to install the two tools; those are split because go install rejects packages from different modules in one invocation. Go sources: modernize -fix output, covering range-over-int, slices.Backward, maps.Copy, strings.CutPrefix and strings.SplitSeq. Behaviour is unchanged.
277 lines
5.0 KiB
Go
277 lines
5.0 KiB
Go
package eventbus
|
|
|
|
import (
|
|
"sync"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/op/go-logging"
|
|
|
|
"github.com/mhsanaei/3x-ui/v3/internal/logger"
|
|
)
|
|
|
|
func TestMain(m *testing.M) {
|
|
logger.InitLogger(logging.ERROR)
|
|
m.Run()
|
|
}
|
|
|
|
func TestBusPublishSubscribe(t *testing.T) {
|
|
b := New(16)
|
|
defer b.Stop()
|
|
|
|
var received Event
|
|
var wg sync.WaitGroup
|
|
wg.Add(1)
|
|
|
|
b.Subscribe("test", func(e Event) {
|
|
received = e
|
|
wg.Done()
|
|
})
|
|
|
|
b.Publish(Event{Type: EventOutboundDown, Source: "my-proxy"})
|
|
|
|
select {
|
|
case <-waitDone(&wg):
|
|
case <-time.After(time.Second):
|
|
t.Fatal("subscriber did not receive event")
|
|
}
|
|
|
|
if received.Type != EventOutboundDown {
|
|
t.Errorf("got type %q, want %q", received.Type, EventOutboundDown)
|
|
}
|
|
if received.Source != "my-proxy" {
|
|
t.Errorf("got source %q, want %q", received.Source, "my-proxy")
|
|
}
|
|
if received.Timestamp.IsZero() {
|
|
t.Error("timestamp not set")
|
|
}
|
|
}
|
|
|
|
func TestBusMultipleSubscribers(t *testing.T) {
|
|
b := New(16)
|
|
defer b.Stop()
|
|
|
|
var count atomic.Int32
|
|
var wg sync.WaitGroup
|
|
wg.Add(2)
|
|
|
|
b.Subscribe("a", func(e Event) {
|
|
count.Add(1)
|
|
wg.Done()
|
|
})
|
|
b.Subscribe("b", func(e Event) {
|
|
count.Add(1)
|
|
wg.Done()
|
|
})
|
|
|
|
b.Publish(Event{Type: EventXrayCrash})
|
|
|
|
select {
|
|
case <-waitDone(&wg):
|
|
case <-time.After(time.Second):
|
|
t.Fatal("subscribers did not receive event")
|
|
}
|
|
|
|
if count.Load() != 2 {
|
|
t.Errorf("got %d calls, want 2", count.Load())
|
|
}
|
|
}
|
|
|
|
func TestBusUnsubscribe(t *testing.T) {
|
|
b := New(16)
|
|
defer b.Stop()
|
|
|
|
var count atomic.Int32
|
|
|
|
b.Subscribe("test", func(e Event) {
|
|
count.Add(1)
|
|
})
|
|
b.Unsubscribe("test")
|
|
|
|
b.Publish(Event{Type: EventOutboundUp})
|
|
time.Sleep(50 * time.Millisecond)
|
|
|
|
if count.Load() != 0 {
|
|
t.Errorf("got %d calls after unsubscribe, want 0", count.Load())
|
|
}
|
|
}
|
|
|
|
func TestBusReplaceSubscriber(t *testing.T) {
|
|
b := New(16)
|
|
defer b.Stop()
|
|
|
|
var last string
|
|
var wg sync.WaitGroup
|
|
wg.Add(1)
|
|
|
|
b.Subscribe("test", func(e Event) {
|
|
last = "old"
|
|
})
|
|
b.Subscribe("test", func(e Event) {
|
|
last = "new"
|
|
wg.Done()
|
|
})
|
|
|
|
b.Publish(Event{Type: EventOutboundDown})
|
|
|
|
select {
|
|
case <-waitDone(&wg):
|
|
case <-time.After(time.Second):
|
|
t.Fatal("subscriber did not receive event")
|
|
}
|
|
|
|
if last != "new" {
|
|
t.Errorf("got %q, want %q", last, "new")
|
|
}
|
|
}
|
|
|
|
func TestBusPanicRecovery(t *testing.T) {
|
|
b := New(16)
|
|
defer b.Stop()
|
|
|
|
var wg sync.WaitGroup
|
|
wg.Add(1)
|
|
|
|
b.Subscribe("panicker", func(e Event) {
|
|
panic("oops")
|
|
})
|
|
b.Subscribe("after", func(e Event) {
|
|
wg.Done()
|
|
})
|
|
|
|
b.Publish(Event{Type: EventOutboundDown})
|
|
|
|
select {
|
|
case <-waitDone(&wg):
|
|
case <-time.After(time.Second):
|
|
t.Fatal("subscriber after panicker did not receive event")
|
|
}
|
|
}
|
|
|
|
func TestBusBlockingSubscriberDoesNotStallOthers(t *testing.T) {
|
|
b := New(16)
|
|
defer b.Stop()
|
|
|
|
release := make(chan struct{})
|
|
b.Subscribe("blocking", func(e Event) {
|
|
<-release
|
|
})
|
|
|
|
fast := make(chan struct{}, 1)
|
|
b.Subscribe("fast", func(e Event) {
|
|
fast <- struct{}{}
|
|
})
|
|
|
|
b.Publish(Event{Type: EventXrayCrash})
|
|
|
|
select {
|
|
case <-fast:
|
|
case <-time.After(time.Second):
|
|
close(release)
|
|
t.Fatal("a blocking subscriber stalled event delivery to another subscriber")
|
|
}
|
|
close(release)
|
|
}
|
|
|
|
func TestBusSubscriberRunsSerially(t *testing.T) {
|
|
b := New(16)
|
|
defer b.Stop()
|
|
|
|
var inFlight atomic.Int32
|
|
var maxSeen atomic.Int32
|
|
var wg sync.WaitGroup
|
|
const n = 8
|
|
wg.Add(n)
|
|
|
|
b.Subscribe("serial", func(Event) {
|
|
cur := inFlight.Add(1)
|
|
for {
|
|
m := maxSeen.Load()
|
|
if cur <= m || maxSeen.CompareAndSwap(m, cur) {
|
|
break
|
|
}
|
|
}
|
|
time.Sleep(5 * time.Millisecond)
|
|
inFlight.Add(-1)
|
|
wg.Done()
|
|
})
|
|
|
|
for range n {
|
|
b.Publish(Event{Type: EventXrayCrash})
|
|
}
|
|
|
|
select {
|
|
case <-waitDone(&wg):
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("subscriber did not process all events")
|
|
}
|
|
if got := maxSeen.Load(); got != 1 {
|
|
t.Fatalf("subscriber ran concurrently with itself: max in-flight = %d, want 1", got)
|
|
}
|
|
}
|
|
|
|
func TestBusBufferFull(t *testing.T) {
|
|
b := New(2)
|
|
defer b.Stop()
|
|
|
|
b.Subscribe("slow", func(e Event) {
|
|
time.Sleep(100 * time.Millisecond)
|
|
})
|
|
|
|
b.Publish(Event{Type: EventOutboundDown})
|
|
b.Publish(Event{Type: EventOutboundUp})
|
|
b.Publish(Event{Type: EventXrayCrash})
|
|
|
|
time.Sleep(50 * time.Millisecond)
|
|
}
|
|
|
|
func TestBusZeroTimestamp(t *testing.T) {
|
|
b := New(16)
|
|
defer b.Stop()
|
|
|
|
var received Event
|
|
var wg sync.WaitGroup
|
|
wg.Add(1)
|
|
|
|
b.Subscribe("test", func(e Event) {
|
|
received = e
|
|
wg.Done()
|
|
})
|
|
|
|
b.Publish(Event{Type: EventOutboundDown})
|
|
|
|
select {
|
|
case <-waitDone(&wg):
|
|
case <-time.After(time.Second):
|
|
t.Fatal("subscriber did not receive event")
|
|
}
|
|
|
|
if received.Timestamp.IsZero() {
|
|
t.Error("timestamp should be set automatically")
|
|
}
|
|
}
|
|
|
|
func waitDone(wg *sync.WaitGroup) <-chan struct{} {
|
|
ch := make(chan struct{})
|
|
go func() {
|
|
wg.Wait()
|
|
close(ch)
|
|
}()
|
|
return ch
|
|
}
|
|
|
|
func TestBusSubscribeAfterStopIsNoop(t *testing.T) {
|
|
b := New(4)
|
|
b.Stop()
|
|
|
|
b.Subscribe("late", func(Event) {})
|
|
|
|
b.mu.RLock()
|
|
n := len(b.subs)
|
|
b.mu.RUnlock()
|
|
if n != 0 {
|
|
t.Fatalf("Subscribe after Stop registered %d subscriber(s), want 0 (a stopped bus must not accept new subscribers, and must not call wg.Add after wg.Wait has been entered)", n)
|
|
}
|
|
}
|