124 lines
3.7 KiB
Rust
124 lines
3.7 KiB
Rust
mod client;
|
|
mod error;
|
|
mod model;
|
|
mod sse;
|
|
|
|
use async_trait::async_trait;
|
|
use crank_core::{
|
|
AdapterResponse, ExecutionMode, PreparedRequest, Protocol, ProtocolAdapter,
|
|
ProtocolAdapterError, RestTarget, RuntimeRequestContext, Target, WindowExecutionResult,
|
|
};
|
|
use serde_json::{Map, Value};
|
|
|
|
pub use client::RestAdapter;
|
|
pub use error::RestAdapterError;
|
|
pub use model::{RestRequest, RestResponse, RestWindowRequest, RestWindowResponse};
|
|
|
|
#[async_trait]
|
|
impl ProtocolAdapter for RestAdapter {
|
|
fn protocol(&self) -> Protocol {
|
|
Protocol::Rest
|
|
}
|
|
|
|
fn supports_mode(&self, mode: ExecutionMode) -> bool {
|
|
matches!(mode, ExecutionMode::Unary | ExecutionMode::Window)
|
|
}
|
|
|
|
async fn invoke_unary(
|
|
&self,
|
|
target: &Target,
|
|
prepared: &PreparedRequest,
|
|
_context: &RuntimeRequestContext,
|
|
) -> Result<AdapterResponse, ProtocolAdapterError> {
|
|
let target = rest_target(target)?;
|
|
let request = RestRequest {
|
|
path_params: prepared.path_params.clone(),
|
|
query_params: prepared.query_params.clone(),
|
|
headers: prepared.headers.clone(),
|
|
body: prepared.body.clone(),
|
|
timeout_ms: prepared.timeout_ms,
|
|
};
|
|
let response = self.execute(target, &request).await?;
|
|
|
|
Ok(AdapterResponse {
|
|
status_code: response.status_code,
|
|
headers: response.headers,
|
|
data: response.body.clone(),
|
|
body: response.body,
|
|
})
|
|
}
|
|
|
|
async fn invoke_window(
|
|
&self,
|
|
target: &Target,
|
|
prepared: &PreparedRequest,
|
|
window_duration_ms: u64,
|
|
max_items: Option<u32>,
|
|
_context: &RuntimeRequestContext,
|
|
) -> Result<WindowExecutionResult, ProtocolAdapterError> {
|
|
let target = rest_target(target)?;
|
|
let request = RestWindowRequest {
|
|
request: RestRequest {
|
|
path_params: prepared.path_params.clone(),
|
|
query_params: prepared.query_params.clone(),
|
|
headers: prepared.headers.clone(),
|
|
body: prepared.body.clone(),
|
|
timeout_ms: prepared.timeout_ms,
|
|
},
|
|
window_duration_ms,
|
|
max_items: max_items.map(|value| value.saturating_add(1)),
|
|
};
|
|
let response = self.execute_window(target, &request).await?;
|
|
Ok(window_result_from_response(response.body, max_items))
|
|
}
|
|
}
|
|
|
|
fn rest_target(target: &Target) -> Result<&RestTarget, ProtocolAdapterError> {
|
|
match target {
|
|
Target::Rest(target) => Ok(target),
|
|
other => Err(ProtocolAdapterError::Message(format!(
|
|
"rest adapter cannot handle target {other:?}"
|
|
))),
|
|
}
|
|
}
|
|
|
|
fn window_result_from_response(body: Value, max_items: Option<u32>) -> WindowExecutionResult {
|
|
let done = body.get("done").and_then(Value::as_bool).unwrap_or(true);
|
|
let mut items = body
|
|
.get("items")
|
|
.and_then(Value::as_array)
|
|
.cloned()
|
|
.unwrap_or_default();
|
|
let mut truncated = false;
|
|
let mut has_more = !done;
|
|
|
|
if let Some(max_items) = max_items.map(|value| value as usize) {
|
|
if items.len() > max_items {
|
|
items.truncate(max_items);
|
|
truncated = true;
|
|
has_more = true;
|
|
}
|
|
}
|
|
|
|
let summary = body
|
|
.get("summary")
|
|
.cloned()
|
|
.unwrap_or_else(|| Value::Object(Map::new()));
|
|
let cursor = body.get("cursor").cloned();
|
|
|
|
WindowExecutionResult {
|
|
summary,
|
|
items,
|
|
cursor,
|
|
window_complete: !has_more,
|
|
truncated,
|
|
has_more,
|
|
}
|
|
}
|
|
|
|
impl From<RestAdapterError> for ProtocolAdapterError {
|
|
fn from(value: RestAdapterError) -> Self {
|
|
ProtocolAdapterError::Message(value.to_string())
|
|
}
|
|
}
|