Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions actions/rest-action/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 4 additions & 2 deletions actions/rest-action/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -25,3 +25,5 @@ ring = "0.17"
lupus = "0.0.2"
jsonschema = "0.46"

[dev-dependencies]
tokio = { version = "1", features = ["test-util"] }
90 changes: 69 additions & 21 deletions actions/rest-action/src/admission.rs
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -16,6 +16,7 @@ fn env(key: &str, default: &str) -> String {
#[derive(Clone)]
pub struct Admission {
semaphore: Arc<Semaphore>,
wait_timeout: Duration,
retry_after: Duration,
}

Expand All @@ -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<OwnedSemaphorePermit> {
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<OwnedSemaphorePermit> {
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 {
Expand All @@ -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());
}
}
9 changes: 9 additions & 0 deletions actions/rest-action/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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`/
Expand All @@ -157,6 +165,7 @@ pub async fn run() -> hercules_sdk::Result<()> {
connected,
pending,
execution_timeout,
queue_write_timeout,
registry: Arc::clone(&registry),
admission,
limits,
Expand Down
15 changes: 13 additions & 2 deletions actions/rest-action/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<RwLock<Arc<RouteRegistry>>>,
pub admission: Admission,
pub limits: BodyLimits,
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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(),
})
Expand Down Expand Up @@ -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,
}
Expand All @@ -324,6 +329,7 @@ async fn execute_and_await_response(request: ExecutionRequest) -> Response<Full<
payload,
route,
execution_timeout,
queue_write_timeout,
permit,
retry_after_secs,
} = request;
Expand Down Expand Up @@ -351,7 +357,12 @@ async fn execute_and_await_response(request: ExecutionRequest) -> Response<Full<
// implies).
let _permit = permit;
let result = connected
.execute_flow_with_id(execution_id.clone(), flow_id.to_string(), payload)
.execute_flow_with_id_wait_for_capacity(
execution_id.clone(),
flow_id.to_string(),
payload,
queue_write_timeout,
)
.await;
// If `respond` was already called, it already claimed this
// entry, so this only fires for the "flow finished without
Expand Down