events

package
v0.1.27-dev3 Latest Latest
Warning

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

Go to latest
Published: Jul 15, 2025 License: MIT Imports: 9 Imported by: 0

Documentation

Overview

Package events provides a generic client for subscribing to onchain events via an EventsQueryClient and transforming the received events into the type defined by the EventsReplayClient's generic type parameter.

The EventsQueryClient emits ReplayObservables which are of the type defined by the EventsReplayClient's generic type parameter.

The usage of of ReplayObservables allows the EventsReplayClient to be always provide the latest event data to the caller, even if the connection to the EventsQueryClient is lost and re-established, without the caller having to re-subscribe to the EventsQueryClient.

Index

Examples

Constants

This section is empty.

Variables

View Source
var (
	ErrEventsDial           = sdkerrors.Register(codespace, 1, "dialing for connection failed")
	ErrEventsConnClosed     = sdkerrors.Register(codespace, 2, "connection closed")
	ErrEventsSubscribe      = sdkerrors.Register(codespace, 3, "failed to subscribe to events")
	ErrEventsUnmarshalEvent = sdkerrors.Register(codespace, 4, "failed to unmarshal event bytes")
	ErrEventsConsClosed     = sdkerrors.Register(codespace, 5, "eventsqueryclient connection closed")
)

Functions

func NewEventsReplayClient

func NewEventsReplayClient[T any](
	ctx context.Context,
	deps depinject.Config,
	queryString string,
	newEventFn NewEventsFn[T],
	replayObsBufferSize int,
) (client.EventsReplayClient[T], error)

NewEventsReplayClient creates a new EventsReplayClient from the given dependencies and subscription query string.

  • It requires a decoder function to be provided which decodes event subscription result into the type defined by the EventsReplayClient's generic type parameter.
  • The replayObsBufferSize is the replay buffer size of the replay observable which is notified of new events.

Required dependencies:

  • cometClient: cometbft/rpc/client/http.HTTP
Example
package main

import (
	"context"
	"fmt"

	"cosmossdk.io/depinject"

	coretypes "github.com/cometbft/cometbft/rpc/core/types"
	"github.com/pokt-network/poktroll/pkg/client/events"
)

const (
	// Define a query string to provide to the EventsQueryClient
	// See: https://docs.cosmos.network/v0.47/learn/advanced/events#subscribing-to-events
	// And: https://docs.cosmos.network/v0.47/learn/advanced/events#default-events
	eventQueryString = "message.action='messageActionName'"
	// the amount of events we want before they are emitted
	replayObsBufferSize = 1
)

var _ EventType = (*eventType)(nil)

// Define an interface to represent the onchain event
type EventType interface {
	GetName() string // Illustrative only; arbitrary interfaces are supported.
}

// Define the event type that implements the interface
type eventType struct {
	Name string `json:"name"`
}

// See: https://pkg.go.dev/github.com/pokt-network/poktroll/pkg/client/events/#NewEventsFn
func eventTypeFactory(ctx context.Context) events.NewEventsFn[EventType] {
	// Define a decoder function that can take the raw event bytes
	// received from the EventsQueryClient and convert them into
	// the desired type for the EventsReplayClient
	return func(eventResult *coretypes.ResultEvent) (EventType, error) {
		eventMsg, ok := eventResult.Data.(eventType)
		if !ok {
			return nil, fmt.Errorf("unable to decode event data: %T", eventResult.Data)
		}

		return &eventMsg, nil
	}
}

func (e *eventType) GetName() string { return e.Name }

func main() {
	depConfig := depinject.Supply()

	// Create a context (this should be cancellable to close the EventsReplayClient)
	ctx, cancel := context.WithCancel(context.Background())

	// Create a new instance of the EventsReplayClient
	// See: https://pkg.go.dev/github.com/pokt-network/poktroll/pkg/client/events/#NewEventsReplayClient
	client, err := events.NewEventsReplayClient[EventType](
		ctx,
		depConfig,
		eventQueryString,
		eventTypeFactory(ctx),
		replayObsBufferSize,
	)
	if err != nil {
		panic(fmt.Errorf("unable to create EventsReplayClient %v", err))
	}

	// Assume events the lastest event emitted of type EventType has the name "testEvent"

	// Retrieve the latest emitted event
	lastEventType := client.LastNEvents(ctx, 1)[0]
	fmt.Printf("Last Event: '%s'\n", lastEventType.GetName())

	// Get the latest replay observable from the EventsReplayClient
	// In order to get the latest events from the sequence
	latestEventsObs := client.EventsSequence(ctx)
	// Get the latest events from the sequence
	lastEventType = latestEventsObs.Last(ctx, 1)[0]
	fmt.Printf("Last Event: '%s'\n", lastEventType.GetName())

	// Cancel the context which will call client.Close and close all
	// subscriptions and the EventsQueryClient
	cancel()
	// Output
	// Last Event: 'testEvent'
	// Last Event: 'testEvent'
}

Types

type NewEventsFn

type NewEventsFn[T any] func(*coretypes.ResultEvent) (T, error)

NewEventsFn is a function that converts a ResultEvent into a new generic type T instance

Jump to

Keyboard shortcuts

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