Merge pull request 'node-tools: say whose words a refused event is (hq issue 276)' (#17) from fix/an-event-handler-answers-what-it-did into main
mesh/delivery delivered
mesh/delivery delivered
This commit was merged in pull request #17.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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:
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user