node-tools: say whose words a refused event is (hq issue 276)

"plex did not take radarr.download.completed: Unexpected end of JSON input"
read as the runtime failing to parse the bundle's answer. It was plex's own
handler error, relayed. A bundle's error answer is now a launch.Refused, and
the line says the handler answered an error, quotes it, names the event id
and the delay before the next offer; the runtime's own failures (no answer
in time, bundle exited) are said as before.

The rule, unchanged and now tested: any mesh/event answer that is not an
error takes the event, whatever its result says, including none.
This commit is contained in:
jochen
2026-10-06 18:14:50 +02:00
parent f75780f1d0
commit 1792b0d9ba
4 changed files with 142 additions and 2 deletions
@@ -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)
}
}
+13 -1
View File
@@ -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: