diff --git a/actions/rest-action/Cargo.lock b/actions/rest-action/Cargo.lock index a9c95a3..048e47c 100644 --- a/actions/rest-action/Cargo.lock +++ b/actions/rest-action/Cargo.lock @@ -636,9 +636,9 @@ checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" [[package]] name = "hercules-macros" -version = "1.4.4" +version = "1.4.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "197c2299ba6864ce41bb0d688b3d6ca69b3819dfaaeade8f4575d8ee630666f8" +checksum = "d15dea904df8795e741dcfe0fcf9b28a22270a212dbd6717503db9a38f07db00" dependencies = [ "proc-macro2", "quote", @@ -647,9 +647,9 @@ dependencies = [ [[package]] name = "hercules-sdk" -version = "1.4.4" +version = "1.4.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0f1372cb2ee4f0a3abe0a0a5e9e956f89beeca01f10d7da9a357a8ffae889edb" +checksum = "50699d0ffac1ef25476196207de6f20c92d6c11a7ac95c974f1b437f935ced6f" dependencies = [ "async-trait", "hercules-macros", diff --git a/actions/rest-action/Cargo.toml b/actions/rest-action/Cargo.toml index 428f2a4..26e4e73 100644 --- a/actions/rest-action/Cargo.toml +++ b/actions/rest-action/Cargo.toml @@ -5,9 +5,9 @@ edition = "2024" publish = false [dependencies] -hercules-sdk = { version = "1.4.4" } +hercules-sdk = { version = "1.4.5" } tucana = { version = "0.0.82", default-features = false } -tokio = { version = "1", features = ["rt-multi-thread", "macros", "net", "sync"] } +tokio = { version = "1", features = ["rt-multi-thread", "macros", "net", "sync", "time"] } tokio-stream = "0.1" hyper = "1.8" hyper-util = "0.1" @@ -25,3 +25,5 @@ ring = "0.17" lupus = "0.0.2" jsonschema = "0.46" +[dev-dependencies] +tokio = { version = "1", features = ["test-util"] } diff --git a/actions/rest-action/src/admission.rs b/actions/rest-action/src/admission.rs index 2cc40c6..b97993f 100644 --- a/actions/rest-action/src/admission.rs +++ b/actions/rest-action/src/admission.rs @@ -1,8 +1,8 @@ //! Bounded admission control for workflow execution: caps how many flow //! executions this instance has in flight at once, independent of however -//! many HTTP connections/requests are open. Once saturated, new requests get -//! a `503` with `Retry-After` immediately — no flow input is built and no -//! execution is started. +//! many HTTP connections/requests are open. Once saturated, new requests wait +//! for a bounded, configurable interval and then get a `503` with +//! `Retry-After` if no permit becomes available. use std::sync::Arc; use std::time::Duration; @@ -16,6 +16,7 @@ fn env(key: &str, default: &str) -> String { #[derive(Clone)] pub struct Admission { semaphore: Arc, + wait_timeout: Duration, retry_after: Duration, } @@ -24,24 +25,42 @@ impl Admission { let max_concurrent: usize = env("HERCULES_REST_MAX_CONCURRENT_EXECUTIONS", "256") .parse() .unwrap_or_else(|err| panic!("invalid HERCULES_REST_MAX_CONCURRENT_EXECUTIONS: {err}")); + let wait_timeout_ms: u64 = env("HERCULES_REST_ADMISSION_WAIT_TIMEOUT_MS", "1000") + .parse() + .unwrap_or_else(|err| panic!("invalid HERCULES_REST_ADMISSION_WAIT_TIMEOUT_MS: {err}")); let retry_after_secs: u64 = env("HERCULES_REST_RETRY_AFTER_SECS", "1") .parse() .unwrap_or_else(|err| panic!("invalid HERCULES_REST_RETRY_AFTER_SECS: {err}")); - Self::new(max_concurrent, Duration::from_secs(retry_after_secs)) + Self::new( + max_concurrent, + Duration::from_millis(wait_timeout_ms), + Duration::from_secs(retry_after_secs), + ) } - pub fn new(max_concurrent: usize, retry_after: Duration) -> Self { + pub fn new(max_concurrent: usize, wait_timeout: Duration, retry_after: Duration) -> Self { Self { semaphore: Arc::new(Semaphore::new(max_concurrent)), + wait_timeout, retry_after, } } - /// Non-blocking: `None` means the instance is at its configured - /// concurrency limit right now. - pub fn try_acquire(&self) -> Option { - Arc::clone(&self.semaphore).try_acquire_owned().ok() + /// Waits up to the configured deadline for execution capacity. A zero + /// timeout preserves fail-fast behavior. + pub async fn acquire(&self) -> Option { + if self.wait_timeout.is_zero() { + return Arc::clone(&self.semaphore).try_acquire_owned().ok(); + } + + tokio::time::timeout( + self.wait_timeout, + Arc::clone(&self.semaphore).acquire_owned(), + ) + .await + .ok() + .and_then(Result::ok) } pub fn retry_after_secs(&self) -> u64 { @@ -53,22 +72,51 @@ impl Admission { mod tests { use super::*; - #[test] - fn admits_up_to_the_configured_limit() { - let admission = Admission::new(2, Duration::from_secs(1)); - let first = admission.try_acquire(); - let second = admission.try_acquire(); + #[tokio::test] + async fn admits_up_to_the_configured_limit() { + let admission = Admission::new(2, Duration::ZERO, Duration::from_secs(1)); + let first = admission.acquire().await; + let second = admission.acquire().await; assert!(first.is_some()); assert!(second.is_some()); - assert!(admission.try_acquire().is_none()); + assert!(admission.acquire().await.is_none()); + } + + #[tokio::test] + async fn releasing_a_permit_frees_capacity() { + let admission = Admission::new(1, Duration::ZERO, Duration::from_secs(1)); + let permit = admission.acquire().await.expect("first acquire succeeds"); + assert!(admission.acquire().await.is_none()); + drop(permit); + assert!(admission.acquire().await.is_some()); } - #[test] - fn releasing_a_permit_frees_capacity() { - let admission = Admission::new(1, Duration::from_secs(1)); - let permit = admission.try_acquire().expect("first acquire succeeds"); - assert!(admission.try_acquire().is_none()); + #[tokio::test] + async fn waits_for_a_permit_to_be_released() { + let admission = Admission::new(1, Duration::from_secs(1), Duration::from_secs(1)); + let permit = admission.acquire().await.expect("first acquire succeeds"); + let waiting = tokio::spawn({ + let admission = admission.clone(); + async move { admission.acquire().await } + }); + + tokio::task::yield_now().await; + assert!(!waiting.is_finished()); drop(permit); - assert!(admission.try_acquire().is_some()); + assert!(waiting.await.expect("waiter task panicked").is_some()); + } + + #[tokio::test(start_paused = true)] + async fn returns_none_when_the_wait_timeout_expires() { + let admission = Admission::new(1, Duration::from_millis(100), Duration::from_secs(1)); + let _permit = admission.acquire().await.expect("first acquire succeeds"); + let waiting = tokio::spawn({ + let admission = admission.clone(); + async move { admission.acquire().await } + }); + + tokio::task::yield_now().await; + tokio::time::advance(Duration::from_millis(100)).await; + assert!(waiting.await.expect("waiter task panicked").is_none()); } } diff --git a/actions/rest-action/src/lib.rs b/actions/rest-action/src/lib.rs index 87457ea..cb8981c 100644 --- a/actions/rest-action/src/lib.rs +++ b/actions/rest-action/src/lib.rs @@ -44,11 +44,15 @@ fn env(key: &str, default: &str) -> String { /// construction time, so it's registered by hand below instead (see its /// `manual` attribute). fn build_action(pending: pending::PendingResponses) -> Action { + let request_queue_capacity: usize = env("HERCULES_REQUEST_QUEUE_CAPACITY", "256") + .parse() + .unwrap_or_else(|err| panic!("invalid HERCULES_REQUEST_QUEUE_CAPACITY: {err}")); let mut action = Action::new( env("HERCULES_ACTION_ID", "rest-action"), env("HERCULES_SDK_VERSION", "0.0.0"), ) .aquila_url(env("HERCULES_AQUILA_URL", "127.0.0.1:8081")) + .request_queue_capacity(request_queue_capacity) // Every instance needs to be reachable to serve HTTP traffic, so every // instance gets every flow rather than splitting them up. .scaling(ScalingOption::Disabled) @@ -133,6 +137,10 @@ pub async fn run() -> hercules_sdk::Result<()> { .parse() .unwrap_or_else(|err| panic!("invalid HERCULES_EXECUTION_TIMEOUT_SECS: {err}")); let execution_timeout = std::time::Duration::from_secs(execution_timeout_secs); + let queue_write_timeout_ms: u64 = env("HERCULES_REST_QUEUE_WRITE_TIMEOUT_MS", "1000") + .parse() + .unwrap_or_else(|err| panic!("invalid HERCULES_REST_QUEUE_WRITE_TIMEOUT_MS: {err}")); + let queue_write_timeout = std::time::Duration::from_millis(queue_write_timeout_ms); // Seeded once here from `Connected::flows()` (not per request — see // `registry.rs`), then kept current purely by `FlowUpserted`/ @@ -157,6 +165,7 @@ pub async fn run() -> hercules_sdk::Result<()> { connected, pending, execution_timeout, + queue_write_timeout, registry: Arc::clone(®istry), admission, limits, diff --git a/actions/rest-action/src/server.rs b/actions/rest-action/src/server.rs index ec37c18..ddb9df3 100644 --- a/actions/rest-action/src/server.rs +++ b/actions/rest-action/src/server.rs @@ -37,6 +37,7 @@ pub struct ServerState { pub connected: Connected, pub pending: PendingResponses, pub execution_timeout: Duration, + pub queue_write_timeout: Duration, pub registry: Arc>>, pub admission: Admission, pub limits: BodyLimits, @@ -79,12 +80,14 @@ async fn handle( connected, pending, execution_timeout, + queue_write_timeout, registry, admission, limits, .. } = state; let execution_timeout = *execution_timeout; + let queue_write_timeout = *queue_write_timeout; let (parts, body) = req.into_parts(); let method = parts.method; let path = parts.uri.path().to_string(); @@ -136,7 +139,7 @@ async fn handle( // parsing, schema validation, flow dispatch) — routing and auth above // are cheap and shouldn't be denied just because executions are // saturated. - let Some(permit) = admission.try_acquire() else { + let Some(permit) = admission.acquire().await else { log::warn!( "{method} {path}: flow {} rejected: execution admission saturated", flow.flow_id @@ -243,6 +246,7 @@ async fn handle( payload, route: format!("{method} {path}"), execution_timeout, + queue_write_timeout, permit, retry_after_secs: admission.retry_after_secs(), }) @@ -312,6 +316,7 @@ struct ExecutionRequest { payload: hercules_sdk::PlainValue, route: String, execution_timeout: Duration, + queue_write_timeout: Duration, permit: tokio::sync::OwnedSemaphorePermit, retry_after_secs: u64, } @@ -324,6 +329,7 @@ async fn execute_and_await_response(request: ExecutionRequest) -> Response Response