@@ -14,6 +14,23 @@ import (
1414// disableES disables the built-in eventstream service so tests can bind port 1002 via driver.
1515func disableES (cfg * daemon.Config ) { cfg .DisableEventStream = true }
1616
17+ // startEventServer launches srv.ListenAndServe in a goroutine and waits until
18+ // the port is accepting connections before returning, so callers don't race
19+ // the goroutine scheduler on slow CI runners.
20+ func startEventServer (t * testing.T , srv * eventstream.Server , info * DaemonInfo ) {
21+ t .Helper ()
22+ go srv .ListenAndServe ()
23+ for i := 0 ; i < 40 ; i ++ {
24+ c , err := eventstream .Subscribe (info .Driver , info .Daemon .Addr (), "_probe" )
25+ if err == nil {
26+ c .Close ()
27+ return
28+ }
29+ time .Sleep (25 * time .Millisecond )
30+ }
31+ t .Fatal ("event stream server did not become ready" )
32+ }
33+
1734func TestEventStream (t * testing.T ) {
1835 t .Parallel ()
1936 env := NewTestEnv (t )
@@ -29,7 +46,7 @@ func TestEventStream(t *testing.T) {
2946
3047 // Start broker on A
3148 srv := eventstream .NewServer (a .Driver )
32- go srv . ListenAndServe ( )
49+ startEventServer ( t , srv , a )
3350
3451 // Subscriber on B
3552 sub , err := eventstream .Subscribe (b .Driver , a .Daemon .Addr (), "test-topic" )
@@ -126,7 +143,7 @@ func TestEventStreamWildcard(t *testing.T) {
126143 c := env .AddDaemon (disableES )
127144
128145 srv := eventstream .NewServer (a .Driver )
129- go srv . ListenAndServe ( )
146+ startEventServer ( t , srv , a )
130147
131148 // B subscribes with wildcard
132149 sub , err := eventstream .Subscribe (b .Driver , a .Daemon .Addr (), "*" )
@@ -170,7 +187,7 @@ func TestEventStreamMultipleTopics(t *testing.T) {
170187 c := env .AddDaemon (disableES )
171188
172189 srv := eventstream .NewServer (a .Driver )
173- go srv . ListenAndServe ( )
190+ startEventServer ( t , srv , a )
174191
175192 // B subscribes to "topic-A"
176193 subA , err := eventstream .Subscribe (b .Driver , a .Daemon .Addr (), "topic-A" )
@@ -239,7 +256,7 @@ func TestEventStreamMultipleSubscribers(t *testing.T) {
239256 c := env .AddDaemon (disableES )
240257
241258 srv := eventstream .NewServer (a .Driver )
242- go srv . ListenAndServe ( )
259+ startEventServer ( t , srv , a )
243260
244261 // Both B and C subscribe to same topic
245262 sub1 , err := eventstream .Subscribe (b .Driver , a .Daemon .Addr (), "shared-topic" )
@@ -317,7 +334,7 @@ func TestEventStreamPublisherExclusion(t *testing.T) {
317334 b := env .AddDaemon (disableES )
318335
319336 srv := eventstream .NewServer (a .Driver )
320- go srv . ListenAndServe ( )
337+ startEventServer ( t , srv , a )
321338
322339 // B subscribes and publishes on the same topic
323340 client , err := eventstream .Subscribe (b .Driver , a .Daemon .Addr (), "self-topic" )
@@ -361,7 +378,7 @@ func TestEventStreamSequentialMessages(t *testing.T) {
361378 c := env .AddDaemon (disableES )
362379
363380 srv := eventstream .NewServer (a .Driver )
364- go srv . ListenAndServe ( )
381+ startEventServer ( t , srv , a )
365382
366383 sub , err := eventstream .Subscribe (b .Driver , a .Daemon .Addr (), "seq-topic" )
367384 if err != nil {
@@ -419,7 +436,7 @@ func TestEventStreamSubscriberDisconnect(t *testing.T) {
419436 c := env .AddDaemon (disableES )
420437
421438 srv := eventstream .NewServer (a .Driver )
422- go srv . ListenAndServe ( )
439+ startEventServer ( t , srv , a )
423440
424441 // B subscribes and then disconnects
425442 sub , err := eventstream .Subscribe (b .Driver , a .Daemon .Addr (), "disc-topic" )
0 commit comments