Every runtime announces what it serves in the NATS services protocol; the console discovers by asking the bus (hq ADR 0197)
The Go runtime answers $SRV.PING, $SRV.INFO and $SRV.STATS (and per name and id) in the io.nats.micro.v1 format with what it serves at the moment it is asked: one service per runtime process, since the bus admits one reply per request from each responder, and one endpoint per tool per subject, its metadata saying module, seat, scope, machine, description, schema and whether the module is interchangeable. Serving is unchanged. The console gathers one $SRV.INFO request's answers instead of asking the catalogue's roster and each module's tools, and reads the controller's records as JSON for what should have answered: an assignment with tools that did not announce is named, a module without tools never is. The text parsers of node list and module list are gone. Packages share the test bus: go test -p 1.
This commit is contained in:
@@ -512,6 +512,59 @@ func (c *Conn) PublishAs(module string, env Envelope) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// ServedOn is where a served module's tool is answered right now: the membership's subjects when
|
||||
// issued, the derived shape otherwise — what Handle subscribes, for what announces it (ADR 0197).
|
||||
func (c *Conn) ServedOn(module, tool string) []Served { return c.servedOn(module, tool) }
|
||||
|
||||
// Raw answers one subject with a function of the request, not a tool's reply envelope: the NATS
|
||||
// services protocol's discovery subjects answer in their own format (novox/hq ADR 0197). A nil
|
||||
// answer is no reply — the request was for another service.
|
||||
func (c *Conn) Raw(subject string, answer func(subject string, data []byte) []byte) (func(), error) {
|
||||
sub, err := c.nc.Subscribe(subject, func(msg *nats.Msg) {
|
||||
if body := answer(msg.Subject, msg.Data); body != nil {
|
||||
_ = msg.Respond(body)
|
||||
}
|
||||
})
|
||||
if err != nil {
|
||||
return func() {}, err
|
||||
}
|
||||
c.track(sub)
|
||||
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) {
|
||||
inbox := c.nc.NewRespInbox()
|
||||
sub, err := c.nc.SubscribeSync(inbox)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Flush waits until the bus has every subscription made so far, so what is served is answerable
|
||||
// when this returns.
|
||||
func (c *Conn) Flush() { _ = c.nc.Flush() }
|
||||
|
||||
Reference in New Issue
Block a user