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
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
mod event_bridge_tests {
    use std::sync::Arc;

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

    use super::super::event_bridge::*;

    #[test]
    fn test_event_bridge_creation() {
        let schema = Arc::new(CompiledSchema::new());
        let manager = Arc::new(SubscriptionManager::new(schema));
        let config = EventBridgeConfig::new();

        let bridge = EventBridge::new(manager, config);

        // Verify bridge is created
        assert!(
            bridge.sender().try_reserve().is_ok(),
            "event bridge channel should have capacity for at least one message"
        );
    }

    #[test]
    fn test_event_conversion_insert() {
        let entity_event = EntityEvent::new(
            "Order",
            "order_123",
            SubscriptionOperation::Create,
            serde_json::json!({
                "id": "order_123",
                "status": "pending"
            }),
        );

        let subscription_event = EventBridge::convert_event(entity_event);

        assert_eq!(subscription_event.entity_type, "Order");
        assert_eq!(subscription_event.entity_id, "order_123");
        assert_eq!(subscription_event.operation, SubscriptionOperation::Create);
    }

    #[test]
    fn test_event_conversion_update() {
        let entity_event = EntityEvent::new(
            "Order",
            "order_123",
            SubscriptionOperation::Update,
            serde_json::json!({
                "id": "order_123",
                "status": "shipped"
            }),
        );

        let subscription_event = EventBridge::convert_event(entity_event);

        assert_eq!(subscription_event.operation, SubscriptionOperation::Update);
    }

    #[test]
    fn test_event_conversion_delete() {
        let entity_event = EntityEvent::new(
            "Order",
            "order_123",
            SubscriptionOperation::Delete,
            serde_json::json!({
                "id": "order_123"
            }),
        );

        let subscription_event = EventBridge::convert_event(entity_event);

        assert_eq!(subscription_event.operation, SubscriptionOperation::Delete);
    }

    #[test]
    fn test_event_conversion_with_old_data() {
        let entity_event = EntityEvent::new(
            "Order",
            "order_123",
            SubscriptionOperation::Update,
            serde_json::json!({
                "id": "order_123",
                "status": "shipped"
            }),
        )
        .with_old_data(serde_json::json!({
            "id": "order_123",
            "status": "pending"
        }));

        let subscription_event = EventBridge::convert_event(entity_event);

        assert!(
            subscription_event.old_data.is_some(),
            "update events should carry old_data for delta computation"
        );
    }

    #[test]
    fn convert_event_propagates_change_spine_envelope() {
        use fraiseql_core::runtime::subscription::ChangeSpineEnvelope;

        let envelope = ChangeSpineEnvelope {
            actor_type:     Some("ai_agent".to_string()),
            acting_for:     Some("11111111-1111-1111-1111-111111111111".to_string()),
            schema_version: Some("v3".to_string()),
            tenant_id:      Some("22222222-2222-2222-2222-222222222222".to_string()),
            duration_ms:    Some(7),
            seq:            Some(99),
        };
        let entity_event = EntityEvent::new(
            "Order",
            "order_123",
            SubscriptionOperation::Update,
            serde_json::json!({ "id": "order_123" }),
        )
        .with_change_spine(envelope.clone());

        let subscription_event = EventBridge::convert_event(entity_event);

        assert_eq!(
            subscription_event.change_spine,
            Some(envelope),
            "the Change-Spine envelope must round-trip through convert_event (#425)"
        );
    }

    #[tokio::test]
    async fn test_event_bridge_spawning() {
        let schema = Arc::new(CompiledSchema::new());
        let manager = Arc::new(SubscriptionManager::new(schema));
        let config = EventBridgeConfig::new();

        let bridge = EventBridge::new(manager, config);
        let handle = bridge.spawn();

        // Verify task was spawned
        assert!(!handle.is_finished());

        // Clean up
        handle.abort();
    }

    /// The `spawn` docstring promises callers may `.await` the handle "for a
    /// clean shutdown". `EventBridge` kept its own `sender` alive inside `run`,
    /// so `recv()` never yielded `None` and that await could never return —
    /// `abort()` was the only way to stop it (#1064).
    #[tokio::test]
    async fn awaiting_the_handle_returns_once_every_external_sender_is_dropped() {
        let schema = Arc::new(CompiledSchema::new());
        let manager = Arc::new(SubscriptionManager::new(schema));
        let config = EventBridgeConfig::new();

        let bridge = EventBridge::new(manager, config);
        let sender = bridge.sender();
        let handle = bridge.spawn();

        drop(sender);

        tokio::time::timeout(std::time::Duration::from_secs(5), handle)
            .await
            .expect("the documented clean-shutdown await never returned")
            .expect("bridge task panicked");
    }

    /// Draining before exit: an event already in the channel must still be
    /// delivered when the sender is dropped straight after.
    #[tokio::test]
    async fn a_queued_event_is_still_delivered_before_shutdown() {
        let schema = Arc::new(CompiledSchema::new());
        let manager = Arc::new(SubscriptionManager::new(schema));
        let config = EventBridgeConfig::new();

        let bridge = EventBridge::new(manager, config);
        let sender = bridge.sender();
        let handle = bridge.spawn();

        sender
            .send(EntityEvent::new(
                "Order",
                "order_1",
                SubscriptionOperation::Create,
                serde_json::json!({"id": "order_1"}),
            ))
            .await
            .expect("channel should be open");
        drop(sender);

        tokio::time::timeout(std::time::Duration::from_secs(5), handle)
            .await
            .expect("bridge did not shut down after its last sender was dropped")
            .expect("bridge task panicked");
    }

    #[tokio::test]
    async fn test_event_bridge_end_to_end_forwarding() {
        let schema = Arc::new(CompiledSchema::new());
        let manager = Arc::new(SubscriptionManager::new(schema));
        let config = EventBridgeConfig::new();

        let bridge = EventBridge::new(manager, config);
        let sender = bridge.sender();
        let handle = bridge.spawn();

        // Send multiple events through the channel
        for i in 0..3 {
            let event = EntityEvent::new(
                "Order",
                format!("order_{i}"),
                SubscriptionOperation::Create,
                serde_json::json!({"id": format!("order_{i}"), "total": 99.95}),
            );
            sender.send(event).await.expect("channel should be open");
        }

        // Yield to let the bridge task process events
        tokio::task::yield_now().await;

        // The bridge should still be running (didn't panic processing events)
        assert!(!handle.is_finished(), "bridge should still be running after processing events");

        handle.abort();
    }

    #[tokio::test]
    async fn test_event_bridge_sender_cloning() {
        let schema = Arc::new(CompiledSchema::new());
        let manager = Arc::new(SubscriptionManager::new(schema));
        let config = EventBridgeConfig::new();

        let bridge = EventBridge::new(manager, config);
        let sender1 = bridge.sender();
        let sender2 = bridge.sender();

        // Both senders should be usable (cloned from the same channel)
        assert!(sender1.try_reserve().is_ok());
        assert!(sender2.try_reserve().is_ok());
    }

    // Tenant-aware CDC filtering: `EventBridge::convert_event` must carry the top-level
    // `tenant_id` from the source `EntityEvent` onto the `SubscriptionEvent` — this is the
    // multi-tenant filtering key on the live `/ws` path. (Ported from the deleted
    // `realtime_integration_test.rs`, whose broadcast/presence tests went with Cluster C but
    // whose live-path conversion tests must survive.)
    #[test]
    fn convert_event_preserves_tenant_id() {
        let entity_event = EntityEvent::new(
            "Order",
            "order_1",
            SubscriptionOperation::Create,
            serde_json::json!({"id": "order_1"}),
        )
        .with_tenant_id("org_42");

        let sub_event = EventBridge::convert_event(entity_event);
        assert_eq!(sub_event.tenant_id.as_deref(), Some("org_42"));
        assert_eq!(sub_event.entity_type, "Order");
    }

    // #773: the bridge can no longer fabricate a Create from an unknown operation
    // string — `EntityEvent.operation` is the closed `SubscriptionOperation` enum, so
    // there is no string to mis-parse. The producer-side filtering (observer
    // `EventKind::Custom`, i.e. a Debezium 'r' snapshot/read row, is never forwarded)
    // is asserted in `observers::runtime::tests::custom_events_are_not_forwarded_to_subscribers`.
    #[test]
    fn convert_event_carries_the_producer_decided_operation_verbatim() {
        for op in [
            SubscriptionOperation::Create,
            SubscriptionOperation::Update,
            SubscriptionOperation::Delete,
        ] {
            let entity_event =
                EntityEvent::new("Order", "order_1", op, serde_json::json!({"id": "order_1"}));
            assert_eq!(EventBridge::convert_event(entity_event).operation, op);
        }
    }

    #[test]
    fn convert_event_without_tenant_id_passes_through_as_none() {
        // No source tenant → `None` (event delivered to all subscribers, not scoped).
        let entity_event = EntityEvent::new(
            "Order",
            "order_1",
            SubscriptionOperation::Create,
            serde_json::json!({"id": "order_1"}),
        );

        let sub_event = EventBridge::convert_event(entity_event);
        assert!(sub_event.tenant_id.is_none());
    }
}

mod entity_fanout_tests {
    #![allow(clippy::unwrap_used)] // Reason: test code, panics acceptable
    #![allow(clippy::panic)] // Reason: test code, panics are the failure mechanism

    use fraiseql_core::runtime::subscription::SubscriptionOperation;

    use super::super::event_bridge::*;

    fn event(entity_id: &str) -> EntityEvent {
        EntityEvent::new("Order", entity_id, SubscriptionOperation::Create, serde_json::json!({}))
    }

    /// The property the whole of #1309 rests on, and the one no `EventTransport` has:
    /// two subscribers each receive **every** event, rather than splitting the stream
    /// between them. `InMemoryTransport` hands out one mpsc receiver,
    /// `PostgresNotifyTransport` marks each batch dispatched, and `NatsTransport` shares
    /// one durable consumer name — so on all three, a second reader takes events away
    /// from the first, and from the observer executor.
    #[test]
    fn every_receiver_gets_every_event() {
        let fanout = EntityEventFanout::new(16);
        let mut first = fanout.subscribe();
        let mut second = fanout.subscribe();

        assert_eq!(fanout.publish(&event("a")), 2, "both receivers must be delivered to");
        assert_eq!(fanout.publish(&event("b")), 2);

        for rx in [&mut first, &mut second] {
            assert_eq!(rx.try_recv().unwrap().entity_id, "a");
            assert_eq!(rx.try_recv().unwrap().entity_id, "b");
        }
    }

    /// Nobody has a `/stream` open. That is the ordinary case — the bridge publishes on
    /// every forwarded event regardless — and it must not read as a failure.
    #[test]
    fn publishing_with_no_receivers_is_not_an_error() {
        let fanout = EntityEventFanout::new(16);
        assert_eq!(fanout.receiver_count(), 0);
        assert_eq!(fanout.publish(&event("a")), 0);
    }

    /// A receiver that never reads is told it lagged rather than served stale events.
    ///
    /// This is what makes the handler's lag arm reachable: without it, the `Lagged`
    /// branch would be a thing the code handles and reality never produces, which is
    /// indistinguishable from dead code.
    #[test]
    fn a_receiver_that_falls_a_whole_capacity_behind_is_told_it_lagged() {
        let fanout = EntityEventFanout::new(2);
        let mut rx = fanout.subscribe();

        for i in 0..5 {
            let _ = fanout.publish(&event(&i.to_string()));
        }

        match rx.try_recv() {
            Err(tokio::sync::broadcast::error::TryRecvError::Lagged(skipped)) => {
                assert_eq!(skipped, 3, "5 published, 2 buffered — 3 dropped");
            },
            other => panic!("a receiver that read nothing must be told it lagged, got {other:?}"),
        }
    }

    /// A late subscriber gets what is published *from now on*, not the buffer's
    /// contents. So opening a `/stream` never replays — which is why `Last-Event-ID`
    /// is refused rather than answered from here (#1310).
    #[test]
    fn a_late_subscriber_does_not_receive_earlier_events() {
        let fanout = EntityEventFanout::new(16);
        let _ = fanout.publish(&event("before"));

        let mut rx = fanout.subscribe();
        assert_eq!(fanout.publish(&event("after")), 1);

        assert_eq!(rx.try_recv().unwrap().entity_id, "after");
        assert!(rx.try_recv().is_err(), "nothing published before the subscription");
    }
}

mod lifecycle_tests {
    use super::super::lifecycle::*;

    #[tokio::test]
    async fn noop_lifecycle_accepts_connect() {
        let lifecycle = NoopLifecycle;
        let result = lifecycle.on_connect(&serde_json::json!({}), "conn-1").await;
        assert!(result.is_ok(), "noop lifecycle should accept any connection");
    }

    #[tokio::test]
    async fn noop_lifecycle_accepts_subscribe() {
        let lifecycle = NoopLifecycle;
        let result = lifecycle.on_subscribe("orderCreated", &serde_json::json!({}), "conn-1").await;
        assert!(result.is_ok(), "noop lifecycle should accept any subscription");
    }
}

mod protocol_tests {
    #![allow(clippy::unwrap_used)] // Reason: test code, panics acceptable
    #![allow(clippy::cast_precision_loss)] // Reason: test metrics reporting
    #![allow(clippy::cast_sign_loss)] // Reason: test data uses small positive integers
    #![allow(clippy::cast_possible_truncation)] // Reason: test data values are bounded
    #![allow(clippy::cast_possible_wrap)] // Reason: test data values are bounded
    #![allow(clippy::missing_panics_doc)] // Reason: test helpers
    #![allow(clippy::missing_errors_doc)] // Reason: test helpers
    #![allow(missing_docs)] // Reason: test code
    #![allow(clippy::items_after_statements)] // Reason: test helpers defined near use site

    use fraiseql_core::runtime::protocol::ServerMessage;

    use super::super::protocol::*;

    // ── WsProtocol::from_header ──────────────────────────────────

    #[test]
    fn from_header_transport_ws() {
        assert_eq!(
            WsProtocol::from_header(Some("graphql-transport-ws")),
            Some(WsProtocol::GraphqlTransportWs)
        );
    }

    #[test]
    fn from_header_legacy_ws() {
        assert_eq!(WsProtocol::from_header(Some("graphql-ws")), Some(WsProtocol::GraphqlWs));
    }

    #[test]
    fn from_header_multiple_prefers_first_known() {
        // Client may offer both; we pick the first recognised one.
        assert_eq!(
            WsProtocol::from_header(Some("graphql-ws, graphql-transport-ws")),
            Some(WsProtocol::GraphqlWs)
        );
        assert_eq!(
            WsProtocol::from_header(Some("graphql-transport-ws, graphql-ws")),
            Some(WsProtocol::GraphqlTransportWs)
        );
    }

    #[test]
    fn from_header_unknown_returns_none() {
        assert_eq!(WsProtocol::from_header(Some("unknown-protocol")), None);
    }

    #[test]
    fn from_header_none_returns_none() {
        assert_eq!(WsProtocol::from_header(None), None);
    }

    // ── ProtocolCodec::decode (modern) ───────────────────────────

    #[test]
    fn decode_transport_ws_subscribe() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlTransportWs);
        let raw = r#"{"type":"subscribe","id":"1","payload":{"query":"subscription { x }"}}"#;
        let msg = codec.decode(raw).unwrap();
        assert_eq!(msg.message_type, "subscribe");
        assert_eq!(msg.id, Some("1".to_string()));
    }

    #[test]
    fn decode_transport_ws_invalid_json() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlTransportWs);
        assert!(
            matches!(codec.decode("not json"), Err(ProtocolError::InvalidJson(_))),
            "expected InvalidJson error for malformed input, got: {:?}",
            codec.decode("not json")
        );
    }

    // ── ProtocolCodec::decode (legacy) ───────────────────────────

    #[test]
    fn decode_legacy_start_becomes_subscribe() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlWs);
        let raw = r#"{"type":"start","id":"1","payload":{"query":"subscription { x }"}}"#;
        let msg = codec.decode(raw).unwrap();
        assert_eq!(msg.message_type, "subscribe");
    }

    #[test]
    fn decode_legacy_stop_becomes_complete() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlWs);
        let raw = r#"{"type":"stop","id":"1"}"#;
        let msg = codec.decode(raw).unwrap();
        assert_eq!(msg.message_type, "complete");
    }

    #[test]
    fn decode_legacy_connection_init_unchanged() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlWs);
        let raw = r#"{"type":"connection_init"}"#;
        let msg = codec.decode(raw).unwrap();
        assert_eq!(msg.message_type, "connection_init");
    }

    // ── ProtocolCodec::encode (modern) ───────────────────────────

    #[test]
    fn encode_transport_ws_next() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlTransportWs);
        let msg = ServerMessage::next("1", serde_json::json!({"x": 1}));
        let json = codec.encode(&msg).unwrap().unwrap();
        assert!(json.contains("\"next\""));
    }

    #[test]
    fn encode_transport_ws_ping() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlTransportWs);
        let msg = ServerMessage::ping(None);
        let json = codec.encode(&msg).unwrap().unwrap();
        assert!(json.contains("\"ping\""));
    }

    // ── ProtocolCodec::encode (legacy) ───────────────────────────

    #[test]
    fn encode_legacy_next_becomes_data() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlWs);
        let msg = ServerMessage::next("1", serde_json::json!({"x": 1}));
        let json = codec.encode(&msg).unwrap().unwrap();
        let parsed: serde_json::Value = serde_json::from_str(&json).unwrap();
        assert_eq!(parsed["type"], "data");
    }

    #[test]
    fn encode_legacy_ping_becomes_ka() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlWs);
        let msg = ServerMessage::ping(None);
        let json = codec.encode(&msg).unwrap().unwrap();
        let parsed: serde_json::Value = serde_json::from_str(&json).unwrap();
        assert_eq!(parsed["type"], "ka");
        // ka has no payload or id
        assert!(parsed.get("payload").is_none() || parsed["payload"].is_null());
    }

    #[test]
    fn encode_legacy_pong_is_suppressed() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlWs);
        let msg = ServerMessage::pong(None);
        let result = codec.encode(&msg).unwrap();
        assert!(result.is_none());
    }

    #[test]
    fn encode_legacy_connection_ack_unchanged() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlWs);
        let msg = ServerMessage::connection_ack(None);
        let json = codec.encode(&msg).unwrap().unwrap();
        let parsed: serde_json::Value = serde_json::from_str(&json).unwrap();
        assert_eq!(parsed["type"], "connection_ack");
    }

    #[test]
    fn encode_legacy_error_unchanged() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlWs);
        let msg = ServerMessage::error(
            "1",
            vec![fraiseql_core::runtime::protocol::GraphQLError::new("test")],
        );
        let json = codec.encode(&msg).unwrap().unwrap();
        let parsed: serde_json::Value = serde_json::from_str(&json).unwrap();
        assert_eq!(parsed["type"], "error");
    }

    // ── uses_keepalive ───────────────────────────────────────────

    #[test]
    fn uses_keepalive_legacy() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlWs);
        assert!(codec.uses_keepalive());
    }

    #[test]
    fn uses_keepalive_modern() {
        let codec = ProtocolCodec::new(WsProtocol::GraphqlTransportWs);
        assert!(!codec.uses_keepalive());
    }
}

mod webhook_lifecycle_tests {
    #![allow(clippy::unwrap_used)] // Reason: test code, panics acceptable
    #![allow(clippy::cast_precision_loss)] // Reason: test metrics reporting
    #![allow(clippy::cast_sign_loss)] // Reason: test data uses small positive integers
    #![allow(clippy::cast_possible_truncation)] // Reason: test data values are bounded
    #![allow(clippy::cast_possible_wrap)] // Reason: test data values are bounded
    #![allow(clippy::missing_panics_doc)] // Reason: test helpers
    #![allow(clippy::missing_errors_doc)] // Reason: test helpers
    #![allow(missing_docs)] // Reason: test code
    #![allow(clippy::items_after_statements)] // Reason: test helpers defined near use site

    use std::time::Duration;

    use super::super::webhook_lifecycle::{MAX_WEBHOOK_RESPONSE_BYTES, WebhookLifecycle};

    #[test]
    fn from_schema_json_no_hooks() {
        let json = serde_json::json!({});
        assert!(WebhookLifecycle::from_schema_json(&json).is_none());
    }

    #[test]
    fn from_schema_json_empty_hooks() {
        let json = serde_json::json!({"hooks": {}});
        assert!(WebhookLifecycle::from_schema_json(&json).is_none());
    }

    #[test]
    fn from_schema_json_with_connect_url() {
        let json = serde_json::json!({
            "hooks": {
                "on_connect": "http://localhost:8001/hooks/ws-connect",
                "timeout_ms": 300
            }
        });
        let wh = WebhookLifecycle::from_schema_json(&json).unwrap();
        assert_eq!(wh.on_connect_url, Some("http://localhost:8001/hooks/ws-connect".to_string()));
        assert!(wh.on_disconnect_url.is_none());
        assert!(wh.on_subscribe_url.is_none());
        assert_eq!(wh.timeout, Duration::from_millis(300));
    }

    #[test]
    fn from_schema_json_default_timeout() {
        let json = serde_json::json!({
            "hooks": {
                "on_disconnect": "http://localhost:8001/hooks/ws-disconnect"
            }
        });
        let wh = WebhookLifecycle::from_schema_json(&json).unwrap();
        assert_eq!(wh.timeout, Duration::from_millis(500));
    }

    #[test]
    fn webhook_response_cap_constant_is_reasonable() {
        // 64 KiB: large enough for any human-readable error, small enough to prevent OOM.
        assert_eq!(MAX_WEBHOOK_RESPONSE_BYTES, 64 * 1024);
    }

    #[test]
    fn webhook_response_body_is_capped_at_limit() {
        // Simulate what on_connect / on_subscribe do: bytes → cap → lossy UTF-8.
        let oversized: Vec<u8> = vec![b'x'; MAX_WEBHOOK_RESPONSE_BYTES + 100];
        let capped = &oversized[..oversized.len().min(MAX_WEBHOOK_RESPONSE_BYTES)];
        let text = String::from_utf8_lossy(capped).into_owned();
        assert_eq!(text.len(), MAX_WEBHOOK_RESPONSE_BYTES);
    }
}