Documentation
¶
Index ¶
- Constants
- Variables
- func MatchFilter(filter Filter, event Event) bool
- func RetryID(eventID, destinationID string) string
- type Attempt
- type Config
- type Credentials
- type Data
- type DeliveryMetadata
- type DeliveryTask
- type DeliveryTelemetry
- type Destination
- type Event
- type EventTelemetry
- type Filter
- type LogEntry
- type MapStringString
- type Metadata
- type Tenant
- type Topics
- func (t *Topics) MarshalBinary() ([]byte, error)
- func (t *Topics) MarshalJSON() ([]byte, error)
- func (t *Topics) MatchTopic(eventTopic string, allowWildcards bool) bool
- func (t *Topics) MatchesAll() bool
- func (t Topics) Normalize() Topics
- func (t *Topics) UnmarshalBinary(data []byte) error
- func (t *Topics) UnmarshalJSON(data []byte) error
- func (t *Topics) Validate(availableTopics []string, allowWildcards bool) error
- func (t Topics) WithoutWildcardPatterns() Topics
Constants ¶
const ( AttemptStatusSuccess = "success" AttemptStatusFailed = "failed" )
Variables ¶
var ( ErrInvalidTopics = errors.New("validation failed: invalid topics") ErrInvalidTopicsFormat = errors.New("validation failed: invalid topics format") )
Functions ¶
func MatchFilter ¶ added in v0.13.0
MatchFilter checks if the given event matches the filter. Returns true if no filter is set (nil or empty) or if the event matches the filter.
Types ¶
type Attempt ¶ added in v0.13.0
type Attempt struct {
ID string `json:"id"`
TenantID string `json:"tenant_id"`
EventID string `json:"event_id"`
DestinationID string `json:"destination_id"`
DestinationType string `json:"destination_type"`
AttemptNumber int `json:"attempt_number"`
Manual bool `json:"manual"`
Status string `json:"status"`
Time time.Time `json:"time"`
Code string `json:"code"`
ResponseData map[string]interface{} `json:"response_data"`
// Provider response time in ms, measured around Publish. Nil when the
// attempt never reached the provider or predates latency tracking.
LatencyMs *int64 `json:"latency_ms"`
}
type Config ¶
type Config = MapStringString
type Credentials ¶
type Credentials = MapStringString
type Data ¶
type Data = json.RawMessage
Data holds the event payload as raw JSON bytes so that key order, whitespace and numeric precision survive every serialisation hop.
type DeliveryMetadata ¶ added in v0.8.0
type DeliveryMetadata = MapStringString
type DeliveryTask ¶ added in v0.13.0
type DeliveryTask struct {
Event Event `json:"event"`
DestinationID string `json:"destination_id"`
Attempt int `json:"attempt"`
Manual bool `json:"manual"`
Telemetry *DeliveryTelemetry `json:"telemetry,omitempty"`
}
DeliveryTask represents a task to deliver an event to a destination. This is a message type (no ID) used by: publishmq -> deliverymq, retry -> deliverymq
func NewDeliveryTask ¶ added in v0.13.0
func NewDeliveryTask(event Event, destinationID string) DeliveryTask
NewDeliveryTask creates a new DeliveryTask for an event and destination.
func NewManualDeliveryTask ¶ added in v0.13.0
func NewManualDeliveryTask(event Event, destinationID string, attemptNumber int) DeliveryTask
NewManualDeliveryTask creates a new DeliveryTask for a manual retry. attemptNumber is the 1-indexed attempt number derived from the count of prior attempts.
func (*DeliveryTask) FromMessage ¶ added in v0.13.0
func (t *DeliveryTask) FromMessage(msg *mqs.Message) error
func (*DeliveryTask) IdempotencyKey ¶ added in v0.13.0
func (t *DeliveryTask) IdempotencyKey() string
IdempotencyKey returns the key used for idempotency checks. Uses event_id:destination_id:attempt_number so that:
- Manual and auto retries with the same attempt_number are deduplicated (race protection)
- Each new attempt gets a fresh key (no need to clear on failure)
- MQ redeliveries of the same message are still deduplicated
type DeliveryTelemetry ¶ added in v0.13.0
type Destination ¶
type Destination struct {
ID string `json:"id" redis:"id"`
TenantID string `json:"tenant_id" redis:"-"`
Type string `json:"type" redis:"type"`
Topics Topics `json:"topics" redis:"-"`
Filter Filter `json:"filter,omitempty" redis:"-"`
Config Config `json:"config" redis:"-"`
Credentials Credentials `json:"credentials" redis:"-"`
DeliveryMetadata DeliveryMetadata `json:"delivery_metadata,omitempty" redis:"-"`
Metadata Metadata `json:"metadata,omitempty" redis:"-"`
CreatedAt time.Time `json:"created_at" redis:"created_at"`
UpdatedAt time.Time `json:"updated_at" redis:"updated_at"`
DisabledAt *time.Time `json:"disabled_at" redis:"disabled_at"`
}
func (*Destination) MatchEvent ¶ added in v0.10.0
func (d *Destination) MatchEvent(event Event, allowWildcards bool) bool
MatchEvent checks if the destination matches the given event. Returns true if the destination is enabled, topic matches, and filter matches. Embedded wildcard topic patterns are ignored when allowWildcards is false.
type Event ¶
type Event struct {
ID string `json:"id"`
TenantID string `json:"tenant_id"`
DestinationID string `json:"destination_id"`
MatchedDestinationIDs []string `json:"matched_destination_ids"`
Topic string `json:"topic"`
EligibleForRetry bool `json:"eligible_for_retry"`
Time time.Time `json:"time"`
Metadata Metadata `json:"metadata"`
Data Data `json:"data"`
// Telemetry data, must exist to properly trace events between publish receiver & delivery handler
Telemetry *EventTelemetry `json:"telemetry,omitempty"`
}
func (*Event) ParsedData ¶ added in v0.13.2
ParsedData unmarshals the raw JSON Data into a map[string]any. This is used by code that needs to inspect individual fields (e.g. filters, partition-key extraction) without losing the original byte representation.
type EventTelemetry ¶
type Filter ¶ added in v0.10.0
Filter represents a JSON schema filter for event matching. It uses the simplejsonmatch schema syntax for filtering events.
func (*Filter) MarshalBinary ¶ added in v0.10.0
func (*Filter) UnmarshalBinary ¶ added in v0.10.0
type LogEntry ¶ added in v0.13.0
type LogEntry struct {
Event *Event `json:"event"`
Attempt *Attempt `json:"attempt"`
Destination *Destination `json:"destination,omitempty"` // carried for alert evaluation in logmq; ignored by logstore
}
LogEntry represents a message for the log queue.
IMPORTANT: Both Event and Attempt are REQUIRED. The logstore requires both to exist for proper data consistency. The logmq consumer validates this requirement and rejects entries missing either field.
func (*LogEntry) FromMessage ¶ added in v0.13.0
type MapStringString ¶
func (*MapStringString) MarshalBinary ¶
func (m *MapStringString) MarshalBinary() ([]byte, error)
func (*MapStringString) UnmarshalBinary ¶
func (m *MapStringString) UnmarshalBinary(data []byte) error
func (*MapStringString) UnmarshalJSON ¶
func (m *MapStringString) UnmarshalJSON(data []byte) error
type Metadata ¶
type Metadata = MapStringString
type Tenant ¶
type Tenant struct {
ID string `json:"id" redis:"id"`
DestinationsCount int `json:"destinations_count" redis:"-"`
Topics []string `json:"topics" redis:"-"`
Metadata Metadata `json:"metadata,omitempty" redis:"-"`
CreatedAt time.Time `json:"created_at" redis:"created_at"`
UpdatedAt time.Time `json:"updated_at" redis:"updated_at"`
}
type Topics ¶
type Topics []string
func TopicsFromString ¶
func (*Topics) MarshalBinary ¶
func (*Topics) MarshalJSON ¶
func (*Topics) MatchTopic ¶ added in v0.9.0
MatchTopic reports whether an event topic matches this subscription. The standalone "*" keeps its subscribe-to-all meaning when allowWildcards is false, while embedded wildcard patterns such as "user.*" are ignored.
func (*Topics) MatchesAll ¶
func (Topics) Normalize ¶ added in v1.1.0
Normalize returns a topic set with redundant entries removed, preserving first-seen order. It performs two reductions:
- exact duplicates are collapsed to their first occurrence (["user.created","user.created"] -> ["user.created"]).
- an entry is folded away when a sibling wildcard pattern covers it and the entry does not itself cover that sibling (["user.*","user.created"] -> ["user.*"]).
Two mutually-non-covering patterns are both kept (["*.created","user.*"] is unchanged), and ["*"] is returned as-is. Normalization never changes MatchTopic results: it only drops entries that are already covered by a retained entry.
func (*Topics) UnmarshalBinary ¶
func (*Topics) UnmarshalJSON ¶
func (Topics) WithoutWildcardPatterns ¶ added in v1.4.0
WithoutWildcardPatterns returns the topics with embedded wildcard patterns such as "user.*" removed. The standalone "*" is kept.