diff --git a/internal/link/bus.go b/internal/link/bus.go index 8b767e3..d35b864 100644 --- a/internal/link/bus.go +++ b/internal/link/bus.go @@ -2,6 +2,7 @@ package link import ( "context" + "errors" "fmt" "time" @@ -29,6 +30,11 @@ type Bus interface { // declaration is not an event, and replaying yesterday's is actively harmful // (design 29 §4, the *state* shape). 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 ----------------------------------------------------- @@ -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 { // 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 @@ -71,7 +87,10 @@ func (b OverCurrent) PublishDeclaration(ctx context.Context, node string, body [ // --- NATS, the bus being built ---------------------------------------------------------------- // 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 // 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 } + +// 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 +} diff --git a/internal/link/bus_test.go b/internal/link/bus_test.go index afd5f94..909c4db 100644 --- a/internal/link/bus_test.go +++ b/internal/link/bus_test.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "os" + "strings" "testing" "time" @@ -93,3 +94,61 @@ func TestAnEventsSubjectIsDerivedFromItsSource(t *testing.T) { 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) + } +}