store

package
v1.0.0 Latest Latest
Warning

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

Go to latest
Published: Sep 16, 2026 License: MIT Imports: 17 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func ReceiveFromStream

func ReceiveFromStream[T any](ctx context.Context, stream grpc.ServerStreamingClient[T]) (<-chan *T, <-chan error)

ReceiveFromStream runs a goroutine that Recv()s from the stream and sends each response to the returned channel. The stream's Recv returns (*T, error) for proto messages; the response channel carries *T. When Recv returns an error (including io.EOF), that error is sent on the error channel and both channels are closed. The caller must not close the stream before consuming the error channel; the goroutine does not close the stream.

Types

type ActiveContractStore

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

func NewActiveContractStore

func NewActiveContractStore(
	logger zerolog.Logger,
	updateService apiv2.UpdateServiceClient,
	stateService apiv2.StateServiceClient,
	metrics Metrics,
) *ActiveContractStore

func (*ActiveContractStore) Get

func (*ActiveContractStore) GetByContractId

func (s *ActiveContractStore) GetByContractId(contractId types.CONTRACT_ID) (*apiv2.ActiveContract, bool)

func (*ActiveContractStore) GetByTemplateId

func (s *ActiveContractStore) GetByTemplateId(party types.PARTY, templateId contracts.TemplateID) ([]*apiv2.ActiveContract, bool)

GetByTemplateId returns the active contracts for the given TemplateID and party, if any exist. It accepts TemplateIds in both the #packageName and PackageId syntax.

func (*ActiveContractStore) RegisterTemplates

func (s *ActiveContractStore) RegisterTemplates(templates ...RegisteredTemplate)

RegisterTemplates registers the given templates with the ActiveContractStore. RegisterTemplates must be called before any call to Run(), calling while the store is already running will lead to undefined behavior.

func (*ActiveContractStore) Run

func (s *ActiveContractStore) Run(ctx context.Context, streamConfig StreamConfig, opts ...RunOption) error

type ActiveContractStoreInterface

type ActiveContractStoreInterface interface {
	// Get returns the active contract for the given instance address.
	Get(address contracts.InstanceAddress) (*apiv2.ActiveContract, bool)
	// GetByTemplateId returns the active contracts for the given party and template ID.
	GetByTemplateId(party types.PARTY, templateId contracts.TemplateID) ([]*apiv2.ActiveContract, bool)
	// GetByContractId returns the active contract for the given contract ID.
	GetByContractId(contractId types.CONTRACT_ID) (*apiv2.ActiveContract, bool)
	// RegisterTemplates registers templates to be tracked by the store.
	RegisterTemplates(templates ...RegisteredTemplate)
}

ActiveContractStoreInterface provides read access to active contracts and the ability to register templates.

type BackfilledStream

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

func (*BackfilledStream) Run

func (s *BackfilledStream) Run(
	ctx context.Context,
	filters FiltersByParty,
	streamConfig StreamConfig,
	onActiveContract func(ctx context.Context, activeContract *apiv2.ActiveContract) error,
	onCreatedEvent func(ctx context.Context, transaction *apiv2.Transaction, createdEvent *apiv2.CreatedEvent) error,
	onArchivedEvent func(ctx context.Context, transaction *apiv2.Transaction, archivedEvent *apiv2.ArchivedEvent) error,
	onExercisedEvent func(ctx context.Context, transaction *apiv2.Transaction, exercisedEvent *apiv2.ExercisedEvent) error,
	onFinishedBackfill func(),
) error

type FiltersByParty

type FiltersByParty map[string]*apiv2.Filters

type InstrumentHoldingStore

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

func NewInstrumentHoldingStore

func NewInstrumentHoldingStore(
	logger zerolog.Logger,
	updateService apiv2.UpdateServiceClient,
	stateService apiv2.StateServiceClient,
	metrics Metrics,
) *InstrumentHoldingStore

func (*InstrumentHoldingStore) GetHolding

func (*InstrumentHoldingStore) ListHoldings

func (*InstrumentHoldingStore) RegisterParty

func (s *InstrumentHoldingStore) RegisterParty(parties ...string)

func (*InstrumentHoldingStore) Run

func (s *InstrumentHoldingStore) Run(ctx context.Context, streamConfig StreamConfig, opts ...RunOption) error

type InstrumentHoldingStoreInterface

type InstrumentHoldingStoreInterface interface {
	// GetHolding returns the active holdings for the given party and instrument ID.
	GetHolding(party types.PARTY, instrumentId splice_api_token_holding_v1.InstrumentId) ([]*apiv2.ActiveContract, bool)
	// ListHoldings returns all active holdings for the given instrument ID.
	ListHoldings(instrumentId splice_api_token_holding_v1.InstrumentId) ([]*apiv2.ActiveContract, error)
	// RegisterParty registers parties whose holdings should be tracked by the store.
	RegisterParty(parties ...string)
}

InstrumentHoldingStoreInterface provides read access to instrument holdings and the ability to register parties.

type Metrics

type Metrics interface {
	IncrementStoreSubscriptionUptime(ctx context.Context)
	IncrementStoreUpdatesCounter(ctx context.Context)
	RecordStoreLedgerEndGauge(ctx context.Context, ledgerEnd int64)
}

Metrics captures a subset of metrics that are used for monitoring a ContractStore. // It is implemented by the global common.EDSMetricLabeler

type RegisteredTemplate

type RegisteredTemplate struct {
	TemplateID contracts.TemplateID
	PartyID    string
}

type RunOption

type RunOption func(o *runOptions)

func WithOnBackfillCompleted

func WithOnBackfillCompleted(callback func()) RunOption

type StreamConfig

type StreamConfig struct {
	MaxRetries    int           // 0 means unlimited retries
	BackoffMin    time.Duration // default 100ms
	BackoffMax    time.Duration // default 3s
	BackoffFactor float64       // default 2
}

func DefaultStreamConfig

func DefaultStreamConfig() StreamConfig

Jump to

Keyboard shortcuts

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