micro--go-micro
75e32f4d87
* Initial plan * Implement NATS connection pool with configuration options Co-authored-by: asim <17530+asim@users.noreply.github.com> * Implement connection pool for transport/nats Co-authored-by: asim <17530+asim@users.noreply.github.com> * Fix connection leaks in events/natsjs and config/source/nats Co-authored-by: asim <17530+asim@users.noreply.github.com> * Fix race condition in connection pool lastUsed field access Co-authored-by: asim <17530+asim@users.noreply.github.com> * Remove unused maxIdle field from connection pools Co-authored-by: asim <17530+asim@users.noreply.github.com> --------- Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com> Co-authored-by: asim <17530+asim@users.noreply.github.com>
NATS JetStream
This plugin uses NATS with JetStream to send and receive events.
Create a stream
ev, err := natsjs.NewStream(
natsjs.Address("nats://10.0.1.46:4222"),
natsjs.MaxAge(24*160*time.Minute),
)
Consume a stream
ee, err := events.Consume("test",
events.WithAutoAck(false, time.Second*30),
events.WithGroup("testgroup"),
)
if err != nil {
panic(err)
}
go func() {
for {
msg := <-ee
// Process the message
logger.Info("Received message:", string(msg.Payload))
err := msg.Ack()
if err != nil {
logger.Error("Error acknowledging message:", err)
} else {
logger.Info("Message acknowledged")
}
}
}()
Publish an Event to the stream
err = ev.Publish("test", []byte("hello world"))
if err != nil {
panic(err)
}