fraiseql-server 2.16.0

HTTP server for FraiseQL v2 GraphQL engine
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
//! Tests for the in-memory DLQ cap (#343), retry-failure lifecycle (#343), and
//! the atomic claim primitive (#344).
#![allow(clippy::unwrap_used)] // Reason: test code; lock/await failures should panic to surface bugs.

use std::{
    collections::HashMap,
    sync::{Arc, Barrier},
};

use fraiseql_observers::{ActionConfig, DeadLetterQueue, DlqItem, EntityEvent, EventKind};
use uuid::Uuid;

use super::InMemoryDlq;

fn test_event() -> EntityEvent {
    EntityEvent::new(
        EventKind::Created,
        "TestEntity".to_string(),
        Uuid::new_v4(),
        serde_json::json!({}),
    )
}

fn test_action() -> ActionConfig {
    ActionConfig::Webhook {
        url:                Some("http://localhost/hook".to_string()),
        url_env:            None,
        method:             None,
        headers:            HashMap::new(),
        body_template:      None,
        signing_secret:     None,
        signing_secret_env: None,
    }
}

async fn push(dlq: &InMemoryDlq) -> Uuid {
    dlq.push(test_event(), test_action(), "boom".to_string()).await.unwrap()
}

// ── cap (drop-newest) ───────────────────────────────────────────────────────

#[tokio::test]
async fn unbounded_dlq_grows_without_limit() {
    let dlq = InMemoryDlq::new_with_max(None);
    for _ in 0..5 {
        push(&dlq).await;
    }
    assert_eq!(dlq.count(), 5);
    assert_eq!(dlq.overflow_count(), 0);
}

#[tokio::test]
async fn capped_dlq_drops_newest_at_capacity() {
    let dlq = InMemoryDlq::new_with_max(Some(2));

    let first = push(&dlq).await;
    let second = push(&dlq).await;
    // Third push is at capacity → dropped (drop-newest): the first two remain.
    push(&dlq).await;

    assert_eq!(dlq.count(), 2, "cap should hold the queue at 2 entries");
    assert_eq!(dlq.overflow_count(), 1, "the dropped entry should bump the overflow counter");

    let ids: Vec<Uuid> = dlq.list_all().into_iter().map(|i| i.id).collect();
    assert!(
        ids.contains(&first) && ids.contains(&second),
        "the first two entries are retained"
    );
}

// ── mark_retry_failed keeps the item ────────────────────────────────────────

#[tokio::test]
async fn mark_retry_failed_keeps_item_and_records_failure() {
    let dlq = InMemoryDlq::new_with_max(None);
    let id = push(&dlq).await;

    dlq.mark_retry_failed(id, "second failure").await.unwrap();

    let item = dlq.get(id).expect("item must still be present after a failed retry");
    assert_eq!(item.attempts, 1, "attempts should be incremented");
    assert_eq!(item.error_message, "second failure", "error_message should be updated");
    assert_eq!(dlq.count(), 1);
}

#[tokio::test]
async fn mark_success_removes_item() {
    let dlq = InMemoryDlq::new_with_max(None);
    let id = push(&dlq).await;

    dlq.mark_success(id).await.unwrap();

    assert!(dlq.get(id).is_none(), "a succeeded item should be removed");
    assert_eq!(dlq.count(), 0);
}

// ── #344: atomic claim ──────────────────────────────────────────────────────

#[tokio::test]
async fn try_claim_removes_and_is_idempotent() {
    let dlq = InMemoryDlq::new_with_max(None);
    let id = push(&dlq).await;

    assert!(dlq.try_claim(id).is_some(), "first claim returns the item");
    assert_eq!(dlq.count(), 0, "claim removes the item");
    assert!(dlq.try_claim(id).is_none(), "second claim finds nothing");
}

#[tokio::test]
async fn try_claim_is_atomic_under_concurrency() {
    // Push one item, then race N claimers: exactly one must win. This is the
    // exactly-once gate the retry handlers rely on (#344) — the old
    // get→process→remove path let every racer dispatch the action.
    let dlq = Arc::new(InMemoryDlq::new_with_max(None));
    let id = push(&dlq).await;

    let n = 8;
    let barrier = Arc::new(Barrier::new(n));
    // Spawn ALL threads before joining any — collecting here is intentional so
    // the claimers actually race rather than run serially.
    #[allow(clippy::needless_collect)] // Reason: must spawn all before joining; see comment above.
    let handles: Vec<_> = (0..n)
        .map(|_| {
            let dlq = Arc::clone(&dlq);
            let barrier = Arc::clone(&barrier);
            std::thread::spawn(move || {
                barrier.wait();
                dlq.try_claim(id).is_some()
            })
        })
        .collect();

    let winners = handles.into_iter().map(|h| h.join().unwrap()).filter(|&won| won).count();
    assert_eq!(winners, 1, "exactly one of {n} concurrent claimers should win");
    assert_eq!(dlq.count(), 0);
}

#[tokio::test]
async fn reinsert_bypasses_the_cap() {
    let dlq = InMemoryDlq::new_with_max(Some(1));
    push(&dlq).await; // fills to cap

    // A claimed item whose retry failed must be restored even when the DLQ
    // refilled to capacity during the claim — never silently dropped (#343/#344).
    let claimed = DlqItem {
        id:            Uuid::new_v4(),
        event:         test_event(),
        action:        test_action(),
        error_message: "retry failed".to_string(),
        attempts:      1,
    };
    dlq.reinsert(claimed);

    assert_eq!(dlq.count(), 2, "reinsert must bypass the cap");
    assert_eq!(dlq.overflow_count(), 0, "reinsert is not an overflow");
}

// ── Function-dispatch DLQ ───────────────────────────────────────────────────

mod function_dlq {
    use fraiseql_observers::{DeadLetterQueue, DispatchSource, FunctionDispatchRecord};

    use super::InMemoryDlq;

    fn record(error: &str) -> FunctionDispatchRecord {
        FunctionDispatchRecord::new(
            DispatchSource::AfterMutation,
            "onUserCreated",
            "after:mutation:onUserCreated",
            "0123456789abcdef0123456789abcdef",
            serde_json::json!({ "event_kind": "insert", "new": { "id": "u1" } }),
            error,
            3,
        )
    }

    #[tokio::test]
    async fn exhausted_dispatch_lands_one_row() {
        // A permanently-failing function dispatch that exhausted its retries
        // lands exactly one inspectable DLQ row.
        let dlq = InMemoryDlq::new_with_max(None);

        let id = dlq.push_function(record("upstream 503")).await.unwrap();

        assert_eq!(dlq.function_count(), 1, "one function DLQ row after exhaustion");
        let pending = dlq.get_pending_functions(10).await.unwrap();
        assert_eq!(pending.len(), 1);
        assert_eq!(pending[0].id, id);
        assert_eq!(pending[0].function_name, "onUserCreated");
        assert_eq!(pending[0].attempts, 3);
        assert_eq!(pending[0].error_message, "upstream 503");
    }

    #[tokio::test]
    async fn function_and_observer_entries_are_separate() {
        // The two collections share one store but are counted independently, so a
        // function failure never masks an observer failure or vice versa.
        let dlq = InMemoryDlq::new_with_max(None);

        super::push(&dlq).await; // observer-action failure
        dlq.push_function(record("boom")).await.unwrap();

        assert_eq!(dlq.count(), 1, "observer entries counted separately");
        assert_eq!(dlq.function_count(), 1, "function entries counted separately");
    }

    // #598: the *in-memory* function DLQ is non-durable by design — a dead-lettered
    // dispatch does not survive a restart (constructing a fresh store models a
    // process restart; there is no persistence layer to reload from). This remains a
    // true characterization of the memory store and documents *why* the durable
    // option exists. The Postgres-backed store is the durable counterpart, and its
    // survival is proven in
    // `observers::pg_function_dlq::tests::dead_lettered_dispatch_survives_a_restart`
    // (the #598 flip). `[functions] dlq_store = "postgres"` selects it in production.
    #[tokio::test]
    async fn in_memory_dlq_loses_function_entries_on_restart() {
        // Dead-letter one dispatch.
        let dlq = InMemoryDlq::new_with_max(None);
        dlq.push_function(record("upstream 503")).await.unwrap();
        assert_eq!(dlq.function_count(), 1, "entry present before restart");

        // "Restart": the old store is dropped, a new one is constructed. The only
        // store implementation is in-memory, so the entry is gone — there is no
        // durable table to reload from.
        drop(dlq);
        let after_restart = InMemoryDlq::new_with_max(None);
        assert_eq!(
            after_restart.function_count(),
            0,
            "M-598: the in-memory DLQ loses dead-lettered function dispatches on restart — \
             phase 07's Postgres-backed store must make this survive."
        );
    }

    #[tokio::test]
    async fn capped_function_dlq_drops_newest() {
        let dlq = InMemoryDlq::new_with_max(Some(2));

        dlq.push_function(record("e1")).await.unwrap();
        dlq.push_function(record("e2")).await.unwrap();
        // Third is at capacity → dropped (drop-newest), mirroring `push`.
        dlq.push_function(record("e3")).await.unwrap();

        assert_eq!(dlq.function_count(), 2, "cap holds the function queue at 2");
        assert_eq!(dlq.overflow_count(), 1, "the dropped entry bumps the overflow counter");
    }
}

// ── Listener selection seam (#350) ──────────────────────────────────────────

mod listener_selection {
    use fraiseql_observers::config::TransportKind;

    use super::super::{ListenerSelection, listener_selection};

    #[test]
    fn postgres_uses_the_change_log_listener() {
        assert_eq!(
            listener_selection(TransportKind::Postgres),
            ListenerSelection::PostgresChangeLog,
        );
    }

    #[test]
    fn nats_uses_the_transport_stream_not_the_pg_listener() {
        // The whole point of #350: a NATS selection must NOT fall through to the
        // PostgreSQL listener.
        assert_eq!(listener_selection(TransportKind::Nats), ListenerSelection::TransportStream,);
    }

    #[test]
    fn in_memory_uses_the_transport_stream() {
        assert_eq!(listener_selection(TransportKind::InMemory), ListenerSelection::TransportStream,);
    }
}

// ── Transport boot-fatality predicate (#350) ────────────────────────────────

mod transport_requires_broker {
    use fraiseql_observers::config::TransportKind;

    use super::super::{ObserverRuntime, ObserverRuntimeConfig};

    fn runtime_with(kind: TransportKind) -> ObserverRuntime {
        // A lazy pool never connects, so this needs no database.
        let pool = sqlx::PgPool::connect_lazy("postgres://u:u@127.0.0.1:1/db")
            .expect("lazy pool construction does not connect");
        let mut config = ObserverRuntimeConfig::new(pool);
        config.transport.transport = kind;
        ObserverRuntime::new(config)
    }

    #[tokio::test]
    async fn postgres_start_failure_is_not_boot_fatal() {
        // The default transport keeps the resilient log-and-continue behaviour.
        assert!(!runtime_with(TransportKind::Postgres).transport_requires_broker());
    }

    #[tokio::test]
    async fn nats_start_failure_is_boot_fatal() {
        // A broker-backed transport must take the server down in production if it
        // cannot start, never silently come up on PostgreSQL (#350).
        assert!(runtime_with(TransportKind::Nats).transport_requires_broker());
    }
}

// ── NATS-cannot-run boot gate (#350 acceptance) ─────────────────────────────

// Deliberately NOT gated on `observers-nats`: this must run in the CI `test`
// leg (which compiles `observers` but not `observers-nats`) so it is never a
// false-green. The asserted property — a configured NATS transport that cannot
// run makes `start()` FAIL rather than silently fall back to the PostgreSQL
// listener — holds in both builds, only the failure *reason* differs:
//   • without `observers-nats`: `start_transport_stream` hits the
//     "lacks the observers-nats feature" arm;
//   • with `observers-nats`: `NatsTransport::new` rejects the loopback URL via
//     the transport SSRF guard before any network I/O (a deterministic stand-in
//     for a dead broker, needing neither PG nor NATS infra).
mod nats_unrunnable_gate {
    use fraiseql_observers::config::TransportKind;

    use super::super::{ObserverRuntime, ObserverRuntimeConfig};

    #[tokio::test]
    async fn nats_that_cannot_run_fails_start_with_no_pg_fallback() {
        // A lazy pool never connects, and start_transport_stream builds the
        // transport before touching the database, so this needs no PG.
        let pool = sqlx::PgPool::connect_lazy("postgres://u:u@127.0.0.1:1/db")
            .expect("lazy pool construction does not connect");
        let mut config = ObserverRuntimeConfig::new(pool);
        config.transport.transport = TransportKind::Nats;
        config.transport.nats.url = "nats://127.0.0.1:4222".to_string();

        let mut runtime = ObserverRuntime::new(config);
        let result = runtime.start().await;

        assert!(
            result.is_err(),
            "NATS that cannot run must fail start(), not fall back to PostgreSQL"
        );
        assert!(
            !runtime.is_running(),
            "runtime must not report running after a failed NATS start"
        );
        assert!(
            runtime.transport_requires_broker(),
            "a NATS transport is broker-backed, so the start failure is boot-fatal in production"
        );
    }
}

/// `truncate_log_payload` (#468): small payloads pass through, oversized ones
/// become a bounded marker.
mod log_payload_truncation {
    use super::super::{MAX_LOG_PAYLOAD_BYTES, truncate_log_payload};

    #[test]
    fn small_payload_is_passed_through_unchanged() {
        let data = serde_json::json!({"id": "abc", "status": "new"});
        assert_eq!(truncate_log_payload(&data), data);
    }

    #[test]
    fn oversized_payload_is_replaced_with_a_size_marker() {
        // A string value comfortably larger than the cap.
        let big = "x".repeat(MAX_LOG_PAYLOAD_BYTES + 1_024);
        let data = serde_json::json!({ "blob": big });

        let out = truncate_log_payload(&data);
        assert_eq!(out["_truncated"], serde_json::Value::Bool(true));
        let recorded = out["_original_size_bytes"].as_u64().unwrap();
        assert!(
            recorded > u64::try_from(MAX_LOG_PAYLOAD_BYTES).unwrap(),
            "marker must record the original (oversized) byte length"
        );
        // The original (large) content must not be persisted verbatim.
        assert!(out.get("blob").is_none());
    }
}

/// #773: the subscription forward site must never fabricate subscriber-visible
/// events from non-change rows. `EventKind::Custom` — how a Debezium `'r'`
/// (snapshot/read/no-op) change-log row surfaces here — maps to `None` and is
/// filtered before the `EventBridge`; the three real changes map 1:1. The match
/// inside `subscription_operation_for` is exhaustive over the closed `EventKind`,
/// so an unmapped variant is a compile error, not a runtime default.
mod subscription_forwarding {
    use fraiseql_core::runtime::subscription::SubscriptionOperation;
    use fraiseql_observers::EventKind;

    use super::super::subscription_operation_for;

    #[test]
    fn custom_events_are_not_forwarded_to_subscribers() {
        assert_eq!(
            subscription_operation_for(EventKind::Custom),
            None,
            "a snapshot/read/no-op row must never become a subscriber-visible event (#773)"
        );
    }

    #[test]
    fn real_changes_map_one_to_one() {
        assert_eq!(
            subscription_operation_for(EventKind::Created),
            Some(SubscriptionOperation::Create)
        );
        assert_eq!(
            subscription_operation_for(EventKind::Updated),
            Some(SubscriptionOperation::Update)
        );
        assert_eq!(
            subscription_operation_for(EventKind::Deleted),
            Some(SubscriptionOperation::Delete)
        );
    }
}

/// #772: the change-log → `EventBridge` seam must not silently lose events under
/// backpressure. The forward is a bounded, awaited send: a full channel makes the
/// producer wait (upstream backpressure against the durable change log) instead of
/// dropping the event for every subscriber with only a `warn!`.
mod bridge_backpressure {
    use std::sync::Arc;

    use fraiseql_core::{
        runtime::subscription::{SubscriptionManager, SubscriptionOperation},
        schema::{CompiledSchema, SubscriptionDefinition},
    };

    use super::super::forward_to_bridge;
    use crate::subscriptions::{EntityEvent, EventBridge, EventBridgeConfig};

    /// A burst far beyond the bridge channel capacity: every event must reach the
    /// subscriber. On a current-thread runtime the producer loop gets no yield
    /// points unless the send awaits, so the old `try_send` deterministically
    /// dropped everything past the channel capacity.
    #[tokio::test]
    async fn burst_beyond_bridge_capacity_loses_no_events() {
        const BURST: usize = 200;

        let mut schema = CompiledSchema::new();
        schema.subscriptions.push(SubscriptionDefinition::new("orderChanged", "Order"));
        let manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
        manager
            .subscribe("orderChanged", serde_json::json!({}), serde_json::json!({}), "conn-1")
            .expect("subscribe");
        let mut rx = manager.receiver();

        let bridge = EventBridge::new(
            Arc::clone(&manager),
            EventBridgeConfig::new().with_channel_capacity(2),
        );
        let sender = bridge.sender();
        let handle = bridge.spawn();

        for i in 0..BURST {
            let event = EntityEvent::new(
                "Order",
                format!("order_{i}"),
                SubscriptionOperation::Create,
                serde_json::json!({"id": format!("order_{i}")}),
            );
            forward_to_bridge(&sender, event, &format!("evt-{i}")).await;
        }

        let mut delivered = 0_usize;
        while delivered < BURST {
            match tokio::time::timeout(std::time::Duration::from_secs(2), rx.recv()).await {
                Ok(result) => {
                    let _payload = result.expect("broadcast receiver must stay healthy");
                    delivered += 1;
                },
                Err(_) => break, // no more events coming
            }
        }

        assert_eq!(
            delivered, BURST,
            "every event in the burst must reach the subscriber — a full bridge channel \
             must apply backpressure, never silently drop (#772)"
        );

        handle.abort();
    }
}

/// #932: every status this runtime writes into `tb_observer_log` must be one the
/// shipped `ck_observer_log_status` CHECK accepts. It emitted `"error"`, which
/// the constraint has never permitted, so the INSERT was rejected and the audit
/// row disappeared behind a `warn!` — precisely when a delivery is failing and
/// the record is what an operator needs.
#[test]
fn observer_log_statuses_are_accepted_by_the_shipped_check() {
    use fraiseql_observers::migrations::OBSERVER_LOG_STATUSES;

    for status in [
        super::OBSERVER_LOG_STATUS_SUCCESS,
        super::OBSERVER_LOG_STATUS_FAILED,
    ] {
        assert!(
            OBSERVER_LOG_STATUSES.contains(&status),
            "the runtime writes tb_observer_log.status = {status:?}, which migration 06's \
             ck_observer_log_status rejects (accepts: {OBSERVER_LOG_STATUSES:?}) — the row \
             is silently dropped"
        );
    }
}