webhook

package
v1.25.0 Latest Latest
Warning

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

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

README

Substreams Webhook Sink Examples

This document provides examples of how to use the webhook sink with the new retry functionality.

Basic Usage

# Basic webhook call with default settings
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name

# With custom start block
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name --start-block 12345

Retry Configuration

The webhook sink now supports configurable retry behavior for handling transient failures:

Default Settings
  • Max Retries: 3 attempts (use -1 for infinite retries)
  • Timeout: 30 seconds per request
  • Max Retry Interval: 30 seconds (exponential backoff cap)
Custom Retry Settings
# Increase retry attempts for unreliable endpoints
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name \
  --webhook-max-retries 5

# Shorter timeout for faster failure detection
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name \
  --webhook-timeout 10s

# Longer maximum retry interval for heavily loaded endpoints
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name \
  --webhook-max-retry-interval 60s

# Disable retries completely
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name \
  --webhook-max-retries 0

# Enable infinite retries (retry until success or context cancellation)
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name \
  --webhook-max-retries -1
Combined Configuration
# Production-ready configuration with aggressive retries
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name \
  --webhook-max-retries 5 \
  --webhook-timeout 15s \
  --webhook-max-retry-interval 45s \
  --state-file ./production.cursor

# Production configuration with infinite retries for critical webhooks
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name \
  --webhook-max-retries -1 \
  --webhook-timeout 30s \
  --webhook-max-retry-interval 15s \
  --state-file ./production.cursor

Retry Behavior

What Gets Retried
  • Network errors (connection failures, timeouts)
  • Server errors (5xx HTTP status codes)
  • Temporary service unavailability
What Doesn't Get Retried
  • Client errors (4xx HTTP status codes like 400, 401, 403, 404)
  • Request creation failures (invalid URLs)
  • Context cancellation
Retry Modes
  • Limited retries (default): Retry up to --webhook-max-retries times
  • No retries (--webhook-max-retries 0): Fail immediately on first error
  • Infinite retries (--webhook-max-retries -1): Retry indefinitely until success or context cancellation
Exponential Backoff

The retry mechanism uses exponential backoff with jitter:

  • First retry: ~1 second
  • Second retry: ~2 seconds
  • Third retry: ~4 seconds
  • Maximum interval is capped by --webhook-max-retry-interval
Infinite Retry Behavior

When --webhook-max-retries is set to -1:

  • Webhook calls will retry indefinitely for transient failures
  • Only context cancellation or permanent errors (4xx) will stop retries
  • Useful for critical webhooks that must eventually succeed
  • Consider setting appropriate --webhook-timeout and --webhook-max-retry-interval values
  • Monitor logs for excessive retry patterns that might indicate service issues

State Management

# Custom state file location
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name \
  --state-file ./custom/path/state.cursor

# Disable state persistence (start from scratch each time)
substreams sink webhook https://api.example.com/webhook manifest.yaml module_name \
  --state-file ""

Error Handling

The webhook sink continues processing even if webhook calls fail after all retries. This ensures that:

  • The substreams processing doesn't get blocked by webhook failures
  • Block processing continues uninterrupted
  • Cursor state is still saved for successful blocks

Monitoring and Logging

The webhook sink provides detailed logging for troubleshooting:

INFO calling webhook block=12345 url=https://api.example.com/webhook
WARN webhook call failed for block 12345: webhook returned server error status 503 for block 12345
INFO calling webhook block=12345 url=https://api.example.com/webhook (retry 1/3)
INFO webhook call successful block=12345 url=https://api.example.com/webhook

Best Practices

  1. Set appropriate timeouts: Use shorter timeouts (5-15s) for real-time processing
  2. Configure retries based on endpoint reliability: More retries for less reliable services
  3. Use infinite retries sparingly: Only for critical webhooks where eventual delivery is essential
  4. Monitor webhook endpoint performance: Adjust settings based on observed behavior
  5. Use state files: Always specify a state file for production deployments
  6. Handle failures gracefully: Implement idempotent webhook handlers that can handle duplicate calls
  7. Context management: When using infinite retries, ensure proper context cancellation for shutdown
Retry Strategy Guidelines

For real-time applications:

--webhook-max-retries 2
--webhook-timeout 5s
--webhook-max-retry-interval 10s

For critical data delivery:

--webhook-max-retries -1
--webhook-timeout 30s
--webhook-max-retry-interval 60s

For testing/development:

--webhook-max-retries 0
--webhook-timeout 10s

Webhook Payload Format

The webhook receives JSON payloads in the following format:

{
  "clock": {
    "timestamp": "2024-02-12T22:23:51.000Z",
    "number": 53448530,
    "id": "f843bc26cea0cbd50b09699546a8a97de6a1727646c17a857c5d8d868fc26142"
  },
  "manifest": {
    "moduleName": "module_name",
    "type": "sf.substreams.ethereum.v1.Events"
  },
  "data": {
    // Your module's output data
  }
}

It is loosely based on the format from https://github.com/pinax-network/substreams-sink-webhook

Payload Structure
  • clock: Contains blockchain timing and identification information
    • timestamp: Block timestamp in RFC3339 format
    • number: Block number
    • id: Block hash/ID
  • manifest: Contains module metadata
    • moduleName: Name of the substreams module that generated this data
    • type: Module output type (automatically strips type.googleapis.com/ prefix)
  • data: Contains the actual output from your substreams module
Example with Real Data
{
  "clock": {
    "timestamp": "2024-02-12T22:23:51.000Z",
    "number": 53448530,
    "id": "f843bc26cea0cbd50b09699546a8a97de6a1727646c17a857c5d8d868fc26142"
  },
  "manifest": {
    "moduleName": "filtered_events",
    "type": "sf.substreams.ethereum.v1.Events"
  },
  "data": {
    "events": [
      {
        "address": "0x1234567890abcdef",
        "topics": ["0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef"],
        "data": "0x000000000000000000000000000000000000000000000001158e460913d00000"
      }
    ]
  }
}

"module": "module_name", "block": 12345, "timestamp": "2023-01-01T00:00:00Z", "type": "type.googleapis.com/sf.substreams.ethereum.v1.Events", "payload": { // Your module's output data } }


Make sure your webhook endpoint can handle:
- POST requests with `Content-Type: application/json`
- Potential duplicate calls (implement idempotency)
- Proper HTTP status code responses (2xx for success, 4xx for permanent errors, 5xx for retryable errors)

### Monitoring via Prometheus

* Prometheus metrics are available at `http://localhost:9102` by default. See the --prometheus-addr flag for more details.

Documentation

Index

Constants

View Source
const DefaultAuthHeaderName = "Authorization"

DefaultAuthHeaderName is the header the auth value is sent under when no name is configured.

View Source
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.

View Source
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

View Source
var ErrInvalidDeliveryKind = errors.New("not a valid DeliveryKind")
View Source
var ErrInvalidOnFailure = fmt.Errorf("not a valid OnFailure, try [%s]", strings.Join(_OnFailureNames, ", "))
View Source
var WebhookCallsCounter = sink.Metrics.NewCounter("substreams_sink_webhook_calls", "Number of calls made to the webhook")
View Source
var WebhookLastDeliveredBlock = sink.Metrics.NewGauge("substreams_sink_webhook_last_delivered_block", "Last block number delivered to the webhook")
View Source
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

func VerifySignature(secret []byte, body []byte, header string) (time.Time, error)

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 BlockRef added in v1.24.0

type BlockRef struct {
	Number uint64 `json:"number"`
	ID     string `json:"id"`
}

BlockRef identifies a block in an undo payload.

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 NewClient

func NewClient(config Config, logger *zap.Logger) *Client

NewClient creates a new webhook client with retry configuration.

func (*Client) Call

func (c *Client) Call(ctx context.Context, url string, payload []byte, blockNumber uint64) error

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.

func (*Client) Sign added in v1.24.0

func (c *Client) Sign(body []byte, at time.Time) string

Sign computes the SignatureHeader value for body at the given time. It returns an empty string when no signing secret is configured.

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

type DeliveryError struct {
	URL         string
	BlockNumber uint64
	StatusCode  int
	Attempts    int
	Err         error
}

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 Manifest

type Manifest struct {
	ModuleName string `json:"moduleName"`
	Type       string `json:"type"`
}

Manifest represents the manifest information in the webhook payload

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)

const (
	// OnFailureSkip is a OnFailure of type skip.
	OnFailureSkip OnFailure = "skip"
	// OnFailureExit is a OnFailure of type exit.
	OnFailureExit OnFailure = "exit"
)

func ParseOnFailure added in v1.24.0

func ParseOnFailure(name string) (OnFailure, error)

ParseOnFailure attempts to convert a string to a OnFailure.

func (OnFailure) IsValid added in v1.24.0

func (x OnFailure) IsValid() bool

IsValid provides a quick way to determine if the typed value is part of the allowed enumerated values

func (OnFailure) String added in v1.24.0

func (x OnFailure) String() string

String implements the Stringer interface.

type Sink

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

Sink represents a webhook sink that sends substream data to HTTP endpoints

func NewSink

func NewSink(config SinkConfig) (*Sink, error)

NewSink creates a new webhook sink

func (*Sink) PrintStats

func (s *Sink) PrintStats()

PrintStats prints final statistics

func (*Sink) Run

func (s *Sink) Run(ctx context.Context) error

Run delivers the pending block left by a previous run, if any, then streams from the cursor. A Substreams stream is never opened while a pending block is undelivered, so retrying against a dead endpoint costs no egress.

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

Jump to

Keyboard shortcuts

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