|
@@ -1,355 +1,40 @@
|
|
|
-// process.go is the shared alert-processing pipeline used by
|
|
|
|
|
-// both the HTTP handler (cmd/ingestd/http.go) and the MQTT
|
|
|
|
|
-// subscriber (cmd/ingestd/mqtt.go). It does the M0→M4 SPEC §22
|
|
|
|
|
-// work in one place: parse → validate → look up source → HMAC
|
|
|
|
|
-// verify → rate-limit (per-source, per-company) → dedupe → stamp
|
|
|
|
|
-// → publish to NATS → return a result struct the caller maps
|
|
|
|
|
-// to its own transport-level response.
|
|
|
|
|
-//
|
|
|
|
|
-// The function is pure (no global state, no transport types).
|
|
|
|
|
-// It returns Result{Accepted/Rejected/Error + reason + alert_id
|
|
|
|
|
-// + dedupe_count} so the HTTP handler can map to a status code
|
|
|
|
|
-// and the MQTT handler can map to a log line.
|
|
|
|
|
|
|
+// process.go is a thin re-export and compatibility shim for the shared
|
|
|
|
|
+// pipeline (internal/pipeline/). All the actual logic lives in pipeline/.
|
|
|
|
|
+// This file exists so that existing callers in package main (http.go,
|
|
|
|
|
+// mqtt.go, ws.go, main.go) that reference types by their package-main
|
|
|
|
|
+// names don't need to change.
|
|
|
package main
|
|
package main
|
|
|
|
|
|
|
|
import (
|
|
import (
|
|
|
"context"
|
|
"context"
|
|
|
- "encoding/json"
|
|
|
|
|
- "errors"
|
|
|
|
|
- "fmt"
|
|
|
|
|
- "log/slog"
|
|
|
|
|
- "strconv"
|
|
|
|
|
- "time"
|
|
|
|
|
|
|
|
|
|
- "git3.techno-world.net/lrosales/broad-announce/internal/alert"
|
|
|
|
|
- "git3.techno-world.net/lrosales/broad-announce/internal/broker"
|
|
|
|
|
- "git3.techno-world.net/lrosales/broad-announce/internal/circuitbreaker"
|
|
|
|
|
- "git3.techno-world.net/lrosales/broad-announce/internal/dedupe"
|
|
|
|
|
- "git3.techno-world.net/lrosales/broad-announce/internal/observability"
|
|
|
|
|
- "git3.techno-world.net/lrosales/broad-announce/internal/quarantine"
|
|
|
|
|
- "git3.techno-world.net/lrosales/broad-announce/internal/ratelimit"
|
|
|
|
|
- "git3.techno-world.net/lrosales/broad-announce/internal/tailhub"
|
|
|
|
|
- "github.com/nats-io/nats.go"
|
|
|
|
|
|
|
+ pipeline "git3.techno-world.net/lrosales/broad-announce/internal/pipeline"
|
|
|
)
|
|
)
|
|
|
|
|
|
|
|
-// Result is the outcome of a ProcessAlert call. The caller
|
|
|
|
|
-// (HTTP/MQTT) maps it to its transport's response shape.
|
|
|
|
|
-type Result struct {
|
|
|
|
|
- // Accepted is true if the alert passed all checks and was
|
|
|
|
|
- // published to NATS.
|
|
|
|
|
- Accepted bool
|
|
|
|
|
- // AlertID is the server-assigned id (empty on error paths).
|
|
|
|
|
- AlertID string
|
|
|
|
|
- // DedupeCount is the count returned by the dedupe layer.
|
|
|
|
|
- DedupeCount uint32
|
|
|
|
|
- // IsNew is true for the first alert in a dedupe window.
|
|
|
|
|
- IsNew bool
|
|
|
|
|
- // RejectReason is the SPEC §22 / auth reason; one of:
|
|
|
|
|
- // "payload_too_large", "bad_request", "invalid_json",
|
|
|
|
|
- // "invalid", "unknown_source", "bad_signature",
|
|
|
|
|
- // "rate_limited_source", "rate_limited_company",
|
|
|
|
|
- // "marshal_failed", "broker_unavailable"
|
|
|
|
|
- RejectReason string
|
|
|
|
|
- // HTTPStatus is the suggested HTTP status code (0 for
|
|
|
|
|
- // accepted / 202).
|
|
|
|
|
- HTTPStatus int
|
|
|
|
|
- // Detail is the free-form string the caller can show in
|
|
|
|
|
- // a response body or a log message.
|
|
|
|
|
- Detail string
|
|
|
|
|
-}
|
|
|
|
|
|
|
+// SourceConfig is re-exported from pipeline so callers in package main
|
|
|
|
|
+// (ws.go, mqtt.go, loadSourcesFromEnv) can use it without an import.
|
|
|
|
|
+type SourceConfig = pipeline.SourceConfig
|
|
|
|
|
+
|
|
|
|
|
+// Result is re-exported for backward compatibility.
|
|
|
|
|
+type Result = pipeline.Result
|
|
|
|
|
|
|
|
-// Accept is the canonical "ok" result.
|
|
|
|
|
|
|
+// Accept is re-exported for backward compatibility.
|
|
|
func Accept(id string, count uint32, isNew bool) Result {
|
|
func Accept(id string, count uint32, isNew bool) Result {
|
|
|
- return Result{Accepted: true, AlertID: id, DedupeCount: count, IsNew: isNew, HTTPStatus: 202}
|
|
|
|
|
|
|
+ return pipeline.Accept(id, count, isNew)
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
-// Reject is the canonical "no" result.
|
|
|
|
|
|
|
+// Reject is re-exported for backward compatibility.
|
|
|
func Reject(reason string, status int, detail string) Result {
|
|
func Reject(reason string, status int, detail string) Result {
|
|
|
- return Result{RejectReason: reason, HTTPStatus: status, Detail: detail}
|
|
|
|
|
|
|
+ return pipeline.Reject(reason, status, detail)
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
-// processDeps is the process-pipeline dependency set. Smaller
|
|
|
|
|
-// than httpDeps — no HTTP-specific fields. Both httpDeps and
|
|
|
|
|
-// the MQTT subscriber construct one and call ProcessAlert.
|
|
|
|
|
-type processDeps struct {
|
|
|
|
|
- Logger *slog.Logger
|
|
|
|
|
- Metrics *observability.IngestdMetrics
|
|
|
|
|
- Limiter *ratelimit.Limiter
|
|
|
|
|
- Deduper *dedupe.Deduper
|
|
|
|
|
- JetStream natsPublisher
|
|
|
|
|
- // Sources is the (company_id, source_id) → SourceConfig map.
|
|
|
|
|
- // M0 reads it from env; M2 from Postgres. The MQTT subscriber
|
|
|
|
|
- // uses the same map keyed on the topic-parsed (co, src).
|
|
|
|
|
- Sources map[string]SourceConfig
|
|
|
|
|
- // Per-company default rate cap. HTTP and MQTT use the same
|
|
|
|
|
- // constant; once the per-company cap lives in DB, both
|
|
|
|
|
- // transports read it.
|
|
|
|
|
- CompanyRatePerSec int
|
|
|
|
|
- // Tail is an optional M5 live-tail hub. When non-nil, every
|
|
|
|
|
- // accepted alert is also published to in-process tail
|
|
|
|
|
- // subscribers. nil is fine (HTTP/MQTT tests don't need it).
|
|
|
|
|
- Tail *tailhub.Hub
|
|
|
|
|
- // Transport is the per-process transport label used in
|
|
|
|
|
- // structured log lines ("http" | "mqtt" | "ws"). The MQTT
|
|
|
|
|
- // path overrides this on the receiver's scoped copy.
|
|
|
|
|
- Transport string
|
|
|
|
|
- // Now is overridable in tests.
|
|
|
|
|
- Now func() time.Time
|
|
|
|
|
- // MaxSeen is the M6 per-source monotonic max tracker for
|
|
|
|
|
- // dedupe_count. processDeps owns one so all transports
|
|
|
|
|
- // (HTTP, MQTT, WS) share the same in-process state and
|
|
|
|
|
- // the same max-observed gauge.
|
|
|
|
|
- MaxSeen *observability.MaxSeen
|
|
|
|
|
- // CircuitBreaker is the M9 layer-6 per-component circuit
|
|
|
|
|
- // breaker wrapping the NATS publish call. Nil is fine (falls
|
|
|
|
|
- // back to direct publish without circuit protection).
|
|
|
|
|
- CircuitBreaker *circuitbreaker.Breaker
|
|
|
|
|
- // Quarantine is the M9 layer-7 per-source error-rate limiter.
|
|
|
|
|
- // Nil is fine (no quarantine enforcement).
|
|
|
|
|
- Quarantine *quarantine.Manager
|
|
|
|
|
-}
|
|
|
|
|
|
|
+// processDeps embeds pipeline.Deps and adds ProcessAlert as a thin
|
|
|
|
|
+// compatibility wrapper so existing callers (http.go, mqtt.go, ws.go)
|
|
|
|
|
+// don't need to change their call sites.
|
|
|
|
|
+type processDeps struct{ pipeline.Deps }
|
|
|
|
|
|
|
|
-// ProcessAlert runs the full SPEC §22 protection chain on one
|
|
|
|
|
-// alert body. It is the single source of truth for the ingest
|
|
|
|
|
-// pipeline; both the HTTP POST handler and the MQTT subscriber
|
|
|
|
|
-// call it.
|
|
|
|
|
-//
|
|
|
|
|
-// Layer order (matches SPEC §22):
|
|
|
|
|
-// 1. payload-size cap (caller does this — http.go via
|
|
|
|
|
-// MaxBytesReader; mqtt.go via the
|
|
|
|
|
-// broker-side max-inflight setting)
|
|
|
|
|
-// 2. quarantine check (M9 layer 7) — per-source ban if error
|
|
|
|
|
-// rate exceeds threshold
|
|
|
|
|
-// 3. per-source rate limit
|
|
|
|
|
-// 4. per-company rate limit
|
|
|
|
|
-// 5. schema validate
|
|
|
|
|
-// 5b. parse JSON
|
|
|
|
|
-// (auth) HMAC verify — see verifyHMAC
|
|
|
|
|
-// 5c. dedupe
|
|
|
|
|
-// publish to NATS (M9 layer 6 circuit breaker wraps this)
|
|
|
|
|
|
|
+// ProcessAlert is the legacy entry point. It delegates to the embedded
|
|
|
|
|
+// pipeline.Deps.Process, which has the same signature.
|
|
|
func (d *processDeps) ProcessAlert(ctx context.Context, body []byte, sigHeader string) Result {
|
|
func (d *processDeps) ProcessAlert(ctx context.Context, body []byte, sigHeader string) Result {
|
|
|
- now := d.now()
|
|
|
|
|
-
|
|
|
|
|
- // 5. Parse + validate. We treat any parse failure as invalid.
|
|
|
|
|
- var a alert.Alert
|
|
|
|
|
- if err := json.Unmarshal(body, &a); err != nil {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
|
|
|
|
|
- return Reject("invalid_json", 400, err.Error())
|
|
|
|
|
- }
|
|
|
|
|
- if err := a.Validate(); err != nil {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
|
|
|
|
|
- return Reject("invalid", 400, err.Error())
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // Look up source. M0: in-memory map. M2: DB.
|
|
|
|
|
- src, ok := d.Sources[a.CompanyID+":"+a.SourceID]
|
|
|
|
|
- if !ok {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
|
|
|
|
|
- return Reject("unknown_source", 401,
|
|
|
|
|
- fmt.Sprintf("no such source %s/%s", a.CompanyID, a.SourceID))
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // M9 layer 7: quarantine check. Check before we spend any CPU
|
|
|
|
|
- // on a known-bad source.
|
|
|
|
|
- if d.Quarantine != nil {
|
|
|
|
|
- if banned, remaining, err := d.Quarantine.IsBanned(ctx, a.SourceID); err == nil && banned {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("quarantined").Inc()
|
|
|
|
|
- d.Logger.Warn("source quarantined",
|
|
|
|
|
- "source_id", a.SourceID,
|
|
|
|
|
- "company_id", a.CompanyID,
|
|
|
|
|
- "remaining", remaining,
|
|
|
|
|
- )
|
|
|
|
|
- return Reject("quarantined", 429,
|
|
|
|
|
- fmt.Sprintf("source quarantined for %v; retry after", remaining.Round(time.Second)))
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // Auth. Stripe-style: X-BA-Signature: t=<unix>,v1=<hex>.
|
|
|
|
|
- // For HTTP it's a header; for MQTT it's a top-level field
|
|
|
|
|
- // on the envelope — both call sites pass the same string.
|
|
|
|
|
- if !verifyHMAC(sigHeader, src.HMACSecret, body, now) {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
|
|
|
|
|
- return Reject("bad_signature", 401, "")
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // M9: quarantine hit tracking. Any rejection after the HMAC
|
|
|
|
|
- // check (where we know the source is real) indicates a
|
|
|
|
|
- // problematic source. We use a closure so the defer pattern
|
|
|
|
|
- // catches every return path without boilerplate at each one.
|
|
|
|
|
- var hitRecorded bool
|
|
|
|
|
- defer func() {
|
|
|
|
|
- if !hitRecorded && d.Quarantine != nil {
|
|
|
|
|
- _ = d.Quarantine.RecordHit(context.Background(), a.SourceID)
|
|
|
|
|
- }
|
|
|
|
|
- }()
|
|
|
|
|
- // recordHit records a quarantine error hit for a.SourceID.
|
|
|
|
|
- // Call it before any rejection return.
|
|
|
|
|
- recordHit := func() {
|
|
|
|
|
- if d.Quarantine != nil && !hitRecorded {
|
|
|
|
|
- hitRecorded = true
|
|
|
|
|
- _ = d.Quarantine.RecordHit(context.Background(), a.SourceID)
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // M6: Dedupe BEFORE rate limit. A duplicate (isNew=false)
|
|
|
|
|
- // is a Redis INCR + JSON marshal + NATS publish — it does
|
|
|
|
|
- // not warrant burning a rate-limit token. The rate limit
|
|
|
|
|
- // exists to backpressure "new alert" volume; the dedupe
|
|
|
|
|
- // itself is the canonical "do less work for repeats"
|
|
|
|
|
- // mechanism. We only burn a token on the first arrival
|
|
|
|
|
- // in a sliding window; the next 999 dupes pass through
|
|
|
|
|
- // the rate limit gates for free, the recipient sees one
|
|
|
|
|
- // consolidated message with `(×N)` appended.
|
|
|
|
|
- isNew, count, err := d.Deduper.Check(ctx, a.SourceID, a.DedupeKey)
|
|
|
|
|
- if err != nil {
|
|
|
|
|
- d.Logger.Warn("dedupe redis error (failing open)", "err", err)
|
|
|
|
|
- isNew, count = true, 1
|
|
|
|
|
- }
|
|
|
|
|
- if !isNew {
|
|
|
|
|
- // M6 metrics: tick the per-source collapse counter and
|
|
|
|
|
- // bump the max-observed gauge if this hit set a new
|
|
|
|
|
- // peak. We use a tiny in-process max tracker (Prom's
|
|
|
|
|
- // Gauge doesn't expose Get() — the canonical pattern
|
|
|
|
|
- // is to read via .Gauges() and pick out the value, but
|
|
|
|
|
- // that's a 2-step write+read; cleaner to just remember
|
|
|
|
|
- // the max in our own map and Set the gauge on growth).
|
|
|
|
|
- d.Metrics.DedupeCollapsed.WithLabelValues(a.SourceID).Inc()
|
|
|
|
|
- d.MaxSeen.RecordAndExport(a.SourceID, count,
|
|
|
|
|
- func(s string, v float64) {
|
|
|
|
|
- d.Metrics.DedupeCountMax.WithLabelValues(s).Set(v)
|
|
|
|
|
- })
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // 3. Per-source rate limit (only charged for new alerts).
|
|
|
|
|
- if isNew {
|
|
|
|
|
- if ok, ttl, err := d.Limiter.Allow(ctx, "source:"+a.CompanyID+":"+a.SourceID, src.RateLimitPerSec); err != nil {
|
|
|
|
|
- d.Logger.Warn("ratelimit redis error (failing open)", "err", err, "scope", "source")
|
|
|
|
|
- } else if !ok {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("rate_limited").Inc()
|
|
|
|
|
- d.Metrics.RateLimitHits.WithLabelValues("source").Inc()
|
|
|
|
|
- _ = ttl
|
|
|
|
|
- recordHit()
|
|
|
|
|
- return Reject("rate_limited_source", 429, strconv.Itoa(int(ttl.Seconds())))
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // 4. Per-company rate limit (only charged for new alerts).
|
|
|
|
|
- if isNew {
|
|
|
|
|
- if ok, ttl, _ := d.Limiter.Allow(ctx, "company:"+a.CompanyID, d.CompanyRatePerSec); !ok {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("rate_limited").Inc()
|
|
|
|
|
- d.Metrics.RateLimitHits.WithLabelValues("company").Inc()
|
|
|
|
|
- _ = ttl
|
|
|
|
|
- recordHit()
|
|
|
|
|
- return Reject("rate_limited_company", 429, strconv.Itoa(int(ttl.Seconds())))
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- // Stamp server-side fields.
|
|
|
|
|
- a.ID = alert.NewID()
|
|
|
|
|
- a.ReceivedAt = now.UTC()
|
|
|
|
|
- a.DedupeCount = count
|
|
|
|
|
-
|
|
|
|
|
- // Publish to NATS. M9 layer 6: circuit breaker wraps the
|
|
|
|
|
- // publish call so a sick NATS server doesn't take down ingestd.
|
|
|
|
|
- subject := broker.AlertsSubject(a.CompanyID)
|
|
|
|
|
- payload, err := json.Marshal(a)
|
|
|
|
|
- if err != nil {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("invalid").Inc()
|
|
|
|
|
- return Reject("marshal_failed", 500, err.Error())
|
|
|
|
|
- }
|
|
|
|
|
- start := time.Now()
|
|
|
|
|
- var publishErr error
|
|
|
|
|
- if d.CircuitBreaker != nil {
|
|
|
|
|
- // Wrap the synchronous publish in the circuit breaker.
|
|
|
|
|
- // We use Publish (sync) so Do() gets immediate feedback.
|
|
|
|
|
- publishErr = d.CircuitBreaker.Do(ctx, func() error {
|
|
|
|
|
- return d.JetStream.Publish(subject, payload)
|
|
|
|
|
- })
|
|
|
|
|
- } else {
|
|
|
|
|
- publishErr = d.JetStream.Publish(subject, payload)
|
|
|
|
|
- }
|
|
|
|
|
- if publishErr != nil {
|
|
|
|
|
- if errors.Is(publishErr, circuitbreaker.ErrCircuitOpen) {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("circuit_open").Inc()
|
|
|
|
|
- d.Metrics.CBState.WithLabelValues("nats").Set(circuitbreaker.StateOpen)
|
|
|
|
|
- d.Logger.Warn("circuit breaker open",
|
|
|
|
|
- "subject", subject,
|
|
|
|
|
- "alert_id", a.ID,
|
|
|
|
|
- "company_id", a.CompanyID,
|
|
|
|
|
- )
|
|
|
|
|
- recordHit()
|
|
|
|
|
- return Reject("circuit_open", 503, "broker circuit breaker open")
|
|
|
|
|
- }
|
|
|
|
|
- // Real publish error (network, auth, etc.).
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("broker_unavailable").Inc()
|
|
|
|
|
- d.Logger.Error("nats publish", "err", publishErr, "subject", subject)
|
|
|
|
|
- recordHit()
|
|
|
|
|
- return Reject("broker_unavailable", 503, publishErr.Error())
|
|
|
|
|
- }
|
|
|
|
|
- d.Metrics.PublishLatency.WithLabelValues(a.SourceID).Observe(time.Since(start).Seconds())
|
|
|
|
|
- d.Metrics.PayloadBytes.Observe(float64(len(payload)))
|
|
|
|
|
-
|
|
|
|
|
- if isNew {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("accepted").Inc()
|
|
|
|
|
- } else {
|
|
|
|
|
- d.Metrics.AlertsReceived.WithLabelValues("deduped").Inc()
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- d.Logger.Info("alert accepted",
|
|
|
|
|
- "alert_id", a.ID,
|
|
|
|
|
- "company_id", a.CompanyID,
|
|
|
|
|
- "source_id", a.SourceID,
|
|
|
|
|
- "severity", string(a.Severity),
|
|
|
|
|
- "transport", d.Transport,
|
|
|
|
|
- "dedupe_count", count,
|
|
|
|
|
- )
|
|
|
|
|
- // M5: fan out to the live-tail hub (if configured). This
|
|
|
|
|
- // is best-effort and never blocks the producer — the hub's
|
|
|
|
|
- // Publish drops on slow consumers.
|
|
|
|
|
- if d.Tail != nil {
|
|
|
|
|
- ev := tailhub.FromAlert(&a, d.Transport)
|
|
|
|
|
- d.Tail.Publish(ev)
|
|
|
|
|
- }
|
|
|
|
|
- // Accepted: mark that we did not get a rejection so the
|
|
|
|
|
- // defer does not record a spurious quarantine hit.
|
|
|
|
|
- hitRecorded = true
|
|
|
|
|
- return Accept(a.ID, count, isNew)
|
|
|
|
|
-}
|
|
|
|
|
-
|
|
|
|
|
-func (d *processDeps) now() time.Time {
|
|
|
|
|
- if d.Now != nil {
|
|
|
|
|
- return d.Now()
|
|
|
|
|
- }
|
|
|
|
|
- return time.Now()
|
|
|
|
|
-}
|
|
|
|
|
-
|
|
|
|
|
-// natsPublisher is the minimal NATS interface. The HTTP and
|
|
|
|
|
-// MQTT paths share it; tests can swap in a fake.
|
|
|
|
|
-type natsPublisher interface {
|
|
|
|
|
- PublishAsync(subj string, data []byte) error
|
|
|
|
|
- // Publish is synchronous. The circuit breaker uses this
|
|
|
|
|
- // to get immediate success/failure feedback.
|
|
|
|
|
- Publish(subj string, data []byte) error
|
|
|
|
|
-}
|
|
|
|
|
-
|
|
|
|
|
-// jsPublisher adapts a nats.JetStreamContext to the natsPublisher interface.
|
|
|
|
|
-type jsPublisher struct {
|
|
|
|
|
- js nats.JetStreamContext
|
|
|
|
|
-}
|
|
|
|
|
-
|
|
|
|
|
-// Publish is synchronous (blocks until server ack or timeout). Used
|
|
|
|
|
-// by the circuit breaker which needs immediate success/failure feedback.
|
|
|
|
|
-func (j *jsPublisher) Publish(subj string, data []byte) error {
|
|
|
|
|
- _, err := j.js.Publish(subj, data)
|
|
|
|
|
- return err
|
|
|
|
|
-}
|
|
|
|
|
-
|
|
|
|
|
-func (j *jsPublisher) PublishAsync(subj string, data []byte) error {
|
|
|
|
|
- _, err := j.js.PublishAsync(subj, data)
|
|
|
|
|
- return err
|
|
|
|
|
-}
|
|
|
|
|
-
|
|
|
|
|
-// newNatsPublisher is the constructor used by main.
|
|
|
|
|
-func newNatsPublisher(js nats.JetStreamContext) natsPublisher {
|
|
|
|
|
- return &jsPublisher{js: js}
|
|
|
|
|
|
|
+ return d.Process(ctx, body, sigHeader)
|
|
|
}
|
|
}
|