From 28a1057635fa7ea48da8555d31097594a1f2244b Mon Sep 17 00:00:00 2001 From: Trezy Date: Mon, 9 Mar 2026 21:08:50 -0500 Subject: [PATCH] feat: process Tap events concurrently --- src/tap.rs | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/src/tap.rs b/src/tap.rs index 07c023d..6ac0192 100644 --- a/src/tap.rs +++ b/src/tap.rs @@ -1,7 +1,9 @@ +use std::sync::Arc; + use futures_util::{SinkExt, StreamExt}; use serde::{Deserialize, Serialize}; use serde_json::Value; -use tokio::sync::watch; +use tokio::sync::{Semaphore, watch}; use tokio_tungstenite::tungstenite::Message; use tokio_tungstenite::tungstenite::client::IntoClientRequest; @@ -321,6 +323,9 @@ async fn run( let (mut write, mut read) = ws.split(); + // Allow up to 50 record events to be processed concurrently. + let semaphore = Arc::new(Semaphore::new(50)); + loop { tokio::select! { msg = read.next() => { @@ -369,7 +374,12 @@ async fn run( rkey = %record.rkey, "received record event from tap" ); - handle_record_event(state, &record).await; + let sem = semaphore.clone(); + let state = state.clone(); + tokio::spawn(async move { + let _permit = sem.acquire().await.unwrap(); + handle_record_event(&state, &record).await; + }); } } "identity" => { -- 2.51.2