Wait for every runtime discovery hears from, and name the ones it misses
The console gathered $SRV.INFO answers for a fixed 750 ms. The laptop's
runtime (48 modules, 341 endpoints, 164 kB) answers last every time: its
answer crosses to the broker on another machine and back, a median of
365 ms on a quiet link and 813 ms in one of 25 rounds measured. When it
missed the window the console said its modules ran nowhere ("nothing in
the mesh is called slack") or only on another machine.
Discovery now asks $SRV.PING alongside $SRV.INFO and waits, past the
window and up to 5 s, for every instance that answered PING. A runtime
that never sends what it serves, or answered before and not now, is
named in every answer that might concern it instead of the module being
called missing. An answer still too large after first-line descriptions
drops them, and says so in its metadata.
mesh_runtimes reports, per runtime, its machine, how long its answer
took, its size, its modules and tools, whether it was shortened and when
it was last heard, and the runtimes and machines not heard.
This commit is contained in:
@@ -649,37 +649,51 @@ func (c *Conn) Raw(subject string, answer func(subject string, data []byte) []by
|
||||
return func() { _ = sub.Unsubscribe() }, nil
|
||||
}
|
||||
|
||||
// Gather publishes one request and collects every answer that arrives within the window: a
|
||||
// discovery request every service instance answers (ADR 0197). It never stops early — how many will
|
||||
// answer is what it is finding out.
|
||||
func (c *Conn) Gather(subject string, body []byte, window time.Duration) ([][]byte, error) {
|
||||
// Arrival is one answer to a request, and when it arrived.
|
||||
type Arrival struct {
|
||||
Data []byte
|
||||
At time.Time
|
||||
}
|
||||
|
||||
// Collect publishes one request and hands every answer to the channel as it arrives, until stopped:
|
||||
// for a caller that decides for itself when it has heard enough — discovery waits for every service
|
||||
// that said it is there, not for a fixed moment (2026-10-05). Answers beyond what the caller has
|
||||
// taken are kept up to the channel's room; an empty answer is no answer.
|
||||
func (c *Conn) Collect(subject string, body []byte) (<-chan Arrival, func(), error) {
|
||||
inbox := c.nc.NewRespInbox()
|
||||
sub, err := c.nc.SubscribeSync(inbox)
|
||||
out := make(chan Arrival, 1024)
|
||||
var mu sync.Mutex
|
||||
stopped := false
|
||||
sub, err := c.nc.Subscribe(inbox, func(msg *nats.Msg) {
|
||||
if len(msg.Data) == 0 {
|
||||
return
|
||||
}
|
||||
a := Arrival{Data: msg.Data, At: time.Now()}
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if stopped {
|
||||
return
|
||||
}
|
||||
select {
|
||||
case out <- a:
|
||||
default:
|
||||
c.Logf("[mesh-tools] more answers to %s than are kept; one (%d bytes) dropped", subject, len(msg.Data))
|
||||
}
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return nil, func() {}, err
|
||||
}
|
||||
stop := func() {
|
||||
_ = sub.Unsubscribe()
|
||||
mu.Lock()
|
||||
stopped = true
|
||||
mu.Unlock()
|
||||
}
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
if err := c.nc.PublishRequest(subject, inbox, body); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out [][]byte
|
||||
deadline := time.Now().Add(window)
|
||||
for {
|
||||
left := time.Until(deadline)
|
||||
if left <= 0 {
|
||||
return out, nil
|
||||
}
|
||||
msg, err := sub.NextMsg(left)
|
||||
if err != nil {
|
||||
if errors.Is(err, nats.ErrTimeout) {
|
||||
return out, nil
|
||||
}
|
||||
return out, err
|
||||
}
|
||||
if len(msg.Data) > 0 {
|
||||
out = append(out, msg.Data)
|
||||
}
|
||||
stop()
|
||||
return nil, func() {}, err
|
||||
}
|
||||
return out, stop, nil
|
||||
}
|
||||
|
||||
// MaxPayload is the largest message the bus carries, as the server told this connection.
|
||||
|
||||
Reference in New Issue
Block a user