fraiseql-server 2.14.1

HTTP server for FraiseQL v2 GraphQL engine
Documentation
//! Tests for the source poller.
//!
//! - [`source_payload_carries_the_trigger_context`] is pure.
//! - [`build_host_binds_both_the_cursor_and_the_executor`] proves the poller's novel composition —
//!   that one firing's host reaches *both* the durable cursor (vs real PostgreSQL) *and* the
//!   `run_as` executor — without a V8 guest, so it runs in the PG integration leg. The full Model B
//!   guest-through-poller round-trip (a Deno connector reading its cursor, mutating via
//!   `fraiseql_query`, advancing) is a local-only V8 test that ships with the runnable example.
#![allow(clippy::unwrap_used)] // Reason: test module
#![allow(clippy::print_stderr)] // Reason: skip diagnostic when no backing Postgres
#![allow(clippy::large_futures)] // Reason: a test future holds a poller + runtime; stack size is irrelevant in a #[tokio::test]

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 * * * *");
}

/// The per-firing idempotency token is signed with the server HMAC subkey when one
/// is configured — the poller threads `idempotency_key` through rather than
/// hard-coding the unsigned digest, so a source's token is as unforgeable as every
/// other dispatch path's. A lazy pool never connects (this derives tokens, it does
/// not query) — `#[tokio::test]` only because `connect_lazy` needs a runtime handle.
#[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);

    // The key changes the token: it is threaded through, not ignored.
    assert_ne!(unsigned, signed);
    // The signed form is exactly the keyed HMAC over the firing's stable identity.
    let expected = fraiseql_observers::derive_idempotency_token(
        Some(&secret),
        fraiseql_observers::DispatchSource::Source,
        "connector",
        &payload.trigger_type,
        &payload.data,
    );
    assert_eq!(signed, expected);
    // Both forms honour the 32-hex-char (128-bit) token contract.
    assert_eq!(unsigned.len(), 32);
}

/// A query executor that returns a canned response and records the query it saw, so
/// a test can prove the host reached *an* executor (the poller wired one on).
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))
}

/// Build a poller whose collaborators are all scoped to `source`, with `executor`
/// as its query bridge. The module/observer are inert here — `build_host` does not
/// invoke the guest.
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";
    // Fresh cursor row so re-runs are independent.
    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");

    // The cursor is bound: it round-trips through the host against real Postgres.
    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"
    );

    // The executor is bound: host.query reaches it (the fraiseql_query bridge that
    // production dispatch left unwired).
    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"
    );

    // Durability: the advance persisted beyond the host.
    let snapshot = PostgresSourceCursorStore::new(pool.clone()).load(source).await.unwrap();
    assert_eq!(snapshot.value, Some(json!({ "page": 3 })));
}

/// A minimal Model B connector: read the cursor, mutate via `fraiseql_query`,
/// advance the cursor — the exact loop the #573 issue shows.
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 };
};
"#;

/// The whole Model B slice end-to-end through the poller: a real Deno connector,
/// fired once under the lease, reads its (null) cursor, issues a `fraiseql_query`
/// mutation (reaching the bound executor), and advances the durable cursor.
///
/// LOCAL-ONLY: this invokes a real Deno guest (one V8 isolate), so — like every
/// `runtime-deno` test — it is excluded from CI (embedded V8 SIGSEGVs in the Dagger
/// exec sandbox); the `.dagger` source suite skips it by name. Run locally with
/// `DATABASE_URL` set, one isolate per process.
#[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"
    );

    // The connector's fraiseql_query reached the bound executor. Clone out of the
    // lock so no guard is held across the later await.
    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]);

    // The connector advanced its durable cursor from null → { page: 1 }.
    let snapshot = PostgresSourceCursorStore::new(pool.clone()).load(source).await.unwrap();
    assert_eq!(snapshot.value, Some(json!({ "page": 1 })));
}