diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go index 608dadc..adb896b 100644 --- a/node-tools/internal/bus/bus.go +++ b/node-tools/internal/bus/bus.go @@ -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, +// `.` another's, `seat:.[@]` 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: `_`. +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..event.` is +// `.` — 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 +} diff --git a/node-tools/internal/launch/launch.go b/node-tools/internal/launch/launch.go index 45caf5d..f96a4fb 100644 --- a/node-tools/internal/launch/launch.go +++ b/node-tools/internal/launch/launch.go @@ -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) diff --git a/node-tools/internal/meshtest/meshtest.go b/node-tools/internal/meshtest/meshtest.go index 7342c0d..337d3d4 100644 --- a/node-tools/internal/meshtest/meshtest.go +++ b/node-tools/internal/meshtest/meshtest.go @@ -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 +} diff --git a/node-tools/internal/runtime/events_test.go b/node-tools/internal/runtime/events_test.go new file mode 100644 index 0000000..2cf56c9 --- /dev/null +++ b/node-tools/internal/runtime/events_test.go @@ -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...) } diff --git a/node-tools/internal/runtime/runtime.go b/node-tools/internal/runtime/runtime.go index 9b616e9..e2c3524 100644 --- a/node-tools/internal/runtime/runtime.go +++ b/node-tools/internal/runtime/runtime.go @@ -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 +} diff --git a/node-tools/package.json b/node-tools/package.json index afef641..56340d4 100644 --- a/node-tools/package.json +++ b/node-tools/package.json @@ -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": { diff --git a/node-tools/test/fixtures/flaky.serve.mjs b/node-tools/test/fixtures/flaky.serve.mjs new file mode 100755 index 0000000..8fcd570 --- /dev/null +++ b/node-tools/test/fixtures/flaky.serve.mjs @@ -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", []); diff --git a/node-tools/test/fixtures/watcher.mjs b/node-tools/test/fixtures/watcher.mjs new file mode 100644 index 0000000..94442d0 --- /dev/null +++ b/node-tools/test/fixtures/watcher.mjs @@ -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", {}) }, +]); diff --git a/node-tools/test/fixtures/watcher.serve.mjs b/node-tools/test/fixtures/watcher.serve.mjs new file mode 100755 index 0000000..f97e72d --- /dev/null +++ b/node-tools/test/fixtures/watcher.serve.mjs @@ -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();