Documentation
¶
Index ¶
- Constants
- func Clip[T cmp.Ordered](value, low, high T) T
- func CreateConfigTxFromConfigBlock(block *common.Block) (*servicepb.LoadGenTx, error)
- func CreateLoadGenNamespacesTX(policy *PolicyProfile) (*servicepb.LoadGenTx, error)
- func CreateNamespacesTxFromEndorser(endorser *TxEndorser, metaNamespaceVersion uint64, includeNS ...string) (*applicationpb.Tx, error)
- func CreateOrExtendConfigBlockWithCrypto(policy *PolicyProfile) (*common.Block, error)
- func CreateOrLoadConfigBlockWithCrypto(policy *PolicyProfile) (*common.Block, error)
- func GenerateTransactions(tb testing.TB, p *Profile, count int) []*servicepb.LoadGenTx
- func Map[T, K any](arr []T, mapper func(index int, value T) K) []K
- func MapToCoordinatorBatch(blockNum uint64, txs []*servicepb.LoadGenTx) *servicepb.CoordinatorBatch
- func MapToEnvelopeBatch(_ uint64, txs []*servicepb.LoadGenTx) []*common.Envelope
- func MapToLoadGenBatch(_ uint64, txs []*servicepb.LoadGenTx) *servicepb.LoadGenBatch
- func MapToOrdererBlock(blockNum uint64, txs []*servicepb.LoadGenTx) *common.Block
- func MapToVcBatch(blockNum uint64, txs []*servicepb.LoadGenTx) *servicepb.VcBatch
- func MapToVerifierBatch(blockNum uint64, txs []*servicepb.LoadGenTx) *servicepb.VerifierBatch
- func NewPolicyEndorserFromMsp(tb testing.TB, artifactsPath string) *testsig.NsEndorser
- func SumInt[T constraints.Integer](arr []T) int64
- type BlockProfile
- type ConsumeParameters
- type ConsumerRateController
- func (l *ConsumerRateController[T]) Consume(ctx context.Context, p ConsumeParameters) []T
- func (l *ConsumerRateController[T]) InstantiateWorker() *ConsumerRateController[T]
- func (l *ConsumerRateController[T]) Next(ctx context.Context) T
- func (l *ConsumerRateController[T]) Rate() uint64
- func (l *ConsumerRateController[T]) SetRate(rate uint64)
- type IndependentTxGenerator
- type KeyPath
- type KeyStats
- type Policy
- type PolicyProfile
- type Probability
- type Profile
- type StreamOptions
- type StreamWithSetup
- type TransactionProfile
- type TxBuilder
- type TxCounter
- type TxEndorser
- type TxGeneratorWithSetup
- type TxStream
Constants ¶
const ( PolicySchemeMSP = "MSP" PolicySchemeDefault = PolicySchemeMSP PolicySchemeUnspecified = "" )
Defines Policy.Scheme.
const DefaultGeneratedNamespaceID = "0"
DefaultGeneratedNamespaceID for now we're only generating transactions for a single namespace.
Variables ¶
This section is empty.
Functions ¶
func CreateConfigTxFromConfigBlock ¶ added in v0.1.8
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
GenerateTransactions is used for benchmarking and tests.
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
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
MapToOrdererBlock creates a Fabric's Orderer output block.
func MapToVcBatch ¶ added in v0.1.6
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
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 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 ¶
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
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
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.
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 ¶
func (g *TxGeneratorWithSetup) Next(ctx context.Context, p BlockProfile) []*servicepb.LoadGenTx
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 ¶
AppendBatch appends a batch to the stream.
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.