Versions in this module Expand all Collapse all v0 v0.1.0 Aug 19, 2026 Changes in this version + func EnsureDLQStream(ctx context.Context, js jetstream.JetStream, maxBytes int64) error + func NewRouter(deps Dependencies) http.Handler + func RequireAdmin(store *policy.Store, logger *slog.Logger) func(http.Handler) http.Handler + type BootState struct + func NewBootState(initialErr error) *BootState + func (b *BootState) Err() error + func (b *BootState) Set(err error) + type DLQHandler struct + JS jetstream.JetStream + Logger *slog.Logger + func NewDLQHandler(js jetstream.JetStream, logger *slog.Logger) *DLQHandler + func (h *DLQHandler) Stats(w http.ResponseWriter, r *http.Request) + type Dependencies struct + AuthMW func(http.Handler) http.Handler + CORSOrigins []string + DLQ *DLQHandler + Health *HealthHandler + Ingest *IngestHandler + JS jetstream.JetStream + Logger *slog.Logger + MetricsHandler http.Handler + MetricsPath string + Pipes *PipesHandler + Policy *PolicyHandler + PolicyStore *policy.Store + Query *QueryHandler + SSE *StreamHandler + Schema *SchemaHandler + StructuredQuery *StructuredQueryHandler + Version *VersionHandler + type HealthHandler struct + Boot *BootState + CHConn driver.Conn + func NewHealthHandler(chConn driver.Conn) *HealthHandler + func (h *HealthHandler) Liveness(w http.ResponseWriter, _ *http.Request) + func (h *HealthHandler) Online(w http.ResponseWriter, _ *http.Request) + func (h *HealthHandler) Readiness(w http.ResponseWriter, r *http.Request) + type IngestHandler struct + Dedup dedupe.Deduplicator + IDField string + PolicyStore *policy.Store + Publisher mq.Publisher + Registry *discovery.SchemaRegistry + RequireID bool + func NewIngestHandler(registry *discovery.SchemaRegistry, pub mq.Publisher, logger *slog.Logger) *IngestHandler + func (h *IngestHandler) Handle(w http.ResponseWriter, r *http.Request) + type PipesHandler struct + CHConn driver.Conn + Cache cache.Cache + PolicyStore *policy.Store + Store *pipes.Store + func NewPipesHandler(store *pipes.Store, policyStore *policy.Store, conn driver.Conn, c cache.Cache, ...) *PipesHandler + func (h *PipesHandler) Delete(w http.ResponseWriter, r *http.Request) + func (h *PipesHandler) Execute(w http.ResponseWriter, r *http.Request) + func (h *PipesHandler) Get(w http.ResponseWriter, r *http.Request) + func (h *PipesHandler) List(w http.ResponseWriter, _ *http.Request) + func (h *PipesHandler) Put(w http.ResponseWriter, r *http.Request) + type PolicyHandler struct + Store *policy.Store + func NewPolicyHandler(store *policy.Store) *PolicyHandler + func (h *PolicyHandler) Get(w http.ResponseWriter, r *http.Request) + func (h *PolicyHandler) Put(w http.ResponseWriter, r *http.Request) + func (h *PolicyHandler) Validate(w http.ResponseWriter, r *http.Request) + type QueryHandler struct + Database string + Endpoint string + HTTPClient *http.Client + Password string + Username string + func NewQueryHandler(endpoint, username, password, database string, queryTimeout time.Duration) *QueryHandler + func (h *QueryHandler) Handle(w http.ResponseWriter, r *http.Request) + type SchemaHandler struct + Registry *discovery.SchemaRegistry + func NewSchemaHandler(registry *discovery.SchemaRegistry) *SchemaHandler + func (h *SchemaHandler) Get(w http.ResponseWriter, r *http.Request) + func (h *SchemaHandler) List(w http.ResponseWriter, _ *http.Request) + func (h *SchemaHandler) Refresh(w http.ResponseWriter, r *http.Request) + type StreamHandler struct + Heartbeater *stream.Heartbeater + Hub *stream.Hub + JS jetstream.JetStream + Metrics *stream.Metrics + func NewStreamHandler(hub *stream.Hub, js jetstream.JetStream) *StreamHandler + func (h *StreamHandler) Handle(w http.ResponseWriter, r *http.Request) + type StructuredQueryHandler struct + BucketSecs int + CHConn driver.Conn + Cache cache.Cache + PolicyStore *policy.Store + Registry *discovery.SchemaRegistry + func NewStructuredQueryHandler(conn driver.Conn, c cache.Cache, registry *discovery.SchemaRegistry, ...) *StructuredQueryHandler + func (h *StructuredQueryHandler) Handle(w http.ResponseWriter, r *http.Request) + type VersionHandler struct + func NewVersionHandler(version, gitCommit, buildTime string) *VersionHandler + func (h *VersionHandler) Handle(w http.ResponseWriter, _ *http.Request)