apirouter

package
v1.6.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 2, 2026 License: Apache-2.0 Imports: 39 Imported by: 0

Documentation

Index

Constants

View Source
const (

	// Role values
	RoleAdmin  = "admin"
	RoleTenant = "tenant"
)
View Source
const (
	// MaxRequestBodySize limits the size of request bodies we'll buffer for logging
	// Set to 10KB to prevent excessive log volumes
	MaxRequestBodySize = 10 * 1024
	// SensitiveFieldMask is used to replace sensitive values
	SensitiveFieldMask = "[REDACTED]"
)

Variables

View Source
var (
	ErrMissingAuthHeader  = errors.New("missing authorization header")
	ErrInvalidBearerToken = errors.New("invalid bearer token format")
)
View Source
var (
	ErrInvalidToken = errors.New("invalid token")
)
View Source
var JWT = jsonwebtoken{}

Functions

func AbortWithError

func AbortWithError(c *gin.Context, code int, err error)

func AbortWithValidationError

func AbortWithValidationError(c *gin.Context, err error)

func AuthMiddleware added in v0.13.0

func AuthMiddleware(apiKey, jwtSecret string, tenantRetriever TenantRetriever, opts AuthOptions) gin.HandlerFunc

AuthMiddleware returns a single gin.HandlerFunc that handles authentication, authorization, and tenant resolution for every route.

Flow:

  1. VPC mode (apiKey=""): grant admin, resolve tenant if RequireTenant, done.
  2. Validate auth header → 401 if missing/malformed.
  3. token == apiKey → admin, resolve tenant if RequireTenant, done.
  4. JWT.Extract(token) → 401 if invalid.
  5. AdminOnly? → 403.
  6. :tenant_id param mismatch? → 403.
  7. Set tenantID + RoleTenant, always resolve tenant for JWT → 401 if missing/deleted.

func CursorToPtr added in v0.13.0

func CursorToPtr(cursor string) *string

CursorToPtr converts an empty cursor string to nil, or returns a pointer to the string. Returns null for empty cursors instead of empty string.

func ErrorHandlerMiddleware

func ErrorHandlerMiddleware() gin.HandlerFunc

func GetRequestLatency

func GetRequestLatency(c *gin.Context) time.Duration

GetRequestLatency returns the request latency from the context

func LatencyMiddleware

func LatencyMiddleware() gin.HandlerFunc

func LoggerMiddleware

func LoggerMiddleware(logger *logging.Logger) gin.HandlerFunc

func LoggerMiddlewareWithSanitizer

func LoggerMiddlewareWithSanitizer(logger *logging.Logger, sanitizer *RequestBodySanitizer) gin.HandlerFunc

func MetricsMiddleware

func MetricsMiddleware() gin.HandlerFunc

func NewRouter

func NewRouter(cfg RouterConfig, deps RouterDeps) http.Handler

func ParseArrayQueryParam added in v0.14.0

func ParseArrayQueryParam(c *gin.Context, fieldName string) []string

ParseArrayQueryParam parses array query parameters. Supports three formats:

  • Indexed bracket notation: field[0]=a&field[1]=b
  • Unindexed bracket notation: field[]=a&field[]=b
  • Plain repeated params: field=a&field=b (fallback when no bracket notation present)

Returns nil if the parameter is absent entirely.

func ParseCursors added in v0.13.0

func ParseCursors(c *gin.Context) (CursorParams, *ErrorResponse)

ParseCursors parses the "next" and "prev" query parameters. Returns an error if both are provided (mutually exclusive).

func ParseDateFilter added in v0.13.0

func ParseDateFilter(c *gin.Context, fieldName string) (*DateFilterResult, *ErrorResponse)

ParseDateFilter parses bracket notation date filters (e.g., field[gte], field[lte], field[gt], field[lt]). Returns 400 for unsupported operators (any). Accepts both RFC3339 format (2024-01-01T00:00:00Z) and date-only format (2024-01-01).

func RecoveryMiddleware added in v1.6.0

func RecoveryMiddleware() gin.HandlerFunc

RecoveryMiddleware recovers from a panic and answers with a JSON 500.

Types

type APIAttempt added in v0.13.0

type APIAttempt struct {
	ID              string                 `json:"id"`
	TenantID        string                 `json:"tenant_id"`
	Status          string                 `json:"status"`
	Time            time.Time              `json:"time"`
	Code            string                 `json:"code,omitempty"`
	ResponseData    map[string]interface{} `json:"response_data,omitempty"`
	AttemptNumber   int                    `json:"attempt_number"`
	Manual          bool                   `json:"manual"`
	DestinationType string                 `json:"destination_type"`
	LatencyMs       *int64                 `json:"latency_ms"`

	EventID       string      `json:"event_id"`
	DestinationID string      `json:"destination_id"`
	Event         interface{} `json:"event,omitempty"`
	Destination   interface{} `json:"destination,omitempty"`
}

APIAttempt is the API response for an attempt

type APIAttemptDestination added in v1.6.0

type APIAttemptDestination struct {
	*destregistry.DestinationDisplay
	Credentials *struct{} `json:"credentials,omitempty"` // shadows the embedded credentials
}

APIAttemptDestination is the destination object when include=destination. Credentials are only returned by the destination endpoints.

type APIEvent added in v0.13.0

type APIEvent struct {
	ID                    string            `json:"id"`
	TenantID              string            `json:"tenant_id"`
	MatchedDestinationIDs []string          `json:"matched_destination_ids"`
	Topic                 string            `json:"topic"`
	Time                  time.Time         `json:"time"`
	EligibleForRetry      bool              `json:"eligible_for_retry"`
	Metadata              map[string]string `json:"metadata,omitempty"`
	Data                  json.RawMessage   `json:"data,omitempty"`
}

APIEvent is the API response for retrieving a single event

type APIEventFull added in v0.13.0

type APIEventFull struct {
	ID                    string            `json:"id"`
	TenantID              string            `json:"tenant_id"`
	MatchedDestinationIDs []string          `json:"matched_destination_ids,omitempty"`
	Topic                 string            `json:"topic"`
	Time                  time.Time         `json:"time"`
	EligibleForRetry      bool              `json:"eligible_for_retry"`
	Metadata              map[string]string `json:"metadata,omitempty"`
	Data                  json.RawMessage   `json:"data,omitempty"`
}

APIEventFull is the event object when expand=event.data

type APIEventSummary added in v0.13.0

type APIEventSummary struct {
	ID                    string            `json:"id"`
	TenantID              string            `json:"tenant_id"`
	MatchedDestinationIDs []string          `json:"matched_destination_ids,omitempty"`
	Topic                 string            `json:"topic"`
	Time                  time.Time         `json:"time"`
	EligibleForRetry      bool              `json:"eligible_for_retry"`
	Metadata              map[string]string `json:"metadata,omitempty"`
}

APIEventSummary is the event object when expand=event (without data)

type APIMetricsDataPoint added in v0.15.0

type APIMetricsDataPoint struct {
	TimeBucket *time.Time     `json:"time_bucket,omitempty"`
	Dimensions map[string]any `json:"dimensions"`
	Metrics    map[string]any `json:"metrics"`
}

type APIMetricsMetadata added in v0.15.0

type APIMetricsMetadata struct {
	Granularity *string `json:"granularity,omitempty"`
	QueryTimeMs int64   `json:"query_time_ms"`
	RowCount    int     `json:"row_count"`
	RowLimit    int     `json:"row_limit"`
	Truncated   bool    `json:"truncated"`
}

type APIMetricsResponse added in v0.15.0

type APIMetricsResponse struct {
	Data     []APIMetricsDataPoint `json:"data"`
	Metadata APIMetricsMetadata    `json:"metadata"`
}

type AttemptPaginatedResult added in v0.13.0

type AttemptPaginatedResult struct {
	Models     []APIAttempt   `json:"models"`
	Pagination SeekPagination `json:"pagination"`
}

AttemptPaginatedResult is the paginated response for listing attempts.

type AuthOptions added in v0.13.0

type AuthOptions struct {
	AdminOnly     bool
	RequireTenant bool
}

AuthOptions configures the behaviour of AuthMiddleware.

type BufferedReader

type BufferedReader struct {
	// contains filtered or unexported fields
}

BufferedReader creates a reader that can be used multiple times

func NewBufferedReader

func NewBufferedReader(r io.Reader) (*BufferedReader, error)

NewBufferedReader creates a new BufferedReader from an io.Reader

func (*BufferedReader) Bytes

func (br *BufferedReader) Bytes() []byte

Bytes returns a copy of the buffered content

func (*BufferedReader) NewReadCloser

func (br *BufferedReader) NewReadCloser() io.ReadCloser

func (*BufferedReader) NewReader

func (br *BufferedReader) NewReader() io.Reader

NewReader creates a new reader from the buffered content

func (*BufferedReader) Read

func (br *BufferedReader) Read(p []byte) (n int, err error)

Read implements io.Reader

type CreateDestinationRequest

type CreateDestinationRequest struct {
	ID               string                  `json:"id" binding:"-"`
	Type             string                  `json:"type" binding:"required"`
	Topics           models.Topics           `json:"topics" binding:"required"`
	Filter           models.Filter           `json:"filter,omitempty" binding:"-"`
	Config           models.Config           `json:"config" binding:"-"`
	Credentials      models.Credentials      `json:"credentials" binding:"-"`
	DeliveryMetadata models.DeliveryMetadata `json:"delivery_metadata,omitempty" binding:"-"`
	Metadata         models.Metadata         `json:"metadata,omitempty" binding:"-"`
	CreatedAt        *time.Time              `json:"created_at,omitempty" binding:"-"`
	UpdatedAt        *time.Time              `json:"updated_at,omitempty" binding:"-"`
	DisabledAt       *time.Time              `json:"disabled_at,omitempty" binding:"-"`
}

func (*CreateDestinationRequest) ToDestination

func (r *CreateDestinationRequest) ToDestination(tenantID string) models.Destination

type CursorParams added in v0.13.0

type CursorParams struct {
	Next string
	Prev string
}

CursorParams holds the parsed cursor values for pagination.

type DateFilterResult added in v0.13.0

type DateFilterResult struct {
	GTE *time.Time // Greater than or equal (>=)
	LTE *time.Time // Less than or equal (<=)
	GT  *time.Time // Greater than (>)
	LT  *time.Time // Less than (<)
}

DateFilterResult holds the parsed date filter values.

type DestinationHandlers

type DestinationHandlers struct {
	// contains filtered or unexported fields
}

func NewDestinationHandlers

func NewDestinationHandlers(logger *logging.Logger, telemetry telemetry.Telemetry, tenantStore tenantstore.TenantStore, emitter SubscriptionEmitter, topics []string, topicsAllowWildcards bool, registry destregistry.Registry, displayer *destinationDisplayer) *DestinationHandlers

func (*DestinationHandlers) Create

func (h *DestinationHandlers) Create(c *gin.Context)

func (*DestinationHandlers) Delete

func (h *DestinationHandlers) Delete(c *gin.Context)

func (*DestinationHandlers) Disable

func (h *DestinationHandlers) Disable(c *gin.Context)

func (*DestinationHandlers) Enable

func (h *DestinationHandlers) Enable(c *gin.Context)

func (*DestinationHandlers) List

func (h *DestinationHandlers) List(c *gin.Context)

func (*DestinationHandlers) ListProviderMetadata

func (h *DestinationHandlers) ListProviderMetadata(c *gin.Context)

func (*DestinationHandlers) Retrieve

func (h *DestinationHandlers) Retrieve(c *gin.Context)

func (*DestinationHandlers) RetrieveProviderMetadata

func (h *DestinationHandlers) RetrieveProviderMetadata(c *gin.Context)

func (*DestinationHandlers) Update

func (h *DestinationHandlers) Update(c *gin.Context)

type ErrorResponse

type ErrorResponse struct {
	Err     error       `json:"-"`
	Code    int         `json:"-"`
	Status  int         `json:"status"`
	Message string      `json:"message"`
	Data    interface{} `json:"data,omitempty"`
}

func NewErrBadRequest

func NewErrBadRequest(err error) ErrorResponse

func NewErrForbidden added in v1.6.0

func NewErrForbidden() ErrorResponse

func NewErrInternalServer

func NewErrInternalServer(err error) ErrorResponse

func NewErrNotFound

func NewErrNotFound(resource string) ErrorResponse

func NewErrUnauthorized added in v1.6.0

func NewErrUnauthorized() ErrorResponse

func ParseDir added in v0.13.0

func ParseDir(c *gin.Context) (string, *ErrorResponse)

ParseDir parses the "dir" query parameter for sort direction. Returns the direction (empty string if not provided) and any validation error response. Valid values are "asc" and "desc". Caller/store should apply default if empty.

func ParseLimit added in v1.6.0

func ParseLimit(c *gin.Context, maxLimit int) (int, *ErrorResponse)

ParseLimit parses the "limit" query parameter. Returns 0 if not provided; caller/store should apply the default. Returns an error response if the value is not an integer between 1 and maxLimit.

func ParseOrderBy added in v0.13.0

func ParseOrderBy(c *gin.Context, allowedValues []string) (string, *ErrorResponse)

ParseOrderBy parses the "order_by" query parameter and validates against allowed values. Returns the order_by value (empty string if not provided) and any validation error response. Caller/store should apply default if empty.

func (ErrorResponse) Error

func (e ErrorResponse) Error() string

func (*ErrorResponse) Parse

func (e *ErrorResponse) Parse(err error)

type EventPaginatedResult added in v0.13.0

type EventPaginatedResult struct {
	Models     []APIEvent     `json:"models"`
	Pagination SeekPagination `json:"pagination"`
}

EventPaginatedResult is the paginated response for listing events.

type IncludeOptions added in v0.13.0

type IncludeOptions struct {
	Event        bool
	EventData    bool
	ResponseData bool
	Destination  bool
}

IncludeOptions represents which fields to include in the response

type JWTClaims added in v0.13.0

type JWTClaims struct {
	TenantID     string
	DeploymentID string
}

JWTClaims contains the custom claims for JWT tokens

type LogHandlers

type LogHandlers struct {
	// contains filtered or unexported fields
}

func NewLogHandlers

func NewLogHandlers(
	logger *logging.Logger,
	logStore logstore.LogStore,
	tenantStore tenantstore.TenantStore,
	displayer *destinationDisplayer,
) *LogHandlers

func (*LogHandlers) ListAttempts added in v0.13.0

func (h *LogHandlers) ListAttempts(c *gin.Context)

ListAttempts handles GET /attempts Query params: tenant_id[], event_id[], destination_id[], status, topic[], time[gte], time[lte], time[gt], time[lt], limit, next, prev, include, order_by, dir

func (*LogHandlers) ListDestinationAttempts added in v0.13.0

func (h *LogHandlers) ListDestinationAttempts(c *gin.Context)

ListDestinationAttempts handles GET /:tenant_id/destinations/:destination_id/attempts Same as ListAttempts but scoped to a specific destination via URL param.

func (*LogHandlers) ListEvents added in v0.13.0

func (h *LogHandlers) ListEvents(c *gin.Context)

ListEvents handles GET /events Query params: tenant_id[], id[], destination_id, topic[], time[gte], time[lte], time[gt], time[lt], limit, next, prev, order_by, dir

func (*LogHandlers) RetrieveAttempt added in v0.13.0

func (h *LogHandlers) RetrieveAttempt(c *gin.Context)

RetrieveAttempt handles GET /attempts/:attempt_id and GET /:tenant_id/destinations/:destination_id/attempts/:attempt_id

func (*LogHandlers) RetrieveEvent

func (h *LogHandlers) RetrieveEvent(c *gin.Context)

RetrieveEvent handles GET /events/:event_id

type MetricsHandlers added in v0.15.0

type MetricsHandlers struct {
	// contains filtered or unexported fields
}

func NewMetricsHandlers added in v0.15.0

func NewMetricsHandlers(logger *logging.Logger, metricsStore logMetricsStore) *MetricsHandlers

func (*MetricsHandlers) MetricsAttempts added in v0.15.0

func (h *MetricsHandlers) MetricsAttempts(c *gin.Context)

MetricsAttempts handles GET /metrics/attempts

func (*MetricsHandlers) MetricsEvents added in v0.15.0

func (h *MetricsHandlers) MetricsEvents(c *gin.Context)

MetricsEvents handles GET /metrics/events

type PublishHandlers

type PublishHandlers struct {
	// contains filtered or unexported fields
}

func NewPublishHandlers

func NewPublishHandlers(
	logger *logging.Logger,
	eventHandler eventHandler,
) *PublishHandlers

func (*PublishHandlers) Ingest

func (h *PublishHandlers) Ingest(c *gin.Context)

type PublishedEvent

type PublishedEvent struct {
	ID               string            `json:"id"`
	TenantID         string            `json:"tenant_id" binding:"required"`
	DestinationID    string            `json:"destination_id"`
	Topic            string            `json:"topic"`
	EligibleForRetry *bool             `json:"eligible_for_retry"`
	Time             time.Time         `json:"time"`
	Metadata         map[string]string `json:"metadata"`
	Data             json.RawMessage   `json:"data" binding:"required"`
}

type RequestBodySanitizer

type RequestBodySanitizer struct {
	// contains filtered or unexported fields
}

RequestBodySanitizer handles sanitization of request bodies for logging

func NewRequestBodySanitizer

func NewRequestBodySanitizer(registry destregistry.Registry) *RequestBodySanitizer

NewRequestBodySanitizer creates a new sanitizer instance

func (*RequestBodySanitizer) SanitizeRequestBody

func (s *RequestBodySanitizer) SanitizeRequestBody(body io.Reader) ([]byte, error)

SanitizeRequestBody reads and sanitizes a request body for safe logging

type RetryHandlers

type RetryHandlers struct {
	// contains filtered or unexported fields
}

func NewRetryHandlers

func NewRetryHandlers(
	logger *logging.Logger,
	tenantStore tenantstore.TenantStore,
	logStore logstore.LogStore,
	deliveryPublisher deliveryPublisher,
	topicsAllowWildcards bool,
) *RetryHandlers

func (*RetryHandlers) Retry

func (h *RetryHandlers) Retry(c *gin.Context)

Retry handles POST /retry Accepts { event_id, destination_id } in body. Looks up the event, verifies the destination exists, is enabled, matches the event and has already been sent it, then publishes a manual delivery task.

type RouteDefinition

type RouteDefinition struct {
	Method        string
	Path          string
	Handler       gin.HandlerFunc
	AdminOnly     bool
	RequireTenant bool
	Middlewares   []gin.HandlerFunc
}

type RouterConfig

type RouterConfig struct {
	ServiceName          string
	APIKey               string
	JWTSecret            string
	DeploymentID         string
	Topics               []string
	TopicsAllowWildcards bool
	Registry             destregistry.Registry
	PortalConfig         portal.PortalConfig
	GinMode              string
}

type RouterDeps added in v0.13.0

type RouterDeps struct {
	TenantStore         tenantstore.TenantStore
	LogStore            logstore.LogStore
	Logger              *logging.Logger
	DeliveryPublisher   deliveryPublisher
	EventHandler        eventHandler
	Telemetry           telemetry.Telemetry
	SubscriptionEmitter SubscriptionEmitter // optional — emits tenant.subscription.updated on destination mutations
}

type SeekPagination added in v0.13.0

type SeekPagination struct {
	OrderBy string  `json:"order_by"`
	Dir     string  `json:"dir"`
	Limit   int     `json:"limit"`
	Next    *string `json:"next"`
	Prev    *string `json:"prev"`
}

SeekPagination represents cursor-based pagination metadata for list responses.

type SubscriptionEmitter added in v0.16.0

type SubscriptionEmitter interface {
	Emit(ctx context.Context, ev opevents.Event) error
}

SubscriptionEmitter emits operator events for subscription changes. Satisfied by opevents.Emitter.

type TenantHandlers

type TenantHandlers struct {
	// contains filtered or unexported fields
}

func NewTenantHandlers

func NewTenantHandlers(
	logger *logging.Logger,
	telemetry telemetry.Telemetry,
	jwtSecret string,
	deploymentID string,
	tenantStore tenantstore.TenantStore,
	topicsAllowWildcards bool,
) *TenantHandlers

func (*TenantHandlers) Delete

func (h *TenantHandlers) Delete(c *gin.Context)

func (*TenantHandlers) List added in v0.12.0

func (h *TenantHandlers) List(c *gin.Context)

func (*TenantHandlers) Retrieve

func (h *TenantHandlers) Retrieve(c *gin.Context)

func (*TenantHandlers) RetrievePortal

func (h *TenantHandlers) RetrievePortal(c *gin.Context)

func (*TenantHandlers) RetrieveToken

func (h *TenantHandlers) RetrieveToken(c *gin.Context)

func (*TenantHandlers) Upsert

func (h *TenantHandlers) Upsert(c *gin.Context)

type TenantRetriever added in v0.13.0

type TenantRetriever interface {
	RetrieveTenant(ctx context.Context, tenantID string) (*models.Tenant, error)
}

TenantRetriever is satisfied by tenantstore.TenantStore. Defined here to avoid coupling the router/middleware to the full store interface.

type TopicHandlers

type TopicHandlers struct {
	// contains filtered or unexported fields
}

func NewTopicHandlers

func NewTopicHandlers(logger *logging.Logger, topics []string) *TopicHandlers

func (*TopicHandlers) List

func (h *TopicHandlers) List(c *gin.Context)

type UpdateDestinationRequest

type UpdateDestinationRequest struct {
	Type             string          `json:"type" binding:"-"`
	Topics           models.Topics   `json:"topics" binding:"-"`
	Filter           json.RawMessage `json:"filter" binding:"-"`
	Config           json.RawMessage `json:"config" binding:"-"`
	Credentials      json.RawMessage `json:"credentials" binding:"-"`
	DeliveryMetadata json.RawMessage `json:"delivery_metadata" binding:"-"`
	Metadata         json.RawMessage `json:"metadata" binding:"-"`
	DisabledAt       json.RawMessage `json:"disabled_at" binding:"-"`
}

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL