Asking a tool goes through the seam, and loses two problems
Step 3.4. On the new bus there is no reply queue to declare and no correlation to check: each account is granted one inbox prefix and no other, so an answer cannot reach the wrong asker. That settles a cost build.go records having paid — on a shared reply exchange every asker saw every result, which is why the correlation was checked rather than assumed. And a tool nobody serves says so at once rather than after the whole wait. The difference between "that module is down" and "that tool is slow" is the first thing a person asking wants, and both tests are against a real server because both are claims about what the server does, not about this code. RequestBuild stays as it is, and is a different shape on the new bus rather than the same one: a build takes minutes, so it is work submitted to a queue with the outcome returning to a reply subject the request carries — the pattern design 25 §2 already sets for anything crossing a stream. It touches the builder too, so it goes with that conversion.
This commit is contained in:
+47
-1
@@ -2,6 +2,7 @@ package link
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -29,6 +30,11 @@ type Bus interface {
|
|||||||
// declaration is not an event, and replaying yesterday's is actively harmful
|
// declaration is not an event, and replaying yesterday's is actively harmful
|
||||||
// (design 29 §4, the *state* shape).
|
// (design 29 §4, the *state* shape).
|
||||||
PublishDeclaration(ctx context.Context, node string, body []byte) error
|
PublishDeclaration(ctx context.Context, node string, body []byte) error
|
||||||
|
|
||||||
|
// AskTool sends one question to a module's tool and awaits one answer. A tool nobody serves
|
||||||
|
// must say so **at once** rather than after the whole wait: the difference between "that
|
||||||
|
// module is down" and "that tool is slow" is the first thing a person asking wants.
|
||||||
|
AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- The bus the mesh runs on today -----------------------------------------------------
|
// --- The bus the mesh runs on today -----------------------------------------------------
|
||||||
@@ -57,6 +63,16 @@ func (b OverCurrent) PublishEvent(ctx context.Context, key, source, node string,
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// AskTool is implemented over the existing reply-queue machinery in ask.go; this seam does not
|
||||||
|
// change how it works today.
|
||||||
|
func (b OverCurrent) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
|
||||||
|
answer, err := Ask(ctx, b.Channel, module, tool, args, timeout)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return answer.Result, nil
|
||||||
|
}
|
||||||
|
|
||||||
func (b OverCurrent) PublishDeclaration(ctx context.Context, node string, body []byte) error {
|
func (b OverCurrent) PublishDeclaration(ctx context.Context, node string, body []byte) error {
|
||||||
// To the queue directly rather than through an exchange: a declaration is for one node, and
|
// To the queue directly rather than through an exchange: a declaration is for one node, and
|
||||||
// routing it by name through a shared exchange would mean a binding per node that nothing
|
// routing it by name through a shared exchange would mean a binding per node that nothing
|
||||||
@@ -71,7 +87,10 @@ func (b OverCurrent) PublishDeclaration(ctx context.Context, node string, body [
|
|||||||
// --- NATS, the bus being built ----------------------------------------------------------------
|
// --- NATS, the bus being built ----------------------------------------------------------------
|
||||||
|
|
||||||
// OverNATS is the bus as a JetStream context.
|
// OverNATS is the bus as a JetStream context.
|
||||||
type OverNATS struct{ JS nats.JetStreamContext }
|
type OverNATS struct {
|
||||||
|
Conn *nats.Conn
|
||||||
|
JS nats.JetStreamContext
|
||||||
|
}
|
||||||
|
|
||||||
// EventSubject is where a module's event lands. Derived from the emitter, never taken from the
|
// EventSubject is where a module's event lands. Derived from the emitter, never taken from the
|
||||||
// caller: a source that could differ from the subject is an envelope that can lie about its
|
// caller: a source that could differ from the subject is an envelope that can lie about its
|
||||||
@@ -118,3 +137,30 @@ func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []by
|
|||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ToolSubject is where a module answers. Derived from the module and the tool, so a caller names
|
||||||
|
// what it wants rather than where it lives.
|
||||||
|
func ToolSubject(module, tool string) string { return "mesh.mod." + module + ".tool." + tool }
|
||||||
|
|
||||||
|
func (b OverNATS) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
|
||||||
|
if len(args) == 0 {
|
||||||
|
args = []byte(`{}`)
|
||||||
|
}
|
||||||
|
ask, cancel := context.WithTimeout(ctx, timeout)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
// **No reply queue, and no correlation to check.** The caller's inbox is its own — each
|
||||||
|
// account is granted one prefix and no other (design 25 §4) — so an answer cannot reach the
|
||||||
|
// wrong asker and there is nothing to correlate against. That also settles a cost recorded
|
||||||
|
// in build.go: on a shared reply exchange every asker saw every result.
|
||||||
|
msg, err := b.Conn.RequestWithContext(ask, ToolSubject(module, tool), args)
|
||||||
|
if err != nil {
|
||||||
|
if errors.Is(err, nats.ErrNoResponders) {
|
||||||
|
// Said at once rather than after the whole wait: nothing is subscribed to that
|
||||||
|
// subject, which is a different fact from a tool being slow.
|
||||||
|
return nil, fmt.Errorf("nothing serves %s.%s", module, tool)
|
||||||
|
}
|
||||||
|
return nil, fmt.Errorf("asking %s.%s: %w", module, tool, err)
|
||||||
|
}
|
||||||
|
return msg.Data, nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"os"
|
"os"
|
||||||
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -93,3 +94,61 @@ func TestAnEventsSubjectIsDerivedFromItsSource(t *testing.T) {
|
|||||||
t.Fatal("two modules share an event subject, so neither owns its own name")
|
t.Fatal("two modules share an event subject, so neither owns its own name")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// A tool nobody serves says so at once. The difference between "that module is down" and "that
|
||||||
|
// tool is slow" is the first thing a person asking wants, and waiting out the timeout to say it
|
||||||
|
// is how a fast answer becomes a slow non-answer.
|
||||||
|
func TestAskingAToolNobodyServesFailsAtOnce(t *testing.T) {
|
||||||
|
url := os.Getenv("MESH_TEST_NATS")
|
||||||
|
if url == "" {
|
||||||
|
t.Skip("MESH_TEST_NATS unset")
|
||||||
|
}
|
||||||
|
conn, err := nats.Connect(url)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer conn.Close()
|
||||||
|
|
||||||
|
bus := OverNATS{Conn: conn}
|
||||||
|
began := time.Now()
|
||||||
|
_, err = bus.AskTool(context.Background(), "nobody", "home", nil, 30*time.Second)
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("asking a tool nothing serves succeeded")
|
||||||
|
}
|
||||||
|
if took := time.Since(began); took > 2*time.Second {
|
||||||
|
t.Errorf("took %v to say nothing serves it; the caller waited out the timeout", took)
|
||||||
|
}
|
||||||
|
if !strings.Contains(err.Error(), "nothing serves") {
|
||||||
|
t.Errorf("the refusal does not say nobody is there: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// And a served tool answers, with no reply queue to declare and no correlation to check.
|
||||||
|
func TestAskingAServedToolAnswers(t *testing.T) {
|
||||||
|
url := os.Getenv("MESH_TEST_NATS")
|
||||||
|
if url == "" {
|
||||||
|
t.Skip("MESH_TEST_NATS unset")
|
||||||
|
}
|
||||||
|
conn, err := nats.Connect(url)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer conn.Close()
|
||||||
|
|
||||||
|
sub, err := conn.Subscribe(ToolSubject("shop", "price"), func(m *nats.Msg) {
|
||||||
|
_ = m.Respond([]byte(`{"total":12}`))
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
defer sub.Unsubscribe()
|
||||||
|
|
||||||
|
got, err := OverNATS{Conn: conn}.AskTool(context.Background(), "shop", "price",
|
||||||
|
[]byte(`{"qty":4}`), 5*time.Second)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if string(got) != `{"total":12}` {
|
||||||
|
t.Fatalf("the answer came back as %s", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user