Files
LocalAI/core/services/messaging/client.go
T
mudler-agentandEttore Di Giacinto e4fa051ee3 refactor(distributed): put the NATS-only paths behind interfaces (#12395)
* feat(messaging): add shared subject rules

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* test(messaging): cover BroadcastRoots, ControlRoots and SubjectRoot

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* feat(messaging): add Broadcaster and enforce subject rules in every carrier

Broadcaster is the fan-out half of MessagingClient. The NATS client and
the in-memory FakeBus now refuse a subject outside the served roots and any
wildcard other than a whole single token, and FakeBus shares MatchSubject
instead of its own copy. FakeBus Unsubscribe now removes its own
subscription instead of the first one with the same subject.

A shared conformance suite in messagingtest runs against both carriers.
The distributed e2e specs that used invented test.* subjects, and the one
that subscribed with a > filter, now use subjects from subjects.go.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: depend on Broadcaster where only publish and subscribe are used

Narrowed to messaging.Broadcaster: nodes/staging_progress.go,
nodes/install_progress_publisher.go, galleryop/operation.go,
galleryop/service.go, agentpool/user_services.go, agentpool/agent_jobs.go,
openresponses/store.go, openresponses/sync.go, syncstate/syncstate.go,
finetune/service.go, quantization/service.go and
failover/distsync/distsync.go. SubscribeJSON now takes a Broadcaster
because it only calls Subscribe, which lets the narrowed consumers use it.

Stayed wide: worker/supervisor.go, because its client field also serves
the SubscribeReply handlers in worker/lifecycle.go. The request/reply,
queue and wiring files (nodes/unloader.go, nodes/file_stager_s3.go,
jobs/dispatcher.go, agents/dispatcher.go, agents/events.go,
worker/file_staging.go, cli/agent_worker.go, http/app.go) are unchanged by
design.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): name the no-route condition and confine the carrier error

Consumers matched nats.ErrNoResponders, which names an absence, to demote a node. They now match ErrNoRoute, the control path maps the carrier's failure onto it, and timeouts and worker refusals are pinned as not being no-route.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs(nodes): state which FileStager implementations return ErrNoRoute

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): build backend clients through one node-aware seam

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: describe the distributed transport seams

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: correct comments that overclaim after the seams refactor

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* fix(agent-worker): refuse an unserved LOCALAI_AGENT_SUBJECT at startup

The messaging client now refuses a subject whose root no carrier serves.
An agent worker started with a custom LOCALAI_AGENT_SUBJECT such as
tenant-a.agent.execute used to start and then wait on a subject the
frontend never publishes to. After the subject rules landed it exited
at subscribe time with an error that did not name the setting.

Behaviour change: the worker now checks LOCALAI_AGENT_SUBJECT before it
registers or connects, and exits with an error that names the variable
and says to use a served subject under the agent root, for example
agent.execute. The served roots are not widened: a custom root was
never delivered by the frontend, and a wider set would reopen the
drift the subject rules exist to close. The flag help and the agent
worker docs state the constraint.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* test(nodes): pin the reactions to ErrNoRoute

Three callers react to ErrNoRoute and had no spec: the reconciler's
upgrade drain falls back to the legacy forced install, the reconciler
marks the node unhealthy when a pending op has no route, and the
backend-op fan-out marks the node unhealthy. Each spec drives the real
caller with a scripted no-responders reply and reads the result from
the registry or the recorded requests.

A fourth spec pins the other side: a pending op that times out leaves
the node healthy and only counts the attempt, so mapping timeouts onto
ErrNoRoute would fail here.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* test(messaging): pin client subject checks and fail the carrier suite in CI

Add specs that call Publish, Request, Subscribe, QueueSubscribe,
SubscribeReply and QueueSubscribeReply on a client with no connection.
Each call must return ErrUnservedSubject for bogus.thing and
ErrUnsupportedWildcard for jobs.>. This proves that the subject check
runs before the connection is used, and needs no server.

The NATS conformance suite is the only check that runs the subject rules
against a real carrier. Before this change it skipped without output
when Docker was missing. Now it fails when CI is set, so a Linux runner
without Docker cannot hide it. It still skips on local runs and on macOS
CI, which has no Docker.

Add SubjectNodeBackendInstallProgress to the list of constructors that
must build served subjects, and ask contributors to extend the list.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: state what ErrNoRoute may change, and group the distributed guides

The seams note said MarkUnhealthy was the only state change allowed
on ErrNoRoute. A pending backend op still records the failed attempt,
counts toward the reconciler's retry limit and is dead-lettered after
the maximum attempts. The note now says that MarkUnhealthy is the only
change to the node's own state, and that the per-op accounting is not
a verdict about the node.

The note also documents that the NATS conformance run fails under CI
when Docker is missing. The distributed-seams row moves next to the
distributed-state row in the topics table. The liveness ping spec
header now says no route is a reason to skip the worker, not proof
that the worker is gone.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): give the backend client factory the node id

Mechanical: the method gains a nodeID parameter and the eight test fakes are updated. No behaviour change.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): drop the optional node-aware factory

The node id is now in the main method, so the optional interface and its helper had no behaviour of their own.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): dial backend probes through the client factory

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(nodes): dial workers' file servers through a per-node dialer

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(http): proxy backend logs through the per-node worker dialer

The admin backend-logs proxy (list, lines and the WebSocket stream) now reaches a worker through the same per-node dialer as the HTTP file stager, so every frontend-to-worker dial goes through one seam. The shared direct dialer keeps alive for 15s where the proxy used 30s. Harmless for requests bounded at 15s.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* fix(http): keep the backend-logs proxy independent of the admin connection

The proxy request had no context before the dialer change and is bounded only by its 15s timeout. Keep it that way.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: move the worker control payloads to workerctl

Mechanical move of the request and reply structs, the install progress event and the file payloads out of messaging. The verbs no longer belong to one carrier. No alias is left behind.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(worker): serve the lifecycle verbs through a controlServer

The worker registers one handler per verb and a NATS server maps each verb to its subject. Registration errors now name the verb. node.stop is served with SubscribeReply, which is identical on the wire because the handler never replies.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(worker): report install progress through the control sink

Install and upgrade now emit download progress through the sink the control server hands them. The debounce and the terminal flush stay in the handler path, built over that sink by the new nodes.NewDebouncedInstallProgressSink, which replaces NewDebouncedInstallProgressPublisher. The subject and payload on the wire are unchanged. The supervisor no longer holds the bus, and installFn and upgradeFn let specs drive both verbs without a gallery.

The malformed-request log lines are restored for install, upgrade, backend.delete, model.unload, model.stop and model.delete, with the reply bytes unchanged. The signal adapter is renamed noReply, which also lets worker.go import os/signal without an alias again.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(worker): serve the file-staging verbs through a controlServer

An empty list-dir answer is now {} rather than {"files":null}, because the typed reply omits an empty Files slice. The frontend decodes both to a nil slice in nodes/file_stager_s3.go.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* feat(messaging): add WorkQueue and the NATS producer

This is the producer side of the competing-consumer seam. The work kinds map one to one to today's subjects and queue groups: task to jobs.new and mcp-ci to jobs.mcp-ci.new (both in group workers), agent-run to agent.execute (group agent-workers). Enqueue publishes the payload as Publish does today, with one JSON marshal. FakeBus now records queue groups and keeps reply handlers so later specs can pin and drive them.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* feat(messaging): add the NATS WorkConsumer

An in-flight limit of one runs the handler inline on the delivery goroutine, as the MCP CI consumer does today. Any other limit spawns per delivery, as the agent consumer does. Queue groups are unchanged.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: publish queued work through WorkQueue

The job dispatcher, the agent pool and the agent scheduler enqueue through messaging.WorkQueue; the NATS implementation publishes to the same subjects as before. DistributedServices builds the queue next to the NATS client and hands it to the dispatcher and the agent pool, whose distributed mode switch now reads a non-nil WorkQueue. The unused AgentPoolService.SetNATSClient is removed.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: consume queued work through WorkConsumer

The agent dispatcher and the MCP CI consumer register through messaging.WorkConsumer. The NATS implementation keeps the inline one-at-a-time model for MCP CI and the per-delivery model for agent runs. handleMCPCIJob reports on the events publisher the carrier hands it instead of a captured client.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: delete the consumers nothing in production reached

jobs.new has a producer and no production consumer, and the agent dispatcher's Dispatch was only called from tests. Publishing jobs.new is unchanged.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(mcp): send MCP requests to agent workers through AgentControl

Timeouts still honour only the deadline, not cancellation, exactly as today. The NATS no-responders error maps to ErrNoRoute and a timeout does not.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(agent-worker): serve MCP requests and backend.stop through agentRPCServer

The agent worker's MCP tool and discovery reply subscriptions and its backend stop listener move behind an unexported agentRPCServer interface, served on NATS by nodes.NATSAgentRPCServer. The handlers become typed mcp.ToolHandler and mcp.DiscoveryHandler values that answer every failure with a reply carrying Error.

Queue group (agent-workers), inline execution on the delivery goroutine, the background handler context, the unmarshal error reply texts and the reply-less backend stop subscription are unchanged. The backend stop handler takes the decoded backend name, so it can still close that backend's MCP sessions.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor(messaging): remove helpers that only tests used

BroadcastRoots, ControlRoots and SubjectRoot had no production caller. The roots spec now asserts every served root through ValidateSubject instead. MatchSubject moves back into the test support package, the only place that used it, with its table. NATSAgentRPCServer drops the subscription list it stored and never read, and NewNATSAgentRPCServer gets a doc comment.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* test(mcp): round trip the agent RPC server over a real NATS server

One spec sends a tool request and a discovery request through NATSAgentControl to NATSAgentRPCServer and checks that the handlers see the decoded requests and the replies come back. It also puts an undecodable body on the tool subject and checks the server answers with an unmarshal error instead of leaving the requester to time out.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: describe the distributed transport seams

The developer note now lists the final seams: fan-out, queues, both halves of the control verbs and of agent RPC, and the dial. It records the open items a second carrier has to handle.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* test: pin the in-flight limit each queue consumer asks for

The work queue specs pin what Consume does for a given limit, but nothing
pinned which limit each production consumer passes. Changing the agent
worker's MCP CI limit from 1 to 0 would have let MCP CI jobs run
concurrently on each worker with every test green.

Move the MCP CI Consume call into startMCPCIConsumer with the same wiring
and pin that it asks for (WorkMCPCI, 1). Pin that NATSDispatcher.Start asks
for (WorkAgentRun, maxConcurrent) for several limits.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* refactor: remove helpers the branch left without a caller

SubjectJobCancelWildcard lost its last subscriber when the frontend
stopped listening on jobs.*.cancel; the NATS permissions and conformance
suite spell the subject out, so nothing reads the constant.

decodeBackendStopRequest returned a stopAll flag that production dropped
and only a test read. decodeBackendStop is now the single decoder with the
same semantics: an empty body is stop-all, an empty Backend is stop-all,
malformed JSON is an error. stopBackends still derives stop-all from
Backend, so no reply changes.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* fix(messaging): keep an explicitly empty agent queue a plain subscription

Before the work queue seam the agent worker passed LOCALAI_AGENT_QUEUE
straight to QueueSubscribe, so an explicitly empty value made a plain
subscription and every agent worker ran every agent run.
WithAgentRunRoute replaced an empty queue with agent-workers, which
silently changed that.

Keep the queue as given once the option is applied. An empty subject still
falls back to agent.execute, since it never had a meaning of its own. The
flag default stays agent-workers, so only an explicitly empty value
reaches this.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* docs: correct comments and record the PR B notes

Fix the recordingFactory comment (it also records the parallel flag),
document that a negative maxInFlight is unbounded and that Unsubscribe from
a handler deadlocks, and say a permanently undecodable payload returns nil.

Record controlHandler's undecodable return as a kept exception, and add the
second carrier notes to the developer note: the reconciler has no
ClientFactory option, the logs proxy honours HTTP_PROXY, verbs one carrier
serves need an opt-out, terminal replies come from the result event, and
agent runs publish through the NATS-bound EventBridge, which is not an
additive change.

Assisted-by: Claude:claude-sonnet-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

---------

Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
2026-10-02 23:58:41 +02:00

334 lines
12 KiB
Go

package messaging
import (
"encoding/json"
"errors"
"fmt"
"strings"
"sync"
"time"
"github.com/mudler/LocalAI/pkg/sanitize"
"github.com/mudler/xlog"
"github.com/nats-io/nats.go"
"github.com/nats-io/nkeys"
)
// subscribeConfirmTimeout bounds the server round-trip used to detect whether a
// subscription was rejected (e.g. by JWT permissions) before returning to the caller.
const subscribeConfirmTimeout = 5 * time.Second
// Client wraps a NATS connection and provides helpers for pub/sub and queue subscriptions.
type Client struct {
conn *nats.Conn
mu sync.RWMutex
// reconnectCbs are invoked after the underlying connection is
// re-established. nats.go transparently resubscribes existing
// subscriptions on reconnect, but it cannot know that a consumer kept
// derived in-memory state (e.g. syncstate.SyncedMap) that may have drifted
// while the link was down — these callbacks let such consumers re-hydrate.
cbMu sync.Mutex
reconnectCbs []func()
}
// New creates a new NATS client with auto-reconnect.
func New(url string, opts ...Option) (*Client, error) {
var cfg connectConfig
for _, o := range opts {
o(&cfg)
}
// Allocate the client up front so the reconnect handler closure can reach
// it; conn is populated after nats.Connect succeeds below.
c := &Client{}
natsOpts := []nats.Option{
nats.RetryOnFailedConnect(true),
nats.MaxReconnects(-1),
nats.DisconnectErrHandler(func(_ *nats.Conn, err error) {
if err != nil {
xlog.Warn("NATS disconnected", "error", err)
}
}),
nats.ReconnectHandler(func(_ *nats.Conn) {
xlog.Info("NATS reconnected")
c.runReconnectCallbacks()
}),
nats.ClosedHandler(func(_ *nats.Conn) {
xlog.Info("NATS connection closed")
}),
// Surface async errors (notably permission violations) that NATS would
// otherwise deliver silently. A subscription the server rejects for a
// JWT permission means the worker never receives those messages, so make
// it loud rather than letting the feature fail invisibly.
nats.ErrorHandler(func(_ *nats.Conn, sub *nats.Subscription, err error) {
subject := ""
if sub != nil {
subject = sub.Subject
}
if errors.Is(err, nats.ErrPermissionViolation) {
xlog.Error("NATS permission violation — check JWT pub/sub allow lists", "subject", subject, "error", err)
return
}
xlog.Warn("NATS async error", "subject", subject, "error", err)
}),
}
switch {
case cfg.jwtProvider != nil:
// Fetch creds on every (re)connect so a refresh loop can rotate the JWT
// before expiry; the server expiring the old JWT triggers a reconnect
// that transparently picks up the new one.
natsOpts = append(natsOpts, nats.UserJWT(
func() (string, error) {
jwt, _ := cfg.jwtProvider()
if jwt == "" {
return "", fmt.Errorf("no NATS user JWT available")
}
return jwt, nil
},
func(nonce []byte) ([]byte, error) {
_, seed := cfg.jwtProvider()
kp, err := nkeys.FromSeed([]byte(seed))
if err != nil {
return nil, fmt.Errorf("loading NATS user seed: %w", err)
}
defer kp.Wipe()
return kp.Sign(nonce)
},
))
case cfg.userJWT != "" && cfg.userSeed != "":
natsOpts = append(natsOpts, nats.UserJWTAndSeed(cfg.userJWT, cfg.userSeed))
}
if cfg.tls.Enabled() {
if err := cfg.tls.Validate(); err != nil {
return nil, err
}
tlsOpts, err := cfg.tls.natsOptions()
if err != nil {
return nil, err
}
natsOpts = append(natsOpts, tlsOpts...)
}
nc, err := nats.Connect(url, natsOpts...)
if err != nil {
return nil, fmt.Errorf("connecting to NATS at %s: %w", sanitize.URL(url), err)
}
c.conn = nc
return c, nil
}
// OnReconnect registers a callback invoked after the NATS connection is
// re-established. It is consumed via an optional interface type-assertion
// (interface{ OnReconnect(func()) }) rather than being added to MessagingClient,
// so the messaging abstraction stays minimal and standalone/test clients are not
// forced to implement reconnect semantics. A nil callback is ignored.
func (c *Client) OnReconnect(cb func()) {
if cb == nil {
return
}
c.cbMu.Lock()
c.reconnectCbs = append(c.reconnectCbs, cb)
c.cbMu.Unlock()
}
// runReconnectCallbacks invokes registered reconnect callbacks. It copies the
// slice under the lock so a callback that (re)registers cannot deadlock.
func (c *Client) runReconnectCallbacks() {
c.cbMu.Lock()
cbs := append([]func(){}, c.reconnectCbs...)
c.cbMu.Unlock()
for _, cb := range cbs {
cb()
}
}
// Publish marshals data as JSON and publishes it to the given subject.
func (c *Client) Publish(subject string, data any) error {
if err := ValidateSubject(subject); err != nil {
return err
}
payload, err := json.Marshal(data)
if err != nil {
return fmt.Errorf("marshalling message for %s: %w", subject, err)
}
c.mu.RLock()
defer c.mu.RUnlock()
return c.conn.Publish(subject, payload)
}
// Subscribe creates a subscription on the given subject. All subscribers receive every message.
func (c *Client) Subscribe(subject string, handler func([]byte)) (Subscription, error) {
return c.confirmSubscription(subject, func(conn *nats.Conn) (*nats.Subscription, error) {
return conn.Subscribe(subject, func(msg *nats.Msg) {
handler(msg.Data)
})
})
}
// QueueSubscribe creates a queue subscription. Within the same queue group,
// only one subscriber receives each message (load-balanced).
func (c *Client) QueueSubscribe(subject, queue string, handler func([]byte)) (Subscription, error) {
return c.confirmSubscription(subject, func(conn *nats.Conn) (*nats.Subscription, error) {
return conn.QueueSubscribe(subject, queue, func(msg *nats.Msg) {
handler(msg.Data)
})
})
}
// confirmSubscription creates a subscription via mk and forces a server
// round-trip so that a permissions violation — which NATS otherwise reports
// only asynchronously — is returned to the caller synchronously. The server
// emits the "-ERR Permissions Violation" for a rejected SUB before the PONG
// that satisfies the flush, so by the time FlushTimeout returns the violation
// is recorded as the connection's last error. Without this, a worker whose JWT
// lacks a subject gets a non-nil subscription that never receives a message,
// turning a permission misconfiguration into a silent failure.
func (c *Client) confirmSubscription(subject string, mk func(*nats.Conn) (*nats.Subscription, error)) (Subscription, error) {
if err := ValidateSubject(subject); err != nil {
return nil, err
}
c.mu.RLock()
conn := c.conn
c.mu.RUnlock()
if conn == nil {
return nil, fmt.Errorf("subscribe to %s: nil NATS connection", subject)
}
sub, err := mk(conn)
if err != nil {
return nil, err
}
// A failed flush here means we could not round-trip to the server (not yet
// connected, reconnecting, slow link). RetryOnFailedConnect intentionally
// buffers subscriptions across that gap, so do NOT fail — keep the
// subscription and let it replay on (re)connect; a later permission
// violation is still logged by the async error handler in New.
if err := conn.FlushTimeout(subscribeConfirmTimeout); err != nil {
xlog.Debug("Could not confirm NATS subscription (will replay on connect)", "subject", subject, "error", err)
return sub, nil
}
// Flush succeeded, so any permission violation for this SUB has already been
// recorded as the connection's last error (the server emits it before the
// PONG). LastError is per-connection; match the exact quoted subject the
// server echoes ("Subscription to \"<subject>\"") so a stale violation for
// another subject can't be mis-attributed here.
if lerr := conn.LastError(); lerr != nil &&
errors.Is(lerr, nats.ErrPermissionViolation) &&
strings.Contains(lerr.Error(), `Subscription to "`+subject+`"`) {
_ = sub.Unsubscribe()
return nil, fmt.Errorf("subscription to %s denied by NATS server (check JWT sub allow list): %w", subject, lerr)
}
return sub, nil
}
// Request sends a request and waits for a reply (request-reply pattern).
// Returns the raw reply data.
func (c *Client) Request(subject string, data []byte, timeout time.Duration) ([]byte, error) {
if err := ValidateSubject(subject); err != nil {
return nil, err
}
c.mu.RLock()
defer c.mu.RUnlock()
msg, err := c.conn.Request(subject, data, timeout)
if err != nil {
return nil, fmt.Errorf("request to %s: %w", subject, err)
}
return msg.Data, nil
}
// SubscribeReply creates a subscription that supports replying to requests.
// The handler receives the raw request data and the reply subject.
func (c *Client) SubscribeReply(subject string, handler func(data []byte, reply func([]byte))) (Subscription, error) {
return c.confirmSubscription(subject, func(conn *nats.Conn) (*nats.Subscription, error) {
return conn.Subscribe(subject, func(msg *nats.Msg) {
handler(msg.Data, func(replyData []byte) {
if msg.Reply != "" {
if err := msg.Respond(replyData); err != nil {
xlog.Warn("Failed to send NATS reply", "subject", subject, "error", err)
}
}
})
})
})
}
// QueueSubscribeReply creates a queue subscription that supports replying to requests.
// Load-balanced across subscribers in the same queue group, with request-reply support.
func (c *Client) QueueSubscribeReply(subject, queue string, handler func(data []byte, reply func([]byte))) (Subscription, error) {
return c.confirmSubscription(subject, func(conn *nats.Conn) (*nats.Subscription, error) {
return conn.QueueSubscribe(subject, queue, func(msg *nats.Msg) {
handler(msg.Data, func(replyData []byte) {
if msg.Reply != "" {
if err := msg.Respond(replyData); err != nil {
xlog.Warn("Failed to send NATS reply", "subject", subject, "error", err)
}
}
})
})
})
}
// SubscribeJSON creates a subscription that automatically unmarshals JSON messages.
// Invalid JSON messages are logged and skipped.
func SubscribeJSON[T any](c Broadcaster, subject string, handler func(T)) (Subscription, error) {
return c.Subscribe(subject, func(data []byte) {
var evt T
if err := json.Unmarshal(data, &evt); err != nil {
xlog.Warn("Failed to unmarshal NATS message", "subject", subject, "error", err)
return
}
handler(evt)
})
}
// RequestJSON sends a JSON request-reply via NATS, marshaling the request and
// unmarshaling the reply. This eliminates the repeated marshal/request/unmarshal
// boilerplate across all NATS request-reply call sites.
func RequestJSON[Req, Reply any](c MessagingClient, subject string, req Req, timeout time.Duration) (*Reply, error) {
data, err := json.Marshal(req)
if err != nil {
return nil, fmt.Errorf("marshaling request: %w", err)
}
replyData, err := c.Request(subject, data, timeout)
if err != nil {
return nil, fmt.Errorf("NATS request to %s: %w", subject, err)
}
var reply Reply
if err := json.Unmarshal(replyData, &reply); err != nil {
return nil, fmt.Errorf("unmarshaling reply from %s: %w", subject, err)
}
return &reply, nil
}
// Conn returns the underlying NATS connection for advanced usage.
//
// Deprecated: Prefer using the MessagingClient interface methods (Publish, Subscribe, etc.)
// instead of accessing the raw NATS connection. This method couples callers to the
// concrete Client type and bypasses the abstraction layer.
func (c *Client) Conn() *nats.Conn {
c.mu.RLock()
defer c.mu.RUnlock()
return c.conn
}
// IsConnected returns true if the client is currently connected to a NATS server.
func (c *Client) IsConnected() bool {
c.mu.RLock()
defer c.mu.RUnlock()
return c.conn != nil && c.conn.IsConnected()
}
// Close drains and closes the NATS connection, waiting for in-flight messages.
func (c *Client) Close() {
c.mu.Lock()
defer c.mu.Unlock()
if c.conn != nil {
c.conn.Drain()
c.conn.FlushTimeout(5 * time.Second)
}
}