@@ -11,6 +11,7 @@ use bytesize::ByteSize;
1111use futures:: StreamExt ;
1212use indicatif:: { ProgressBar , ProgressStyle } ;
1313use objectstore_client:: { ExpirationPolicy , Usecase } ;
14+ use objectstore_metrics:: { count, gauge, record, timer} ;
1415use sketches_ddsketch:: DDSketch ;
1516use tokio:: sync:: Semaphore ;
1617use yansi:: Paint ;
@@ -258,6 +259,7 @@ async fn run_workload(
258259
259260 let workload = Arc :: new ( Mutex :: new ( workload) ) ;
260261 let metrics = Arc :: new ( Mutex :: new ( WorkloadMetrics :: default ( ) ) ) ;
262+ let workload_name = workload. lock ( ) . unwrap ( ) . name . clone ( ) ;
261263
262264 // See <https://docs.rs/tokio/latest/tokio/time/struct.Sleep.html#examples>
263265 let sleep = tokio:: time:: sleep_until ( deadline) ;
@@ -272,6 +274,7 @@ async fn run_workload(
272274 let workload = Arc :: clone( & workload) ;
273275 let remote = Arc :: clone( & remote) ;
274276 let metrics = Arc :: clone( & metrics) ;
277+ let workload_name = workload_name. clone( ) ;
275278
276279 let action = loop {
277280 if let Some ( action) = workload. lock( ) . unwrap( ) . next_action( ) {
@@ -284,9 +287,11 @@ async fn run_workload(
284287
285288 let multipart_config = multipart_config. clone( ) ;
286289 let task = async move {
290+ gauge!( "stresstest.concurrency" += 1usize , workload = workload_name. clone( ) ) ;
287291 let start = Instant :: now( ) ;
288292 match action {
289293 Action :: Write ( internal_id, payload) => {
294+ let timer = timer!( "stresstest.write.duration" , workload = workload_name. clone( ) ) ;
290295 let file_size = payload. len;
291296 let organization_id = workload. lock( ) . unwrap( ) . next_organization_id( ) ;
292297
@@ -304,6 +309,9 @@ async fn run_workload(
304309
305310 match result {
306311 Ok ( object_key) => {
312+ timer. record( ) ;
313+ count!( "stresstest.bytes_written" += file_size, workload = workload_name. clone( ) ) ;
314+ record!( "stresstest.write.size" = file_size, workload = workload_name. clone( ) ) ;
307315 let external_id = ( usecase, organization_id, object_key) ;
308316 workload. lock( ) . unwrap( ) . push_file( internal_id, external_id) ;
309317 let mut metrics = metrics. lock( ) . unwrap( ) ;
@@ -312,38 +320,49 @@ async fn run_workload(
312320 metrics. bytes_written += file_size;
313321 }
314322 Err ( err) => {
323+ drop( timer) ;
315324 print_error( "writing object" , & err) ;
316325 let mut metrics = metrics. lock( ) . unwrap( ) ;
317326 metrics. write_failures += 1 ;
318327 }
319328 }
320329 }
321330 Action :: Read ( internal_id, external_id, payload) => {
331+ let timer = timer!( "stresstest.read.duration" , workload = workload_name. clone( ) ) ;
322332 let file_size = payload. len;
323333 let ( usecase, organization_id, object_key) = & external_id;
324334 match remote. read( usecase, * organization_id, object_key, payload) . await {
325335 Ok ( _) => {
336+ timer. record( ) ;
337+ count!( "stresstest.bytes_read" += file_size, workload = workload_name. clone( ) ) ;
326338 workload. lock( ) . unwrap( ) . push_file( internal_id, external_id) ;
327339 let mut metrics = metrics. lock( ) . unwrap( ) ;
328340 metrics. read_timing. add( start. elapsed( ) . as_secs_f64( ) ) ;
329341 metrics. bytes_read += file_size;
330342 }
331343 Err ( err) => {
344+ drop( timer) ;
332345 print_error( "reading object" , & err) ;
333346 let mut metrics = metrics. lock( ) . unwrap( ) ;
334347 metrics. read_failures += 1 ;
335348 }
336349 }
337350 }
338351 Action :: Delete ( external_id) => {
352+ let timer = timer!( "stresstest.delete.duration" , workload = workload_name. clone( ) ) ;
339353 let ( usecase, organization_id, object_key) = & external_id;
340- if let Err ( err) = remote. delete( usecase, * organization_id, object_key) . await {
341- print_error( "deleting object" , & err) ;
354+ match remote. delete( usecase, * organization_id, object_key) . await {
355+ Ok ( ( ) ) => timer. record( ) ,
356+ Err ( err) => {
357+ drop( timer) ;
358+ print_error( "deleting object" , & err) ;
359+ }
342360 }
343361 let mut metrics = metrics. lock( ) . unwrap( ) ;
344362 metrics. delete_timing. add( start. elapsed( ) . as_secs_f64( ) ) ;
345363 }
346364 }
365+ gauge!( "stresstest.concurrency" -= 1usize , workload = workload_name) ;
347366 drop( permit) ;
348367 } ;
349368 tokio:: spawn( task) ;
@@ -388,6 +407,7 @@ async fn run_batch_workload(
388407
389408 let workload = Arc :: new ( Mutex :: new ( workload) ) ;
390409 let metrics = Arc :: new ( Mutex :: new ( WorkloadMetrics :: default ( ) ) ) ;
410+ let workload_name = workload. lock ( ) . unwrap ( ) . name . clone ( ) ;
391411
392412 let sleep = tokio:: time:: sleep_until ( deadline) ;
393413 tokio:: pin!( sleep) ;
@@ -401,6 +421,7 @@ async fn run_batch_workload(
401421 let workload = Arc :: clone( & workload) ;
402422 let remote = Arc :: clone( & remote) ;
403423 let metrics = Arc :: clone( & metrics) ;
424+ let workload_name = workload_name. clone( ) ;
404425
405426 let ( payloads, usecase, org_id) = {
406427 let mut wl = workload. lock( ) . unwrap( ) ;
@@ -414,6 +435,7 @@ async fn run_batch_workload(
414435 } ;
415436
416437 let task = async move {
438+ gauge!( "stresstest.concurrency" += 1usize , workload = workload_name. clone( ) ) ;
417439 let session = remote. session( & usecase, org_id) ;
418440 let mut many = session. many( ) ;
419441
@@ -446,15 +468,18 @@ async fn run_batch_workload(
446468
447469 metrics. lock( ) . unwrap( ) . many_requests += 1 ;
448470
449- let batch_start = Instant :: now ( ) ;
471+ let batch_timer = timer! ( "stresstest.batch.duration" , workload = workload_name . clone ( ) ) ;
450472 let mut results = many. send( ) . await ;
473+ let mut batch_had_errors = false ;
451474
452475 while let Some ( result) = results. next( ) . await {
453476 match result {
454477 objectstore_client:: OperationResult :: Put ( key, Ok ( _) ) => {
455478 if let Some ( ( internal_id, file_size) ) =
456479 payload_info. remove( & key)
457480 {
481+ count!( "stresstest.bytes_written" += file_size, workload = workload_name. clone( ) ) ;
482+ record!( "stresstest.write.size" = file_size, workload = workload_name. clone( ) ) ;
458483 let external_id = ( usecase. clone( ) , org_id, key) ;
459484 workload. lock( ) . unwrap( ) . push_file( internal_id, external_id) ;
460485
@@ -466,10 +491,12 @@ async fn run_batch_workload(
466491 }
467492 }
468493 objectstore_client:: OperationResult :: Put ( _, Err ( err) ) => {
494+ batch_had_errors = true ;
469495 print_error( "batch write" , & err. into( ) ) ;
470496 metrics. lock( ) . unwrap( ) . write_failures += 1 ;
471497 }
472498 objectstore_client:: OperationResult :: Error ( err) => {
499+ batch_had_errors = true ;
473500 print_error( "batch request" , & err. into( ) ) ;
474501 metrics. lock( ) . unwrap( ) . write_failures += 1 ;
475502 }
@@ -479,11 +506,18 @@ async fn run_batch_workload(
479506 }
480507 }
481508
509+ let batch_elapsed = batch_timer. elapsed( ) ;
510+ if batch_had_errors {
511+ drop( batch_timer) ;
512+ } else {
513+ batch_timer. record( ) ;
514+ }
515+
482516 metrics
483517 . lock( )
484518 . unwrap( )
485519 . batch_timing
486- . add( batch_start . elapsed ( ) . as_secs_f64( ) ) ;
520+ . add( batch_elapsed . as_secs_f64( ) ) ;
487521
488522 // Remove the temp files now that the batch has been uploaded.
489523 if let Err ( err) = tokio:: fs:: remove_dir_all( & temp_dir) . await {
@@ -493,6 +527,7 @@ async fn run_batch_workload(
493527 ) ;
494528 }
495529
530+ gauge!( "stresstest.concurrency" -= 1usize , workload = workload_name) ;
496531 drop( permit) ;
497532 } ;
498533 tokio:: spawn( task) ;
0 commit comments