diff --git a/node-tools/internal/launch/event_answer_test.go b/node-tools/internal/launch/event_answer_test.go new file mode 100644 index 0000000..a102b0f --- /dev/null +++ b/node-tools/internal/launch/event_answer_test.go @@ -0,0 +1,89 @@ +package launch + +import ( + "encoding/json" + "errors" + "os" + "path/filepath" + "sync" + "testing" + "time" +) + +// A bundle that answers `mesh/event` as $ANSWER says: an error, a result, or an answer with neither. +const answeringBundle = `#!/bin/sh +while IFS= read -r line; do + id=$(printf '%s' "$line" | sed -n 's/.*"id":\([0-9]*\)[,}].*/\1/p') + case "$line" in + *'"method":"initialize"'*) printf '{"jsonrpc":"2.0","id":%s,"result":{}}\n' "$id" ;; + *'"method":"tools/list"'*) + printf '{"jsonrpc":"2.0","id":%s,"result":{"tools":[]}}\n' "$id" + printf '{"jsonrpc":"2.0","id":"s1","method":"mesh/subscribe","params":{"pattern":"#"}}\n' ;; + *'"method":"mesh/event"'*) + case "$ANSWER" in + error) printf '{"jsonrpc":"2.0","id":%s,"error":{"code":-32000,"message":"Unexpected end of JSON input"}}\n' "$id" ;; + result) printf '{"jsonrpc":"2.0","id":%s,"result":{}}\n' "$id" ;; + bare) printf '{"jsonrpc":"2.0","id":%s}\n' "$id" ;; + esac ;; + esac +done +` + +type subscribing struct { + mu sync.Mutex + deliver func(json.RawMessage) error +} + +func (b *subscribing) Publish(json.RawMessage) error { return nil } +func (b *subscribing) Ask(json.RawMessage) (json.RawMessage, error) { return nil, nil } +func (b *subscribing) State(string, json.RawMessage) (json.RawMessage, error) { return nil, nil } +func (b *subscribing) Watch(json.RawMessage, func(json.RawMessage) error) (func(), error) { + return func() {}, nil +} +func (b *subscribing) Subscribe(d func(json.RawMessage) error) error { + b.mu.Lock() + defer b.mu.Unlock() + b.deliver = d + return nil +} + +func deliverTo(t *testing.T, answer string) error { + t.Helper() + entry := filepath.Join(t.TempDir(), "bundle") + if err := os.WriteFile(entry, []byte(answeringBundle), 0o755); err != nil { + t.Fatal(err) + } + b := &subscribing{} + l, err := Start("plex", entry, append(os.Environ(), "ANSWER="+answer), b, t.Logf) + if err != nil { + t.Fatal(err) + } + t.Cleanup(l.Stop) + for i := 0; ; i++ { + b.mu.Lock() + d := b.deliver + b.mu.Unlock() + if d != nil { + return d(json.RawMessage(`{"key":"radarr.download.completed","body":{"title":"x"}}`)) + } + if i > 100 { + t.Fatal("the bundle never subscribed") + } + time.Sleep(20 * time.Millisecond) + } +} + +// An event is taken by any answer that is not an error, and not taken by an error, which reaches the +// runtime as the bundle's own words — not as something the runtime failed to read (hq issue 276). +func TestAnEventIsTakenByAnyAnswerThatIsNotAnError(t *testing.T) { + for _, answer := range []string{"result", "bare"} { + if err := deliverTo(t, answer); err != nil { + t.Errorf("an answer %q did not take the event: %v", answer, err) + } + } + err := deliverTo(t, "error") + var refused *Refused + if !errors.As(err, &refused) || refused.Method != "mesh/event" || refused.Message != "Unexpected end of JSON input" { + t.Fatalf("an error answer was not the bundle's refusal of mesh/event: %#v", err) + } +} diff --git a/node-tools/internal/launch/launch.go b/node-tools/internal/launch/launch.go index 9b626ab..a691111 100644 --- a/node-tools/internal/launch/launch.go +++ b/node-tools/internal/launch/launch.go @@ -65,6 +65,17 @@ type Bus interface { // EventTimeout bounds how long a bundle has to handle one event before it is offered again. var EventTimeout = 2 * time.Minute +// Refused is a bundle's error answer to what the runtime asked it: the bundle's own words, as opposed +// to the runtime not reaching it (not running, exited, no answer in time). For `mesh/event` it is the +// one way a bundle says it did not take an event: any answer that is not an error — `{}`, `null`, no +// result at all — takes it (mesh-sdk, "Taking an event"; novox/hq issue 276). +type Refused struct { + Method string + Message string +} + +func (r *Refused) Error() string { return r.Message } + // Restart backoff: a bundle that exits is started again at once, then after growing pauses while it // keeps exiting, back to at once once it has run a while. var ( @@ -144,6 +155,7 @@ func Start(module, entry string, env []string, mesh Bus, logf func(string, ...an if c == nil { return errors.New(module + "'s bundle is not running to take its events") } + // Taken unless the answer is an error; what a result says is not read (see Refused). _, err := c.ask(module, "mesh/event", map[string]any{"envelope": envelope}, EventTimeout) return err } @@ -528,7 +540,7 @@ func (c *child) ask(module, method string, params any, timeout time.Duration) (j if msg == "" { msg = "the bundle refused the request" } - return nil, errors.New(msg) + return nil, &Refused{Method: method, Message: msg} } return m.Result, nil case <-c.dead: diff --git a/node-tools/internal/runtime/runtime.go b/node-tools/internal/runtime/runtime.go index 5293f82..d778471 100644 --- a/node-tools/internal/runtime/runtime.go +++ b/node-tools/internal/runtime/runtime.go @@ -7,6 +7,7 @@ package runtime import ( "encoding/json" + "errors" "fmt" "os" "path/filepath" @@ -565,7 +566,8 @@ func (b *moduleBus) Subscribe(deliver func(json.RawMessage) error) error { c.mu.Unlock() for _, d := range targets { if err := (*d)(raw); err != nil { - c.logf("[mesh-tools] %s did not take %s: %v; offered again", module, env.Key, err) + c.logf("[mesh-tools] %s did not take %s%s: %s; offered again in %s", module, env.Key, + eventRef(env), whyNotTaken(err), bus.NakDelay) return err } } @@ -579,6 +581,26 @@ func (b *moduleBus) Subscribe(deliver func(json.RawMessage) error) error { return nil } +// whyNotTaken says why a bundle did not take an event, so the line names who failed: the module's +// handler, in its own words — an error it answered — or the runtime not reaching it. A bundle's +// handler error read bare ("Unexpected end of JSON input") looked like the runtime's (issue 276). +func whyNotTaken(err error) string { + var refused *launch.Refused + if errors.As(err, &refused) { + return fmt.Sprintf("its handler answered an error, in its own words: %q (an event is taken by any answer that is not an error)", refused.Message) + } + return err.Error() +} + +// eventRef names the event a line is about by its id, when it has one, so its offers can be told +// apart from another event's. +func eventRef(env bus.Envelope) string { + if id := env.Headers["x-event-id"]; id != "" { + return " (event " + id + ")" + } + return "" +} + // stateAsked is what a bundle names when it reaches its state (ADR 0201): the state by the name its // module uses, a key, and for a put the value. type stateAsked struct { diff --git a/node-tools/internal/runtime/runtime_test.go b/node-tools/internal/runtime/runtime_test.go index 0cb991b..1c1040f 100644 --- a/node-tools/internal/runtime/runtime_test.go +++ b/node-tools/internal/runtime/runtime_test.go @@ -2,12 +2,14 @@ package runtime import ( "encoding/json" + "errors" "os" "strings" "testing" "time" "github.com/novox/mesh-tools/node-tools/internal/bus" + "github.com/novox/mesh-tools/node-tools/internal/launch" mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" ) @@ -249,3 +251,18 @@ func TestToolModulesNamesOtherModulesOnly(t *testing.T) { } } } + +// A line about an event a bundle did not take says who failed: the handler in its own words, or the +// runtime not reaching it (hq issue 276). +func TestAnUntakenEventSaysWhoFailed(t *testing.T) { + handler := whyNotTaken(&launch.Refused{Method: "mesh/event", Message: "Unexpected end of JSON input"}) + if !strings.Contains(handler, `its handler answered an error, in its own words: "Unexpected end of JSON input"`) { + t.Errorf("a handler's error is not said as the handler's: %s", handler) + } + if got := whyNotTaken(errors.New("plex's bundle did not answer mesh/event in 120s")); got != "plex's bundle did not answer mesh/event in 120s" { + t.Errorf("the runtime's own failure is not said as it is: %s", got) + } + if got := eventRef(bus.Envelope{Headers: map[string]string{"x-event-id": "e1"}}); got != " (event e1)" { + t.Errorf("the event is not named by its id: %q", got) + } +}