项目文件夹

文件
wehub-resource-sync bf9395e022
CI / license-header (push) Has been skipped
CI / e2e-dry-run (push) Has been skipped
CI / fast-gate (push) Failing after 0s
Test PR Label Logic / test-pr-labels (push) Failing after 1s
Skill Format Check / check-format (push) Failing after 2s
CI / security (push) Failing after 5s
CI / unit-test (push) Has been skipped
CI / lint (push) Has been skipped
CI / script-test (push) Has been skipped
CI / deterministic-gate (push) Has been skipped
CI / coverage (push) Has been skipped
CI / results (push) Has been cancelled
CI / deadcode (push) Has been cancelled
CI / e2e-live (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 12:22:54 +08:00

91 行
1.9 KiB
Go

此文件含有模棱两可的 Unicode 字符
此文件含有可能会与其他字符混淆的 Unicode 字符。 如果您是想特意这样的,可以安全地忽略该警告。 使用 Escape 按钮显示他们。
// Copyright (c) 2026 Lark Technologies Pte. Ltd.
// SPDX-License-Identifier: MIT
package bus
import (
"io"
"log"
"net"
"testing"
"time"
)
// Reproduces Run × onClose re-entrant deadlock if b.mu is held across Close.
func TestRunShutdownWithMultipleConns(t *testing.T) {
logger := log.New(io.Discard, "", 0)
hub := NewHub()
b := &Bus{
hub: hub,
logger: logger,
conns: make(map[*Conn]struct{}),
}
const N = 3
pipes := make([]net.Conn, 0, N*2)
t.Cleanup(func() {
for _, p := range pipes {
p.Close()
}
})
for i := 0; i < N; i++ {
server, client := net.Pipe()
pipes = append(pipes, server, client)
bc := NewConn(server, nil, "im.msg", []string{"im.message.receive_v1"}, 1000+i, "")
bc.SetLogger(logger)
hub.RegisterAndIsFirst(bc)
bc.SetOnClose(func(c *Conn) {
b.hub.UnregisterAndIsLast(c)
b.mu.Lock()
delete(b.conns, c)
b.mu.Unlock()
})
b.mu.Lock()
b.conns[bc] = struct{}{}
b.mu.Unlock()
}
done := make(chan struct{})
go func() {
shutdownConns(b)
close(done)
}()
select {
case <-done:
case <-time.After(2 * time.Second):
t.Fatal("shutdownConns deadlocked: did not complete within 2s")
}
if got := hub.ConnCount(); got != 0 {
t.Errorf("expected 0 subscribers in hub after shutdown, got %d", got)
}
b.mu.Lock()
remaining := len(b.conns)
b.mu.Unlock()
if remaining != 0 {
t.Errorf("expected 0 conns in Bus after shutdown, got %d", remaining)
}
}
// shutdownCh must be buffered so a signal sent before Run's select loop is still delivered.
func TestShutdownSignalNotDroppedBeforeRunSelects(t *testing.T) {
b := NewBus("test-app", "test-secret", "", nil, log.New(io.Discard, "", 0))
select {
case b.shutdownCh <- struct{}{}:
default:
t.Fatal("handleShutdown's send took default branch — signal would be lost")
}
select {
case <-b.shutdownCh:
case <-time.After(200 * time.Millisecond):
t.Fatal("shutdown signal was not latched")
}
}