use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use async_trait::async_trait;
use platform_core::automation::http_client::AsyncHttpClientService;
use platform_core::automation::{self, get_event_http_target, EventApiService};
use platform_core::{
overrides, resources, AppConfigReader, AppError, ComposableFunction, EventEnvelope,
FunctionOptions, Platform, PostOffice,
};
use rmpv::Value;
struct SaveGet {
store: Mutex<Option<Value>>,
}
#[async_trait]
impl ComposableFunction for SaveGet {
async fn handle_event(
&self,
headers: HashMap<String, String>,
input: EventEnvelope,
_instance: usize,
) -> Result<EventEnvelope, AppError> {
match headers.get("type").map(String::as_str) {
Some("save") => {
*self.store.lock().expect("store") = Some(input.body().clone());
EventEnvelope::new().set_body("saved")
}
Some("get") => {
let stored = self
.store
.lock()
.expect("store")
.clone()
.unwrap_or(Value::Nil);
Ok(EventEnvelope::new().set_raw_body(stored))
}
_ => Err(AppError::new(400, "unknown type")),
}
}
}
struct CaptureCallback {
tx: tokio::sync::mpsc::UnboundedSender<Value>,
}
#[async_trait]
impl ComposableFunction for CaptureCallback {
async fn handle_event(
&self,
_headers: HashMap<String, String>,
input: EventEnvelope,
_instance: usize,
) -> Result<EventEnvelope, AppError> {
let _ = self.tx.send(input.body().clone());
Ok(EventEnvelope::new())
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn declarative_event_over_http() {
resources::prepend_resource_root("tests/resources");
let holding = std::env::temp_dir().join(format!("mercury-eoh-decl-{}", std::process::id()));
overrides::set("transient.data.store", &holding.display().to_string());
overrides::set("rest.server.port", "0");
overrides::set("yaml.rest.automation", "classpath:/event-rest.yaml");
overrides::set("yaml.event.over.http", "classpath:/event-over-http.yaml");
let _ = AppConfigReader::get_instance();
let platform = Platform::new();
platform
.register_with_options(
automation::EVENT_API_SERVICE,
Arc::new(EventApiService::new(&platform)),
10,
FunctionOptions {
zero_traced: false,
interceptor: true,
private: true,
},
)
.unwrap();
platform
.register(
"event.save.get",
Arc::new(SaveGet {
store: Mutex::new(None),
}),
2,
)
.unwrap();
platform
.register_with_options(
automation::ASYNC_HTTP_REQUEST,
Arc::new(AsyncHttpClientService::new(&platform)),
10,
FunctionOptions {
interceptor: true,
private: true,
..FunctionOptions::default()
},
)
.unwrap();
let addr = automation::start_http_server(&platform).await.unwrap();
overrides::set("server.port", &addr.port().to_string());
let entry = get_event_http_target("event.http.test").expect("configured route");
assert_eq!(
entry.target,
format!("http://127.0.0.1:{}/api/event", addr.port())
);
assert_eq!(
entry.headers.get("authorization").map(String::as_str),
Some("demo")
);
assert!(get_event_http_target("event.save.get@2").is_some());
assert!(get_event_http_target("no.such.route").is_none());
let po = PostOffice::new(&platform);
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
platform
.register_private("blocking.event.wait", Arc::new(CaptureCallback { tx }), 1)
.unwrap();
po.send(
EventEnvelope::new()
.set_to("event.save.get")
.set_header("type", "save")
.set_reply_to("blocking.event.wait")
.set_body("hello")
.unwrap(),
)
.await
.expect("declarative send");
let callback_body = tokio::time::timeout(Duration::from_secs(5), rx.recv())
.await
.expect("callback within 5s")
.expect("callback value");
assert_eq!(callback_body.as_str(), Some("saved"));
let response = po
.request(
EventEnvelope::new()
.set_to("event.save.get")
.set_header("type", "get"),
Duration::from_secs(10),
)
.await
.expect("declarative request");
assert_eq!(response.body().as_str(), Some("hello"));
}