#![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",
"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)),
fraiseql_core::security::SecurityContext::system_job(
"test-source",
"test-request",
vec![],
vec![],
None,
),
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,
source,
CronSchedule::parse("*/5 * * * *").unwrap(),
FunctionModule::from_source("noop".to_string(), String::new(), RuntimeType::Deno),
Arc::new(FunctionObserver::new()),
PostgresSourceCursorStore::new(pool.clone()),
executor,
fraiseql_core::security::SecurityContext::system_job(
"test-source",
"test-request",
vec![],
vec![],
None,
),
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,
source,
CronSchedule::parse("*/5 * * * *").unwrap(),
module,
Arc::new(observer),
PostgresSourceCursorStore::new(pool.clone()),
executor.clone(),
fraiseql_core::security::SecurityContext::system_job(
"test-source",
"test-request",
vec![],
vec![],
None,
),
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 })));
}
async fn wait_for_cursor_page(
pool: &PgPool,
source: &str,
target: i64,
deadline: std::time::Duration,
) -> i64 {
let store = PostgresSourceCursorStore::new(pool.clone());
let started = std::time::Instant::now();
loop {
let page = store
.load(source)
.await
.unwrap()
.value
.and_then(|v| v.get("page").and_then(Value::as_i64))
.unwrap_or(0);
if page >= target {
return page;
}
assert!(
started.elapsed() < deadline,
"cursor page stuck at {page} (target {target}) after {:?} — the scheduler stopped \
firing (the #796 failure shape)",
started.elapsed()
);
tokio::time::sleep(std::time::Duration::from_secs(2)).await;
}
}
#[tokio::test]
#[ignore = "local-only #573 gate: real V8 guest + real multi-minute schedule windows (~4 min)"]
async fn ingests_across_schedule_windows_and_a_restart() {
let Some((pool, _svc)) = connect_pool().await else {
eprintln!("SKIP ingests_across_schedule_windows_and_a_restart: no postgres");
return;
};
let source = "test-poller-multi-window";
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 build = || {
let mut observer = FunctionObserver::new();
observer
.register_runtime(RuntimeType::Deno, DenoRuntime::new(&DenoConfig::default()).unwrap());
SourcePoller::new(
source,
source,
CronSchedule::parse("* * * * *").unwrap(),
FunctionModule::from_source(
"connector".to_string(),
CONNECTOR_TS.to_string(),
RuntimeType::Deno,
),
Arc::new(observer),
PostgresSourceCursorStore::new(pool.clone()),
executor.clone(),
fraiseql_core::security::SecurityContext::system_job(
"test-source",
"test-request",
vec![],
vec![],
None,
),
LeaseGuardedRunner::in_process(source),
HostContextConfig::default(),
ResourceLimits::default(),
None,
false,
)
};
let first = tokio::spawn(build().run_forever());
wait_for_cursor_page(&pool, source, 2, std::time::Duration::from_mins(4)).await;
first.abort();
assert!(first.await.unwrap_err().is_cancelled());
let before = wait_for_cursor_page(&pool, source, 2, std::time::Duration::from_secs(5)).await;
let seen_at_restart = executor.seen.lock().unwrap().clone();
assert!(
seen_at_restart.len() >= 2,
"each window's firing issued its mutation: {seen_at_restart:?}"
);
assert!(
seen_at_restart[0].contains("(page: 1)"),
"the first window ingested from a fresh cursor: {}",
seen_at_restart[0]
);
let second = tokio::spawn(build().run_forever());
wait_for_cursor_page(&pool, source, before + 1, std::time::Duration::from_mins(4)).await;
second.abort();
assert!(second.await.unwrap_err().is_cancelled());
let seen = executor.seen.lock().unwrap().clone();
let first_after_restart = &seen[seen_at_restart.len()];
assert!(
first_after_restart.contains(&format!("(page: {})", before + 1)),
"after the restart the connector resumed from the durable cursor (expected page {}, \
got: {first_after_restart})",
before + 1
);
}
#[tokio::test]
async fn the_poller_advances_the_declared_cursor_not_the_source_name() {
let Some((pool, _svc)) = connect_pool().await else {
eprintln!("SKIP the_poller_advances_the_declared_cursor_not_the_source_name: no postgres");
return;
};
let source_name = "test-poller-renamed-source";
let declared_cursor = "test-poller-original-cursor";
let store = PostgresSourceCursorStore::new(pool.clone());
store.init().await.unwrap();
for key in [source_name, declared_cursor] {
sqlx::query("DELETE FROM _fraiseql_source_cursor WHERE source_name = $1")
.bind(key)
.execute(&pool)
.await
.unwrap();
}
let poller = SourcePoller::new(
source_name,
declared_cursor,
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)),
fraiseql_core::security::SecurityContext::system_job(
"test-source",
"test-request",
vec![],
vec![],
None,
),
LeaseGuardedRunner::in_process(source_name),
HostContextConfig::default(),
ResourceLimits::default(),
None,
false,
);
let host =
poller.build_host(build_source_payload(source_name, "*/5 * * * *", Utc::now()), "tok");
host.advance_cursor(json!({ "page": 7 })).await.unwrap();
assert_eq!(
store.load(declared_cursor).await.unwrap().value,
Some(json!({ "page": 7 })),
"the watermark must land under the declared `cursor` override"
);
assert_eq!(
store.load(source_name).await.unwrap().value,
None,
"nothing may be written under the source name when an override is declared — that is \
the row a rename-with-cursor-preservation was trying to avoid creating"
);
}