Documentation
¶
Index ¶
- Constants
- Variables
- func OnFailureNames() []string
- func VerifySignature(secret []byte, body []byte, header string) (time.Time, error)
- type BatchPayload
- type BlockEntry
- type BlockRef
- type Client
- type Clock
- type Config
- type DeliveryError
- type DeliveryFailedError
- type DeliveryKind
- type Manifest
- type OnFailure
- type Sink
- type SinkConfig
- type UndoManifest
- type UndoPayload
- type WebhookPayload
Constants ¶
const DefaultAuthHeaderName = "Authorization"
DefaultAuthHeaderName is the header the auth value is sent under when no name is configured.
const ExitCodeDeliveryFailed = 75
ExitCodeDeliveryFailed is the process exit status for a delivery failure in OnFailureExit mode. It is EX_TEMPFAIL from sysexits: the input was fine, try again later.
const SignatureHeader = "X-Substreams-Signature"
SignatureHeader carries the HMAC signature of the request body when a signing secret is configured. Its value has the form "t=<unix seconds>,v1=<hex>", where <hex> is HMAC-SHA256 over the string "<t>.<body>" keyed with the signing secret. Receivers recompute the HMAC over the raw body bytes they read and compare with a constant-time comparison; the timestamp lets them reject replays older than they accept.
Variables ¶
var ErrInvalidDeliveryKind = errors.New("not a valid DeliveryKind")
var ErrInvalidOnFailure = fmt.Errorf("not a valid OnFailure, try [%s]", strings.Join(_OnFailureNames, ", "))
var WebhookCallsCounter = sink.Metrics.NewCounter("substreams_sink_webhook_calls", "Number of calls made to the webhook")
var WebhookLastDeliveredBlock = sink.Metrics.NewGauge("substreams_sink_webhook_last_delivered_block", "Last block number delivered to the webhook")
var WebhookSizeBytes = sink.Metrics.NewCounter("substreams_sink_webhook_bytes_sent", "Number of bytes sent via webhook")
Functions ¶
func OnFailureNames ¶ added in v1.24.0
func OnFailureNames() []string
OnFailureNames returns a list of possible string values of OnFailure.
func VerifySignature ¶ added in v1.24.0
VerifySignature checks a SignatureHeader value against body using secret. It returns the timestamp embedded in the header so callers can apply their own tolerance window. Exposed for receivers written in Go and for tests.
Types ¶
type BatchPayload ¶ added in v1.24.0
type BatchPayload struct {
Manifest Manifest `json:"manifest"`
Blocks []BlockEntry `json:"blocks"`
}
BatchPayload is sent instead of WebhookPayload when batching is on. The manifest is the same for every block so it is carried once; blocks are in ascending order. Every call in batch mode uses this shape, a batch of one included, so a receiver only ever parses one format.
func NewBatchPayload ¶ added in v1.24.0
func NewBatchPayload(moduleName string, msgType string) *BatchPayload
NewBatchPayload creates an empty batch for the given module.
func (*BatchPayload) ToJSON ¶ added in v1.24.0
func (p *BatchPayload) ToJSON() ([]byte, error)
ToJSON serializes the batch payload to JSON bytes
type BlockEntry ¶ added in v1.24.0
type BlockEntry struct {
Clock Clock `json:"clock"`
Data json.RawMessage `json:"data"`
}
BlockEntry is one block inside a BatchPayload.
type Client ¶
type Client struct {
// contains filtered or unexported fields
}
Client holds the HTTP client and retry configuration for webhook calls. It implements exponential backoff retry logic for transient failures while avoiding retries for permanent client errors (4xx status codes).
func (*Client) Call ¶
Call makes a webhook call with exponential backoff retry logic. It retries on network errors and server errors (5xx), but not on client errors (4xx). The function respects context cancellation and implements proper error classification. When every attempt fails the returned error is a *DeliveryError.
type Clock ¶
type Clock struct {
Timestamp string `json:"timestamp"`
Number uint64 `json:"number"`
ID string `json:"id"`
}
Clock represents the clock information in the webhook payload
type Config ¶
type Config struct {
Timeout time.Duration // HTTP request timeout for individual calls
MaxRetries int // Maximum number of retry attempts for transient failures (-1 for infinite retries)
MaxInterval time.Duration // Maximum interval between retries (exponential backoff cap)
// AuthHeaderName and AuthHeaderValue are sent verbatim on every request
// when AuthHeaderValue is not empty. AuthHeaderName defaults to
// DefaultAuthHeaderName.
AuthHeaderName string
AuthHeaderValue string
// SigningSecret, when not empty, adds a SignatureHeader to every request.
SigningSecret string
}
Config holds configuration for the webhook client
type DeliveryError ¶ added in v1.24.0
DeliveryError is returned by Client.Call once every attempt for a payload has failed. It carries what a receiver's operator needs to see: the last HTTP status (0 when no response was received), how many attempts were made, and the last underlying error.
func (*DeliveryError) Error ¶ added in v1.24.0
func (e *DeliveryError) Error() string
func (*DeliveryError) Unwrap ¶ added in v1.24.0
func (e *DeliveryError) Unwrap() error
type DeliveryFailedError ¶ added in v1.24.0
type DeliveryFailedError struct {
Delivery *DeliveryError
Kind DeliveryKind
FirstAttemptAt time.Time
}
DeliveryFailedError is returned from Sink.Run in OnFailureExit mode. The pending block stays on disk for the next start.
func (*DeliveryFailedError) Error ¶ added in v1.24.0
func (e *DeliveryFailedError) Error() string
func (*DeliveryFailedError) TerminationReason ¶ added in v1.24.0
func (e *DeliveryFailedError) TerminationReason() []byte
TerminationReason is the one-line JSON written to the termination log so an orchestrator can read why the process stopped and since when.
func (*DeliveryFailedError) Unwrap ¶ added in v1.24.0
func (e *DeliveryFailedError) Unwrap() error
type DeliveryKind ¶ added in v1.24.0
type DeliveryKind string
DeliveryKind tells a block payload from a reorg notification. ENUM(block, undo)
const ( // DeliveryKindBlock is a DeliveryKind of type block. DeliveryKindBlock DeliveryKind = "block" // DeliveryKindUndo is a DeliveryKind of type undo. DeliveryKindUndo DeliveryKind = "undo" )
func ParseDeliveryKind ¶ added in v1.24.0
func ParseDeliveryKind(name string) (DeliveryKind, error)
ParseDeliveryKind attempts to convert a string to a DeliveryKind.
func (DeliveryKind) IsValid ¶ added in v1.24.0
func (x DeliveryKind) IsValid() bool
IsValid provides a quick way to determine if the typed value is part of the allowed enumerated values
func (DeliveryKind) String ¶ added in v1.24.0
func (x DeliveryKind) String() string
String implements the Stringer interface.
type OnFailure ¶ added in v1.24.0
type OnFailure string
OnFailure selects what the sink does once every attempt to deliver a block has failed.
OnFailureSkip logs the failure, drops the block and moves on to the next one. The cursor is not advanced past the dropped block, so a later restart replays it. An undo notification is never dropped: it is retried until it goes through.
OnFailureExit keeps the block in the pending file, writes a termination reason and stops the sink with ExitCodeDeliveryFailed. The next start delivers the pending block before it opens a Substreams stream.
ENUM(skip, exit)
func ParseOnFailure ¶ added in v1.24.0
ParseOnFailure attempts to convert a string to a OnFailure.
type Sink ¶
type Sink struct {
// contains filtered or unexported fields
}
Sink represents a webhook sink that sends substream data to HTTP endpoints
type SinkConfig ¶
type SinkConfig struct {
WebhookURL string
// UndoURL receives an UndoPayload for every undo signal the sink sees.
// Empty disables the notification; the cursor still moves back so the
// blocks that replace the undone ones are delivered as usual.
UndoURL string
// StateFile holds the cursor of the last delivered block. The pending
// block lives next to it in "<StateFile>.pending". Empty disables both.
StateFile string
OnFailure OnFailure
SinkerConfig *sink.SinkerConfig
ClientConfig Config
// BatchMaxBlocks above zero switches every call to the BatchPayload shape
// and sends up to that many blocks per call. Zero sends one WebhookPayload
// per block.
BatchMaxBlocks int
// BatchMaxBytes above zero bounds the size of a batch payload: a batch is
// sent before the next block would take it past this many bytes. A block
// larger than the limit on its own is sent alone. Zero means no limit.
BatchMaxBytes int
// BatchMaxWait bounds how long a batch waits for more blocks. It is
// checked when the next block arrives. Defaults to one second.
BatchMaxWait time.Duration
// TerminationLogPath receives the reason for a delivery-failure exit. It
// is written only when the file already exists, which is the case under
// Kubernetes, so the default of /dev/termination-log is safe elsewhere.
TerminationLogPath string
Logger *zap.Logger
}
SinkConfig holds configuration for the webhook sink
type UndoManifest ¶ added in v1.24.0
type UndoManifest struct {
ModuleName string `json:"moduleName"`
}
UndoManifest names the module whose blocks are being undone.
type UndoPayload ¶ added in v1.24.0
type UndoPayload struct {
LastValidBlock BlockRef `json:"lastValidBlock"`
Manifest UndoManifest `json:"manifest"`
}
UndoPayload is sent to the undo URL when the chain reorganizes. Every block the receiver got with a number above LastValidBlock is no longer on the chain; the blocks that replace them follow as regular payloads.
func NewUndoPayload ¶ added in v1.24.0
func NewUndoPayload(moduleName string, lastValidBlock *pbsubstreams.BlockRef) *UndoPayload
NewUndoPayload creates the payload for an undo signal.
func (*UndoPayload) ToJSON ¶ added in v1.24.0
func (p *UndoPayload) ToJSON() ([]byte, error)
ToJSON serializes the undo payload to JSON bytes
type WebhookPayload ¶
type WebhookPayload struct {
Clock Clock `json:"clock"`
Manifest Manifest `json:"manifest"`
Data json.RawMessage `json:"data"`
}
WebhookPayload represents the payload structure sent to webhook endpoints
func NewWebhookPayload ¶
func NewWebhookPayload(moduleName string, clock *pbsubstreams.Clock, msgType string, data json.RawMessage) (*WebhookPayload, error)
NewWebhookPayload creates a new webhook payload with the desired format
func (*WebhookPayload) ToJSON ¶
func (p *WebhookPayload) ToJSON() ([]byte, error)
ToJSON serializes the webhook payload to JSON bytes