Documentation
¶
Index ¶
- func ReceiveFromStream[T any](ctx context.Context, stream grpc.ServerStreamingClient[T]) (<-chan *T, <-chan error)
- type ActiveContractStore
- func (s *ActiveContractStore) Get(address contracts.InstanceAddress) (*apiv2.ActiveContract, bool)
- func (s *ActiveContractStore) GetByContractId(contractId types.CONTRACT_ID) (*apiv2.ActiveContract, bool)
- func (s *ActiveContractStore) GetByTemplateId(party types.PARTY, templateId contracts.TemplateID) ([]*apiv2.ActiveContract, bool)
- func (s *ActiveContractStore) RegisterTemplates(templates ...RegisteredTemplate)
- func (s *ActiveContractStore) Run(ctx context.Context, streamConfig StreamConfig, opts ...RunOption) error
- type ActiveContractStoreInterface
- type BackfilledStream
- type FiltersByParty
- type InstrumentHoldingStore
- func (s *InstrumentHoldingStore) GetHolding(party types.PARTY, instrumentId splice_api_token_holding_v1.InstrumentId) ([]*apiv2.ActiveContract, bool)
- func (s *InstrumentHoldingStore) ListHoldings(instrumentId splice_api_token_holding_v1.InstrumentId) ([]*apiv2.ActiveContract, error)
- func (s *InstrumentHoldingStore) RegisterParty(parties ...string)
- func (s *InstrumentHoldingStore) Run(ctx context.Context, streamConfig StreamConfig, opts ...RunOption) error
- type InstrumentHoldingStoreInterface
- type Metrics
- type RegisteredTemplate
- type RunOption
- type StreamConfig
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 (s *ActiveContractStore) Get(address contracts.InstanceAddress) (*apiv2.ActiveContract, bool)
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 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 (s *InstrumentHoldingStore) GetHolding(party types.PARTY, instrumentId splice_api_token_holding_v1.InstrumentId) ([]*apiv2.ActiveContract, bool)
func (*InstrumentHoldingStore) ListHoldings ¶
func (s *InstrumentHoldingStore) ListHoldings(instrumentId splice_api_token_holding_v1.InstrumentId) ([]*apiv2.ActiveContract, error)
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