use std::collections::HashMap;
use std::net::IpAddr;
use std::sync::Arc;
use std::time::Instant;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::{TcpListener, TcpStream};
use tropel_sdk::types::{Body, Method, Request as TropelRequest, ResponseType};
use tropel_sdk::TropelError;
const RATE_LIMIT_PER_SEC: u64 = 200;
#[derive(Default)]
struct ScriptSendLog {
sends: std::sync::Mutex<Vec<serde_json::Value>>,
}
impl ScriptSendLog {
fn record(&self, entry: serde_json::Value) {
let mut g = match self.sends.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
};
g.push(entry);
}
fn drain(&self) -> Vec<serde_json::Value> {
let g = match self.sends.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
};
g.clone()
}
}
struct RecordingHttpClient {
inner: Arc<dyn tropel_sdk::traits::DriverHttpClient>,
log: Arc<ScriptSendLog>,
}
#[async_trait::async_trait]
impl tropel_sdk::traits::DriverHttpClient for RecordingHttpClient {
async fn execute(
&self,
req: &tropel_sdk::Request,
) -> tropel_sdk::Result<tropel_sdk::types::Response> {
let started = std::time::Instant::now();
let out = self.inner.execute(req).await;
let (status, error) = match &out {
Ok(res) => (res.status_code, None),
Err(e) => (0u16, Some(e.to_string())),
};
self.log.record(serde_json::json!({
"source": "sendRequest",
"kind": "ad-hoc",
"method": req.method.to_string(),
"url": req.url,
"status": status,
"responseTime": started.elapsed().as_millis() as u64,
"error": error,
}));
out
}
}
#[derive(Default)]
struct ScriptCookies {
jar: std::sync::Mutex<Vec<serde_json::Value>>,
ops: std::sync::Mutex<Vec<serde_json::Value>>,
}
impl ScriptCookies {
fn lock_jar(&self) -> std::sync::MutexGuard<'_, Vec<serde_json::Value>> {
match self.jar.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
}
}
fn lock_ops(&self) -> std::sync::MutexGuard<'_, Vec<serde_json::Value>> {
match self.ops.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
}
}
fn record(&self, op: serde_json::Value) {
self.lock_ops().push(op);
}
}
fn cookie_matches_url(cookie: &serde_json::Value, url: &str) -> bool {
let host = url
.split("://")
.nth(1)
.unwrap_or(url)
.split('/')
.next()
.unwrap_or("")
.split(':')
.next()
.unwrap_or("");
let path = {
let after = url.split("://").nth(1).unwrap_or(url);
match after.find('/') {
Some(i) => after[i..].split('?').next().unwrap_or("/").to_string(),
None => "/".to_string(),
}
};
let domain = cookie
.get("domain")
.and_then(|d| d.as_str())
.unwrap_or("")
.trim_start_matches('.');
let c_path = cookie.get("path").and_then(|p| p.as_str()).unwrap_or("/");
let domain_ok = domain.is_empty() || host == domain || host.ends_with(&format!(".{domain}"));
domain_ok && path.starts_with(c_path)
}
#[derive(Clone, serde::Serialize)]
struct PendingHostCall {
#[serde(rename = "callId")]
call_id: u64,
kind: &'static str,
path: String,
}
#[derive(Default)]
struct RunCallbacks {
pending: std::sync::Mutex<std::collections::VecDeque<PendingHostCall>>,
ready: tokio::sync::Notify,
replies: std::sync::Mutex<HashMap<u64, std::sync::mpsc::Sender<String>>>,
next_id: std::sync::atomic::AtomicU64,
done: std::sync::atomic::AtomicBool,
}
impl RunCallbacks {
fn finish(&self) {
self.done.store(true, std::sync::atomic::Ordering::SeqCst);
self.ready.notify_waiters();
}
}
struct AgentState {
token: Option<String>,
client: tropel_http::HttpClient,
allowed_origins: Vec<String>,
runs: std::sync::Mutex<HashMap<String, Arc<RunCallbacks>>>,
}
fn cors_headers(state: &AgentState, origin: Option<&str>, is_preflight: bool) -> Option<String> {
let origin = origin?;
if !state.allowed_origins.iter().any(|o| o == origin) {
return None;
}
let mut h = format!("Access-Control-Allow-Origin: {origin}\r\n");
h.push_str("Vary: Origin\r\n");
if is_preflight {
h.push_str("Access-Control-Allow-Methods: GET, POST, OPTIONS\r\n");
h.push_str("Access-Control-Allow-Headers: Authorization, Content-Type\r\n");
h.push_str("Access-Control-Max-Age: 600\r\n");
}
Some(h)
}
pub async fn run_agent(
port: u16,
bind: &str,
token: Option<&str>,
allowed_origins: &[String],
exit_with_parent: bool,
) -> tropel_sdk::Result<()> {
let ip: IpAddr = bind
.parse()
.map_err(|_| TropelError::Other(format!("invalid bind address: {bind}")))?;
if !ip.is_loopback() {
return Err(TropelError::Other(format!(
"refusing to bind {bind}: the agent is a localhost-only execution endpoint (TR-405)"
)));
}
if exit_with_parent {
std::thread::spawn(|| {
use std::io::Read;
let mut stdin = std::io::stdin();
let mut buf = [0u8; 256];
loop {
match stdin.read(&mut buf) {
Ok(0) | Err(_) => {
tracing::info!(
"tropel agent exiting: stdin reached EOF — either the process that \
spawned it is gone, or stdin was not a pipe held open by it \
(--exit-with-parent requires one)"
);
std::process::exit(0);
}
Ok(_) => {}
}
}
});
}
let addr = format!("{bind}:{port}");
let listener = TcpListener::bind(&addr).await.map_err(TropelError::Io)?;
tracing::info!(
"tropel agent listening on http://{addr} (token auth {})",
if token.is_some() { "on" } else { "off" }
);
let http_config = tropel_http::HttpConfig::default();
let client = tropel_http::HttpClient::new(&http_config)
.map_err(|e| TropelError::Other(format!("http client init failed: {e}")))?;
let state = Arc::new(AgentState {
token: token.map(str::to_string),
client,
allowed_origins: allowed_origins.to_vec(),
runs: std::sync::Mutex::new(HashMap::new()),
});
loop {
let (mut sock, peer) = listener.accept().await.map_err(TropelError::Io)?;
tracing::debug!("agent: connection from {peer}");
let state = state.clone();
tokio::spawn(async move {
if let Err(e) = handle_connection(&mut sock, state).await {
tracing::debug!("agent: connection error: {e}");
}
});
}
}
async fn handle_connection(sock: &mut TcpStream, state: Arc<AgentState>) -> tropel_sdk::Result<()> {
let mut buf = vec![0u8; 64 * 1024];
let n = sock.read(&mut buf).await.map_err(TropelError::Io)?;
let raw = String::from_utf8_lossy(&buf[..n]);
let prefetched_body: Vec<u8> = raw
.find("\r\n\r\n")
.map(|i| buf[i + 4..n].to_vec())
.unwrap_or_default();
let mut lines = raw.lines();
let request_line = lines.next().unwrap_or("");
let mut parts = request_line.split_whitespace();
let method = parts.next().unwrap_or("");
let path = parts.next().unwrap_or("/");
let mut content_length = 0usize;
let mut auth_header = String::new();
let mut origin: Option<String> = None;
let mut wants_private_network = false;
for line in lines {
if line.is_empty() {
break;
}
if let Some((k, v)) = line.split_once(':') {
let key = k.trim().to_ascii_lowercase();
if key == "content-length" {
content_length = v.trim().parse().unwrap_or(0);
} else if key == "authorization" {
auth_header = v.trim().to_string();
} else if key == "origin" {
origin = Some(v.trim().to_string());
} else if key == "access-control-request-private-network" {
wants_private_network = v.trim().eq_ignore_ascii_case("true");
}
}
}
let cors = cors_headers(&state, origin.as_deref(), false);
{
let mut limiter = RateLimiter::new();
if limiter.allow().is_err() {
return respond_raw_cors(sock, cors.as_deref(), 429, "rate limit exceeded").await;
}
}
if method == "OPTIONS" {
let Some(mut headers) = cors_headers(&state, origin.as_deref(), true) else {
return respond(
sock,
403,
&error_body(&format!(
"origin {:?} is not allowed to reach this agent \u{2014} start it with --allow-origin {}",
origin.as_deref().unwrap_or("(none)"),
origin.as_deref().unwrap_or("<origin>")
)),
)
.await;
};
if wants_private_network {
headers.push_str("Access-Control-Allow-Private-Network: true\r\n");
}
return respond_raw(sock, 204, "", Some(&headers)).await;
}
if let Some(expected) = &state.token {
if auth_header != format!("Bearer {expected}") {
return respond_raw_cors(sock, cors.as_deref(), 401, r#"{"error":"unauthorized"}"#)
.await;
}
}
match (method, path) {
("GET", "/version") => {
let body = format!(r#"{{"version":"{}"}}"#, env!("CARGO_PKG_VERSION"));
respond_raw_cors(sock, cors.as_deref(), 200, &body).await
}
("POST", "/resolve/batch") => {
let Some(payload) =
read_json_body(sock, content_length, 8 * 1024 * 1024, &prefetched_body).await?
else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await;
};
let vars: HashMap<String, String> = payload
.get("variables")
.and_then(|v| serde_json::from_value(v.clone()).ok())
.unwrap_or_default();
let items: Vec<BatchResolveItem> =
match serde_json::from_value(payload.get("items").cloned().unwrap_or_default()) {
Ok(v) => v,
Err(e) => {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
&error_body(&format!("invalid items: {e}")),
)
.await
}
};
let scope = tropel_variables::VariableScope {
env: vars.clone(),
..Default::default()
};
let resolver = tropel_variables::VariableResolver::new();
let mut out = Vec::with_capacity(items.len());
for item in &items {
let mode = item.mode.as_deref().unwrap_or("plain");
if item.deep == Some(false) {
out.push(serde_json::json!({
"error": "deep: false is not supported by POST /resolve/batch \u{2014} the batched reply reports hitCap/unresolved, which only a chain resolved to settlement has; use POST /resolve for a shallow resolve"
}));
continue;
}
match resolver.resolve_reporting(
&item.template,
&scope,
tropel_variables::MAX_VARIABLE_RESOLUTION_PASSES,
mode,
) {
Ok(outcome) => out.push(serde_json::json!({
"value": outcome.value,
"hitCap": outcome.hit_cap,
"unresolved": outcome.unresolved,
})),
Err(why) => out.push(serde_json::json!({ "error": why })),
}
}
respond_raw_cors(
sock,
cors.as_deref(),
200,
&serde_json::json!({ "items": out }).to_string(),
)
.await
}
("POST", "/resolve") => {
let Some(payload) =
read_json_body(sock, content_length, 1024 * 1024, &prefetched_body).await?
else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await;
};
let template = payload
.get("template")
.and_then(|t| t.as_str())
.unwrap_or("");
let vars: HashMap<String, String> = payload
.get("variables")
.and_then(|v| serde_json::from_value(v.clone()).ok())
.unwrap_or_default();
let mode = payload
.get("mode")
.and_then(|m| m.as_str())
.unwrap_or("none");
let deep = payload
.get("deep")
.and_then(|d| d.as_bool())
.unwrap_or(true);
match tropel_variables::resolve_template_for_host(template, &vars, mode, deep) {
Ok(resolved) => {
respond(
sock,
200,
&serde_json::json!({ "resolved": resolved }).to_string(),
)
.await
}
Err(why) => respond_raw_cors(sock, cors.as_deref(), 400, &error_body(&why)).await,
}
}
("POST", "/assert") => {
let Some(payload) =
read_json_body(sock, content_length, 8 * 1024 * 1024, &prefetched_body).await?
else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await;
};
let target: tropel_variables::assertions::AssertionTarget = match serde_json::from_value(
payload.get("response").cloned().unwrap_or_default(),
) {
Ok(t) => t,
Err(e) => {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
&error_body(&format!("invalid response: {e}")),
)
.await
}
};
let specs: Vec<AgentAssertionSpec> = match serde_json::from_value(
payload.get("assertions").cloned().unwrap_or_default(),
) {
Ok(v) => v,
Err(e) => {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
&error_body(&format!("invalid assertions: {e}")),
)
.await
}
};
let outcomes: Vec<_> = specs
.iter()
.map(|spec| {
let name = spec
.name
.clone()
.unwrap_or_else(|| format!("{} {}", spec.target, spec.operator));
match tropel_variables::assertions::resolve_assertion_target(
&spec.target,
&target,
) {
Ok(actual) => tropel_variables::assertions::assert_evaluate(
&name,
&spec.target,
&actual,
&spec.operator,
&spec.expected,
None,
),
Err(why) => tropel_variables::assertions::AssertionOutcome {
name,
passed: false,
unsupported: Some(why),
message: None,
},
}
})
.collect();
respond(
sock,
200,
&serde_json::to_string(&outcomes).unwrap_or_default(),
)
.await
}
("POST", "/variables/dynamic/batch") => {
let Some(payload) =
read_json_body(sock, content_length, 8 * 1024 * 1024, &prefetched_body).await?
else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await;
};
let items: Vec<BatchResolveItem> = payload
.get("items")
.and_then(|v| serde_json::from_value(v.clone()).ok())
.unwrap_or_default();
let catalog = tropel_variables::DynamicCatalog::new();
let mut out = Vec::with_capacity(items.len());
for item in &items {
match catalog.resolve(&item.template) {
Ok(value) => out.push(serde_json::json!({ "value": value })),
Err(why) => out.push(serde_json::json!({ "error": why })),
}
}
respond_raw_cors(
sock,
cors.as_deref(),
200,
&serde_json::json!({ "items": out }).to_string(),
)
.await
}
("POST", "/variables/dynamic") => {
let Some(payload) =
read_json_body(sock, content_length, 1024 * 1024, &prefetched_body).await?
else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await;
};
let template = payload
.get("template")
.and_then(|t| t.as_str())
.unwrap_or("");
let catalog = tropel_variables::DynamicCatalog::new();
match catalog.resolve(template) {
Ok(value) => {
respond(
sock,
200,
&serde_json::json!({ "value": value }).to_string(),
)
.await
}
Err(why) => respond_raw_cors(sock, cors.as_deref(), 400, &error_body(&why)).await,
}
}
("GET", "/constants") => {
let variables: Vec<serde_json::Value> = tropel_variables::PREDEFINED_VARIABLE_META
.iter()
.map(|m| serde_json::json!({ "name": m.name, "description": m.description }))
.collect();
let body = serde_json::json!({
"maxVariableResolutionPasses": tropel_variables::MAX_VARIABLE_RESOLUTION_PASSES,
"predefinedVariables": variables,
});
respond_raw_cors(sock, cors.as_deref(), 200, &body.to_string()).await
}
("GET", "/operators") => {
let body = serde_json::to_string(tropel_variables::assertions::ASSERTION_OPERATORS)
.unwrap_or_default();
respond_raw_cors(sock, cors.as_deref(), 200, &body).await
}
("POST", "/auth/sign") => {
let Some(payload) =
read_json_body(sock, content_length, 8 * 1024 * 1024, &prefetched_body).await?
else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await;
};
let scheme = payload.get("scheme").and_then(|s| s.as_str()).unwrap_or("");
let params = payload.get("params").cloned().unwrap_or_default();
match sign_with_scheme(scheme, ¶ms) {
Ok(headers) => {
respond(
sock,
200,
&serde_json::to_string(&headers).unwrap_or_default(),
)
.await
}
Err(why) => respond_raw_cors(sock, cors.as_deref(), 400, &error_body(&why)).await,
}
}
("POST", "/script") => {
let Some(payload) =
read_json_body(sock, content_length, 4 * 1024 * 1024, &prefetched_body).await?
else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await;
};
let Some(code) = payload.get("code").and_then(|c| c.as_str()) else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"`code` is required and must be a string"}"#,
)
.await;
};
let sandbox_cfg = payload
.get("sandbox")
.map(|v| tropel_sandbox::config::SandboxConfig {
namespace: v
.get("namespace")
.and_then(|n| n.as_str())
.unwrap_or("trp")
.to_string(),
aliases: v
.get("aliases")
.and_then(|a| a.as_array())
.map(|a| {
a.iter()
.filter_map(|x| x.as_str().map(str::to_string))
.collect()
})
.unwrap_or_default(),
})
.unwrap_or_default();
let scopes = ScriptScopes::from_payload(&payload);
let script_request: Option<TropelRequest> = payload
.get("request")
.and_then(|v| serde_json::from_value(v.clone()).ok());
let script_response: Option<tropel_sdk::types::Response> = payload
.get("response")
.and_then(|v| serde_json::from_value(v.clone()).ok());
let run_id = payload
.get("runId")
.and_then(|v| v.as_str())
.map(str::to_string);
let cookies = payload
.get("cookies")
.and_then(|c| c.as_array())
.map(|arr| {
let sc = ScriptCookies::default();
*sc.lock_jar() = arr.clone();
Arc::new(sc)
});
let callbacks = run_id.as_ref().map(|id| {
let cb = Arc::new(RunCallbacks::default());
let mut runs = match state.runs.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
};
runs.insert(id.clone(), cb.clone());
cb
});
let outcome = run_script_once(
code,
scopes,
script_request,
script_response,
sandbox_cfg,
ScriptHost {
http: state.client.clone(),
callbacks: callbacks.clone(),
cookies: cookies.clone(),
},
)
.await;
if let (Some(id), Some(cb)) = (run_id, callbacks) {
cb.finish();
let mut runs = match state.runs.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
};
runs.remove(&id);
}
match outcome {
Ok(out) => respond_raw_cors(sock, cors.as_deref(), 200, &out.to_string()).await,
Err(why) => respond_raw_cors(sock, cors.as_deref(), 500, &error_body(&why)).await,
}
}
("POST", "/script/callback/next") => {
let Some(payload) =
read_json_body(sock, content_length, 64 * 1024, &prefetched_body).await?
else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await;
};
let Some(run_id) = payload.get("runId").and_then(|v| v.as_str()) else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"`runId` is required"}"#,
)
.await;
};
let cb = {
let runs = match state.runs.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
};
runs.get(run_id).cloned()
};
let Some(cb) = cb else {
return respond_raw_cors(sock, cors.as_deref(), 204, "").await;
};
let deadline = std::time::Duration::from_secs(20);
let started = std::time::Instant::now();
loop {
let waiting = cb.ready.notified();
if let Some(call) = {
let mut pending = match cb.pending.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
};
pending.pop_front()
} {
let body = serde_json::to_string(&call).unwrap_or_default();
return respond_raw_cors(sock, cors.as_deref(), 200, &body).await;
}
if cb.done.load(std::sync::atomic::Ordering::SeqCst) {
return respond_raw_cors(sock, cors.as_deref(), 204, "").await;
}
let left = deadline.saturating_sub(started.elapsed());
if left.is_zero() {
return respond_raw_cors(sock, cors.as_deref(), 204, "").await;
}
tokio::select! {
_ = waiting => {}
_ = tokio::time::sleep(left) => {}
}
}
}
("POST", "/script/callback/reply") => {
let Some(payload) =
read_json_body(sock, content_length, 8 * 1024 * 1024, &prefetched_body).await?
else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await;
};
let run_id = payload.get("runId").and_then(|v| v.as_str()).unwrap_or("");
let call_id = payload.get("callId").and_then(|v| v.as_u64());
let Some(call_id) = call_id else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"`callId` is required and must be a number"}"#,
)
.await;
};
let cb = {
let runs = match state.runs.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
};
runs.get(run_id).cloned()
};
let Some(cb) = cb else {
return respond_raw_cors(
sock,
cors.as_deref(),
404,
r#"{"error":"no such run — it may have already finished"}"#,
)
.await;
};
let tx = {
let mut replies = match cb.replies.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
};
replies.remove(&call_id)
};
let Some(tx) = tx else {
return respond_raw_cors(
sock,
cors.as_deref(),
404,
r#"{"error":"no such callId — already answered, or it timed out"}"#,
)
.await;
};
let result = payload
.get("result")
.cloned()
.unwrap_or(serde_json::Value::Null);
let _ = tx.send(result.to_string());
respond_raw_cors(sock, cors.as_deref(), 200, r#"{"delivered":true}"#).await
}
("POST", "/auth/oauth2") => {
let Some(payload) =
read_json_body(sock, content_length, 1024 * 1024, &prefetched_body).await?
else {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await;
};
let op = payload.get("op").and_then(|o| o.as_str()).unwrap_or("");
let params = payload.get("params").cloned().unwrap_or_default();
match oauth2_dispatch(op, ¶ms) {
Ok(out) => respond_raw_cors(sock, cors.as_deref(), 200, &out.to_string()).await,
Err(why) => respond_raw_cors(sock, cors.as_deref(), 400, &error_body(&why)).await,
}
}
("POST", "/execute") => {
let body_buf = read_body(sock, content_length, 64 * 1024, &prefetched_body).await?;
let req: serde_json::Value = match serde_json::from_slice(&body_buf) {
Ok(v) => v,
Err(_) => {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await
}
};
let out = execute_single(&state, &req).await;
respond_raw_cors(sock, cors.as_deref(), 200, &out.to_string()).await
}
("POST", "/run") => {
let raw_lower = raw.to_ascii_lowercase();
if raw_lower.contains("x-tropel-relay")
|| raw_lower.contains("x-knockport-relay")
|| raw_lower.contains("via: relay")
|| raw_lower.contains("x-relay-transport")
{
return respond(
sock,
403,
r#"{"error":"relay is not a load transport — POST /run refused (TR-411); use the desktop tauri transport or the native CLI"}"#,
)
.await;
}
let body_buf =
read_body(sock, content_length, 4 * 1024 * 1024, &prefetched_body).await?;
let payload: serde_json::Value = match serde_json::from_slice(&body_buf) {
Ok(v) => v,
Err(_) => {
return respond_raw_cors(
sock,
cors.as_deref(),
400,
r#"{"error":"invalid JSON body"}"#,
)
.await
}
};
let scenario_json = payload
.get("scenario")
.and_then(|s| s.as_str())
.unwrap_or("");
let iterations = payload
.get("iterations")
.and_then(|i| i.as_u64())
.unwrap_or(1)
.min(1000); let scenario: tropel_sdk::scenario::Scenario = match serde_json::from_str(scenario_json)
{
Ok(s) => s,
Err(e) => {
return respond(
sock,
400,
&format!(r#"{{"error":"invalid scenario: {e}"}}"#),
)
.await
}
};
let thresholds: std::collections::HashMap<String, String> = payload
.get("thresholds")
.and_then(|t| t.as_object())
.map(|o| {
o.iter()
.map(|(k, v)| (k.clone(), v.as_str().unwrap_or("").to_string()))
.collect()
})
.unwrap_or_default();
let stream = payload
.get("stream")
.and_then(|s| s.as_bool())
.unwrap_or(false);
if stream {
return run_load_streaming(sock, &state, &scenario, iterations, &thresholds).await;
}
let out = run_load(&state, &scenario, iterations, &thresholds).await;
respond_raw_cors(sock, cors.as_deref(), 200, &out.to_string()).await
}
_ => respond_raw_cors(sock, cors.as_deref(), 404, r#"{"error":"not found"}"#).await,
}
}
fn eval_threshold(expr: &str, reqs: u64, failed: u64) -> Result<bool, String> {
let parts: Vec<&str> = expr.split_whitespace().collect();
if parts.len() != 3 {
return Err(format!(
"invalid threshold '{expr}': expected '<metric> <op> <value>'"
));
}
let actual = match parts[0] {
"http_reqs" => reqs as f64,
"http_req_failed" => {
if reqs == 0 {
0.0
} else {
failed as f64 / reqs as f64
}
}
other => return Err(format!("unsupported threshold metric '{other}'")),
};
let threshold: f64 = parts[2]
.parse()
.map_err(|_| format!("invalid threshold value '{}'", parts[2]))?;
let passed = match parts[1] {
"<" => actual < threshold,
"<=" => actual <= threshold,
">" => actual > threshold,
">=" => actual >= threshold,
"==" | "===" => (actual - threshold).abs() < f64::EPSILON,
"!=" => (actual - threshold).abs() > f64::EPSILON,
other => return Err(format!("unknown operator '{other}'")),
};
Ok(passed)
}
async fn run_load(
state: &AgentState,
scenario: &tropel_sdk::scenario::Scenario,
iterations: u64,
thresholds: &std::collections::HashMap<String, String>,
) -> serde_json::Value {
let mut samples: Vec<serde_json::Value> = Vec::new();
let mut total_failures = 0u64;
let mut unsupported_errors: Vec<String> = Vec::new();
for it in 0..iterations {
for item in &scenario.items {
let Some(request) = item.request.as_ref() else {
continue;
};
let signer_opt = match &request.auth {
Some(auth) => match state.client.get_signer(auth) {
Ok(s) => s,
Err(e) => {
let msg = e.to_string();
if !unsupported_errors.contains(&msg) {
unsupported_errors.push(msg.clone());
}
total_failures += 1;
let elapsed_ms = 0.0;
samples.push(serde_json::json!({
"metric": "http_reqs",
"iteration": it,
"url": request.url,
"status": 0,
"duration_ms": elapsed_ms,
"error": msg,
}));
continue;
}
},
None => None,
};
let start = Instant::now();
let result = state.client.execute(request, signer_opt.as_deref()).await;
let elapsed_ms = start.elapsed().as_millis() as f64;
let (status, ok) = match &result {
Ok(resp) => (resp.status_code, (200..400).contains(&resp.status_code)),
Err(_) => (0, false),
};
if !ok {
total_failures += 1;
}
let mut sample = serde_json::json!({
"metric": "http_reqs",
"iteration": it,
"url": request.url,
"status": status,
"duration_ms": elapsed_ms,
});
if let Err(e) = &result {
sample["error"] = serde_json::Value::String(e.to_string());
}
samples.push(sample);
}
}
let mut out = serde_json::json!({
"iterations": iterations,
"samples": samples,
"failures": total_failures,
"has_failures": total_failures > 0,
"thresholds": threshold_verdict(thresholds, iterations, total_failures),
});
if !unsupported_errors.is_empty() {
out["unsupported_auth"] = serde_json::Value::Array(
unsupported_errors
.into_iter()
.map(serde_json::Value::String)
.collect(),
);
}
out
}
async fn run_load_streaming(
sock: &mut TcpStream,
state: &AgentState,
scenario: &tropel_sdk::scenario::Scenario,
iterations: u64,
thresholds: &std::collections::HashMap<String, String>,
) -> tropel_sdk::Result<()> {
let head = "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nTransfer-Encoding: chunked\r\nConnection: close\r\n\r\n";
sock.write_all(head.as_bytes())
.await
.map_err(TropelError::Io)?;
let mut total_failures = 0u64;
for it in 0..iterations {
let mut batch: Vec<serde_json::Value> = Vec::new();
for item in &scenario.items {
let Some(request) = item.request.as_ref() else {
continue;
};
let signer_opt = match &request.auth {
Some(auth) => match state.client.get_signer(auth) {
Ok(s) => s,
Err(e) => {
total_failures += 1;
batch.push(serde_json::json!({
"metric": "http_reqs",
"iteration": it,
"url": request.url,
"status": 0,
"duration_ms": 0.0,
"error": e.to_string(),
}));
continue;
}
},
None => None,
};
let start = Instant::now();
let result = state.client.execute(request, signer_opt.as_deref()).await;
let elapsed_ms = start.elapsed().as_millis() as f64;
let (status, ok) = match &result {
Ok(resp) => (resp.status_code, (200..400).contains(&resp.status_code)),
Err(_) => (0, false),
};
if !ok {
total_failures += 1;
}
let mut sample = serde_json::json!({
"metric": "http_reqs",
"iteration": it,
"url": request.url,
"status": status,
"duration_ms": elapsed_ms,
});
if let Err(e) = &result {
sample["error"] = serde_json::Value::String(e.to_string());
}
batch.push(sample);
}
let chunk = serde_json::json!({
"iteration": it,
"samples": batch,
"failures": total_failures,
})
.to_string();
write_chunk(sock, &chunk).await?;
}
let verdict = serde_json::json!({
"done": true,
"iterations": iterations,
"failures": total_failures,
"has_failures": total_failures > 0,
"thresholds": threshold_verdict(thresholds, iterations, total_failures),
})
.to_string();
write_chunk(sock, &verdict).await?;
sock.write_all(b"0\r\n\r\n").await.map_err(TropelError::Io)
}
async fn write_chunk(sock: &mut TcpStream, data: &str) -> tropel_sdk::Result<()> {
let size = format!("{:x}\r\n", data.len());
sock.write_all(size.as_bytes())
.await
.map_err(TropelError::Io)?;
sock.write_all(data.as_bytes())
.await
.map_err(TropelError::Io)?;
sock.write_all(b"\r\n").await.map_err(TropelError::Io)
}
fn threshold_verdict(
thresholds: &std::collections::HashMap<String, String>,
iterations: u64,
failures: u64,
) -> serde_json::Value {
let mut results: Vec<serde_json::Value> = Vec::new();
let mut all_passed = true;
for (name, expr) in thresholds {
let passed = match eval_threshold(expr, iterations, failures) {
Ok(p) => p,
Err(e) => {
all_passed = false;
results.push(serde_json::json!({
"name": name, "expression": expr, "passed": false,
"error": e, "actual": null, "threshold": null,
}));
continue;
}
};
if !passed {
all_passed = false;
}
let actual = if expr.starts_with("http_req_failed") && iterations > 0 {
failures as f64 / iterations as f64
} else {
iterations as f64
};
results.push(serde_json::json!({
"name": name, "expression": expr, "passed": passed,
"actual": actual,
"threshold": expr.split_whitespace().nth(2).and_then(|v| v.parse::<f64>().ok()),
}));
}
results.sort_by(|a, b| a["name"].as_str().cmp(&b["name"].as_str()));
serde_json::json!({
"results": results,
"passed": all_passed,
})
}
async fn respond(sock: &mut TcpStream, status: u16, body: &str) -> tropel_sdk::Result<()> {
respond_raw(sock, status, body, None).await
}
async fn respond_raw_cors(
sock: &mut TcpStream,
cors: Option<&str>,
status: u16,
body: &str,
) -> tropel_sdk::Result<()> {
respond_raw(sock, status, body, cors).await
}
async fn respond_raw(
sock: &mut TcpStream,
status: u16,
body: &str,
extra_headers: Option<&str>,
) -> tropel_sdk::Result<()> {
let status_text = match status {
200 => "OK",
204 => "No Content",
400 => "Bad Request",
401 => "Unauthorized",
403 => "Forbidden",
404 => "Not Found",
429 => "Too Many Requests",
_ => "Error",
};
let extra = extra_headers.unwrap_or("");
let resp = format!(
"HTTP/1.1 {status} {status_text}\r\nContent-Type: application/json\r\nContent-Length: {}\r\n{extra}Connection: close\r\n\r\n{body}",
body.len()
);
sock.write_all(resp.as_bytes())
.await
.map_err(TropelError::Io)
}
#[derive(serde::Deserialize)]
struct BatchResolveItem {
template: String,
#[serde(default)]
mode: Option<String>,
#[serde(default)]
deep: Option<bool>,
}
#[derive(serde::Deserialize)]
struct AgentAssertionSpec {
#[serde(default)]
name: Option<String>,
target: String,
operator: String,
#[serde(default)]
expected: serde_json::Value,
}
async fn read_json_body(
sock: &mut TcpStream,
content_length: usize,
max: usize,
prefetched: &[u8],
) -> Result<Option<serde_json::Value>, TropelError> {
Ok(serde_json::from_slice(&read_body(sock, content_length, max, prefetched).await?).ok())
}
async fn read_body(
sock: &mut TcpStream,
content_length: usize,
max: usize,
prefetched: &[u8],
) -> Result<Vec<u8>, TropelError> {
let want = content_length.min(max);
let mut body = prefetched.to_vec();
body.truncate(want);
if body.len() < want {
let mut rest = vec![0u8; want - body.len()];
sock.read_exact(&mut rest).await.map_err(TropelError::Io)?;
body.extend_from_slice(&rest);
}
Ok(body)
}
fn error_body(message: &str) -> String {
serde_json::json!({ "error": message }).to_string()
}
fn sign_with_scheme(
scheme: &str,
p: &serde_json::Value,
) -> Result<Vec<tropel_auth::builders::HeaderOut>, String> {
let s = |k: &str| p.get(k).and_then(|v| v.as_str()).unwrap_or("").to_string();
let opt = |k: &str| {
p.get(k)
.and_then(|v| v.as_str())
.filter(|v| !v.is_empty())
.map(str::to_string)
};
let pairs = |k: &str| -> Vec<(String, String)> {
p.get(k)
.and_then(|v| serde_json::from_value(v.clone()).ok())
.unwrap_or_default()
};
match scheme {
"digest" => {
let challenge = s("wwwAuthenticate");
let Some(c) = tropel_auth::builders::find_digest_challenge(&challenge) else {
return Err("the WWW-Authenticate header carries no Digest challenge".to_string());
};
let get = |k: &str| c.get(k).map(String::as_str);
Ok(vec![tropel_auth::builders::digest_build_authorization(
&tropel_auth::builders::DigestBuildParams {
username: &s("username"),
password: &s("password"),
method: &s("method"),
uri: &s("uri"),
realm: get("realm").unwrap_or(""),
nonce: get("nonce").unwrap_or(""),
nc: p.get("nc").and_then(|v| v.as_u64()).unwrap_or(1),
cnonce: &s("cnonce"),
qop: get("qop"),
algorithm: get("algorithm"),
opaque: get("opaque"),
},
)])
}
"hawk" => Ok(vec![tropel_auth::builders::hawk_build_header(
&tropel_auth::builders::HawkBuildParams {
method: &s("method"),
resource: &s("resource"),
host: &s("host"),
port: p.get("port").and_then(|v| v.as_u64()).unwrap_or(443) as u16,
id: &s("id"),
key: &s("key"),
algorithm: opt("algorithm").as_deref(),
ts: &s("ts"),
nonce: &s("nonce"),
ext: &s("ext"),
},
)]),
"awsSigV4" => {
let host = s("host");
let region = opt("region").unwrap_or_else(|| "us-east-1".to_string());
let service =
opt("service").unwrap_or_else(|| tropel_auth::builders::default_service(&host));
let signing_service = tropel_auth::builders::signing_name(&service);
let path = s("path");
let canonical_uri = tropel_auth::builders::sigv4_canonical_uri(&path, &service);
let secret = s("secretKey");
let date_stamp = s("dateStamp");
let key = tropel_auth::builders::derive_signing_key(
&secret,
&date_stamp,
®ion,
signing_service,
);
let body = match opt("bodyBase64") {
Some(b64) => Some(
base64_decode(&b64).map_err(|e| format!("bodyBase64 is not base64: {e}"))?,
),
None => None,
};
let headers = pairs("headers");
let out = tropel_auth::builders::aws_sigv4_build_headers(
&tropel_auth::builders::AwsSigV4BuildParams {
method: &s("method"),
path: &path,
query: &s("query"),
host: &tropel_auth::builders::bracket_host(&host),
headers: &headers,
body: body.as_deref(),
access_key: &s("accessKey"),
secret_key: &secret,
session_token: opt("sessionToken").as_deref(),
region: ®ion,
service: &service,
amz_date: &s("amzDate"),
date_stamp: &date_stamp,
},
&canonical_uri,
signing_service,
&key,
);
Ok(out.headers)
}
"oauth1" => {
let base_uri = tropel_auth::builders::oauth1_base_uri(
&s("scheme"),
&s("host"),
p.get("port").and_then(|v| v.as_u64()).map(|v| v as u16),
&s("path"),
);
let mut params = pairs("queryParams");
if let Some(form) = opt("formBody") {
params.extend(tropel_auth::builders::parse_form(form.as_bytes()));
}
let method = s("signatureMethod");
tropel_auth::builders::oauth1_build_header(&tropel_auth::builders::OAuth1BuildParams {
method: &s("method"),
base_uri: &base_uri,
request_params: ¶ms,
consumer_key: &s("consumerKey"),
consumer_secret: &s("consumerSecret"),
token: opt("token").as_deref(),
token_secret: opt("tokenSecret").as_deref(),
signature_method: &method,
nonce: &s("nonce"),
timestamp: &s("timestamp"),
})
.map(|o| vec![o.header])
.ok_or_else(|| {
format!(
"unsupported OAuth1 signature_method '{method}' — supported: {}",
tropel_auth::builders::OAUTH1_SIGNATURE_METHODS.join(", ")
)
})
}
"akamai-edgegrid" => {
let body = opt("body").map(|b| b.into_bytes());
let headers_to_sign: Vec<String> = p
.get("headersToSign")
.or_else(|| p.get("headers_to_sign"))
.and_then(|v| serde_json::from_value(v.clone()).ok())
.unwrap_or_default();
let params = tropel_auth::edgegrid::EdgeGridBuildParams {
method: s("method"),
url: s("url"),
headers_to_sign,
body,
access_token: s("accessToken"),
client_token: s("clientToken"),
client_secret: s("clientSecret"),
nonce: opt("nonce"),
timestamp: opt("timestamp"),
max_body: p
.get("maxBody")
.or_else(|| p.get("max_body"))
.and_then(|v| v.as_u64())
.map(|n| n as usize)
.unwrap_or(tropel_auth::edgegrid::DEFAULT_MAX_BODY),
};
let signed_headers = pairs("headers");
tropel_auth::edgegrid::edgegrid_build_header(¶ms, &signed_headers)
.map(|value| {
vec![tropel_auth::builders::HeaderOut {
name: "Authorization".to_string(),
value,
}]
})
.map_err(|e| e.to_string())
}
"wsse" => {
let signed = tropel_auth::oauth::sign_wsse(&tropel_auth::oauth::WsseParams {
username: s("username"),
password: s("password"),
nonce: s("nonce"),
created: s("created"),
})
.map_err(|e| e.to_string())?;
Ok(vec![
tropel_auth::builders::HeaderOut {
name: "X-WSSE".to_string(),
value: signed.authorization,
},
tropel_auth::builders::HeaderOut {
name: "Authorization".to_string(),
value: "WSSE profile=\"UsernameToken\"".to_string(),
},
])
}
other => Err(format!(
"unknown auth scheme '{other}' — supported: {}",
AUTH_SIGN_SCHEMES.join(", ")
)),
}
}
pub const AUTH_SIGN_SCHEMES: &[&str] = &[
"digest",
"hawk",
"awsSigV4",
"oauth1",
"akamai-edgegrid",
"wsse",
];
fn base64_encode(bytes: &[u8]) -> String {
use base64::Engine as _;
base64::engine::general_purpose::STANDARD.encode(bytes)
}
fn base64_decode(s: &str) -> Result<Vec<u8>, String> {
use base64::Engine as _;
base64::engine::general_purpose::STANDARD
.decode(s)
.map_err(|e| e.to_string())
}
#[derive(Default)]
struct ScriptScopes {
environment: HashMap<String, String>,
collection: HashMap<String, serde_json::Value>,
globals: HashMap<String, serde_json::Value>,
variables: HashMap<String, serde_json::Value>,
}
impl ScriptScopes {
fn from_payload(payload: &serde_json::Value) -> Self {
fn map_of(payload: &serde_json::Value, key: &str) -> HashMap<String, serde_json::Value> {
payload
.get(key)
.and_then(|v| v.as_object())
.map(|o| o.iter().map(|(k, v)| (k.clone(), v.clone())).collect())
.unwrap_or_default()
}
Self {
environment: payload
.get("environment")
.and_then(|e| serde_json::from_value(e.clone()).ok())
.unwrap_or_default(),
collection: map_of(payload, "collectionVariables"),
globals: map_of(payload, "globals"),
variables: map_of(payload, "variables"),
}
}
}
struct ScriptHost {
http: tropel_http::HttpClient,
callbacks: Option<Arc<RunCallbacks>>,
cookies: Option<Arc<ScriptCookies>>,
}
async fn run_script_once(
code: &str,
scopes: ScriptScopes,
request: Option<TropelRequest>,
response: Option<tropel_sdk::types::Response>,
sandbox: tropel_sandbox::config::SandboxConfig,
host: ScriptHost,
) -> Result<serde_json::Value, String> {
let ScriptHost {
http,
callbacks,
cookies,
} = host;
let request_url = request.as_ref().map(|r| r.url.clone()).unwrap_or_default();
let mut ctx = tropel_js::JsContext::new(None, Some(std::time::Duration::from_secs(10)))
.await
.map_err(|e| format!("js context: {e:?}"))?;
ctx.eval(&sandbox.render_js_preamble())
.await
.map_err(|e| format!("sandbox config preamble: {e:?}"))?;
ctx.eval(include_str!("../js/shared/deep-equal.js"))
.await
.map_err(|e| format!("deep-equal shim: {e:?}"))?;
for entry in crate::js_bootstrap::ShimBundle::default().0 {
ctx.eval(&entry.1)
.await
.map_err(|e| format!("{} shim: {e:?}", entry.0))?;
}
let state = tropel_sandbox::state::SharedPmState::default();
{
let mut st = state.lock().unwrap_or_else(|e| e.into_inner());
st.environment = scopes.environment;
st.collection_vars = std::sync::Arc::new(scopes.collection);
st.globals = std::sync::Arc::new(scopes.globals);
st.local_vars = scopes.variables;
st.request = request;
st.response = response;
}
let sends = Arc::new(ScriptSendLog::default());
let vu_client = tropel_http::VuCookieClient::new(http);
let base_client: std::sync::Arc<dyn tropel_sdk::traits::DriverHttpClient> =
crate::vu_loop::DriverHttpClientImpl::new_arc(vu_client);
let driver_client: std::sync::Arc<dyn tropel_sdk::traits::DriverHttpClient> =
Arc::new(RecordingHttpClient {
inner: base_client,
log: sends.clone(),
});
tropel_sandbox::bindings::trp::TrpBridge::with_http_client(state.clone(), driver_client)
.install(&mut ctx)
.map_err(|e| format!("bridge install: {e:?}"))?;
if let Some(cookies) = cookies.clone() {
let url_for_current = request_url.clone();
let installed: Result<(), String> = ctx.with_ctx(|rq| {
let globals = rq.globals();
let mut fail: Option<String> = None;
let u = url_for_current.clone();
if let Err(e) = globals.set(
"__tropel_cookies_current_url",
rquickjs::function::Func::from(move || -> String { u.clone() }),
) {
fail = Some(e.to_string());
}
let c = cookies.clone();
if let Err(e) = globals.set(
"__tropel_cookies_all",
rquickjs::function::Func::from(move |url: String| -> String {
let jar = c.lock_jar();
let matching: Vec<&serde_json::Value> = jar
.iter()
.filter(|ck| cookie_matches_url(ck, &url))
.collect();
serde_json::to_string(&matching).unwrap_or_else(|_| "[]".to_string())
}),
) {
fail = Some(e.to_string());
}
let c = cookies.clone();
if let Err(e) = globals.set(
"__tropel_cookies_set",
rquickjs::function::Func::from(move |url: String, cookie_json: String| {
let Ok(cookie) = serde_json::from_str::<serde_json::Value>(&cookie_json) else {
return;
};
let name = cookie
.get("key")
.or_else(|| cookie.get("name"))
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
{
let mut jar = c.lock_jar();
jar.retain(|ck| {
ck.get("key").and_then(|v| v.as_str()).unwrap_or("") != name
});
jar.push(cookie.clone());
}
c.record(serde_json::json!({ "op": "set", "url": url, "cookie": cookie }));
}),
) {
fail = Some(e.to_string());
}
let c = cookies.clone();
if let Err(e) = globals.set(
"__tropel_cookies_delete",
rquickjs::function::Func::from(move |url: String, name: String| {
{
let mut jar = c.lock_jar();
jar.retain(|ck| {
ck.get("key").and_then(|v| v.as_str()).unwrap_or("") != name
});
}
c.record(serde_json::json!({ "op": "delete", "url": url, "name": name }));
}),
) {
fail = Some(e.to_string());
}
let c = cookies.clone();
if let Err(e) = globals.set(
"__tropel_cookies_clear",
rquickjs::function::Func::from(move |url: String| {
c.lock_jar().clear();
c.record(serde_json::json!({ "op": "clear", "url": url }));
}),
) {
fail = Some(e.to_string());
}
match fail {
Some(e) => Err(e),
None => Ok(()),
}
});
installed.map_err(|e| format!("cookie bridge: {e}"))?;
}
if let Some(cb) = callbacks.clone() {
let sends_in_cb = sends.clone();
let installed: Result<(), String> = ctx.with_ctx(|rq| {
rq.globals()
.set(
"__tropel_trp_run_request",
rquickjs::function::Func::from(move |path: String| -> String {
let sends_for_run = sends_in_cb.clone();
let requested_path = path.clone();
let call_id = cb
.next_id
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let (tx, rx) = std::sync::mpsc::channel::<String>();
{
let mut replies = match cb.replies.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
};
replies.insert(call_id, tx);
}
{
let mut pending = match cb.pending.lock() {
Ok(g) => g,
Err(e) => e.into_inner(),
};
pending.push_back(PendingHostCall {
call_id,
kind: "runRequest",
path,
});
}
cb.ready.notify_waiters();
let started = std::time::Instant::now();
match tokio::task::block_in_place(|| {
rx.recv_timeout(std::time::Duration::from_secs(30))
}) {
Ok(json) => {
let parsed: serde_json::Value =
serde_json::from_str(&json).unwrap_or(serde_json::Value::Null);
sends_for_run.record(serde_json::json!({
"source": "runRequest",
"kind": "collection",
"requestName": requested_path,
"method": parsed.get("method")
.and_then(|v| v.as_str()).unwrap_or(""),
"url": parsed.get("url")
.and_then(|v| v.as_str()).unwrap_or(""),
"status": parsed.get("status")
.and_then(|v| v.as_u64()).unwrap_or(0),
"responseTime": started.elapsed().as_millis() as u64,
"error": parsed.get("error").cloned(),
}));
json
}
Err(_) => {
if let Ok(mut r) = cb.replies.lock() {
r.remove(&call_id);
}
r#"{"error":"runRequest timed out: the host did not answer within 30s"}"#
.to_string()
}
}
}),
)
.map_err(|e| e.to_string())
});
installed.map_err(|e| format!("run_request bridge: {e}"))?;
}
let wrapped = format!(
"globalThis.__tropel_script_error = null;\n\
(async function () {{ try {{\n{code}\n}} catch (e) {{ \
globalThis.__tropel_script_error = \
(e && e.message ? e.message : String(e)) + \
(e && e.stack ? \"\\n\" + e.stack : \"\"); \
}} }})();"
);
let script_error = match ctx.eval(&wrapped).await {
Ok(_) => {
let raw = ctx
.eval("globalThis.__tropel_script_error === null ? \"\" : String(globalThis.__tropel_script_error)")
.await
.unwrap_or_default();
if raw.is_empty() {
None
} else {
Some(raw)
}
}
Err(e) => Some(format!("{e:?}")),
};
let st = state.lock().unwrap_or_else(|e| e.into_inner());
let tests: Vec<serde_json::Value> = st
.samples
.iter()
.filter(|s| s.metric == "checks")
.map(|s| {
serde_json::json!({
"name": s.tags.get("check").unwrap_or_default(),
"passed": s.value != 0.0,
})
})
.collect();
Ok(serde_json::json!({
"tests": tests,
"assertions": {
"total": st.assertions.total,
"passed": st.assertions.passed,
"failed": st.assertions.failed,
},
"environment": st.environment,
"collectionVariables": *st.collection_vars,
"globals": *st.globals,
"variables": st.local_vars,
"request": st.request,
"scriptSends": sends.drain(),
"cookieOps": cookies
.as_ref()
.map(|c| serde_json::Value::Array(c.lock_ops().clone()))
.unwrap_or(serde_json::Value::Null),
"scriptError": script_error,
}))
}
fn as_json<T: serde::Serialize>(v: &T) -> Result<serde_json::Value, String> {
serde_json::to_value(v).map_err(|e| e.to_string())
}
fn oauth2_dispatch(op: &str, p: &serde_json::Value) -> Result<serde_json::Value, String> {
use tropel_auth::oauth;
let as_str = |k: &str| p.get(k).and_then(|v| v.as_str()).unwrap_or("");
let opt = |k: &str| {
p.get(k)
.and_then(|v| v.as_str())
.filter(|s| !s.is_empty())
.map(str::to_string)
};
match op {
"buildAuthorizeUrl" => {
let params: oauth::AuthorizeParams =
serde_json::from_value(p.clone()).map_err(|e| e.to_string())?;
as_json(&oauth::build_authorize_url(¶ms).map_err(|e| e.to_string())?)
}
"buildTokenRequest" => {
let params: oauth::TokenRequestParams =
serde_json::from_value(p.clone()).map_err(|e| e.to_string())?;
as_json(&oauth::build_token_request(¶ms).map_err(|e| e.to_string())?)
}
"parseTokenResponse" => {
as_json(&oauth::parse_token_response(as_str("body")).map_err(|e| e.to_string())?)
}
"attachToken" => {
let placement = match p.get("placement").and_then(|v| v.as_str()) {
Some("header") | None => oauth::TokenPlacement::Header,
Some("query") => oauth::TokenPlacement::Query,
Some(other) => {
return Err(format!(
"unknown token placement '{other}' — expected header or query"
))
}
};
as_json(&oauth::attach_token(
as_str("token"),
opt("tokenType").as_deref(),
placement,
opt("headerPrefix").as_deref(),
opt("queryKey").as_deref(),
))
}
"decodeJwt" => as_json(&oauth::decode_jwt(as_str("token")).map_err(|e| e.to_string())?),
"jwtExpiresAt" => {
let exp = oauth::jwt_expires_at(as_str("token")).map_err(|e| e.to_string())?;
Ok(serde_json::json!({ "expiresAt": exp }))
}
"signJwt" => {
let algorithm = match p.get("algorithm").and_then(|v| v.as_str()) {
Some("HS256") | None => oauth::JwtAlgorithm::Hs256,
Some("HS384") => oauth::JwtAlgorithm::Hs384,
Some("HS512") => oauth::JwtAlgorithm::Hs512,
Some(other) => {
return Err(format!(
"unsupported JWT algorithm '{other}' — supported: HS256, HS384, HS512"
))
}
};
let payload = p.get("payload").cloned().unwrap_or_default();
let header = p.get("header").cloned().filter(|h| !h.is_null());
let token = oauth::sign_jwt(header.as_ref(), &payload, algorithm, as_str("secret"))
.map_err(|e| e.to_string())?;
Ok(serde_json::json!({ "token": token }))
}
"wsseSign" => {
let params: oauth::WsseParams =
serde_json::from_value(p.clone()).map_err(|e| e.to_string())?;
as_json(&oauth::sign_wsse(¶ms).map_err(|e| e.to_string())?)
}
"codeChallengeS256" => Ok(serde_json::json!({
"codeChallenge": oauth::code_challenge_s256(as_str("verifier")),
"codeChallengeMethod": "S256",
})),
other => Err(format!(
"unknown oauth2 op '{other}' — supported: buildAuthorizeUrl, buildTokenRequest, \
parseTokenResponse, attachToken, decodeJwt, jwtExpiresAt, signJwt, wsseSign, \
codeChallengeS256"
)),
}
}
async fn execute_single(state: &AgentState, req: &serde_json::Value) -> serde_json::Value {
let method = req.get("method").and_then(|m| m.as_str()).unwrap_or("GET");
let url = req.get("url").and_then(|u| u.as_str()).unwrap_or("");
let follow = req
.get("follow_redirects")
.and_then(|f| f.as_bool())
.unwrap_or(true);
let headers: Vec<(String, String)> = match req.get("headers") {
Some(serde_json::Value::Array(rows)) => rows
.iter()
.filter_map(|row| match row {
serde_json::Value::Array(pair) if pair.len() == 2 => Some((
pair[0].as_str()?.to_string(),
pair[1].as_str().unwrap_or("").to_string(),
)),
serde_json::Value::Object(o) => Some((
o.get("name")?.as_str()?.to_string(),
o.get("value")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string(),
)),
_ => None,
})
.collect(),
Some(serde_json::Value::Object(o)) => o
.iter()
.map(|(k, v)| (k.clone(), v.as_str().unwrap_or("").to_string()))
.collect(),
_ => Vec::new(),
};
let certificate: Option<tropel_sdk::types::CertificateConfig> = req
.get("certificate")
.and_then(|v| serde_json::from_value(v.clone()).ok());
let proxy: Option<tropel_sdk::types::ProxyConfig> = req
.get("proxy")
.and_then(|v| serde_json::from_value(v.clone()).ok());
let host: Option<String> = req.get("host").and_then(|h| h.as_str()).map(str::to_string);
let cookies: Vec<tropel_sdk::types::RequestCookie> = req
.get("cookies")
.and_then(|v| serde_json::from_value(v.clone()).ok())
.unwrap_or_default();
let timeout = req
.get("timeout_ms")
.and_then(|t| t.as_u64())
.map(std::time::Duration::from_millis);
let query_params: HashMap<String, String> = req
.get("query_params")
.and_then(|v| serde_json::from_value(v.clone()).ok())
.unwrap_or_default();
let method_parsed = Method::parse(method).unwrap_or(Method::GET);
let auth: Option<tropel_sdk::types::AuthConfig> = req
.get("auth")
.and_then(|v| serde_json::from_value(v.clone()).ok());
let request = TropelRequest {
url: url.to_string(),
method: method_parsed,
headers,
query_params,
body: req
.get("body")
.and_then(|b| b.as_str())
.map(|s| Body::Raw(s.to_string())),
auth: auth.clone(),
certificate,
proxy,
follow_redirects: follow,
host,
cookies,
timeout,
response_type: ResponseType::Text,
};
let signer_opt = match &request.auth {
Some(a) => match state.client.get_signer(a) {
Ok(s) => s,
Err(e) => {
return serde_json::json!({
"status": 0,
"status_text": "Unsupported Auth",
"headers": {},
"body": "",
"timings": {
"blocked": 0.0, "dns": 0.0, "connecting": 0.0, "tls_handshaking": 0.0,
"sending": 0.0, "waiting": 0.0, "receiving": 0.0, "duration": 0.0,
},
"error": e.to_string(),
});
}
},
None => None,
};
let start = Instant::now();
let result = state.client.execute(&request, signer_opt.as_deref()).await;
let elapsed_ms = start.elapsed().as_millis() as f64;
match result {
Ok(resp) => {
let waiting = resp
.timings
.as_ref()
.map(|t| t.waiting.as_millis() as f64)
.unwrap_or(0.0);
let receiving = resp
.timings
.as_ref()
.map(|t| t.receiving.as_millis() as f64)
.unwrap_or(0.0);
let (body, body_encoding) = match std::str::from_utf8(&resp.body) {
Ok(text) => (text.to_string(), "utf8"),
Err(_) => (base64_encode(&resp.body), "base64"),
};
serde_json::json!({
"status": resp.status_code,
"status_text": resp.status_text,
"headers": resp.headers,
"body": body,
"bodyEncoding": body_encoding,
"timings": {
"blocked": 0.0, "dns": 0.0, "connecting": 0.0, "tls_handshaking": 0.0,
"sending": 0.0, "waiting": waiting, "receiving": receiving,
"duration": elapsed_ms,
},
"error": null,
})
}
Err(e) => {
serde_json::json!({
"status": 0,
"status_text": "Transport Error",
"headers": {},
"body": "",
"timings": {
"blocked": 0.0, "dns": 0.0, "connecting": 0.0, "tls_handshaking": 0.0,
"sending": 0.0, "waiting": 0.0, "receiving": 0.0, "duration": elapsed_ms,
},
"error": e.to_string(),
})
}
}
}
struct RateLimiter {
window_start: Instant,
count: u64,
}
impl RateLimiter {
fn new() -> Self {
Self {
window_start: Instant::now(),
count: 0,
}
}
fn allow(&mut self) -> Result<(), ()> {
if self.window_start.elapsed().as_secs() >= 1 {
self.window_start = Instant::now();
self.count = 0;
}
self.count += 1;
if self.count > RATE_LIMIT_PER_SEC {
return Err(());
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn refuses_non_loopback_bind_is_enforced() {
let ip: IpAddr = "0.0.0.0".parse().unwrap();
assert!(!ip.is_loopback(), "0.0.0.0 must be rejected");
let ip2: IpAddr = "127.0.0.1".parse().unwrap();
assert!(ip2.is_loopback(), "127.0.0.1 must be accepted");
}
#[test]
fn threshold_verdict_evaluates_and_verdicts() {
use std::collections::HashMap;
let mut t = HashMap::new();
t.insert("reqs_ok".into(), "http_reqs < 100".into());
t.insert("fail_rate_ok".into(), "http_req_failed <= 0.1".into());
let v = threshold_verdict(&t, 10, 1);
assert!(v["passed"].as_bool().unwrap(), "all must pass: {v}");
let results = v["results"].as_array().unwrap();
assert_eq!(results.len(), 2);
assert!(results.iter().all(|r| r["passed"].as_bool().unwrap()));
let v2 = threshold_verdict(&t, 200, 1);
assert!(!v2["passed"].as_bool().unwrap(), "reqs threshold must fail");
let reqs = v2["results"]
.as_array()
.unwrap()
.iter()
.find(|r| r["name"] == "reqs_ok")
.unwrap();
assert!(!reqs["passed"].as_bool().unwrap());
let mut bad = HashMap::new();
bad.insert("bogus".into(), "http_reqs".into());
let v3 = threshold_verdict(&bad, 10, 0);
assert!(!v3["passed"].as_bool().unwrap(), "malformed must fail");
assert!(
v3["results"][0]["error"].as_str().is_some(),
"the malformed threshold must report its error"
);
}
#[tokio::test]
async fn the_rules_endpoints_answer_over_the_socket() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let port = listener.local_addr().expect("addr").port();
let state = Arc::new(AgentState {
token: None,
client: tropel_http::HttpClient::new(&tropel_http::config::HttpConfig::default())
.expect("http client"),
runs: std::sync::Mutex::new(HashMap::new()),
allowed_origins: vec![],
});
tokio::spawn(async move {
while let Ok((mut sock, _)) = listener.accept().await {
let st = state.clone();
tokio::spawn(async move {
let _ = handle_connection(&mut sock, st).await;
});
}
});
let post = |path: &'static str, body: String| async move {
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
let req = format!(
"POST {path} HTTP/1.1\r\nHost: localhost\r\nContent-Length: {}\r\n\r\n{body}",
body.len()
);
s.write_all(req.as_bytes()).await.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
String::from_utf8_lossy(&out).to_string()
};
let raw = post(
"/resolve",
serde_json::json!({
"template": "{{base}}/v1", "variables": {"base": "https://x.test"},
"mode": "plain", "deep": true
})
.to_string(),
)
.await;
assert!(raw.contains("https://x.test/v1"), "{raw}");
let raw = post(
"/resolve",
serde_json::json!({"template": "x", "variables": {}, "mode": "jsonn"}).to_string(),
)
.await;
assert!(raw.starts_with("HTTP/1.1 400"), "{raw}");
assert!(raw.contains("unknown mode"), "{raw}");
let raw = post(
"/assert",
serde_json::json!({
"response": {
"status": 200, "status_text": "OK",
"headers": [["Content-Type", "application/json"]],
"body": "{\"count\":2}", "response_time": 5.0, "size": 12,
"cookies": []
},
"assertions": [
{"name": "ok", "target": "status", "operator": "eq", "expected": 200},
{"target": "json.count", "operator": "eq", "expected": 99}
]
})
.to_string(),
)
.await;
assert!(raw.contains(r#""name":"ok""#), "{raw}");
assert!(raw.contains(r#""passed":true"#), "{raw}");
assert!(
raw.contains("expected target json.count equals 99"),
"{raw}"
);
let raw = post(
"/assert",
serde_json::json!({
"response": {
"status": 200, "status_text": "OK", "headers": [],
"body": "abc", "response_time": 1.0, "size": 3, "cookies": []
},
"assertions": [{"target": "body", "operator": "matches", "expected": "^a"}]
})
.to_string(),
)
.await;
assert!(raw.contains("regex matcher"), "{raw}");
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
s.write_all(b"GET /operators HTTP/1.1\r\nHost: localhost\r\n\r\n")
.await
.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
let raw = String::from_utf8_lossy(&out).to_string();
assert!(raw.contains(r#""name":"eq""#), "{raw}");
assert!(raw.contains(r#""arity":"unary""#), "{raw}");
let raw = post(
"/variables/dynamic",
serde_json::json!({"template": "{{$guid}}|{{$guid}}"}).to_string(),
)
.await;
let body = raw.split("\r\n\r\n").nth(1).unwrap_or_default().to_string();
let parsed: serde_json::Value = serde_json::from_str(&body).expect("json body");
let value = parsed["value"].as_str().expect("value");
let (first, second) = value.split_once('|').expect("two guids");
assert_ne!(
first, second,
"each occurrence generates a fresh value: {value}"
);
assert_eq!(first.len(), 36, "a v4 GUID, not a placeholder: {value}");
let raw = post(
"/variables/dynamic",
serde_json::json!({"template": "{{base}}/{{$timestamp}}"}).to_string(),
)
.await;
assert!(
raw.contains("{{base}}"),
"a plain variable must survive the dynamic pass untouched: {raw}"
);
let raw = post(
"/variables/dynamic/batch",
serde_json::json!({
"items": [
{"template": "{{$guid}}"},
{"template": "no tokens here"},
{"template": "{{$timestamp}}"}
]
})
.to_string(),
)
.await;
let body = raw.split("\r\n\r\n").nth(1).unwrap_or_default().to_string();
let parsed: serde_json::Value = serde_json::from_str(&body).expect("json body");
let items = parsed["items"].as_array().expect("items");
assert_eq!(items.len(), 3, "one output per input, always: {body}");
assert_eq!(
items[1]["value"], "no tokens here",
"a template with nothing to resolve comes back unchanged: {body}"
);
assert_ne!(
items[0]["value"], items[2]["value"],
"different tokens, different values: {body}"
);
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
s.write_all(b"GET /constants HTTP/1.1\r\nHost: localhost\r\n\r\n")
.await
.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
let raw = String::from_utf8_lossy(&out).to_string();
let body = raw.split("\r\n\r\n").nth(1).unwrap_or_default().to_string();
let parsed: serde_json::Value = serde_json::from_str(&body).expect("json body");
assert_eq!(
parsed["maxVariableResolutionPasses"],
tropel_variables::MAX_VARIABLE_RESOLUTION_PASSES,
"the cap must come from the resolver: {body}"
);
let vars = parsed["predefinedVariables"]
.as_array()
.expect("predefinedVariables array");
assert_eq!(
vars.len(),
tropel_variables::PREDEFINED_VARIABLE_META.len(),
"the whole catalogue, not a subset: {body}"
);
assert!(
vars.iter()
.any(|v| v["name"] == "$guid" && v["description"].is_string()),
"names AND descriptions — the editor renders both: {body}"
);
}
#[tokio::test]
async fn the_resolution_corpus_agrees_over_the_socket() {
const CORPUS: &str = include_str!("../testdata/resolve-corpus.json");
fn generated_vars(kind: &str) -> serde_json::Value {
match kind {
"chain_longer_than_cap" => {
let cap = tropel_variables::MAX_VARIABLE_RESOLUTION_PASSES;
let mut m = serde_json::Map::new();
for i in 0..=cap {
m.insert(
format!("v{i}"),
serde_json::json!(format!("{{{{v{}}}}}", i + 1)),
);
}
m.insert(format!("v{}", cap + 1), serde_json::json!("end"));
serde_json::Value::Object(m)
}
other => panic!("unknown vars_generated kind: {other}"),
}
}
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let port = listener.local_addr().expect("addr").port();
let state = Arc::new(AgentState {
token: None,
client: tropel_http::HttpClient::new(&tropel_http::config::HttpConfig::default())
.expect("http client"),
runs: std::sync::Mutex::new(HashMap::new()),
allowed_origins: vec![],
});
tokio::spawn(async move {
while let Ok((mut sock, _)) = listener.accept().await {
let st = state.clone();
tokio::spawn(async move {
let _ = handle_connection(&mut sock, st).await;
});
}
});
let doc: serde_json::Value = serde_json::from_str(CORPUS).expect("corpus is valid JSON");
let cases = doc["cases"].as_array().expect("cases array");
assert!(!cases.is_empty(), "an empty corpus asserts nothing");
for case in cases {
let name = case["name"].as_str().expect("every case is named");
let template = case["template"].as_str().expect("template");
let mode = case["mode"].as_str().expect("mode");
let vars = match case.get("vars_generated") {
Some(kind) => generated_vars(kind.as_str().unwrap()),
None => case["vars"].clone(),
};
let body = serde_json::json!({
"variables": vars,
"items": [{"template": template, "mode": mode}],
})
.to_string();
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
let req = format!(
"POST /resolve/batch HTTP/1.1\r\nHost: localhost\r\nContent-Length: {}\r\n\r\n{body}",
body.len()
);
s.write_all(req.as_bytes()).await.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
let raw = String::from_utf8_lossy(&out).to_string();
let reply = raw.split("\r\n\r\n").nth(1).unwrap_or_default().to_string();
let parsed: serde_json::Value =
serde_json::from_str(&reply).unwrap_or_else(|e| panic!("{name}: {e} — {reply}"));
let item = &parsed["items"][0];
assert!(
item["error"].is_null(),
"{name}: the agent refused a corpus case: {item}"
);
let value = item["value"]
.as_str()
.unwrap_or_else(|| panic!("{name}: no value"));
if let Some(expected) = case.get("expect").and_then(|v| v.as_str()) {
assert_eq!(value, expected, "{name}");
}
if case.get("parses_as_json").and_then(|v| v.as_bool()) == Some(true) {
serde_json::from_str::<serde_json::Value>(value).unwrap_or_else(|e| {
panic!("{name}: result must stay parseable JSON: {e} — {value}")
});
}
if let Some(expected) = case.get("expect_hit_cap").and_then(|v| v.as_bool()) {
assert_eq!(item["hitCap"], expected, "{name}: hitCap");
}
if let Some(expected) = case.get("expect_unresolved") {
assert_eq!(&item["unresolved"], expected, "{name}: unresolved");
}
if let Some(required) = case
.get("expect_unresolved_contains")
.and_then(|v| v.as_array())
{
let got = item["unresolved"].as_array().unwrap();
for n in required {
assert!(
got.contains(n),
"{name}: unresolved must contain {n} — got {got:?}"
);
}
}
}
}
#[tokio::test]
async fn the_preflight_answers_cors_and_private_network_for_allowed_origins() {
async fn agent_with(origins: Vec<String>) -> u16 {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let port = listener.local_addr().expect("addr").port();
let state = Arc::new(AgentState {
token: Some("s3cret".into()),
client: tropel_http::HttpClient::new(&tropel_http::config::HttpConfig::default())
.expect("http client"),
runs: std::sync::Mutex::new(HashMap::new()),
allowed_origins: origins,
});
tokio::spawn(async move {
while let Ok((mut sock, _)) = listener.accept().await {
let st = state.clone();
tokio::spawn(async move {
let _ = handle_connection(&mut sock, st).await;
});
}
});
port
}
async fn preflight(port: u16, origin: &str, ask_pna: bool) -> String {
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
let pna = if ask_pna {
"Access-Control-Request-Private-Network: true\r\n"
} else {
""
};
let req = format!(
"OPTIONS /resolve/batch HTTP/1.1\r\nHost: localhost\r\nOrigin: {origin}\r\n\
Access-Control-Request-Method: POST\r\n{pna}\r\n"
);
s.write_all(req.as_bytes()).await.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
String::from_utf8_lossy(&out).to_string()
}
let allowed = "https://app.knockport.dev";
let port = agent_with(vec![allowed.to_string()]).await;
let raw = preflight(port, allowed, true).await;
assert!(raw.starts_with("HTTP/1.1 204"), "{raw}");
assert!(
raw.contains("Access-Control-Allow-Private-Network: true"),
"Chrome's PNA preflight must be answered or this breaks in Chrome ONLY: {raw}"
);
assert!(
raw.contains(&format!("Access-Control-Allow-Origin: {allowed}")),
"{raw}"
);
assert!(raw.contains("Vary: Origin"), "{raw}");
let raw = preflight(port, allowed, false).await;
assert!(raw.starts_with("HTTP/1.1 204"), "{raw}");
assert!(
!raw.contains("Access-Control-Allow-Private-Network"),
"the grant must not be volunteered: {raw}"
);
let raw = preflight(port, "https://evil.test", true).await;
assert!(raw.starts_with("HTTP/1.1 403"), "{raw}");
assert!(
raw.contains("--allow-origin"),
"the refusal must say how to fix it: {raw}"
);
assert!(
!raw.contains("Access-Control-Allow-Origin"),
"a refused origin must NOT be handed a grant: {raw}"
);
let closed = agent_with(vec![]).await;
let raw = preflight(closed, allowed, true).await;
assert!(raw.starts_with("HTTP/1.1 403"), "{raw}");
}
#[tokio::test]
async fn the_preflight_does_not_require_the_token() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let port = listener.local_addr().expect("addr").port();
let state = Arc::new(AgentState {
token: Some("s3cret".into()),
client: tropel_http::HttpClient::new(&tropel_http::config::HttpConfig::default())
.expect("http client"),
runs: std::sync::Mutex::new(HashMap::new()),
allowed_origins: vec!["https://app.knockport.dev".into()],
});
tokio::spawn(async move {
while let Ok((mut sock, _)) = listener.accept().await {
let st = state.clone();
tokio::spawn(async move {
let _ = handle_connection(&mut sock, st).await;
});
}
});
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
s.write_all(
b"OPTIONS /resolve HTTP/1.1\r\nHost: localhost\r\nOrigin: https://app.knockport.dev\r\n\r\n",
)
.await
.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
let raw = String::from_utf8_lossy(&out).to_string();
assert!(
raw.starts_with("HTTP/1.1 204"),
"a preflight must not be 401'd: {raw}"
);
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
s.write_all(
b"GET /version HTTP/1.1\r\nHost: localhost\r\nOrigin: https://app.knockport.dev\r\n\r\n",
)
.await
.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
let raw = String::from_utf8_lossy(&out).to_string();
assert!(raw.starts_with("HTTP/1.1 401"), "{raw}");
}
#[tokio::test]
async fn a_test_script_sees_the_response() {
let response = tropel_sdk::types::Response {
url: "https://api.test/v1".into(),
status_code: 201,
status_text: "Created".into(),
protocol: "HTTP/1.1".into(),
headers: HashMap::from([("content-type".to_string(), "application/json".to_string())]),
body: br#"{"id":7}"#.to_vec(),
text_cache: std::sync::OnceLock::new(),
json_cache: std::sync::OnceLock::new(),
response_time: std::time::Duration::from_millis(1),
timings: None,
cookies: Vec::new(),
size: 8,
redirects: Vec::new(),
request_body_size: 0,
};
let out = run_script_once(
"pm.test('status', () => pm.response.code === 201);\n pm.test('body', () => pm.response.json().id === 7);",
ScriptScopes::default(),
None,
Some(response),
tropel_sandbox::config::SandboxConfig::default(),
ScriptHost { http: test_http_client(), callbacks: None, cookies: None },
)
.await
.expect("the realm runs");
assert!(
out.get("scriptError").is_some_and(|e| e.is_null()),
"the script must not error: {out}"
);
let tests = out["tests"].as_array().expect("tests array");
assert_eq!(tests.len(), 2, "both checks must have run: {out}");
for t in tests {
assert_eq!(
t["passed"], true,
"a check asserting against the seeded response must PASS — a \
failing one means the response was not there: {t}"
);
}
}
#[tokio::test]
async fn a_pre_request_script_mutation_comes_back() {
let request = TropelRequest {
url: "https://api.test/v1".into(),
method: Method::GET,
headers: vec![("Accept".into(), "application/json".into())],
query_params: HashMap::new(),
body: None,
auth: None,
certificate: None,
proxy: None,
follow_redirects: true,
host: None,
cookies: Vec::new(),
timeout: None,
response_type: ResponseType::Text,
};
let out = run_script_once(
"pm.request.headers.add({ key: 'X-Trace', value: 'abc' });",
ScriptScopes::default(),
Some(request),
None,
tropel_sandbox::config::SandboxConfig::default(),
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
assert!(
out.get("scriptError").is_some_and(|e| e.is_null()),
"the script must not error: {out}"
);
let headers = out
.get("request")
.and_then(|r| r.get("headers"))
.and_then(|h| h.as_array())
.unwrap_or_else(|| panic!("the mutated request must come back: {out}"));
let rendered = format!("{headers:?}");
assert!(
rendered.contains("X-Trace") && rendered.contains("abc"),
"the header the script added must survive the round trip: {rendered}"
);
assert!(
rendered.contains("Accept"),
"the original headers must survive too: {rendered}"
);
}
fn test_http_client() -> tropel_http::HttpClient {
tropel_http::HttpClient::new(&tropel_http::config::HttpConfig::default())
.expect("test http client")
}
#[tokio::test]
async fn the_embedder_namespace_reaches_the_script_realm() {
let code = "kp.environment.set('viaKp', '2');";
let configured = run_script_once(
code,
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig {
namespace: "kp".into(),
aliases: Vec::new(),
},
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
assert!(
configured
.get("scriptError")
.map(|e| e.is_null())
.unwrap_or(false),
"a kp.* script must run when the caller declares the kp namespace: {configured}"
);
assert_eq!(
configured
.get("environment")
.and_then(|e| e.get("viaKp"))
.and_then(|v| v.as_str()),
Some("2"),
"the effect must come back, not just the absence of an error: {configured}"
);
let stock = run_script_once(
code,
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig::default(),
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
let err = stock
.get("scriptError")
.and_then(|e| e.as_str())
.unwrap_or("");
assert!(
err.contains("kp is not defined"),
"the stock install must NOT bind kp — if it does, this test no \
longer proves the preamble is what carries the namespace: {stock}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_script_calls_back_into_the_host_and_resumes() {
let cb = Arc::new(RunCallbacks::default());
let host = {
let cb = cb.clone();
tokio::spawn(async move {
for _ in 0..200 {
let call = {
let mut p = cb.pending.lock().unwrap();
p.pop_front()
};
if let Some(call) = call {
let tx = cb.replies.lock().unwrap().remove(&call.call_id);
if let Some(tx) = tx {
let _ = tx.send(
serde_json::json!({ "status": 201, "body": call.path }).to_string(),
);
}
return true;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
false
})
};
let out = run_script_once(
"var r = bru.runRequest('Login');\
kp.environment.set('status', String(r.status));\
kp.environment.set('echoed', String(r.body));",
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig {
namespace: "kp".into(),
aliases: Vec::new(),
},
ScriptHost {
http: test_http_client(),
callbacks: Some(cb.clone()),
cookies: None,
},
)
.await
.expect("the realm runs");
assert!(host.await.unwrap_or(false), "the host never saw the call");
assert!(
out.get("scriptError").map(|e| e.is_null()).unwrap_or(false),
"the script must not error: {out}"
);
let env = out.get("environment").expect("environment comes back");
assert_eq!(
env.get("status").and_then(|v| v.as_str()),
Some("201"),
"the script must resume with the HOST's answer, not a placeholder: {out}"
);
assert_eq!(
env.get("echoed").and_then(|v| v.as_str()),
Some("Login"),
"the path the script asked for must reach the host verbatim: {out}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn run_request_refuses_when_no_channel_is_open() {
let out = run_script_once(
"bru.runRequest('Login');",
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig::default(),
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
let err = out
.get("scriptError")
.and_then(|e| e.as_str())
.unwrap_or("");
assert!(
err.contains("not available here"),
"it must refuse by name rather than return undefined: {out}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_failed_script_send_is_still_recorded() {
let out = run_script_once(
"pm.sendRequest('http://127.0.0.1:1/nope', function () {});",
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig {
namespace: "kp".into(),
aliases: Vec::new(),
},
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
let sends = out
.get("scriptSends")
.and_then(|v| v.as_array())
.expect("scriptSends is always an array, never null");
assert_eq!(
sends.len(),
1,
"the send must be recorded even though it failed: {out}"
);
let only = &sends[0];
assert_eq!(
only.get("url").and_then(|v| v.as_str()),
Some("http://127.0.0.1:1/nope"),
"the recorded URL must be the one the script asked for: {out}"
);
assert_eq!(
only.get("status").and_then(|v| v.as_u64()),
Some(0),
"a send that never connected has no status — 0, not a fabricated one: {out}"
);
assert!(
only.get("error").map(|e| !e.is_null()).unwrap_or(false),
"the failure must be NAMED on the record, not implied by status 0: {out}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_script_that_sends_nothing_reports_an_empty_list() {
let out = run_script_once(
"kp.environment.set('x', '1');",
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig::default(),
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
assert_eq!(
out.get("scriptSends")
.and_then(|v| v.as_array())
.map(|a| a.len()),
Some(0),
"an empty array, not null: {out}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_results_answer_and_assertion_results_refuse() {
let out = run_script_once(
"kp.test('first', function () {});\
kp.test('second', function () { throw new Error('no'); });\
var r = bru.getTestResults();\
kp.environment.set('count', String(r.length));\
kp.environment.set('shape', r.map(function (x) { return x.name + ':' + x.passed; }).join(','));",
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig {
namespace: "kp".into(),
aliases: Vec::new(),
},
ScriptHost { http: test_http_client(), callbacks: None, cookies: None },
)
.await
.expect("the realm runs");
let env = out.get("environment").expect("environment comes back");
assert_eq!(
env.get("count").and_then(|v| v.as_str()),
Some("2"),
"both checks must be reported, not just the passing one: {out}"
);
assert_eq!(
env.get("shape").and_then(|v| v.as_str()),
Some("first:true,second:false"),
"each result must carry its own name AND outcome — a count alone \
cannot tell a passing suite from a failing one: {out}"
);
let refused = run_script_once(
"bru.getAssertionResults();",
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig::default(),
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
let err = refused
.get("scriptError")
.and_then(|e| e.as_str())
.unwrap_or("");
assert!(
err.contains("not available here"),
"it must refuse by NAME rather than return an empty list, which would \
read as \"no assertions failed\": {refused}"
);
}
#[tokio::test]
async fn the_cookie_jar_round_trips_and_refuses_when_absent() {
let jar = Arc::new(ScriptCookies::default());
*jar.lock_jar() = vec![serde_json::json!({
"key": "sid", "value": "abc", "domain": "api.example.test", "path": "/"
})];
let out = run_script_once(
"kp.environment.set('seeded', String(bru.cookies.get('sid')));\
bru.cookies.add({ key: 'new', value: 'n1', domain: 'api.example.test', path: '/' });\
kp.environment.set('readback', String(bru.cookies.get('new')));\
bru.cookies.remove('sid');\
kp.environment.set('count', String(bru.cookies.count()));",
ScriptScopes::default(),
Some(TropelRequest {
url: "https://api.example.test/things".to_string(),
..Default::default()
}),
None,
tropel_sandbox::config::SandboxConfig {
namespace: "kp".into(),
aliases: Vec::new(),
},
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: Some(jar.clone()),
},
)
.await
.expect("the realm runs");
let env = out.get("environment").expect("environment comes back");
assert_eq!(
env.get("seeded").and_then(|v| v.as_str()),
Some("abc"),
"the script must read what the caller seeded: {out}"
);
assert_eq!(
env.get("readback").and_then(|v| v.as_str()),
Some("n1"),
"read-your-writes: a cookie the script just set must be visible to it, \
not shadowed by the seeded snapshot: {out}"
);
assert_eq!(
env.get("count").and_then(|v| v.as_str()),
Some("1"),
"the removed cookie must be gone from the script's view too: {out}"
);
let ops = out
.get("cookieOps")
.and_then(|o| o.as_array())
.expect("cookieOps come back");
let kinds: Vec<&str> = ops
.iter()
.filter_map(|o| o.get("op").and_then(|v| v.as_str()))
.collect();
assert_eq!(
kinds,
vec!["set", "delete"],
"the caller replays these onto the real jar, so ORDER is part of the \
contract, not just the set of ops: {out}"
);
}
#[tokio::test]
async fn bru_cookies_refuses_when_no_jar_was_supplied() {
let out = run_script_once(
"bru.cookies.get('sid');",
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig::default(),
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
let err = out
.get("scriptError")
.and_then(|e| e.as_str())
.unwrap_or("");
assert!(
err.contains("not available here"),
"it must refuse by name rather than read as an empty jar: {out}"
);
assert!(
out.get("cookieOps").map(|v| v.is_null()).unwrap_or(false),
"no jar means no ops, not an empty op list: {out}"
);
}
#[tokio::test]
async fn fetch_is_bound_to_the_same_host_client() {
let out = run_script_once(
"(async function () {\
try {\
var r = await fetch('http://127.0.0.1:1/');\
kp.environment.set('status', String(r.status));\
} catch (e) {\
kp.environment.set('err', String(e && e.message));\
}\
})();",
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig {
namespace: "kp".into(),
aliases: Vec::new(),
},
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
let env = out.get("environment").expect("environment comes back");
let err = env.get("err").and_then(|v| v.as_str()).unwrap_or("");
assert!(
!err.contains("not available here"),
"fetch must be bound to a real client, not refuse for lack of one: {out}"
);
assert!(
!err.is_empty() || env.get("status").is_some(),
"the call must have gone somewhere — neither a result nor an error \
means `fetch` never ran: {out}"
);
}
#[tokio::test]
async fn every_variable_scope_round_trips() {
let mut scopes = ScriptScopes::default();
scopes.environment.insert("envIn".into(), "e".into());
scopes
.collection
.insert("colIn".into(), serde_json::json!("c"));
scopes
.globals
.insert("gloIn".into(), serde_json::json!("g"));
scopes
.variables
.insert("varIn".into(), serde_json::json!("v"));
let out = run_script_once(
"kp.environment.set('envSeen', String(kp.environment.get('envIn')));\
kp.collectionVariables.set('colSeen', String(kp.collectionVariables.get('colIn')));\
kp.globals.set('gloSeen', String(kp.globals.get('gloIn')));\
kp.variables.set('varSeen', String(kp.variables.get('varIn')));",
scopes,
None,
None,
tropel_sandbox::config::SandboxConfig {
namespace: "kp".into(),
aliases: Vec::new(),
},
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
assert!(
out.get("scriptError").map(|e| e.is_null()).unwrap_or(false),
"the script must not error: {out}"
);
for (scope, seen, expected) in [
("environment", "envSeen", "e"),
("collectionVariables", "colSeen", "c"),
("globals", "gloSeen", "g"),
("variables", "varSeen", "v"),
] {
let got = out
.get(scope)
.and_then(|m| m.get(seen))
.and_then(|v| v.as_str());
assert_eq!(
got,
Some(expected),
"`{scope}` must be seeded AND returned — the script read \
`{expected}` from it and wrote `{seen}` back: {out}"
);
}
}
#[tokio::test]
async fn send_request_reaches_a_real_http_client() {
let out = run_script_once(
"pm.sendRequest('http://127.0.0.1:1/', function (err, res) {\
kp.environment.set('err', String(err));\
kp.environment.set('code', String(res && res.code));\
});",
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig {
namespace: "kp".into(),
aliases: Vec::new(),
},
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
let err = out
.get("environment")
.and_then(|e| e.get("err"))
.and_then(|v| v.as_str())
.unwrap_or("");
assert!(
!err.contains("unavailable in this build"),
"the send-request bridge must hold a real client, not be installed \
over `http_client: None`: {out}"
);
assert_ne!(
err, "null",
"port 1 cannot have answered — a success here means the call never \
left the realm: {out}"
);
}
#[tokio::test]
async fn the_script_realm_has_no_host_escape() {
fn kp() -> tropel_sandbox::config::SandboxConfig {
tropel_sandbox::config::SandboxConfig {
namespace: "kp".into(),
aliases: Vec::new(),
}
}
let probe = "kp.environment.set('process', typeof process);\
kp.environment.set('require', typeof require);\
kp.environment.set('module', typeof module);\
kp.environment.set('viaFn', typeof Function('return this')().process);";
let out = run_script_once(
probe,
ScriptScopes::default(),
None,
None,
kp(),
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
let env = out.get("environment").expect("the environment comes back");
for name in ["process", "require", "module", "viaFn"] {
assert_eq!(
env.get(name).and_then(|v| v.as_str()),
Some("undefined"),
"`{name}` must not be reachable from a script: {out}"
);
}
for (label, code) in [
("direct", "import('fs');"),
("indirect eval", "var e = eval; e(\"import\" + \"('fs')\");"),
] {
let result = run_script_once(
code,
ScriptScopes::default(),
None,
None,
kp(),
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await;
let refused = match &result {
Err(why) => why.contains("module"),
Ok(out) => out
.get("scriptError")
.and_then(|e| e.as_str())
.is_some_and(|e| e.contains("module")),
};
assert!(
refused,
"a {label} import must not load a host module: {result:?}"
);
}
}
#[tokio::test]
async fn the_script_realm_matches_the_committed_corpus() {
const CORPUS: &str = include_str!("../testdata/script-realm-corpus.json");
let doc: serde_json::Value = serde_json::from_str(CORPUS).expect("corpus is valid JSON");
let probes = doc["probes"].as_array().expect("probes array");
assert!(!probes.is_empty(), "an empty corpus asserts nothing");
let script = probes
.iter()
.map(|p| {
let name = p["name"].as_str().expect("name");
let expr = p["expression"].as_str().expect("expression");
format!("pm.environment.set({name:?}, String({expr}));")
})
.collect::<Vec<_>>()
.join("\n");
let out = run_script_once(
&script,
ScriptScopes::default(),
None,
None,
tropel_sandbox::config::SandboxConfig::default(),
ScriptHost {
http: test_http_client(),
callbacks: None,
cookies: None,
},
)
.await
.expect("the realm runs");
let mut wrong: Vec<String> = Vec::new();
for probe in probes {
let name = probe["name"].as_str().unwrap();
let want = probe["quickJs"].as_str().unwrap();
let got = out
.get("environment")
.and_then(|e| e.get(name))
.and_then(|v| v.as_str())
.unwrap_or("(probe did not run)");
if got != want {
wrong.push(format!(
" {name}: corpus says {want:?}, realm answered {got:?}"
));
}
}
assert!(
wrong.is_empty(),
"the QuickJS realm no longer matches the committed corpus.\n{}\n\n\
If the realm CHANGED on purpose, update \
packages/shims/fixtures/script-realm-corpus.json — and update the \
`why` line too, because KnockPort's half of this corpus asserts \
the same file and a user reads those lines to know what their \
script can use.",
wrong.join("\n")
);
}
#[test]
fn the_execute_wire_carries_a_per_request_proxy() {
fn parse(v: &serde_json::Value) -> Option<tropel_sdk::types::ProxyConfig> {
v.get("proxy")
.and_then(|p| serde_json::from_value(p.clone()).ok())
}
let cfg = parse(&serde_json::json!({
"proxy": {"mode": "fixed", "protocol": "http", "host": "p.internal",
"port": 3128, "username": "u", "password": "p",
"bypass": ["localhost", "*.internal"]}
}))
.expect("a fixed proxy reads");
assert_eq!(cfg.mode, tropel_sdk::types::ProxyMode::Fixed);
assert_eq!(cfg.fixed_url().as_deref(), Some("http://p.internal:3128"));
assert_eq!(cfg.bypass.len(), 2);
let pac = parse(&serde_json::json!({
"proxy": {"mode": "pac", "pacUrl": "http://wpad/proxy.pac"}
}))
.expect("a pac proxy reads");
assert_eq!(pac.pac_url.as_deref(), Some("http://wpad/proxy.pac"));
assert!(parse(&serde_json::json!({"url": "https://x/y"})).is_none());
assert!(parse(&serde_json::json!({"proxy": "http://p:3128"})).is_none());
}
#[test]
fn duplicate_header_names_survive_the_execute_wire_format() {
fn parse(v: &serde_json::Value) -> Vec<(String, String)> {
match v.get("headers") {
Some(serde_json::Value::Array(rows)) => rows
.iter()
.filter_map(|row| match row {
serde_json::Value::Array(pair) if pair.len() == 2 => Some((
pair[0].as_str()?.to_string(),
pair[1].as_str().unwrap_or("").to_string(),
)),
serde_json::Value::Object(o) => Some((
o.get("name")?.as_str()?.to_string(),
o.get("value")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string(),
)),
_ => None,
})
.collect(),
Some(serde_json::Value::Object(o)) => o
.iter()
.map(|(k, v)| (k.clone(), v.as_str().unwrap_or("").to_string()))
.collect(),
_ => Vec::new(),
}
}
let pairs = parse(&serde_json::json!({
"headers": [["Accept", "application/json"], ["Accept", "text/plain"]]
}));
assert_eq!(
pairs,
vec![
("Accept".to_string(), "application/json".to_string()),
("Accept".to_string(), "text/plain".to_string()),
],
"a duplicate header name must reach the wire twice"
);
let obj = parse(&serde_json::json!({"headers": {"Accept": "application/json"}}));
assert_eq!(
obj,
vec![("Accept".to_string(), "application/json".to_string())]
);
let named = parse(&serde_json::json!({
"headers": [{"name": "X-A", "value": "1"}, {"name": "X-A", "value": "2"}]
}));
assert_eq!(named.len(), 2, "{named:?}");
let collapsed: serde_json::Value =
serde_json::from_str(r#"{"headers":{"Accept":"a","Accept":"b"}}"#).unwrap();
assert_eq!(
parse(&collapsed).len(),
1,
"the object form loses one row before the agent ever sees it"
);
}
#[test]
fn a_binary_response_body_is_base64_not_mojibake() {
fn png_header() -> Vec<u8> {
vec![0x89, b'P', b'N', b'G', 0x0D, 0x0A, 0x1A, 0x0A, 0xFF, 0xD8]
}
let png = png_header();
assert!(
std::str::from_utf8(&png).is_err(),
"the fixture must actually be invalid UTF-8 or this test proves nothing"
);
let (body, encoding) = match std::str::from_utf8(&png) {
Ok(text) => (text.to_string(), "utf8"),
Err(_) => (base64_encode(&png), "base64"),
};
assert_eq!(encoding, "base64");
assert_eq!(
base64_decode(&body).expect("round trips"),
png,
"the bytes must come back EXACTLY — lossy is the bug"
);
let lossy = String::from_utf8_lossy(&png);
assert!(
lossy.contains('\u{FFFD}'),
"the old path really did corrupt these bytes: {lossy:?}"
);
assert_ne!(
lossy.as_bytes(),
png,
"and the corruption was unrecoverable — no field said so"
);
fn json_body() -> Vec<u8> {
br#"{"ok":true}"#.to_vec()
}
let text = json_body();
let (body, encoding) = match std::str::from_utf8(&text) {
Ok(t) => (t.to_string(), "utf8"),
Err(_) => (base64_encode(&text), "base64"),
};
assert_eq!(encoding, "utf8");
assert_eq!(body, "{\"ok\":true}");
}
#[tokio::test]
async fn the_auth_sign_endpoint_applies_the_rules_server_side() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let port = listener.local_addr().expect("addr").port();
let state = Arc::new(AgentState {
token: None,
client: tropel_http::HttpClient::new(&tropel_http::config::HttpConfig::default())
.expect("http client"),
runs: std::sync::Mutex::new(HashMap::new()),
allowed_origins: vec![],
});
tokio::spawn(async move {
while let Ok((mut sock, _)) = listener.accept().await {
let st = state.clone();
tokio::spawn(async move {
let _ = handle_connection(&mut sock, st).await;
});
}
});
let sign = |body: String| async move {
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
let req = format!(
"POST /auth/sign HTTP/1.1\r\nHost: localhost\r\nContent-Length: {}\r\n\r\n{body}",
body.len()
);
s.write_all(req.as_bytes()).await.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
String::from_utf8_lossy(&out).to_string()
};
let raw = sign(
serde_json::json!({
"scheme": "awsSigV4",
"params": {
"method": "GET", "host": "examplebucket.s3.amazonaws.com",
"path": "/test.txt", "accessKey": "AKID", "secretKey": "SECRET",
"region": "us-east-1", "amzDate": "20130524T000000Z",
"dateStamp": "20130524"
}
})
.to_string(),
)
.await;
assert!(
raw.contains("/20130524/us-east-1/s3/aws4_request"),
"service must derive to s3, not the bucket: {raw}"
);
assert!(raw.contains("x-amz-date"), "{raw}");
assert!(raw.contains("x-amz-content-sha256"), "{raw}");
let raw = sign(
serde_json::json!({
"scheme": "digest",
"params": {
"wwwAuthenticate": "Basic realm=\"b\", Digest realm=\"r\", qop=\"auth, auth-int\", nonce=\"n\"",
"username": "u", "password": "p", "method": "GET",
"uri": "/dir/index.html", "nc": 1, "cnonce": "0a4f113b"
}
})
.to_string(),
)
.await;
assert!(raw.contains("Digest "), "{raw}");
assert!(raw.contains(r#"realm=\"r\""#), "{raw}");
let raw = sign(
serde_json::json!({
"scheme": "digest",
"params": {"wwwAuthenticate": "Basic realm=\"b\"", "username": "u"}
})
.to_string(),
)
.await;
assert!(raw.starts_with("HTTP/1.1 400"), "{raw}");
assert!(raw.contains("no Digest challenge"), "{raw}");
let raw = sign(
serde_json::json!({
"scheme": "oauth1",
"params": {
"method": "POST", "scheme": "http", "host": "::1", "port": 8080,
"path": "/request", "formBody": "c2=&a3=2+q",
"consumerKey": "ck", "consumerSecret": "cs",
"signatureMethod": "HMAC-SHA1", "nonce": "n", "timestamp": "1"
}
})
.to_string(),
)
.await;
assert!(raw.contains("oauth_signature="), "{raw}");
let raw = sign(
serde_json::json!({
"scheme": "oauth1",
"params": {
"method": "GET", "scheme": "https", "host": "x.test", "path": "/",
"consumerKey": "ck", "consumerSecret": "cs",
"signatureMethod": "RSA-SHA1", "nonce": "n", "timestamp": "1"
}
})
.to_string(),
)
.await;
assert!(raw.starts_with("HTTP/1.1 400"), "{raw}");
assert!(raw.contains("RSA-SHA1"), "{raw}");
let raw = sign(
serde_json::json!({
"scheme": "akamai-edgegrid",
"params": {
"method": "GET", "url": "https://akaa-x.luna.akamaiapis.net/diagnostic/v1/x",
"clientToken": "ct", "accessToken": "at", "clientSecret": "cs",
"nonce": "nnn", "timestamp": "20260909T12:00:00+0000"
}
})
.to_string(),
)
.await;
assert!(raw.contains("EG1-HMAC-SHA256 "), "{raw}");
assert!(raw.contains("client_token=ct"), "{raw}");
assert!(raw.contains("nonce=nnn"), "{raw}");
assert!(
!raw.contains("cs\""),
"the client secret must not echo: {raw}"
);
let raw = sign(
serde_json::json!({
"scheme": "akamai-edgegrid",
"params": {"method": "GET", "url": "https://x/y", "clientToken": "ct"}
})
.to_string(),
)
.await;
assert!(raw.starts_with("HTTP/1.1 400"), "{raw}");
assert!(raw.contains("access_token"), "{raw}");
assert!(raw.contains("client_secret"), "{raw}");
let raw = sign(
serde_json::json!({
"scheme": "wsse",
"params": {"username": "u", "password": "p", "nonce": "n", "created": "2026-01-01T00:00:00Z"}
})
.to_string(),
)
.await;
assert!(raw.contains("X-WSSE"), "{raw}");
assert!(raw.contains("UsernameToken"), "{raw}");
assert!(raw.contains("PasswordDigest"), "{raw}");
let raw = sign(serde_json::json!({"scheme": "ntlm", "params": {}}).to_string()).await;
assert!(raw.starts_with("HTTP/1.1 400"), "{raw}");
assert!(raw.contains("unknown auth scheme"), "{raw}");
for scheme in AUTH_SIGN_SCHEMES {
assert!(
raw.contains(scheme),
"the refusal must name {scheme}: {raw}"
);
}
}
#[tokio::test]
async fn the_script_endpoint_runs_in_the_same_realm_as_a_load_run() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let port = listener.local_addr().expect("addr").port();
let state = Arc::new(AgentState {
token: None,
client: tropel_http::HttpClient::new(&tropel_http::config::HttpConfig::default())
.expect("http client"),
runs: std::sync::Mutex::new(HashMap::new()),
allowed_origins: vec![],
});
tokio::spawn(async move {
while let Ok((mut sock, _)) = listener.accept().await {
let st = state.clone();
tokio::spawn(async move {
let _ = handle_connection(&mut sock, st).await;
});
}
});
let run = |body: String| async move {
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
let req = format!(
"POST /script HTTP/1.1\r\nHost: localhost\r\nContent-Length: {}\r\n\r\n{body}",
body.len()
);
s.write_all(req.as_bytes()).await.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
String::from_utf8_lossy(&out).to_string()
};
let raw = run(serde_json::json!({
"code": r#"
pm.test("passes", function () { pm.expect(1).to.eql(1); });
pm.test("fails", function () { pm.expect(1).to.eql(2); });
"#
})
.to_string())
.await;
assert!(raw.contains(r#""name":"passes""#), "{raw}");
assert!(raw.contains(r#""name":"fails""#), "{raw}");
assert!(raw.contains(r#""passed":1"#), "{raw}");
assert!(raw.contains(r#""failed":1"#), "{raw}");
let raw = run(serde_json::json!({
"code": r#"pm.environment.set("token", "abc123");"#,
"environment": {"seeded": "yes"}
})
.to_string())
.await;
assert!(raw.contains(r#""token":"abc123""#), "{raw}");
assert!(
raw.contains(r#""seeded":"yes""#),
"seed must survive: {raw}"
);
let raw = run(serde_json::json!({
"code": r#"pm.test("ran", function () {}); throw new Error("boom");"#
})
.to_string())
.await;
assert!(raw.starts_with("HTTP/1.1 200"), "{raw}");
assert!(raw.contains("scriptError"), "{raw}");
assert!(raw.contains("boom"), "{raw}");
assert!(
raw.contains(r#""name":"ran""#),
"what ran before the throw survives: {raw}"
);
let _ = run(serde_json::json!({"code": "globalThis.__leak = 1;"}).to_string()).await;
let raw = run(serde_json::json!({
"code": r#"pm.test("isolated", function () {
pm.expect(typeof globalThis.__leak).to.eql("undefined");
});"#
})
.to_string())
.await;
assert!(
raw.contains(r#""passed":1"#),
"realms must not share globals: {raw}"
);
}
#[tokio::test]
async fn the_oauth2_endpoint_serves_the_whole_family() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let port = listener.local_addr().expect("addr").port();
let state = Arc::new(AgentState {
token: None,
client: tropel_http::HttpClient::new(&tropel_http::config::HttpConfig::default())
.expect("http client"),
runs: std::sync::Mutex::new(HashMap::new()),
allowed_origins: vec![],
});
tokio::spawn(async move {
while let Ok((mut sock, _)) = listener.accept().await {
let st = state.clone();
tokio::spawn(async move {
let _ = handle_connection(&mut sock, st).await;
});
}
});
let call = |op: &'static str, params: serde_json::Value| async move {
let body = serde_json::json!({ "op": op, "params": params }).to_string();
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
let req = format!(
"POST /auth/oauth2 HTTP/1.1\r\nHost: localhost\r\nContent-Length: {}\r\n\r\n{body}",
body.len()
);
s.write_all(req.as_bytes()).await.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
String::from_utf8_lossy(&out).to_string()
};
let raw = call(
"codeChallengeS256",
serde_json::json!({"verifier": "abc123"}),
)
.await;
assert!(raw.contains("codeChallenge"), "{raw}");
assert!(raw.contains(r#""codeChallengeMethod":"S256""#), "{raw}");
let raw = call(
"signJwt",
serde_json::json!({"payload": {"sub": "u1"}, "algorithm": "HS512", "secret": "s"}),
)
.await;
assert!(raw.contains(r#""token":"#), "{raw}");
let raw = call(
"signJwt",
serde_json::json!({"payload": {"sub": "u1"}, "algorithm": "RS256", "secret": "s"}),
)
.await;
assert!(raw.starts_with("HTTP/1.1 400"), "{raw}");
assert!(raw.contains("RS256"), "refused by name: {raw}");
let raw = call(
"attachToken",
serde_json::json!({"token": "t", "placement": "cookie"}),
)
.await;
assert!(raw.starts_with("HTTP/1.1 400"), "{raw}");
assert!(raw.contains("unknown token placement"), "{raw}");
let raw = call(
"attachToken",
serde_json::json!({"token": "t", "tokenType": "Bearer", "placement": "header"}),
)
.await;
assert!(raw.contains("Bearer"), "{raw}");
let raw = call(
"wsseSign",
serde_json::json!({"username": "u", "password": "p"}),
)
.await;
assert!(raw.starts_with("HTTP/1.1 200"), "{raw}");
let raw = call("nope", serde_json::json!({})).await;
assert!(raw.starts_with("HTTP/1.1 400"), "{raw}");
assert!(raw.contains("unknown oauth2 op"), "{raw}");
assert!(raw.contains("signJwt"), "it lists the alternatives: {raw}");
}
#[tokio::test]
async fn the_batch_resolve_endpoint_preserves_order_and_isolates_failures() {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let port = listener.local_addr().expect("addr").port();
let state = Arc::new(AgentState {
token: None,
client: tropel_http::HttpClient::new(&tropel_http::config::HttpConfig::default())
.expect("http client"),
runs: std::sync::Mutex::new(HashMap::new()),
allowed_origins: vec![],
});
tokio::spawn(async move {
while let Ok((mut sock, _)) = listener.accept().await {
let st = state.clone();
tokio::spawn(async move {
let _ = handle_connection(&mut sock, st).await;
});
}
});
let post = |body: String| async move {
let mut s = TcpStream::connect(("127.0.0.1", port))
.await
.expect("connect");
let req = format!(
"POST /resolve/batch HTTP/1.1\r\nHost: localhost\r\nContent-Length: {}\r\n\r\n{body}",
body.len()
);
s.write_all(req.as_bytes()).await.expect("write");
let mut out = Vec::new();
s.read_to_end(&mut out).await.expect("read");
String::from_utf8_lossy(&out).to_string()
};
let raw = post(
serde_json::json!({
"variables": {"base": "https://api.test", "tok": "abc", "n": "2"},
"items": [
{"template": "{{base}}/v{{n}}"},
{"template": "Bearer {{tok}}"},
{"template": "{\"t\":\"{{tok}}\"}", "mode": "json"},
{"template": "{{tok}}", "mode": "nope"},
{"template": "{{base}}"}
]
})
.to_string(),
)
.await;
let body = raw.split("\r\n\r\n").nth(1).unwrap_or_default().to_string();
let parsed: serde_json::Value = serde_json::from_str(&body).expect("json body");
let items = parsed["items"].as_array().expect("items array");
assert_eq!(items.len(), 5, "one output per input, always: {body}");
assert_eq!(items[0]["value"], "https://api.test/v2");
assert_eq!(items[1]["value"], "Bearer abc");
assert_eq!(items[2]["value"], "{\"t\":\"abc\"}");
assert!(
items[3]["error"].is_string(),
"item 3 should carry an error: {body}"
);
assert!(items[3]["value"].is_null());
assert_eq!(
items[4]["value"], "https://api.test",
"later items still resolve"
);
let raw = post(
serde_json::json!({
"variables": {"a": "{{b}}", "b": "{{a}}"},
"items": [
{"template": "{{a}}"},
{"template": "{{nosuchvar}}"}
]
})
.to_string(),
)
.await;
let body = raw.split("\r\n\r\n").nth(1).unwrap_or_default().to_string();
let parsed: serde_json::Value = serde_json::from_str(&body).expect("json body");
let items = parsed["items"].as_array().expect("items");
assert_eq!(
items[0]["hitCap"], true,
"a cycle must report hitCap: {body}"
);
assert_eq!(
items[1]["hitCap"], false,
"an unknown name is NOT a cycle: {body}"
);
assert!(
items[1]["unresolved"]
.as_array()
.is_some_and(|u| u.iter().any(|n| n == "nosuchvar")),
"the unknown name must be reported so the user sees their typo: {body}"
);
let raw = post(
serde_json::json!({
"variables": {"a": "{{b}}", "b": "final"},
"items": [{"template": "{{a}}", "deep": false}]
})
.to_string(),
)
.await;
let body = raw.split("\r\n\r\n").nth(1).unwrap_or_default().to_string();
let parsed: serde_json::Value = serde_json::from_str(&body).expect("json body");
let item = &parsed["items"][0];
assert!(
item["value"].is_null(),
"a refused item carries no value: {body}"
);
assert!(
item["error"]
.as_str()
.is_some_and(|e| e.contains("deep: false") && e.contains("/resolve")),
"the refusal must name the field AND the endpoint that serves it: {body}"
);
let raw = post(serde_json::json!({"variables": {}, "items": []}).to_string()).await;
assert!(raw.starts_with("HTTP/1.1 200"), "{raw}");
assert!(raw.contains(r#""items":[]"#), "{raw}");
}
}