workload

package
v1.0.5 Latest Latest
Warning

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

Go to latest
Published: Sep 3, 2026 License: Apache-2.0 Imports: 39 Imported by: 0

Documentation

Index

Constants

View Source
const (
	PolicySchemeMSP         = "MSP"
	PolicySchemeDefault     = PolicySchemeMSP
	PolicySchemeUnspecified = ""
)

Defines Policy.Scheme.

View Source
const DefaultGeneratedNamespaceID = "0"

DefaultGeneratedNamespaceID for now we're only generating transactions for a single namespace.

Variables

This section is empty.

Functions

func Clip

func Clip[T cmp.Ordered](value, low, high T) T

Clip returns the given value, clipped between the given boundaries.

func CreateConfigTxFromConfigBlock added in v0.1.8

func CreateConfigTxFromConfigBlock(block *common.Block) (*servicepb.LoadGenTx, error)

CreateConfigTxFromConfigBlock creates a config TX.

func CreateLoadGenNamespacesTX added in v0.1.6

func CreateLoadGenNamespacesTX(policy *PolicyProfile) (*servicepb.LoadGenTx, error)

CreateLoadGenNamespacesTX creating the transaction containing the requested namespaces into the MetaNamespace.

func CreateNamespacesTxFromEndorser added in v0.1.8

func CreateNamespacesTxFromEndorser(
	endorser *TxEndorser, metaNamespaceVersion uint64, includeNS ...string,
) (*applicationpb.Tx, error)

CreateNamespacesTxFromEndorser creating the transaction containing the requested namespaces into the MetaNamespace.

func CreateOrExtendConfigBlockWithCrypto added in v0.1.9

func CreateOrExtendConfigBlockWithCrypto(policy *PolicyProfile) (*common.Block, error)

CreateOrExtendConfigBlockWithCrypto generates or extends the crypto artifacts for a policy. This will generate a new config block, or overwrite the existing config block if it already exists.

func CreateOrLoadConfigBlockWithCrypto added in v0.1.9

func CreateOrLoadConfigBlockWithCrypto(policy *PolicyProfile) (*common.Block, error)

CreateOrLoadConfigBlockWithCrypto generates the crypto artifacts for a policy. If it is already generated, it will load the config block from the file system.

func GenerateTransactions added in v0.1.4

func GenerateTransactions(tb testing.TB, p *Profile, count int) []*servicepb.LoadGenTx

GenerateTransactions is used for benchmarking and tests.

func Map

func Map[T, K any](arr []T, mapper func(index int, value T) K) []K

Map maps an array to a new array of the same size using a transformation function.

func MapToCoordinatorBatch added in v0.1.6

func MapToCoordinatorBatch(blockNum uint64, txs []*servicepb.LoadGenTx) *servicepb.CoordinatorBatch

MapToCoordinatorBatch creates a Coordinator batch.

func MapToEnvelopeBatch added in v0.1.6

func MapToEnvelopeBatch(_ uint64, txs []*servicepb.LoadGenTx) []*common.Envelope

MapToEnvelopeBatch creates a batch of Fabric's Orderer input envelopes.

func MapToLoadGenBatch added in v0.1.6

func MapToLoadGenBatch(_ uint64, txs []*servicepb.LoadGenTx) *servicepb.LoadGenBatch

MapToLoadGenBatch creates a load-gen batch.

func MapToOrdererBlock added in v0.1.6

func MapToOrdererBlock(blockNum uint64, txs []*servicepb.LoadGenTx) *common.Block

MapToOrdererBlock creates a Fabric's Orderer output block.

func MapToVcBatch added in v0.1.6

func MapToVcBatch(blockNum uint64, txs []*servicepb.LoadGenTx) *servicepb.VcBatch

MapToVcBatch creates a VC batch.

func MapToVerifierBatch added in v0.1.6

func MapToVerifierBatch(blockNum uint64, txs []*servicepb.LoadGenTx) *servicepb.VerifierBatch

MapToVerifierBatch creates a Verifier batch.

func NewPolicyEndorserFromMsp added in v0.1.9

func NewPolicyEndorserFromMsp(tb testing.TB, artifactsPath string) *testsig.NsEndorser

NewPolicyEndorserFromMsp creates an MSP-based endorser and namespace policy from the peer organization crypto artifacts under artifactsPath.

func SumInt

func SumInt[T constraints.Integer](arr []T) int64

SumInt returns the sum of an array of integers.

Types

type BlockProfile

type BlockProfile struct {
	MaxSize       uint64        `mapstructure:"max-size" yaml:"max-size"`
	MinSize       uint64        `mapstructure:"min-size" yaml:"min-size"`
	PreferredRate time.Duration `mapstructure:"preferred-rate" yaml:"preferred-rate"`
}

BlockProfile describes generate block characteristics (when applying load to the VC or Verifier, blocks are translated to batches). The generated block size is aimed to be MaxSize, if the generated TXs rate is sufficient. If the generated TXs rate is too low, the block size might be less than MaxSize, but at least MinSize. In such case, the block is generated at a preferred rate of PreferredRate. Blocks wait up to PreferredRate (default: 1 second) before submission. If a full block is not ready by then, a partial block is submitted if it meets MinSize (default: 1); otherwise, the system waits until MinSize is available. If the MaxSize is less than or equal to MinSize, PreferredRate is ignored.

type ConsumeParameters added in v0.1.8

type ConsumeParameters struct {
	RequestedItems uint64
	MinItems       uint64
	SoftTimeout    time.Duration
}

ConsumeParameters describes the consume request parameters.

type ConsumerRateController added in v0.1.8

type ConsumerRateController[T any] struct {
	// contains filtered or unexported fields
}

ConsumerRateController controls the consumption rate from a inputQueue while ensuring:

  • smooth pacing,
  • bounded waiting,
  • and useful batch sizes.

It is designed for progress-driven workloads (streaming, chunked processing, UI updates, network sends), where callers prefer some progress within a deadline rather than waiting indefinitely for a full request. The rate controller is based on a token bucket with two additional guarantees: - a maximum wait time for partial delivery, - a minimum batch size for efficiency.

Unlike standard rate limiters, it enforces both a latency deadline and a minimum useful batch, making it suitable for block- and batch-driven systems rather than per-request throttling.

func NewConsumerRateController added in v0.1.8

func NewConsumerRateController[T any](rate uint64, inputQueue <-chan []T) *ConsumerRateController[T]

NewConsumerRateController create a new rate controller.

func (*ConsumerRateController[T]) Consume added in v0.1.8

func (l *ConsumerRateController[T]) Consume(ctx context.Context, p ConsumeParameters) []T

Consume requests up to RequestedItems from the producer. The call may block for some time, and may return fewer items than requested depending on MinItems and softTimeout.

The rate controller chooses one of three outcomes:

  • Full request (fast path): If RequestedItems can be made available within SoftTimeout, the call waits just long enough and returns exactly the RequestedItems.
  • Partial batch (deadline path): If waiting SoftTimeout yields a meaningful amount of items (>= MinItems), the call waits SoftTimeout and returns as many items as are available at that time (less than requested).
  • Minimum batch (efficiency path): If even after SoftTimeout fewer than MinItems would be available, the call waits longer until MinItems is ready, then returns exactly MinItems.

Notes:

  • If SoftTimeout is not positive, or MinItems is not less than RequestedItems, the call waits indefinitely for RequestedItems.
  • If the context is canceled before any items are available, the call returns nil.
  • If RequestedItems is less than 1, the call returns nil immediately.

func (*ConsumerRateController[T]) InstantiateWorker added in v0.1.8

func (l *ConsumerRateController[T]) InstantiateWorker() *ConsumerRateController[T]

InstantiateWorker creates a new rate controller instance that shares the same state and inputQueue. The worker instance maintains its own outputBuffer.

func (*ConsumerRateController[T]) Next added in v0.1.8

func (l *ConsumerRateController[T]) Next(ctx context.Context) T

Next consumes a single item from the inputQueue.

func (*ConsumerRateController[T]) Rate added in v0.1.8

func (l *ConsumerRateController[T]) Rate() uint64

Rate returns the current rate limit in tokens per second.

func (*ConsumerRateController[T]) SetRate added in v0.1.8

func (l *ConsumerRateController[T]) SetRate(rate uint64)

SetRate updates the rate limit to the specified value in tokens per second. It resets the timeline, so any accumulated tokens will be erased.

type IndependentTxGenerator

type IndependentTxGenerator struct {
	TransactionProfile

	TxBuilder *TxBuilder
	// contains filtered or unexported fields
}

IndependentTxGenerator builds transactions from a txRandomProcess. The generator only assembles: it embeds the TransactionProfile (the fixed layout set sizes and value sizes), lays out the namespace, optionally stamps a dummy endorsement in place of a real signature, and hands off to the TxBuilder. The random content — keys, values, nonce/TX ID, and metadata — comes from the process and is a pure function of the transaction index, so a transaction is reproducible by index. Generation is deterministic and side-effect free: fetching reads' committed versions is a separate, context-bound stage the stream runs between buildBatch and signBatch.

type KeyPath

type KeyPath struct {
	SigningKey      string `mapstructure:"signing-key" yaml:"signing-key"`
	VerificationKey string `mapstructure:"verification-key" yaml:"verification-key"`
	SignCertificate string `mapstructure:"sign-certificate" yaml:"sign-certificate"`
}

KeyPath describes how to find/generate the signature keys.

type KeyStats added in v1.0.5

type KeyStats struct {
	KeyFrontier         uint64 // keys introduced so far, including ones read or written but never committed
	ReferencedReadKeys  uint64 // read-only slots that reused an existing key
	ReferencedWriteKeys uint64 // write slots that reused an existing key
}

KeyStats is a cumulative, monotone snapshot of the workload's key-generation counts, so its fields can back counter metrics directly.

type Policy

type Policy struct {
	Scheme        signature.Scheme              `mapstructure:"scheme" yaml:"scheme"`
	Seed          int64                         `mapstructure:"seed" yaml:"seed"`
	KeyPath       *KeyPath                      `mapstructure:"key-path" yaml:"key-path"`
	MSPIdentities []*ordererdial.IdentityConfig `mapstructure:"msp-identities" yaml:"msp-identities"`
}

Policy describes how to sign/verify a TX. It supports a signing with a raw signing key, or via a local MSP. Scheme can be a valid signature schemes (NONE, ECDSA, BLS, or EDDSA) or MSP to indicate using a local MSP. When Scheme is not MSP, we generate a key using the given Seed, or loading one if KeyPath is given, ignoring MSPIdentities. When Scheme is MSP, we load the signing identities from MSPIdentities, ignoring Seed and KeyPath. In such case, we use the default rule, which state that all peer organization should sign. If MSPIdentities is not provided, we load the signing identities from ArtifactsPath.

type PolicyProfile

type PolicyProfile struct {
	// NamespacePolicies specifies the namespace policies.
	NamespacePolicies map[string]*Policy `mapstructure:"namespace-policies" yaml:"namespace-policies"`

	// OrdererEndpoints may specify the endpoints to add to the config block.
	// If this field is empty, no endpoints will be configured.
	// If ConfigBlockPath is specified, this value is ignored.
	OrdererEndpoints []*commontypes.OrdererEndpoint `mapstructure:"orderer-endpoints" yaml:"orderer-endpoints"`

	// ArtifactsPath may specify the path to the artifacts generated by CreateOrExtendConfigBlockWithCrypto().
	// If this field is empty, the artifacts will be generated into a temporary folder.
	// If this path does not exist, or it is empty, the artifacts will be generated into it.
	// The config block will be fetched from ArtifactsPath.
	ArtifactsPath string `mapstructure:"artifacts-path" yaml:"artifacts-path"`

	// ChannelID and Identity are used to create the TX envelop.
	ChannelID string                      `mapstructure:"channel-id"`
	Identity  *ordererdial.IdentityConfig `mapstructure:"identity"`

	// PeerOrganizationCount may specify the number of peer organizations to generate if the ArtifactsPath
	// is not provided.
	PeerOrganizationCount uint32 `mapstructure:"peer-organization-count"`
}

PolicyProfile holds the policy information for the load generation.

func (*PolicyProfile) Validate added in v0.1.9

func (p *PolicyProfile) Validate() error

Validate checks that the PolicyProfile does not contain invalid entries. System namespace "_config" must not be provided explicitly as it is not a real namespace. System namespace "_meta" can be derived from the artifacts' path when given. But it can be provided explicitly if desired. If provided explicitly, it must use a MSP rule.

type Probability

type Probability = float64

Probability is a float in the closed interval [0,1].

type Profile

type Profile struct {
	Block       BlockProfile       `mapstructure:"block" yaml:"block"`
	Transaction TransactionProfile `mapstructure:"transaction" yaml:"transaction"`
	Policy      PolicyProfile      `mapstructure:"policy" yaml:"policy"`

	// Seed is the single PRF root for the whole workload. Every generated item (keys, values, nonce,
	// metadata, and the new-vs-existing layout) is a pure function of this seed and the item's global
	// transaction index, so the same Seed reproduces the same items.
	Seed int64 `mapstructure:"seed" yaml:"seed"`

	// Workers is the number of parallel producers. They share one global transaction-index counter and
	// the same Seed, so the multiset of generated transactions is independent of the worker count —
	// workers are pure parallelism. The count therefore does not need to be preserved between runs to
	// reproduce items.
	Workers uint32 `mapstructure:"workers" yaml:"workers"`
}

Profile describes the generated workload characteristics. It only contains parameters that deterministically affect the generated items. The items order, however, might be affected by other parameters.

func DefaultProfile

func DefaultProfile(workers uint32) *Profile

DefaultProfile is used for testing and benchmarking.

type StreamOptions

type StreamOptions struct {
	// GenBatch impacts the rate by batching generated items before inserting then the channel.
	// This helps overcome the inherit rate limitation of Go channels.
	GenBatch uint32 `mapstructure:"gen-batch" yaml:"gen-batch"`
	// BuffersSize impact the rate by masking fluctuation in performance.
	BuffersSize int `mapstructure:"buffers-size" yaml:"buffers-size"`
	// RateLimit directly impacts the rate by limiting it.
	// TXs are released at RateLimit (default: unlimited).
	RateLimit uint64 `mapstructure:"rate-limit" yaml:"rate-limit"`
}

StreamOptions allows adjustment to the stream rate. It only contains parameters that do not affect the produced items. However, these parameters might affect the order of the items.

type StreamWithSetup

type StreamWithSetup struct {
	WorkloadSetupTXs channel.Reader[*servicepb.LoadGenTx]
	TxStream         *TxStream
}

StreamWithSetup implements the TxStream interface.

func (*StreamWithSetup) MakeTxGenerator

func (c *StreamWithSetup) MakeTxGenerator() *TxGeneratorWithSetup

MakeTxGenerator instantiate clientTxGenerator.

type TransactionProfile

type TransactionProfile struct {
	// The byte sizes of the generated key/values/metadata (size=0 => nil), ordered key, value, metadata.
	KeySize             uint32 `mapstructure:"key-size" yaml:"key-size"`
	ReadWriteValueSize  uint32 `mapstructure:"read-write-value-size" yaml:"read-write-value-size"`
	BlindWriteValueSize uint32 `mapstructure:"blind-write-value-size" yaml:"blind-write-value-size"`
	MetadataSize        uint32 `mapstructure:"metadata-size" yaml:"metadata-size"`

	// The number of keys to generate (read ver=nil)
	ReadOnlyCount uint32 `mapstructure:"read-only-count" yaml:"read-only-count"`
	// The number of keys to generate (read ver=nil/write)
	ReadWriteCount uint32 `mapstructure:"read-write-count" yaml:"read-write-count"`
	// The number of keys to generate (write)
	BlindWriteCount uint32 `mapstructure:"write-count" yaml:"write-count"`

	// KeyBackrefRate is the average number of BACKWARD REFERENCES (reused keys) per transaction: slots that
	// reference a key created by an earlier transaction instead of creating a fresh one, so keys are reused
	// and the coordinator sees commit-time contention. The remaining slots create fresh keys, so the
	// average number of NEW keys per transaction is the total slot count minus this rate.
	// 0 (default) means no references — every slot gets a fresh unique key, the historical contention-free
	// workload. Below 1 it acts as a probability: the fraction of transactions that get a single reference
	// (e.g. 0.3 ≈ 30% of transactions reference one existing key). 1 or more is roughly a fixed count of
	// references per transaction, and a fractional part adds that sub-1 probability on top (e.g. 2.5 ≈ two
	// or three references). It must not exceed the total slot count (read-only + read-write + write); at the
	// maximum every slot is a reference (a fully static working set, no new keys). References fill slots
	// after the fresh keys, in layout order: read-write, then blind-write, then read-only.
	KeyBackrefRate float64 `mapstructure:"key-backref-rate" yaml:"key-backref-rate" validate:"gte=0"`
	// TxReferenceGap and KeyLookbackWindow shape where backward references point; both are optional, both
	// default to 0, and both are irrelevant when key-backref-rate is 0. A reference is drawn from the
	// KeyLookbackWindow newest keys that existed TxReferenceGap transactions ago.
	//
	// TxReferenceGap is how far back, in transactions, references reach. 0 (default) draws the newest keys,
	// which may still be in flight — so conflicting transactions can land in the same block (a live
	// dependency); larger values draw older, already-committed keys.
	TxReferenceGap uint64 `mapstructure:"tx-reference-gap" yaml:"tx-reference-gap"`
	// KeyLookbackWindow is how many of the newest keys a reference is drawn from. 0 (default) means no
	// window: references step straight back from the gap position (the newest key, then the one before it,
	// and so on) — a fixed, most-contended pattern. A positive value spreads references across that many
	// keys (larger = less contention); it has no minimum, and when a transaction needs more distinct
	// references than the window holds, the surplus simply steps back beyond the window.
	KeyLookbackWindow uint64 `mapstructure:"key-lookback-window" yaml:"key-lookback-window"`

	// InvalidSignatures is the probability [0,1] that a transaction is stamped with a bad signature
	// (default: 0). The decision is derived deterministically from the transaction index.
	InvalidSignatures Probability `mapstructure:"invalid-signatures" yaml:"invalid-signatures" validate:"gte=0,lte=1"`
}

TransactionProfile describes generate TX characteristics.

func (*TransactionProfile) Validate added in v1.0.5

func (p *TransactionProfile) Validate() error

Validate checks the split configuration. key-backref-rate must not exceed the total slot count, so the per-transaction reference count never exceeds the transaction's slots; 0 (the default) keeps the historical fresh-key workload. tx-reference-gap and key-lookback-window are always optional and have no minimum — a zero or small window never fails, as references step back beyond it as needed.

type TxBuilder added in v0.1.6

type TxBuilder struct {
	ChannelID   string
	EnvSigner   identity.Signer
	TxEndorser  *TxEndorser
	EnvCreator  []byte
	NonceSource io.Reader
}

TxBuilder is a convenient way to create an enveloped TX.

func NewTxBuilderFromPolicy added in v0.1.6

func NewTxBuilderFromPolicy(policy *PolicyProfile, nonceSource io.Reader) (*TxBuilder, error)

NewTxBuilderFromPolicy instantiate a TxBuilder from a given policy profile.

func (*TxBuilder) MakeTx added in v0.1.6

func (txb *TxBuilder) MakeTx(tx *applicationpb.Tx) *servicepb.LoadGenTx

MakeTx makes an enveloped TX with the builder's properties.

func (*TxBuilder) MakeTxWithID added in v0.1.6

func (txb *TxBuilder) MakeTxWithID(txID string, tx *applicationpb.Tx) *servicepb.LoadGenTx

MakeTxWithID makes an enveloped TX with the builder's properties. It uses the given TX-ID instead of generating a valid one.

func (*TxBuilder) MakeTxWithNonce added in v1.0.5

func (txb *TxBuilder) MakeTxWithNonce(nonce []byte, tx *applicationpb.Tx) *servicepb.LoadGenTx

MakeTxWithNonce makes an enveloped TX using the given nonce instead of reading one from the builder's nonce source. The derived TX-ID (ComputeTxID(nonce, creator)) is therefore a deterministic function of the nonce, letting the caller make it addressable by index.

type TxCounter added in v1.0.5

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

TxCounter is the shared transaction-index counter. All workers draw their index ranges from this one counter, which is what makes the generated transaction multiset independent of the worker count. It also holds the profile so KeyStats can report the key-generation counts without reaching into the stream.

func NewTxCounter added in v1.0.5

func NewTxCounter(profile TransactionProfile) *TxCounter

NewTxCounter creates the shared transaction counter.

func (*TxCounter) KeyStats added in v1.0.5

func (c *TxCounter) KeyStats() KeyStats

KeyStats computes the current key-generation counts from the counter value.

type TxEndorser added in v0.1.8

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

TxEndorser supports endorsing a TX.

func NewTxEndorser added in v0.1.8

func NewTxEndorser(policy *PolicyProfile) *TxEndorser

NewTxEndorser creates a new TxEndorser given a workload profile.

func (*TxEndorser) Endorse added in v0.1.8

func (e *TxEndorser) Endorse(txID string, tx *applicationpb.Tx)

Endorse a TX.

func (*TxEndorser) VerificationPolicies added in v0.1.8

func (e *TxEndorser) VerificationPolicies() map[string]*applicationpb.NamespacePolicy

VerificationPolicies returns the verification policies.

type TxGeneratorWithSetup

type TxGeneratorWithSetup struct {
	WorkloadSetupTXs channel.Reader[*servicepb.LoadGenTx]
	TxGen            *ConsumerRateController[*servicepb.LoadGenTx]
}

TxGeneratorWithSetup is a TX generator that first submit TXs from the WorkloadSetupTXs, and blocks until indicated that it was committed. Then, it submits transactions from the tx stream.

func (*TxGeneratorWithSetup) Next

Next generates the next TX batch.

type TxStream

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

TxStream yields transactions from the stream.

func NewTxStream

func NewTxStream(profile *Profile, options *StreamOptions, counter *TxCounter) (*TxStream, error)

NewTxStream creates a stream that generates transactions in batches into a queue. The counter is the shared transaction-index counter (created before the stream and shared with the metrics); the stream's workers reserve their index ranges from it via the generators.

func (*TxStream) AppendBatch

func (s *TxStream) AppendBatch(ctx context.Context, batch []*servicepb.LoadGenTx)

AppendBatch appends a batch to the stream.

func (*TxStream) GetRate added in v0.1.8

func (s *TxStream) GetRate() uint64

GetRate reads the stream limit.

func (*TxStream) MakeGenerator

func (s *TxStream) MakeGenerator() *ConsumerRateController[*servicepb.LoadGenTx]

MakeGenerator creates a new generator that consumes from the stream. Each generator must be used from a single goroutine, but different generators from the same Stream can be used concurrently.

func (*TxStream) Run

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

Run starts the stream workers.

func (*TxStream) SetRate added in v0.1.8

func (s *TxStream) SetRate(rate uint64)

SetRate sets the stream limit.

Jump to

Keyboard shortcuts

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