The node's runtime is its modules' bus: it binds each module's consumer and hands its events to the bundle (hq ADR 0198)
A launched bundle's mesh/subscribe binds the module's own durable consumer — EVENTS, <node>_<module>, by name as the module's own runtime bound it, so nothing is lost or replayed in the move — and every event goes to each child of the module that subscribed as mesh/event, acknowledged only when all answered, negatively acknowledged after a short delay when one failed or died, terminated when it is not an event. mesh/ask calls a tool as the module. Every launched bundle is started again when it exits, with backoff, since long-running code waits for no call. Requires SDK 0.1.6.
This commit is contained in:
@@ -8,6 +8,7 @@
|
||||
package bus
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
@@ -22,6 +23,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
"github.com/nats-io/nats.go/jetstream"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/wire"
|
||||
)
|
||||
@@ -619,3 +621,115 @@ func SeatToolSubject(seat, verb, scope, node string) string {
|
||||
}
|
||||
return base
|
||||
}
|
||||
|
||||
// AskAs calls a tool on a module's behalf (novox/hq ADR 0198): a bare key is that module's own tool,
|
||||
// `<module>.<tool>` another's, `seat:<seat>.<verb>[@<node>]` a role's — resolved through what the
|
||||
// mesh issued that module to reach, as its own runtime resolved it, else the derived shape.
|
||||
func (c *Conn) AskAs(module, key string, body any) (Answered, error) {
|
||||
name, wanted, _ := strings.Cut(key, "@")
|
||||
subject := ""
|
||||
if m := c.Membership(module); m != nil {
|
||||
if reach := m.Reaches[name]; len(reach) > 0 {
|
||||
subject = reach[0]
|
||||
if wanted != "" {
|
||||
subject = ""
|
||||
for _, s := range reach {
|
||||
if strings.HasSuffix(s, "."+wanted) {
|
||||
subject = s
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if subject == "" {
|
||||
s, err := ToolSubject(key, module)
|
||||
if err != nil {
|
||||
return Answered{}, err
|
||||
}
|
||||
subject = s
|
||||
}
|
||||
return c.Ask(key, body, subject)
|
||||
}
|
||||
|
||||
// EventsStream is where every module's events land, and ConsumerOf the durable consumer the
|
||||
// controller makes for a module on a machine: the same names the module's own runtime bound
|
||||
// (node-tools/src/broker-nats.ts), so moving the module into the node's runtime neither loses an
|
||||
// event nor sees one twice.
|
||||
const EventsStream = "EVENTS"
|
||||
|
||||
// ConsumerOf is a module's durable consumer on a machine: `<node>_<module>`.
|
||||
func ConsumerOf(node, module string) string {
|
||||
if node == "" {
|
||||
node = "?"
|
||||
}
|
||||
return node + "_" + module
|
||||
}
|
||||
|
||||
// NakDelay is how long an event a handler failed waits before it is offered again: a transient cause
|
||||
// gets another attempt, a permanent one exhausts the consumer's max-deliver rather than spinning.
|
||||
var NakDelay = 5 * time.Second
|
||||
|
||||
// ConsumeAs reads a module's durable consumer and hands each event to deliver (novox/hq ADR 0198):
|
||||
// acknowledged when deliver returns nil, negatively acknowledged after NakDelay when it returns an
|
||||
// error, terminated when it is not an event at all. The consumer is the controller's to create; this
|
||||
// binds to it and never makes one. Runs until stopped.
|
||||
func (c *Conn) ConsumeAs(module string, deliver func(Envelope) error) (func(), error) {
|
||||
durable := ConsumerOf(c.node, module)
|
||||
js, err := jetstream.New(c.nc)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
// Bound by name, never created and never matched against a subject: the consumer and its
|
||||
// filters are the controller's (design 29 §3), exactly as the module's own runtime bound it.
|
||||
consumer, err := js.Consumer(ctx, EventsStream, durable)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("binding %s's consumer %s: %w", module, durable, err)
|
||||
}
|
||||
reading, err := consumer.Consume(func(msg jetstream.Msg) {
|
||||
env, ok := envelopeOf(msg.Subject(), msg.Headers(), msg.Data())
|
||||
if !ok {
|
||||
// Unparseable: redelivering bytes no version of a handler can read is a loop.
|
||||
_ = msg.Term()
|
||||
return
|
||||
}
|
||||
if err := deliver(env); err != nil {
|
||||
_ = msg.NakWithDelay(NakDelay)
|
||||
return
|
||||
}
|
||||
_ = msg.Ack()
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("reading %s's consumer %s: %w", module, durable, err)
|
||||
}
|
||||
var once sync.Once
|
||||
return func() { once.Do(reading.Stop) }, nil
|
||||
}
|
||||
|
||||
// envelopeOf rebuilds the envelope a module sees from a delivered event: its key is the emitter and
|
||||
// the event, recovered from the subject; its metadata rides as headers (novox/hq ADR 0042).
|
||||
func envelopeOf(subject string, header nats.Header, data []byte) (Envelope, bool) {
|
||||
if !json.Valid(data) {
|
||||
return Envelope{}, false
|
||||
}
|
||||
headers := map[string]string{}
|
||||
for k := range header {
|
||||
headers[k] = header.Get(k)
|
||||
}
|
||||
return Envelope{Key: keyOf(subject), Node: headers["x-node"], Body: json.RawMessage(data), Headers: headers}, true
|
||||
}
|
||||
|
||||
// keyOf is the key a module sees for an event subject: `mesh.mod.<emitter>.event.<event>` is
|
||||
// `<emitter>.<event>` — the vocabulary its manifest names what it consumes in.
|
||||
func keyOf(subject string) string {
|
||||
before, event, found := strings.Cut(subject, ".event.")
|
||||
if !found {
|
||||
return subject
|
||||
}
|
||||
parts := strings.Split(before, ".")
|
||||
if emitter := parts[len(parts)-1]; emitter != "" {
|
||||
return emitter + "." + event
|
||||
}
|
||||
return event
|
||||
}
|
||||
|
||||
@@ -46,8 +46,26 @@ type Registration struct {
|
||||
Tools []Tool
|
||||
}
|
||||
|
||||
// Publisher publishes an event a bundle asked the runtime to emit, as the bundle's module.
|
||||
type Publisher func(params json.RawMessage) error
|
||||
// Bus is what a launched bundle reaches the mesh through (novox/hq ADR 0193, ADR 0198): the runtime
|
||||
// publishes, asks and subscribes on the module's behalf. Subscribe is asked each time the module's
|
||||
// code subscribes; the runtime binds the module's consumer once and calls deliver for every event,
|
||||
// acknowledging it on the bus when deliver returns nil.
|
||||
type Bus interface {
|
||||
Publish(params json.RawMessage) error
|
||||
Ask(params json.RawMessage) (json.RawMessage, error)
|
||||
Subscribe(deliver func(envelope json.RawMessage) error) error
|
||||
}
|
||||
|
||||
// EventTimeout bounds how long a bundle has to handle one event before it is offered again.
|
||||
var EventTimeout = 2 * time.Minute
|
||||
|
||||
// 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 (
|
||||
RestartFirst = 500 * time.Millisecond
|
||||
RestartMost = 30 * time.Second
|
||||
RestartSettle = time.Minute
|
||||
)
|
||||
|
||||
// Launched is a running bundle: what it registered, and how to stop it.
|
||||
type Launched struct {
|
||||
@@ -103,10 +121,26 @@ var stackLine = regexp.MustCompile(`^\s+at\s`)
|
||||
|
||||
// Start launches one bundle and learns its tools. It fails when the child cannot be started or does
|
||||
// not complete the handshake. A child that exits later is started again on its next call.
|
||||
func Start(module, entry string, env []string, publish Publisher, logf func(string, ...any)) (*Launched, error) {
|
||||
func Start(module, entry string, env []string, mesh Bus, logf func(string, ...any)) (*Launched, error) {
|
||||
var mu sync.Mutex
|
||||
var current *child
|
||||
// The child that last subscribed is the one events are handed to: a restarted child subscribes
|
||||
// again as its code is imported, and from then on the module's events are its.
|
||||
var subscriber *child
|
||||
stopped := false
|
||||
deliver := func(envelope json.RawMessage) error {
|
||||
mu.Lock()
|
||||
c := subscriber
|
||||
mu.Unlock()
|
||||
if c == nil {
|
||||
return errors.New(module + "'s bundle is not running to take its events")
|
||||
}
|
||||
_, err := c.ask(module, "mesh/event", map[string]any{"envelope": envelope}, EventTimeout)
|
||||
return err
|
||||
}
|
||||
var restart func(after time.Duration)
|
||||
var startedAt time.Time
|
||||
pause := RestartFirst
|
||||
|
||||
start := func() (*child, error) {
|
||||
cmd := exec.Command(entry)
|
||||
@@ -163,18 +197,32 @@ func Start(module, entry string, env []string, publish Publisher, logf func(stri
|
||||
logf("[mesh-tools] %s's bundle said something that is not a reply: %s", module, line)
|
||||
continue
|
||||
}
|
||||
// The bundle asks the runtime to emit (ADR 0193): published as this module, answered
|
||||
// once the bus has accepted it. Nothing else a bundle may ask.
|
||||
// The bundle asks the runtime (ADR 0193, ADR 0198): to emit, to call a tool, to hand it the
|
||||
// module's events. Each on the module's behalf, answered when done; nothing else.
|
||||
if m.Method != "" {
|
||||
go func(m message) {
|
||||
if m.Method != "mesh/publish" {
|
||||
var result any = map[string]any{}
|
||||
var err error
|
||||
switch m.Method {
|
||||
case "mesh/publish":
|
||||
err = mesh.Publish(m.Params)
|
||||
case "mesh/ask":
|
||||
var asked json.RawMessage
|
||||
if asked, err = mesh.Ask(m.Params); err == nil {
|
||||
result = asked
|
||||
}
|
||||
case "mesh/subscribe":
|
||||
mu.Lock()
|
||||
subscriber = c
|
||||
mu.Unlock()
|
||||
err = mesh.Subscribe(deliver)
|
||||
default:
|
||||
if len(m.ID) > 0 {
|
||||
_ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID,
|
||||
"error": map[string]any{"code": -32601, "message": "the runtime answers no " + m.Method + " from a bundle"}})
|
||||
}
|
||||
return
|
||||
}
|
||||
err := publish(m.Params)
|
||||
if len(m.ID) == 0 {
|
||||
return
|
||||
}
|
||||
@@ -183,7 +231,7 @@ func Start(module, entry string, env []string, publish Publisher, logf func(stri
|
||||
"error": map[string]any{"code": -32000, "message": err.Error()}})
|
||||
return
|
||||
}
|
||||
_ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, "result": map[string]any{}})
|
||||
_ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, "result": result})
|
||||
}(m)
|
||||
continue
|
||||
}
|
||||
@@ -220,13 +268,32 @@ func Start(module, entry string, env []string, publish Publisher, logf func(stri
|
||||
c.mu.Unlock()
|
||||
close(c.dead)
|
||||
mu.Lock()
|
||||
// Only a child that was serving is brought back; one that died in its own handshake was
|
||||
// never accepted, and whoever started it was told why.
|
||||
wasServing := current == c
|
||||
if current == c {
|
||||
current = nil
|
||||
}
|
||||
if subscriber == c {
|
||||
subscriber = nil
|
||||
}
|
||||
wasStopped := stopped
|
||||
// Started again at once if it ran a while; after a growing pause while it keeps exiting.
|
||||
wait := RestartFirst
|
||||
if time.Since(startedAt) < RestartSettle {
|
||||
wait = pause
|
||||
if pause *= 2; pause > RestartMost {
|
||||
pause = RestartMost
|
||||
}
|
||||
} else {
|
||||
pause = RestartFirst
|
||||
}
|
||||
mu.Unlock()
|
||||
if !wasStopped {
|
||||
logf("[mesh-tools] %s; started again on its next call", why)
|
||||
if !wasStopped && wasServing {
|
||||
// Every bundle stays up (ADR 0198): code that runs long — a handler, a provisioner —
|
||||
// is not waiting for a call to bring it back, and a tool bundle back early costs nothing.
|
||||
logf("[mesh-tools] %s; started again in %s", why, wait.Round(100*time.Millisecond))
|
||||
restart(wait)
|
||||
}
|
||||
}()
|
||||
if _, err := c.ask(module, "initialize", map[string]any{"protocolVersion": Protocol, "capabilities": map[string]any{},
|
||||
@@ -248,19 +315,62 @@ func Start(module, entry string, env []string, publish Publisher, logf func(stri
|
||||
return nil, err
|
||||
}
|
||||
mu.Lock()
|
||||
current = fresh
|
||||
if current == nil {
|
||||
current = fresh
|
||||
startedAt = time.Now()
|
||||
} else {
|
||||
_ = fresh.cmd.Process.Signal(syscall.SIGTERM)
|
||||
fresh = current
|
||||
}
|
||||
c = fresh
|
||||
mu.Unlock()
|
||||
}
|
||||
return c.ask(module, method, params, timeout)
|
||||
}
|
||||
|
||||
restart = func(after time.Duration) {
|
||||
go func() {
|
||||
time.Sleep(after)
|
||||
mu.Lock()
|
||||
if stopped || current != nil {
|
||||
mu.Unlock()
|
||||
return
|
||||
}
|
||||
mu.Unlock()
|
||||
fresh, err := start()
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if err != nil {
|
||||
if !stopped {
|
||||
wait := pause
|
||||
if pause *= 2; pause > RestartMost {
|
||||
pause = RestartMost
|
||||
}
|
||||
logf("[mesh-tools] %s's bundle did not start again: %v; trying in %s", module, err, wait)
|
||||
go restart(wait)
|
||||
}
|
||||
return
|
||||
}
|
||||
if stopped {
|
||||
_ = fresh.cmd.Process.Signal(syscall.SIGTERM)
|
||||
return
|
||||
}
|
||||
if current == nil {
|
||||
current = fresh
|
||||
startedAt = time.Now()
|
||||
} else {
|
||||
_ = fresh.cmd.Process.Signal(syscall.SIGTERM)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
first, err := start()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
mu.Lock()
|
||||
current = first
|
||||
startedAt = time.Now()
|
||||
mu.Unlock()
|
||||
|
||||
raw, err := asking("tools/list", map[string]any{}, HandshakeTimeout)
|
||||
|
||||
@@ -145,3 +145,40 @@ func (l *Logs) Has(fragments ...string) bool {
|
||||
|
||||
// All is every line.
|
||||
func (l *Logs) All() string { return strings.Join(l.lines, "\n") }
|
||||
|
||||
// Consumer makes a module's durable consumer on a machine as the controller does: pull, on the
|
||||
// EVENTS stream, filtered to what the module consumes, with a short ack wait so a test sees a
|
||||
// redelivery in seconds rather than the mesh's minutes.
|
||||
func (m *Mesh) Consumer(t *testing.T, node, module string, filters []string, ackWait time.Duration) {
|
||||
t.Helper()
|
||||
cfg := &nats.ConsumerConfig{Durable: node + "_" + module, AckPolicy: nats.AckExplicitPolicy,
|
||||
AckWait: ackWait, MaxDeliver: 10, DeliverPolicy: nats.DeliverNewPolicy}
|
||||
if len(filters) == 1 {
|
||||
cfg.FilterSubject = filters[0]
|
||||
} else {
|
||||
cfg.FilterSubjects = filters
|
||||
}
|
||||
if _, err := m.js.AddConsumer("EVENTS", cfg); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
// Emit publishes an event as a module would, into the EVENTS stream.
|
||||
func (m *Mesh) Emit(t *testing.T, subject string, body any) {
|
||||
t.Helper()
|
||||
data, _ := json.Marshal(body)
|
||||
if _, err := m.js.Publish(subject, data); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
// Pending is what a module's consumer still holds: delivered and not acknowledged, and not yet
|
||||
// delivered.
|
||||
func (m *Mesh) Pending(t *testing.T, node, module string) (ackPending, notDelivered uint64) {
|
||||
t.Helper()
|
||||
info, err := m.js.ConsumerInfo("EVENTS", node+"_"+module)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return uint64(info.NumAckPending), info.NumPending
|
||||
}
|
||||
|
||||
@@ -0,0 +1,157 @@
|
||||
package runtime
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"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"
|
||||
)
|
||||
|
||||
// A logger safe to write from the runtime's goroutines.
|
||||
type lines struct {
|
||||
mu sync.Mutex
|
||||
l []string
|
||||
}
|
||||
|
||||
func (s *lines) logf(format string, args ...any) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.l = append(s.l, strings.TrimSpace(sprintf(format, args...)))
|
||||
}
|
||||
|
||||
func (s *lines) all() string {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return strings.Join(s.l, "\n")
|
||||
}
|
||||
|
||||
func sprintf(format string, args ...any) string { return fmtSprintf(format, args...) }
|
||||
|
||||
func countLines(t *testing.T, path, want string) int {
|
||||
t.Helper()
|
||||
b, _ := os.ReadFile(path)
|
||||
n := 0
|
||||
for _, l := range strings.Split(string(b), "\n") {
|
||||
if l == want {
|
||||
n++
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
// novox/hq ADR 0198: a module's long-running code is launched by the node's runtime, and the runtime
|
||||
// is its bus — it binds the module's own consumer, hands each event to the bundle, and acknowledges it
|
||||
// only when the bundle has taken it; one the bundle failed or died on is offered again.
|
||||
func TestTheRuntimeHandsAModulesEventsToItsBundleAndAcknowledgesThemOnlyWhenTaken(t *testing.T) {
|
||||
mesh := mt.New(t)
|
||||
restore := bus.NakDelay
|
||||
bus.NakDelay = 300 * time.Millisecond
|
||||
t.Cleanup(func() { bus.NakDelay = restore })
|
||||
firstPause := launch.RestartFirst
|
||||
launch.RestartFirst = 100 * time.Millisecond
|
||||
t.Cleanup(func() { launch.RestartFirst = firstPause })
|
||||
|
||||
mesh.Issue(t, mt.MembershipOf("watcher", "anchor", false, nil))
|
||||
mesh.Issue(t, mt.MembershipOf("beta", "anchor", true, nil))
|
||||
// The controller's consumer for the module, as its own runtime bound it: anchor_watcher on EVENTS.
|
||||
mesh.Consumer(t, "anchor", "watcher", []string{"mesh.mod.alpha.event.>"}, 2*time.Second)
|
||||
|
||||
nodeTools := connect(t, "node-tools", "anchor")
|
||||
asker := connect(t, "console", "workstation")
|
||||
dir := t.TempDir()
|
||||
log := filepath.Join(dir, "watch.log")
|
||||
said := &lines{}
|
||||
stop, err := Run(nodeTools, []Served{
|
||||
{Module: "watcher", Entrypoints: []string{mt.Fixture("watcher.serve.mjs")}},
|
||||
{Module: "beta", Entrypoints: []string{mt.Fixture("many-beta.serve.mjs")}},
|
||||
}, map[string]map[string]string{"watcher": {"WATCH_LOG": log}}, said.logf)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stop)
|
||||
if !strings.Contains(said.all(), "watcher's events arrive on its consumer anchor_watcher") {
|
||||
t.Fatalf("the module's consumer was not bound:\n%s", said.all())
|
||||
}
|
||||
|
||||
// Handled once, acknowledged once.
|
||||
mesh.Emit(t, "mesh.mod.alpha.event.happened", map[string]any{"n": 1})
|
||||
mt.Until(t, func() error {
|
||||
if countLines(t, log, "handled 1") != 1 {
|
||||
return errorf("event 1 not handled yet")
|
||||
}
|
||||
return nil
|
||||
})
|
||||
// The handler fails the first time: not acknowledged, offered again, then handled.
|
||||
mesh.Emit(t, "mesh.mod.alpha.event.happened", map[string]any{"n": 2, "fail": true})
|
||||
// The bundle dies on the first offer: not acknowledged, the bundle is started again and the
|
||||
// event is offered again once its ack wait has passed.
|
||||
mesh.Emit(t, "mesh.mod.alpha.event.happened", map[string]any{"n": 3, "die": true})
|
||||
deadline := time.Now().Add(20 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
if countLines(t, log, "handled 2") == 1 && countLines(t, log, "handled 3") == 1 {
|
||||
break
|
||||
}
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
}
|
||||
if countLines(t, log, "handled 2") != 1 || countLines(t, log, "handled 3") != 1 {
|
||||
b, _ := os.ReadFile(log)
|
||||
t.Fatalf("a failed or interrupted event was not offered again and handled once:\n%s\nruntime:\n%s", b, said.all())
|
||||
}
|
||||
if countLines(t, log, "started") < 2 {
|
||||
t.Errorf("the bundle that died was not started again: %s", said.all())
|
||||
}
|
||||
mt.Until(t, func() error {
|
||||
ack, notYet := mesh.Pending(t, "anchor", "watcher")
|
||||
if ack != 0 || notYet != 0 {
|
||||
return errorf("still pending: %d unacknowledged, %d undelivered", ack, notYet)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if countLines(t, log, "handled 1") != 1 {
|
||||
t.Errorf("event 1 was handled more than once")
|
||||
}
|
||||
|
||||
// A tool asks another module's tool through the runtime, as the module.
|
||||
got, err := call(t, asker, "watcher.relay@anchor", map[string]any{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
same(t, got, `{"beta":3}`)
|
||||
}
|
||||
|
||||
// A long-running bundle that exits is started again, without waiting for a call (ADR 0198).
|
||||
func TestALongRunningBundleThatExitsIsStartedAgain(t *testing.T) {
|
||||
mesh := mt.New(t)
|
||||
firstPause := launch.RestartFirst
|
||||
launch.RestartFirst = 100 * time.Millisecond
|
||||
t.Cleanup(func() { launch.RestartFirst = firstPause })
|
||||
mesh.Issue(t, mt.MembershipOf("flaky", "anchor", false, nil))
|
||||
nodeTools := connect(t, "node-tools", "anchor")
|
||||
log := filepath.Join(t.TempDir(), "flaky.log")
|
||||
said := &lines{}
|
||||
stop, err := Run(nodeTools, []Served{{Module: "flaky", Entrypoints: []string{mt.Fixture("flaky.serve.mjs")}}},
|
||||
map[string]map[string]string{"flaky": {"FLAKY_LOG": log}}, said.logf)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stop)
|
||||
mt.Until(t, func() error {
|
||||
if countLines(t, log, "started") < 3 {
|
||||
return errorf("started %d time(s)", countLines(t, log, "started"))
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if !strings.Contains(said.all(), "flaky's bundle exited (3)") || !strings.Contains(said.all(), "started again in") {
|
||||
t.Errorf("the restart was not said:\n%s", said.all())
|
||||
}
|
||||
}
|
||||
|
||||
func fmtSprintf(format string, args ...any) string { return fmt.Sprintf(format, args...) }
|
||||
func errorf(format string, args ...any) error { return fmt.Errorf(format, args...) }
|
||||
@@ -157,6 +157,10 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log
|
||||
return out
|
||||
}
|
||||
|
||||
// The module's events, for every child of it that subscribes (ADR 0198): one consumer per module,
|
||||
// bound the first time any of its children subscribes, each event handed to every child that did.
|
||||
events := &consumers{conn: conn, logf: logf, of: map[string]*moduleEvents{}}
|
||||
|
||||
failed := map[string]string{}
|
||||
var registrations []registration
|
||||
var stops []func()
|
||||
@@ -172,13 +176,7 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log
|
||||
fail(path + " is not executable; a bundle the runtime serves is started, never imported, and its build makes it executable (novox/hq ADR 0193)")
|
||||
continue
|
||||
}
|
||||
child, err := launch.Start(module, path, envFor(module), func(params json.RawMessage) error {
|
||||
var env bus.Envelope
|
||||
if err := json.Unmarshal(params, &env); err != nil {
|
||||
return fmt.Errorf("not an event envelope: %w", err)
|
||||
}
|
||||
return conn.PublishAs(module, env)
|
||||
}, logf)
|
||||
child, err := launch.Start(module, path, envFor(module), events.forModule(module), logf)
|
||||
if err != nil {
|
||||
fail(err.Error())
|
||||
continue
|
||||
@@ -209,6 +207,7 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log
|
||||
}
|
||||
}
|
||||
stopAll := func() {
|
||||
events.stopAll()
|
||||
for i := len(stops) - 1; i >= 0; i-- {
|
||||
stops[i]()
|
||||
}
|
||||
@@ -450,3 +449,132 @@ func serveSeats(conn *bus.Conn, modules []string, registrations []registration,
|
||||
stops = nil
|
||||
}
|
||||
}
|
||||
|
||||
// consumers holds, per module the runtime serves, the one durable consumer its events arrive on and
|
||||
// the children its events are handed to (novox/hq ADR 0198).
|
||||
type consumers struct {
|
||||
conn *bus.Conn
|
||||
logf func(string, ...any)
|
||||
mu sync.Mutex
|
||||
of map[string]*moduleEvents
|
||||
}
|
||||
|
||||
type moduleEvents struct {
|
||||
stop func()
|
||||
delivers []*func(json.RawMessage) error
|
||||
}
|
||||
|
||||
// stopAll unbinds every module's consumer.
|
||||
func (c *consumers) stopAll() {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
for _, m := range c.of {
|
||||
if m.stop != nil {
|
||||
m.stop()
|
||||
}
|
||||
}
|
||||
c.of = map[string]*moduleEvents{}
|
||||
}
|
||||
|
||||
// forModule is the bus one launched bundle of a module reaches the mesh through.
|
||||
func (c *consumers) forModule(module string) launch.Bus {
|
||||
return &moduleBus{all: c, module: module}
|
||||
}
|
||||
|
||||
type moduleBus struct {
|
||||
all *consumers
|
||||
module string
|
||||
mu sync.Mutex
|
||||
deliver *func(json.RawMessage) error
|
||||
}
|
||||
|
||||
// Publish emits an event as the module (ADR 0193).
|
||||
func (b *moduleBus) Publish(params json.RawMessage) error {
|
||||
var env bus.Envelope
|
||||
if err := json.Unmarshal(params, &env); err != nil {
|
||||
return fmt.Errorf("not an event envelope: %w", err)
|
||||
}
|
||||
return b.all.conn.PublishAs(b.module, env)
|
||||
}
|
||||
|
||||
// Ask calls a tool as the module: `{key, body}`, answered with the tool's result (ADR 0198).
|
||||
func (b *moduleBus) Ask(params json.RawMessage) (json.RawMessage, error) {
|
||||
var asked struct {
|
||||
Key string `json:"key"`
|
||||
Body json.RawMessage `json:"body"`
|
||||
}
|
||||
if err := json.Unmarshal(params, &asked); err != nil || asked.Key == "" {
|
||||
return nil, fmt.Errorf("mesh/ask names no tool: {key, body}")
|
||||
}
|
||||
body := any(asked.Body)
|
||||
if len(asked.Body) == 0 {
|
||||
body = map[string]any{}
|
||||
}
|
||||
answered, err := b.all.conn.AskAs(b.module, asked.Key, body)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(answered.Result) == 0 {
|
||||
return json.RawMessage("null"), nil
|
||||
}
|
||||
return answered.Result, nil
|
||||
}
|
||||
|
||||
// Subscribe hands this bundle the module's events. The consumer is bound once per module; each of
|
||||
// the module's children that subscribed is handed every event, and the event is acknowledged only
|
||||
// when all of them took it — one consumer split between two readers would give each half.
|
||||
func (b *moduleBus) Subscribe(deliver func(json.RawMessage) error) error {
|
||||
b.mu.Lock()
|
||||
if b.deliver == nil {
|
||||
d := deliver
|
||||
b.deliver = &d
|
||||
} else {
|
||||
*b.deliver = deliver
|
||||
}
|
||||
mine := b.deliver
|
||||
b.mu.Unlock()
|
||||
|
||||
c := b.all
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
m := c.of[b.module]
|
||||
if m == nil {
|
||||
m = &moduleEvents{}
|
||||
c.of[b.module] = m
|
||||
}
|
||||
listed := false
|
||||
for _, d := range m.delivers {
|
||||
if d == mine {
|
||||
listed = true
|
||||
}
|
||||
}
|
||||
if !listed {
|
||||
m.delivers = append(m.delivers, mine)
|
||||
}
|
||||
if m.stop != nil {
|
||||
return nil
|
||||
}
|
||||
module := b.module
|
||||
stop, err := c.conn.ConsumeAs(module, func(env bus.Envelope) error {
|
||||
raw, err := json.Marshal(env)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
c.mu.Lock()
|
||||
targets := append([]*func(json.RawMessage) error(nil), c.of[module].delivers...)
|
||||
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)
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
m.stop = stop
|
||||
c.logf("[mesh-tools] %s's events arrive on its consumer %s", module, bus.ConsumerOf(c.conn.Node(), module))
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -13,7 +13,7 @@
|
||||
"test": "node --test --test-concurrency=1 --experimental-strip-types 'test/*.test.ts'"
|
||||
},
|
||||
"dependencies": {
|
||||
"@novox/mesh-sdk": "^0.1.5",
|
||||
"@novox/mesh-sdk": "^0.1.6",
|
||||
"nats": "^2.29.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
|
||||
+7
@@ -0,0 +1,7 @@
|
||||
#!/usr/bin/env node
|
||||
// A long-running bundle that keeps exiting (novox/hq ADR 0198): the runtime starts it again each time.
|
||||
import { appendFileSync } from "node:fs";
|
||||
import { serveStdio } from "@novox/mesh-sdk/stdio";
|
||||
appendFileSync(process.env.FLAKY_LOG, "started\n");
|
||||
setTimeout(() => process.exit(3), 300);
|
||||
await serveStdio("flaky", []);
|
||||
Vendored
+24
@@ -0,0 +1,24 @@
|
||||
// A module whose code runs long (novox/hq ADR 0198): it subscribes as it is imported, handles each
|
||||
// event once it can, fails the first time it is asked to, dies the first time it is told to, and one
|
||||
// tool asks another module's tool through the runtime. What it did is written to WATCH_LOG.
|
||||
import { appendFileSync, existsSync, writeFileSync } from "node:fs";
|
||||
import { on } from "@novox/mesh-sdk/events";
|
||||
import { broker } from "@novox/mesh-sdk/messaging";
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
const log = process.env.WATCH_LOG;
|
||||
const once = (mark) => {
|
||||
const file = `${log}.${mark}`;
|
||||
if (existsSync(file)) return false;
|
||||
writeFileSync(file, "1");
|
||||
return true;
|
||||
};
|
||||
appendFileSync(log, "started\n");
|
||||
await on("alpha.happened", async (e) => {
|
||||
if (e.body.fail && once(`fail-${e.body.n}`)) throw new Error(`not yet ${e.body.n}`);
|
||||
if (e.body.die && once(`die-${e.body.n}`)) process.exit(7);
|
||||
appendFileSync(log, `handled ${e.body.n}\n`);
|
||||
});
|
||||
registerModuleTools("watcher", () => [
|
||||
{ name: "relay", description: "asks beta", input: {}, run: async () => broker().request("beta.three", {}) },
|
||||
]);
|
||||
+5
@@ -0,0 +1,5 @@
|
||||
#!/usr/bin/env node
|
||||
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for watcher.mjs.
|
||||
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
|
||||
await import("./watcher.mjs");
|
||||
await serveRegisteredOverStdio();
|
||||
Reference in New Issue
Block a user