#![allow(clippy::unwrap_used)] #![allow(clippy::print_stderr)] #![allow(clippy::large_futures)]
use std::{future::Future, pin::Pin, sync::Arc};
use chrono::Utc;
use fraiseql_functions::{
FunctionModule, FunctionObserver, ResourceLimits, RuntimeType,
host::live::{HostContextConfig, QueryExecutor},
runtime::deno::{DenoConfig, DenoRuntime},
triggers::CronSchedule,
};
use fraiseql_observers::{
LeaseGuardedRunner, PostgresSourceCursorStore, RunOutcome, SourceCursorStore,
};
use serde_json::{Value, json};
use sqlx::PgPool;
use super::{SourcePoller, build_source_payload};
#[test]
fn source_payload_carries_the_trigger_context() {
let payload = build_source_payload("orders", "*/5 * * * *", Utc::now());
assert_eq!(payload.trigger_type, "source:orders");
assert_eq!(payload.entity, "source");
assert_eq!(payload.event_kind, "scheduled");
assert_eq!(payload.data["source"], "orders");
assert_eq!(payload.data["schedule"], "*/5 * * * *");
}
#[tokio::test]
async fn idempotency_token_is_signed_when_a_key_is_configured() {
let pool = PgPool::connect_lazy("postgres://localhost/unused").unwrap();
let build = |key: Option<Arc<[u8]>>| {
SourcePoller::new(
"orders",
CronSchedule::parse("*/5 * * * *").unwrap(),
FunctionModule::from_source("connector".to_string(), String::new(), RuntimeType::Deno),
Arc::new(FunctionObserver::new()),
PostgresSourceCursorStore::new(pool.clone()),
StubExecutor::new(json!(null)),
LeaseGuardedRunner::in_process("orders"),
HostContextConfig::default(),
ResourceLimits::default(),
key,
false,
)
};
let payload = build_source_payload("orders", "*/5 * * * *", Utc::now());
let unsigned = build(None).idempotency_token(&payload);
let secret: Arc<[u8]> = Arc::from(b"server-hmac-secret".as_slice());
let signed = build(Some(Arc::clone(&secret))).idempotency_token(&payload);
assert_ne!(unsigned, signed);
let expected = fraiseql_observers::derive_idempotency_token(
Some(&secret),
fraiseql_observers::DispatchSource::Source,
"connector",
&payload.trigger_type,
&payload.data,
);
assert_eq!(signed, expected);
assert_eq!(unsigned.len(), 32);
}
struct StubExecutor {
response: Value,
seen: std::sync::Mutex<Vec<String>>,
}
impl StubExecutor {
fn new(response: Value) -> Arc<Self> {
Arc::new(Self {
response,
seen: std::sync::Mutex::new(Vec::new()),
})
}
}
impl QueryExecutor for StubExecutor {
fn execute_query(
&self,
query: &str,
_variables: Option<&Value>,
) -> Pin<Box<dyn Future<Output = fraiseql_error::Result<Value>> + Send + '_>> {
self.seen.lock().unwrap().push(query.to_string());
let response = self.response.clone();
Box::pin(async move { Ok(response) })
}
}
async fn connect_pool() -> Option<(PgPool, fraiseql_test_support::Service)> {
let svc = fraiseql_test_support::postgres().await?;
let pool = PgPool::connect(svc.url()).await.unwrap();
Some((pool, svc))
}
fn poller(pool: &PgPool, source: &str, executor: Arc<dyn QueryExecutor>) -> SourcePoller {
SourcePoller::new(
source,
CronSchedule::parse("*/5 * * * *").unwrap(),
FunctionModule::from_source("noop".to_string(), String::new(), RuntimeType::Deno),
Arc::new(FunctionObserver::new()),
PostgresSourceCursorStore::new(pool.clone()),
executor,
LeaseGuardedRunner::in_process(source),
HostContextConfig::default(),
ResourceLimits::default(),
None,
false,
)
}
#[tokio::test]
async fn build_host_binds_both_the_cursor_and_the_executor() {
let Some((pool, _svc)) = connect_pool().await else {
eprintln!("SKIP build_host_binds_both_the_cursor_and_the_executor: no postgres");
return;
};
let source = "test-poller-build-host";
PostgresSourceCursorStore::new(pool.clone()).init().await.unwrap();
sqlx::query("DELETE FROM _fraiseql_source_cursor WHERE source_name = $1")
.bind(source)
.execute(&pool)
.await
.unwrap();
let executor = StubExecutor::new(json!({ "data": { "createOrder": { "status": "ok" } } }));
let poller = poller(&pool, source, executor.clone());
let host = poller.build_host(build_source_payload(source, "*/5 * * * *", Utc::now()), "tok");
assert!(host.cursor().await.unwrap().is_none(), "a fresh source has no cursor");
host.advance_cursor(json!({ "page": 3 })).await.unwrap();
assert_eq!(
host.cursor().await.unwrap(),
Some(json!({ "page": 3 })),
"the host reads back what it advanced"
);
let result = host.query("mutation { createOrder }", json!({})).await.unwrap();
assert_eq!(result, json!({ "data": { "createOrder": { "status": "ok" } } }));
assert_eq!(
executor.seen.lock().unwrap().as_slice(),
["mutation { createOrder }"],
"the query reached the bound executor"
);
let snapshot = PostgresSourceCursorStore::new(pool.clone()).load(source).await.unwrap();
assert_eq!(snapshot.value, Some(json!({ "page": 3 })));
}
const CONNECTOR_TS: &str = r#"
export default async () => {
const before = JSON.parse(await Deno.core.ops.fraiseql_cursor_get());
const page = (before && before.page ? before.page : 0) + 1;
await Deno.core.ops.fraiseql_query(
"mutation { createOrder(page: " + page + ") { id } }",
"{}"
);
await Deno.core.ops.fraiseql_cursor_advance(JSON.stringify({ page: page }));
return { page: page };
};
"#;
#[tokio::test]
async fn fires_a_model_b_connector_end_to_end() {
let Some((pool, _svc)) = connect_pool().await else {
eprintln!("SKIP fires_a_model_b_connector_end_to_end: no postgres");
return;
};
let source = "test-poller-e2e";
PostgresSourceCursorStore::new(pool.clone()).init().await.unwrap();
sqlx::query("DELETE FROM _fraiseql_source_cursor WHERE source_name = $1")
.bind(source)
.execute(&pool)
.await
.unwrap();
let executor = StubExecutor::new(json!({ "data": { "createOrder": { "id": "1" } } }));
let mut observer = FunctionObserver::new();
observer.register_runtime(RuntimeType::Deno, DenoRuntime::new(&DenoConfig::default()).unwrap());
let module = FunctionModule::from_source(
"connector".to_string(),
CONNECTOR_TS.to_string(),
RuntimeType::Deno,
);
let poller = SourcePoller::new(
source,
CronSchedule::parse("*/5 * * * *").unwrap(),
module,
Arc::new(observer),
PostgresSourceCursorStore::new(pool.clone()),
executor.clone(),
LeaseGuardedRunner::in_process(source),
HostContextConfig::default(),
ResourceLimits::default(),
None,
false,
);
let outcome = poller.fire_once(chrono::Utc::now()).await;
assert!(
matches!(outcome, RunOutcome::Ran(Ok(_))),
"the connector fired and ran to completion under the lease"
);
let seen = executor.seen.lock().unwrap().clone();
assert_eq!(seen.len(), 1, "the guest issued exactly one query");
assert!(seen[0].contains("createOrder"), "it was the connector's mutation: {}", seen[0]);
let snapshot = PostgresSourceCursorStore::new(pool.clone()).load(source).await.unwrap();
assert_eq!(snapshot.value, Some(json!({ "page": 1 })));
}