@@ -33,14 +33,16 @@ public sealed class DirectErrorQueueWriter
3333 private readonly Counter < long > _bypassCounter ;
3434
3535 private CancellationTokenSource ? _cts ;
36- private Task ? _loop ;
36+ private Task [ ] ? _loops ;
3737 private long _errorsWritten ;
38+ private long _errorsFailed ;
3839 private double _currentRate ;
3940 private string ? _activeScenario ;
4041 private DateTimeOffset _startedAt ;
4142
4243 public bool IsRunning => _cts is not null ;
4344 public long ErrorsWritten => Interlocked . Read ( ref _errorsWritten ) ;
45+ public long ErrorsFailed => Interlocked . Read ( ref _errorsFailed ) ;
4446 public double CurrentRate => _currentRate ;
4547 public string ? ActiveScenario => _activeScenario ;
4648
@@ -62,7 +64,7 @@ public DirectErrorQueueWriter(
6264 }
6365
6466 /// <summary>Starts writing failed-message envelopes directly to the error queue.</summary>
65- public bool TryStart ( string scenarioName , double rate , TimeSpan ? duration , out string ? error )
67+ public bool TryStart ( string scenarioName , double rate , TimeSpan ? duration , int ? parallelism , out string ? error )
6668 {
6769 var scenario = _registry . Get ( scenarioName ) ;
6870 if ( scenario is null )
@@ -83,6 +85,11 @@ public bool TryStart(string scenarioName, double rate, TimeSpan? duration, out s
8385 return false ;
8486 }
8587
88+ // Default to ProcessorCount workers. Each worker runs its own timer at rate/N so the
89+ // aggregate approaches the target. More workers = more concurrent sends = higher
90+ // throughput when individual sends have latency.
91+ var workerCount = parallelism is { } p and > 0 ? p : Environment . ProcessorCount ;
92+
8693 _activeScenario = scenarioName ;
8794 _currentRate = rate ;
8895 _startedAt = DateTimeOffset . UtcNow ;
@@ -92,10 +99,15 @@ public bool TryStart(string scenarioName, double rate, TimeSpan? duration, out s
9299 : new CancellationTokenSource ( ) ;
93100 _cts = cts ;
94101
95- _loop = Task . Run ( ( ) => WriteLoop ( scenario , rate , cts . Token ) ) ;
102+ _loops = new Task [ workerCount ] ;
103+ for ( var i = 0 ; i < workerCount ; i ++ )
104+ {
105+ var workerIndex = i ;
106+ _loops [ i ] = Task . Run ( ( ) => WriteLoop ( scenario , rate / workerCount , workerIndex , cts . Token ) ) ;
107+ }
96108
97- _logger . LogInformation ( "Started bypass writer for scenario {Scenario} at {Rate:F1} msg/s{Duration}" ,
98- scenarioName , rate , duration is null ? "" : $ " for { duration . Value } ") ;
109+ _logger . LogInformation ( "Started bypass writer for scenario {Scenario} at {Rate:F1} msg/s{Duration} across {Workers} workers " ,
110+ scenarioName , rate , duration is null ? "" : $ " for { duration . Value } ", workerCount ) ;
99111
100112 error = null ;
101113 return true ;
@@ -107,10 +119,27 @@ public void Stop()
107119 if ( _cts is null ) return ;
108120
109121 _cts . Cancel ( ) ;
122+
123+ // Await all worker loops before disposing the token so we don't dispose a CTS that's still
124+ // in flight inside _session.Send. Swallow the expected cancellation/timeout.
125+ var loops = _loops ;
126+ if ( loops is not null )
127+ {
128+ try
129+ {
130+ Task . WaitAll ( loops , TimeSpan . FromSeconds ( 5 ) ) ;
131+ }
132+ catch
133+ {
134+ // Loops were cancelled or timed out — expected during stop.
135+ }
136+ }
137+
110138 _cts . Dispose ( ) ;
111139 _cts = null ;
140+ _loops = null ;
112141
113- _logger . LogInformation ( "Stopped bypass writer after {Errors} errors written" , ErrorsWritten ) ;
142+ _logger . LogInformation ( "Stopped bypass writer: {Errors} written, {Failed} failed " , ErrorsWritten , ErrorsFailed ) ;
114143
115144 _activeScenario = null ;
116145 _currentRate = 0 ;
@@ -123,14 +152,16 @@ public void Stop()
123152 Scenario = _activeScenario ,
124153 Rate = _currentRate ,
125154 ErrorsWritten = ErrorsWritten ,
155+ ErrorsFailed = ErrorsFailed ,
126156 StartedAt = IsRunning ? _startedAt . ToString ( "O" ) : null
127157 } ;
128158
129159 /// <summary>
130160 /// The load generation loop: sends <see cref="LoadMessage"/> directly to the error queue
131- /// with failure headers at the target rate until cancelled.
161+ /// with failure headers at the target rate until cancelled. Multiple instances run in
162+ /// parallel, each handling a fraction of the total rate.
132163 /// </summary>
133- private async Task WriteLoop ( IScenario scenario , double rate , CancellationToken ct )
164+ private async Task WriteLoop ( IScenario scenario , double rate , int workerIndex , CancellationToken ct )
134165 {
135166 var interval = TimeSpan . FromSeconds ( 1.0 / rate ) ;
136167 using var timer = new PeriodicTimer ( interval ) ;
@@ -184,9 +215,20 @@ private async Task WriteLoop(IScenario scenario, double rate, CancellationToken
184215 _metrics . AddBypassErrorsWritten ( 1 ) ;
185216 _bypassCounter . Add ( 1 , new KeyValuePair < string , object ? > ( "scenario" , scenario . Name ) ) ;
186217 }
187- catch ( Exception ex ) when ( ex is not OperationCanceledException )
218+ catch ( OperationCanceledException )
188219 {
189- _logger . LogDebug ( ex , "Bypass send failed for scenario {Scenario} seq {Seq}" , scenario . Name , seq ) ;
220+ // Cancellation is expected on stop/timeout — let it propagate to the outer handler.
221+ throw ;
222+ }
223+ catch ( Exception ex )
224+ {
225+ Interlocked . Increment ( ref _errorsFailed ) ;
226+ _metrics . AddBypassErrorsFailed ( 1 ) ;
227+
228+ // Log at Warning so send failures are visible in default logging configs.
229+ // Previously this was LogDebug, which silently swallowed transport/broker
230+ // failures and made the bypass appear idle when sends were actually failing.
231+ _logger . LogWarning ( ex , "Bypass send failed for scenario {Scenario} worker {Worker} seq {Seq}" , scenario . Name , workerIndex , seq ) ;
190232 }
191233 }
192234 }
0 commit comments