Files
agentmon/internal/queue/nats/nats.go
T
William Valentin 256b841cbf feat: scaffold agentmon services and k8s deploy
Adds Go microservices (ingest-gateway, event-processor, query-api, web-ui), NATS+Postgres wiring, initial schema/init job, ingress manifests for LAN+tailnet, and a multi-arch image build script.
2026-01-17 01:06:57 -08:00

35 lines
615 B
Go

package nats
import (
"context"
"time"
gnats "github.com/nats-io/nats.go"
)
type Publisher struct {
conn *gnats.Conn
topic string
timeout time.Duration
}
func NewPublisher(url, topic string) (*Publisher, error) {
conn, err := gnats.Connect(url)
if err != nil {
return nil, err
}
return &Publisher{conn: conn, topic: topic, timeout: 5 * time.Second}, nil
}
func (p *Publisher) Close() {
p.conn.Close()
}
func (p *Publisher) Publish(ctx context.Context, data []byte) error {
ctx, cancel := context.WithTimeout(ctx, p.timeout)
defer cancel()
_ = ctx
return p.conn.Publish(p.topic, data)
}