1use std::sync::OnceLock;
30
31use serde::Serialize;
32use serde::de::DeserializeOwned;
33use snafu::ResultExt;
34
35use crate::cassettes::discovery::Discovery;
36use crate::core::contract::{self, core, ops};
37use crate::core::models::params::ContractParams;
38use crate::core::models::{
39 RawTurnHeaderItem, RawTurnListParams, RawTurnListResponse, SeedDemoRequest, SeedResult,
40 SessionDetailResponse, SessionItem, SessionListParams, SessionListResponse,
41 SessionTracesParams, SessionTracesResponse, SessionUpdateRequest, SpanItem,
42 StandaloneTraceDetail, StatsParams, StatsResponse, TraceListParams, TraceListResponse,
43 TraceParams,
44};
45use crate::decode;
46use crate::error::{Result, error};
47use crate::page;
48use crate::transport::{StreamingTransport, TapesTransport, WireRequest};
49
50#[derive(Debug, Clone, Copy)]
52pub struct CoreClient<T> {
53 transport: T,
54}
55
56impl<T> CoreClient<T> {
57 #[must_use]
59 pub fn new(transport: T) -> Self {
60 Self { transport }
61 }
62
63 #[must_use]
65 pub fn transport(&self) -> &T {
66 &self.transport
67 }
68
69 #[must_use]
71 pub fn into_transport(self) -> T {
72 self.transport
73 }
74}
75
76impl<T: TapesTransport> CoreClient<T> {
77 pub async fn call<R: DeserializeOwned>(
89 &self,
90 operation_id: &str,
91 values: Vec<(&str, String)>,
92 ) -> Result<R> {
93 self.call_with_body(operation_id, values, None).await
94 }
95
96 pub async fn call_with_body<R: DeserializeOwned>(
109 &self,
110 operation_id: &str,
111 values: Vec<(&str, String)>,
112 body: Option<String>,
113 ) -> Result<R> {
114 self.call_shaped(operation_id, values, &[], body).await
115 }
116
117 pub async fn call_with_claimed<R: DeserializeOwned>(
140 &self,
141 operation_id: &str,
142 values: Vec<(&str, String)>,
143 claimed: &[(String, String)],
144 ) -> Result<R> {
145 self.call_shaped(operation_id, values, claimed, None).await
146 }
147
148 async fn call_shaped<R: DeserializeOwned>(
172 &self,
173 operation_id: &str,
174 values: Vec<(&str, String)>,
175 claimed: &[(String, String)],
176 body: Option<String>,
177 ) -> Result<R> {
178 if !claimed.is_empty() && !ops::CLAIM_BEARING_OPS.contains(&operation_id) {
179 return error::ContractClaimsSnafu {
180 operation: operation_id,
181 }
182 .fail();
183 }
184 let method = core()?.method(operation_id)?;
185 let mut request = contract::call_for_with_body(method, values, body)?;
186 request.query.extend(claimed.iter().cloned());
187 let response = self
188 .transport
189 .send(&request)
190 .await
191 .context(error::TransportSnafu)?;
192 decode::json_typed(&response)
193 }
194
195 pub fn request_for(
205 &self,
206 operation_id: &str,
207 values: Vec<(&str, String)>,
208 ) -> Result<WireRequest<'static>> {
209 contract::call_for(core()?.method(operation_id)?, values)
210 }
211
212 async fn with_params<P: ContractParams, R: DeserializeOwned>(&self, params: &P) -> Result<R> {
214 self.call(P::OPERATION, params.values()).await
215 }
216
217 async fn with_params_at<P: ContractParams, R: DeserializeOwned>(
219 &self,
220 params: &P,
221 path: Vec<(&str, String)>,
222 ) -> Result<R> {
223 let mut values: Vec<(&str, String)> = params.values();
224 values.extend(path);
225 self.call(P::OPERATION, values).await
226 }
227
228 async fn with_body<B: Serialize, R: DeserializeOwned>(
230 &self,
231 operation_id: &str,
232 values: Vec<(&str, String)>,
233 body: &B,
234 ) -> Result<R> {
235 let rendered = serde_json::to_string(body).context(error::RenderBodySnafu)?;
236 self.call_with_body(operation_id, values, Some(rendered))
237 .await
238 }
239
240 pub async fn list_sessions(&self, params: &SessionListParams) -> Result<SessionListResponse> {
246 self.call_shaped(ops::LIST_SESSIONS, params.values(), ¶ms.claimed, None)
247 .await
248 }
249
250 pub async fn list_all_sessions(&self, params: &SessionListParams) -> Result<Vec<SessionItem>> {
260 page::walk(|cursor| {
261 let mut params = params.clone();
262 params.cursor = cursor;
263 async move { Ok(self.list_sessions(¶ms).await?.into_page()) }
264 })
265 .await
266 }
267
268 pub async fn get_session(&self, id: &str) -> Result<SessionDetailResponse> {
274 self.call(ops::GET_SESSION, vec![("id", id.to_owned())])
275 .await
276 }
277
278 pub async fn update_session(
284 &self,
285 id: &str,
286 body: &SessionUpdateRequest,
287 ) -> Result<SessionDetailResponse> {
288 self.with_body(ops::UPDATE_SESSION, vec![("id", id.to_owned())], body)
289 .await
290 }
291
292 pub async fn delete_session(&self, id: &str) -> Result<()> {
298 self.call(ops::DELETE_SESSION, vec![("id", id.to_owned())])
299 .await
300 }
301
302 pub async fn get_session_traces(
315 &self,
316 id: &str,
317 params: &SessionTracesParams,
318 ) -> Result<SessionTracesResponse> {
319 self.with_params_at(params, vec![("id", id.to_owned())])
320 .await
321 }
322
323 pub async fn get_whole_session_traces(
338 &self,
339 id: &str,
340 params: &SessionTracesParams,
341 ) -> Result<SessionTracesResponse> {
342 let envelope: OnceLock<SessionTracesResponse> = OnceLock::new();
348 let traces = page::walk(|cursor| {
349 let mut params = params.clone();
350 params.cursor = cursor;
351 let envelope = &envelope;
352 async move {
353 let mut response = self.get_session_traces(id, ¶ms).await?;
354 let page = response.take_page();
355 let _ = envelope.set(response);
356 Ok(page)
357 }
358 })
359 .await?;
360 let mut whole = envelope.into_inner().unwrap_or_default();
363 whole.traces = traces;
364 Ok(whole)
365 }
366
367 pub async fn list_raw_turns(
378 &self,
379 id: &str,
380 params: &RawTurnListParams,
381 ) -> Result<RawTurnListResponse> {
382 self.with_params_at(params, vec![("id", id.to_owned())])
383 .await
384 }
385
386 pub async fn list_all_raw_turns(
398 &self,
399 id: &str,
400 params: &RawTurnListParams,
401 ) -> Result<Vec<RawTurnHeaderItem>> {
402 page::walk(|cursor| {
403 let mut params = params.clone();
404 params.cursor = cursor;
405 async move { Ok(self.list_raw_turns(id, ¶ms).await?.into_page()) }
406 })
407 .await
408 }
409
410 pub async fn list_traces(&self, params: &TraceListParams) -> Result<TraceListResponse> {
416 self.with_params(params).await
417 }
418
419 pub async fn get_trace(
431 &self,
432 trace_id: &str,
433 params: &TraceParams,
434 ) -> Result<StandaloneTraceDetail> {
435 self.with_params_at(params, vec![("trace_id", trace_id.to_owned())])
436 .await
437 }
438
439 pub async fn get_whole_trace(
453 &self,
454 trace_id: &str,
455 params: &TraceParams,
456 ) -> Result<StandaloneTraceDetail> {
457 let envelope: OnceLock<StandaloneTraceDetail> = OnceLock::new();
459 let spans = page::walk(|cursor| {
460 let mut params = params.clone();
461 params.cursor = cursor;
462 let envelope = &envelope;
463 async move {
464 let mut response = self.get_trace(trace_id, ¶ms).await?;
465 let page = response.take_page();
466 let _ = envelope.set(response);
467 Ok(page)
468 }
469 })
470 .await?;
471 let mut whole = envelope.into_inner().unwrap_or_default();
472 whole.spans = spans;
473 Ok(whole)
474 }
475
476 pub async fn get_span(&self, trace_id: &str, span_id: &str) -> Result<SpanItem> {
482 self.call(
483 ops::GET_SPAN,
484 vec![
485 ("trace_id", trace_id.to_owned()),
486 ("span_id", span_id.to_owned()),
487 ],
488 )
489 .await
490 }
491
492 pub async fn get_stats(&self, params: &StatsParams) -> Result<StatsResponse> {
498 self.with_params(params).await
499 }
500
501 pub async fn list_cassettes(&self) -> Result<Discovery> {
512 self.call(ops::LIST_CASSETTES, Vec::new()).await
513 }
514
515 pub async fn seed_demo(&self, body: &SeedDemoRequest) -> Result<SeedResult> {
521 self.with_body(ops::SEED_DEMO, Vec::new(), body).await
522 }
523}
524
525impl<T: StreamingTransport> CoreClient<T> {
526 pub async fn stream(&self, operation_id: &str, values: Vec<(&str, String)>) -> Result<T::Body> {
539 let method = core()?.method(operation_id)?;
540 let request = contract::call_for(method, values)?;
541 self.transport.send_stream(&request).await
542 }
543}
544
545#[cfg(test)]
546#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
547mod tests {
548 use super::*;
549 use crate::core::models::params::PayloadDetail;
550 use crate::path::{PathMode, call_url};
551 use crate::transport::{TransportError, WireResponse};
552 use serde::Deserialize;
553 use serde_json::Value;
554 use std::cell::RefCell;
555 use url::Url;
556
557 struct Recorder {
565 base: Url,
566 responses: RefCell<Vec<Value>>,
567 seen: RefCell<Vec<String>>,
568 bodies: RefCell<Vec<Option<String>>>,
569 }
570
571 impl Recorder {
572 fn new(base: &str, responses: Vec<Value>) -> Self {
573 Self {
574 base: Url::parse(base).unwrap(),
575 responses: RefCell::new(responses),
576 seen: RefCell::new(Vec::new()),
577 bodies: RefCell::new(Vec::new()),
578 }
579 }
580 }
581
582 impl TapesTransport for Recorder {
583 async fn send(
584 &self,
585 request: &WireRequest<'_>,
586 ) -> std::result::Result<WireResponse, TransportError> {
587 let url = call_url(&self.base, request, PathMode::UnderBase)
588 .map_err(|error| TransportError::new(error.to_string()))?;
589 self.seen.borrow_mut().push(url.to_string());
590 self.bodies.borrow_mut().push(request.body.clone());
591 let mut responses = self.responses.borrow_mut();
592 let body = if responses.len() > 1 {
593 responses.remove(0)
594 } else {
595 responses.first().cloned().unwrap_or(Value::Null)
596 };
597 Ok(WireResponse::new(
598 200,
599 url.to_string(),
600 Vec::new(),
601 body.to_string().into_bytes(),
602 ))
603 }
604 }
605
606 fn client(base: &str, response: Value) -> CoreClient<Recorder> {
607 CoreClient::new(Recorder::new(base, vec![response]))
608 }
609
610 #[tokio::test]
611 async fn an_operation_is_routed_through_the_contract_and_the_transport() {
612 let client = client(
613 "https://acme.example/primary/tapes/",
614 serde_json::json!({"traces": []}),
615 );
616 let _ = client
617 .get_session_traces("s-1", &SessionTracesParams::default())
618 .await
619 .unwrap();
620
621 assert_eq!(
622 client.transport().seen.borrow()[0],
623 "https://acme.example/primary/tapes/v1/sessions/s-1/traces",
624 );
625 }
626
627 #[tokio::test]
628 async fn a_typed_method_decodes_the_contracts_own_shape() {
629 let client = client(
632 "http://127.0.0.1:8081",
633 serde_json::json!({
634 "items": [{"id": "s1", "rollup": {"turn_count": 3}}],
635 "next_cursor": "abc",
636 }),
637 );
638 let listing = client
639 .list_sessions(&SessionListParams::default())
640 .await
641 .unwrap();
642
643 assert_eq!(listing.items[0].id, "s1");
644 assert_eq!(listing.items[0].rollup.turn_count, 3);
645 assert_eq!(listing.next_cursor, "abc");
646 }
647
648 #[tokio::test]
649 async fn a_typed_method_survives_a_field_it_has_never_heard_of() {
650 let client = client(
653 "http://127.0.0.1:8081",
654 serde_json::json!({"items": [{"id": "s1", "a_field_from_the_future": 7}]}),
655 );
656 let listing = client
657 .list_sessions(&SessionListParams::default())
658 .await
659 .unwrap();
660 assert_eq!(listing.items[0].id, "s1");
661 }
662
663 #[tokio::test]
664 async fn the_generic_seam_still_decodes_into_a_callers_own_type() {
665 #[derive(Debug, Deserialize)]
668 struct Listing {
669 next_cursor: String,
670 }
671
672 let client = client(
673 "http://127.0.0.1:8081",
674 serde_json::json!({"items": [], "next_cursor": "abc"}),
675 );
676 let got: Listing = client.call(ops::LIST_SESSIONS, Vec::new()).await.unwrap();
677 assert_eq!(got.next_cursor, "abc");
678
679 let raw: Value = client.call(ops::LIST_SESSIONS, Vec::new()).await.unwrap();
680 assert_eq!(raw["next_cursor"], "abc");
681 }
682
683 #[tokio::test]
684 async fn a_typed_parameter_travels_under_the_contracts_own_name() {
685 let client = client("http://127.0.0.1:8081", serde_json::json!({"traces": []}));
686 let _ = client
687 .get_session_traces(
688 "s-1",
689 &SessionTracesParams {
690 payload: Some(PayloadDetail::Preview),
691 ..Default::default()
692 },
693 )
694 .await
695 .unwrap();
696 assert!(
697 client.transport().seen.borrow()[0].ends_with("/traces?payload=preview"),
698 "got: {:?}",
699 client.transport().seen.borrow(),
700 );
701 }
702
703 #[tokio::test]
704 async fn a_page_request_travels_under_the_contracts_own_names_on_every_paged_read() {
705 let client = client("http://127.0.0.1:8081", serde_json::json!({}));
710 let _ = client
711 .get_session_traces(
712 "s-1",
713 &SessionTracesParams {
714 payload: Some(PayloadDetail::Full),
715 limit: Some(50),
716 cursor: Some("c1".to_owned()),
717 },
718 )
719 .await
720 .unwrap();
721 let _ = client
722 .get_trace(
723 "t-1",
724 &TraceParams {
725 payload: None,
726 limit: Some(200),
727 cursor: Some("c2".to_owned()),
728 },
729 )
730 .await
731 .unwrap_err(); let _ = client
733 .list_raw_turns(
734 "s-1",
735 &RawTurnListParams {
736 limit: Some(1000),
737 cursor: Some("c3".to_owned()),
738 },
739 )
740 .await
741 .unwrap();
742
743 let seen = client.transport().seen.borrow();
744 assert_eq!(
745 *seen,
746 vec![
747 "http://127.0.0.1:8081/v1/sessions/s-1/traces?payload=full&limit=50&cursor=c1",
748 "http://127.0.0.1:8081/v1/traces/t-1?limit=200&cursor=c2",
749 "http://127.0.0.1:8081/v1/sessions/s-1/raw_turns?limit=1000&cursor=c3",
750 ],
751 );
752 }
753
754 #[tokio::test]
755 async fn a_standalone_trace_page_without_its_session_id_is_refused() {
756 let client = client(
760 "http://127.0.0.1:8081",
761 serde_json::json!({"trace": {"trace_id": "t-1"}, "spans": []}),
762 );
763 let err = client
764 .get_trace("t-1", &TraceParams::default())
765 .await
766 .unwrap_err();
767 let crate::Error::Decode { source } = &err else {
768 panic!("expected a decode failure, got: {err}");
769 };
770 assert!(
771 source.to_string().contains("session_id"),
772 "the refusal must name the missing property: {source}",
773 );
774 }
775
776 #[tokio::test]
777 async fn a_null_cursor_ends_a_whole_read_instead_of_failing_it() {
778 let traces = client(
782 "http://127.0.0.1:8081",
783 serde_json::json!({"traces": [{"trace": {"trace_id": "t-1"}, "spans": []}], "next_cursor": null}),
784 );
785 let whole = traces
786 .get_whole_session_traces("s-1", &SessionTracesParams::default())
787 .await
788 .unwrap();
789 assert_eq!(whole.traces.len(), 1);
790 assert_eq!(whole.next_cursor, "");
791
792 let trace = client(
793 "http://127.0.0.1:8081",
794 serde_json::json!({"session_id": "s-1", "trace": {"trace_id": "t-1"}, "spans": [{"span_id": "sp-1"}], "next_cursor": null}),
795 );
796 let whole = trace
797 .get_whole_trace("t-1", &TraceParams::default())
798 .await
799 .unwrap();
800 assert_eq!(whole.spans.len(), 1);
801 assert_eq!(whole.next_cursor, "");
802
803 let turns = client(
804 "http://127.0.0.1:8081",
805 serde_json::json!({"items": [{"id": 7}], "next_cursor": null}),
806 );
807 let all = turns
808 .list_all_raw_turns("s-1", &RawTurnListParams::default())
809 .await
810 .unwrap();
811 assert_eq!(all.len(), 1);
812 }
813
814 #[tokio::test]
815 async fn a_listing_walk_follows_the_cursor_to_the_end() {
816 let client = CoreClient::new(Recorder::new(
819 "http://127.0.0.1:8081",
820 vec![
821 serde_json::json!({"items": [{"id": "s1"}], "next_cursor": "c1"}),
822 serde_json::json!({"items": [{"id": "s2"}], "next_cursor": ""}),
823 ],
824 ));
825 let sessions = client
826 .list_all_sessions(&SessionListParams::default())
827 .await
828 .unwrap();
829
830 assert_eq!(
831 sessions.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(),
832 vec!["s1", "s2"],
833 );
834 assert!(
835 client.transport().seen.borrow()[1].contains("cursor=c1"),
836 "got: {:?}",
837 client.transport().seen.borrow(),
838 );
839 }
840
841 #[tokio::test]
842 async fn a_session_traces_walk_keeps_the_first_envelope_and_appends_every_page() {
843 let client = CoreClient::new(Recorder::new(
848 "http://127.0.0.1:8081",
849 vec![
850 serde_json::json!({
851 "schema": "20260615",
852 "session": {"id": "s-1"},
853 "links": [{"from_span_id": "a", "to_span_id": "b"}],
854 "traces": [{"trace": {"trace_id": "t-1"}}],
855 "next_cursor": "c1",
856 }),
857 serde_json::json!({
858 "schema": "20260615",
859 "session": {"id": "s-1"},
860 "links": [{"from_span_id": "a", "to_span_id": "b"}],
861 "traces": [{"trace": {"trace_id": "t-2"}}],
862 }),
863 ],
864 ));
865 let whole = client
866 .get_whole_session_traces(
867 "s-1",
868 &SessionTracesParams {
869 payload: Some(PayloadDetail::Preview),
870 ..Default::default()
871 },
872 )
873 .await
874 .unwrap();
875
876 assert_eq!(whole.session.id, "s-1");
877 assert_eq!(whole.schema, "20260615");
878 assert_eq!(
879 whole.links.len(),
880 1,
881 "the envelope is kept once, not per page"
882 );
883 assert_eq!(
884 whole
885 .traces
886 .iter()
887 .map(|t| t.trace.trace_id.as_str())
888 .collect::<Vec<_>>(),
889 vec!["t-1", "t-2"],
890 );
891 assert_eq!(whole.next_cursor, "", "a whole response has no next page");
892 let seen = client.transport().seen.borrow();
893 assert_eq!(seen.len(), 2);
894 assert!(seen[1].contains("cursor=c1"), "got: {seen:?}");
895 assert!(
896 seen.iter().all(|url| url.contains("payload=preview")),
897 "every page of a walk must carry the caller's parameters: {seen:?}",
898 );
899 }
900
901 #[tokio::test]
902 async fn a_trace_walk_keeps_the_first_envelope_and_appends_every_span_page() {
903 let client = CoreClient::new(Recorder::new(
904 "http://127.0.0.1:8081",
905 vec![
906 serde_json::json!({
907 "session_id": "s-1",
908 "trace": {"trace_id": "t-1"},
909 "links": [{"from_span_id": "a", "to_span_id": "b"}],
910 "spans": [{"span_id": "sp-1"}, {"span_id": "sp-2"}],
911 "next_cursor": "c1",
912 }),
913 serde_json::json!({
914 "session_id": "s-1",
915 "trace": {"trace_id": "t-1"},
916 "links": [{"from_span_id": "a", "to_span_id": "b"}],
917 "spans": [{"span_id": "sp-3"}],
918 "next_cursor": "",
919 }),
920 ],
921 ));
922 let whole = client
923 .get_whole_trace("t-1", &TraceParams::default())
924 .await
925 .unwrap();
926
927 assert_eq!(whole.session_id, "s-1");
928 assert_eq!(whole.trace.trace_id, "t-1");
929 assert_eq!(
930 whole.links.len(),
931 1,
932 "the envelope is kept once, not per page"
933 );
934 assert_eq!(
935 whole
936 .spans
937 .iter()
938 .map(|s| s.span_id.as_str())
939 .collect::<Vec<_>>(),
940 vec!["sp-1", "sp-2", "sp-3"],
941 );
942 assert_eq!(whole.next_cursor, "");
943 assert!(
944 client.transport().seen.borrow()[1].ends_with("/v1/traces/t-1?cursor=c1"),
945 "got: {:?}",
946 client.transport().seen.borrow(),
947 );
948 }
949
950 #[tokio::test]
951 async fn a_raw_turn_walk_follows_the_cursor_to_the_end() {
952 let client = CoreClient::new(Recorder::new(
953 "http://127.0.0.1:8081",
954 vec![
955 serde_json::json!({"items": [{"id": 1}], "next_cursor": "c1"}),
956 serde_json::json!({"items": [{"id": 2, "raw_response_dropped": true}]}),
957 ],
958 ));
959 let turns = client
960 .list_all_raw_turns("s-1", &RawTurnListParams::default())
961 .await
962 .unwrap();
963
964 assert_eq!(turns.iter().map(|t| t.id).collect::<Vec<_>>(), vec![1, 2]);
965 assert!(turns[1].raw_response_dropped);
966 assert!(
967 client.transport().seen.borrow()[1].ends_with("/v1/sessions/s-1/raw_turns?cursor=c1"),
968 "got: {:?}",
969 client.transport().seen.borrow(),
970 );
971 }
972
973 #[tokio::test]
974 async fn a_whole_read_stops_on_a_repeated_cursor_rather_than_hanging() {
975 let client = client(
979 "http://127.0.0.1:8081",
980 serde_json::json!({
981 "session_id": "s-1",
982 "trace": {"trace_id": "t-1"},
983 "spans": [{"span_id": "sp-1"}],
984 "next_cursor": "stuck",
985 }),
986 );
987 let whole = client
988 .get_whole_trace("t-1", &TraceParams::default())
989 .await
990 .unwrap();
991 assert_eq!(client.transport().seen.borrow().len(), 2);
992 assert_eq!(whole.spans.len(), 2);
993 }
994
995 #[tokio::test]
996 async fn an_undeclared_parameter_is_refused_before_the_transport_is_reached() {
997 let client = client("http://127.0.0.1:8081", Value::Null);
998 let err = client
999 .call::<Value>(ops::GET_SESSION, vec![("payolad", "full".to_owned())])
1000 .await
1001 .unwrap_err();
1002 assert!(err.to_string().contains("payolad"), "got: {err}");
1003 assert!(
1004 client.transport().seen.borrow().is_empty(),
1005 "nothing may be sent for a call the contract refused",
1006 );
1007 }
1008
1009 #[tokio::test]
1010 async fn a_typed_body_reaches_the_transport_as_the_contracts_own_json() {
1011 let client = client(
1015 "http://127.0.0.1:8081",
1016 serde_json::json!({"session": {"id": "s-1"}}),
1017 );
1018 let updated = client
1019 .update_session(
1020 "s-1",
1021 &SessionUpdateRequest {
1022 display_name: Some("gum glow charm".to_owned()),
1023 },
1024 )
1025 .await
1026 .unwrap();
1027
1028 assert_eq!(updated.session.id, "s-1");
1029 let bodies = client.transport().bodies.borrow();
1030 let sent: Value = serde_json::from_str(bodies[0].as_deref().unwrap()).unwrap();
1031 assert_eq!(sent["display_name"], "gum glow charm");
1032 }
1033
1034 #[tokio::test]
1035 async fn the_bodyless_facade_refuses_an_operation_that_requires_a_body() {
1036 let client = client("http://127.0.0.1:8081", Value::Null);
1040 let err = client
1041 .call::<Value>(ops::UPDATE_SESSION, vec![("id", "s-1".to_owned())])
1042 .await
1043 .unwrap_err();
1044
1045 assert!(
1046 err.to_string().contains("requires a request body"),
1047 "got: {err}",
1048 );
1049 assert!(client.transport().seen.borrow().is_empty());
1050 }
1051
1052 #[tokio::test]
1053 async fn the_facade_refuses_a_body_on_an_operation_that_declares_none() {
1054 let client = client("http://127.0.0.1:8081", Value::Null);
1055 let err = client
1056 .call_with_body::<Value>(
1057 ops::GET_SESSION,
1058 vec![("id", "s-1".to_owned())],
1059 Some("{}".to_owned()),
1060 )
1061 .await
1062 .unwrap_err();
1063
1064 assert!(
1065 err.to_string().contains("declares no request body"),
1066 "got: {err}",
1067 );
1068 assert!(client.transport().seen.borrow().is_empty());
1069 }
1070
1071 #[tokio::test]
1072 async fn the_bodyless_facade_still_sends_no_body_for_an_ordinary_read() {
1073 let client = client("http://127.0.0.1:8081", serde_json::json!({"items": []}));
1076 let _ = client
1077 .list_sessions(&SessionListParams::default())
1078 .await
1079 .unwrap();
1080 assert_eq!(client.transport().bodies.borrow().as_slice(), [None]);
1081 }
1082
1083 impl crate::transport::StreamingTransport for Recorder {
1084 type Body = Vec<u8>;
1085
1086 async fn send_stream(&self, request: &WireRequest<'_>) -> Result<Self::Body> {
1087 let url = call_url(&self.base, request, PathMode::UnderBase).map_err(|error| {
1088 crate::Error::Transport {
1089 source: TransportError::new(error.to_string()),
1090 }
1091 })?;
1092 self.seen.borrow_mut().push(url.to_string());
1093 Ok(Vec::new())
1094 }
1095 }
1096
1097 #[tokio::test]
1098 async fn the_stream_escape_hatch_builds_the_contract_url() {
1099 let client = CoreClient::new(Recorder::new(
1103 "http://127.0.0.1:8081",
1104 vec![serde_json::json!({})],
1105 ));
1106 let _ = client
1107 .stream(ops::LIST_RAW_TURNS, vec![("id", "s-1".to_owned())])
1108 .await
1109 .unwrap();
1110 let seen = client.transport().seen.borrow();
1111 assert_eq!(seen[0], "http://127.0.0.1:8081/v1/sessions/s-1/raw_turns");
1112 }
1113
1114 #[tokio::test]
1115 async fn a_named_method_and_its_operation_id_build_the_same_request() {
1116 let named = client("http://127.0.0.1:8081", serde_json::json!({}));
1120 let _ = named.get_span("t-1", "sp-1").await.unwrap();
1121
1122 let raw = client("http://127.0.0.1:8081", serde_json::json!({}));
1123 let _: Value = raw
1124 .call(
1125 ops::GET_SPAN,
1126 vec![
1127 ("trace_id", "t-1".to_owned()),
1128 ("span_id", "sp-1".to_owned()),
1129 ],
1130 )
1131 .await
1132 .unwrap();
1133
1134 assert_eq!(
1135 *named.transport().seen.borrow(),
1136 *raw.transport().seen.borrow()
1137 );
1138 }
1139
1140 #[tokio::test]
1141 async fn claimed_params_append_to_the_query_in_order() {
1142 let client = client("http://127.0.0.1:8081", serde_json::json!({"items": []}));
1148 let _ = client
1149 .list_sessions(&SessionListParams {
1150 limit: Some(25),
1151 claimed: vec![
1152 ("flavor".to_owned(), "grape".to_owned()),
1153 ("flavor".to_owned(), "sour cherry".to_owned()),
1154 ("vintage".to_owned(), "1998".to_owned()),
1155 ],
1156 ..Default::default()
1157 })
1158 .await
1159 .unwrap();
1160 assert_eq!(
1161 client.transport().seen.borrow()[0],
1162 "http://127.0.0.1:8081/v1/sessions?limit=25&flavor=grape&flavor=sour+cherry&vintage=1998",
1163 );
1164 }
1165
1166 #[tokio::test]
1167 async fn claimed_values_are_percent_encoded_and_nothing_more() {
1168 let client = client("http://127.0.0.1:8081", serde_json::json!({"items": []}));
1173 let _ = client
1174 .list_sessions(&SessionListParams {
1175 claimed: vec![("flavor".to_owned(), "Grüße 🍇".to_owned())],
1176 ..Default::default()
1177 })
1178 .await
1179 .unwrap();
1180 assert_eq!(
1181 client.transport().seen.borrow()[0],
1182 "http://127.0.0.1:8081/v1/sessions?flavor=Gr%C3%BC%C3%9Fe+%F0%9F%8D%87",
1183 );
1184 }
1185
1186 #[tokio::test]
1187 async fn an_empty_claimed_set_leaves_the_request_as_it_always_was() {
1188 let client = client("http://127.0.0.1:8081", serde_json::json!({"items": []}));
1191 let _ = client
1192 .list_sessions(&SessionListParams::default())
1193 .await
1194 .unwrap();
1195 assert_eq!(
1196 client.transport().seen.borrow()[0],
1197 "http://127.0.0.1:8081/v1/sessions",
1198 );
1199 }
1200
1201 #[tokio::test]
1202 async fn the_page_walk_carries_claimed_params_onto_every_page() {
1203 let client = CoreClient::new(Recorder::new(
1207 "http://127.0.0.1:8081",
1208 vec![
1209 serde_json::json!({"items": [{"id": "s1"}], "next_cursor": "c1"}),
1210 serde_json::json!({"items": [{"id": "s2"}], "next_cursor": ""}),
1211 ],
1212 ));
1213 let _ = client
1214 .list_all_sessions(&SessionListParams {
1215 claimed: vec![("flavor".to_owned(), "grape".to_owned())],
1216 ..Default::default()
1217 })
1218 .await
1219 .unwrap();
1220 let seen = client.transport().seen.borrow();
1221 assert_eq!(seen.len(), 2);
1222 assert!(seen[1].contains("cursor=c1"), "got: {seen:?}");
1223 assert!(
1224 seen.iter().all(|url| url.contains("flavor=grape")),
1225 "every page of a filtered walk must carry the claimed pairs: {seen:?}",
1226 );
1227 }
1228
1229 #[tokio::test]
1230 async fn the_typed_and_untyped_claimed_spellings_build_the_same_request() {
1231 let named = client("http://127.0.0.1:8081", serde_json::json!({"items": []}));
1236 let _ = named
1237 .list_sessions(&SessionListParams {
1238 limit: Some(1),
1239 claimed: vec![("flavor".to_owned(), "grape".to_owned())],
1240 ..Default::default()
1241 })
1242 .await
1243 .unwrap();
1244
1245 let raw = client("http://127.0.0.1:8081", serde_json::json!({"items": []}));
1246 let _: Value = raw
1247 .call_with_claimed(
1248 ops::LIST_SESSIONS,
1249 vec![("limit", "1".to_owned())],
1250 &[("flavor".to_owned(), "grape".to_owned())],
1251 )
1252 .await
1253 .unwrap();
1254
1255 assert_eq!(
1256 *named.transport().seen.borrow(),
1257 *raw.transport().seen.borrow()
1258 );
1259 }
1260
1261 #[tokio::test]
1262 async fn claimed_params_do_not_loosen_the_declared_parameter_refusal() {
1263 let client = client("http://127.0.0.1:8081", Value::Null);
1267 let err = client
1268 .call_with_claimed::<Value>(
1269 ops::LIST_SESSIONS,
1270 vec![("limt", "25".to_owned())],
1271 &[("flavor".to_owned(), "grape".to_owned())],
1272 )
1273 .await
1274 .unwrap_err();
1275 assert!(err.to_string().contains("limt"), "got: {err}");
1276 assert!(
1277 client.transport().seen.borrow().is_empty(),
1278 "nothing may be sent for a call the contract refused",
1279 );
1280 }
1281
1282 #[tokio::test]
1283 async fn claimed_pairs_on_a_non_claim_bearing_operation_are_refused() {
1284 let client = client("http://127.0.0.1:8081", Value::Null);
1290 let err = client
1291 .call_with_claimed::<Value>(
1292 ops::GET_SESSION,
1293 vec![("id", "s-1".to_owned())],
1294 &[("flavor".to_owned(), "grape".to_owned())],
1295 )
1296 .await
1297 .unwrap_err();
1298 assert!(
1299 err.to_string().contains("no claimed filter params"),
1300 "got: {err}",
1301 );
1302 assert!(
1303 client.transport().seen.borrow().is_empty(),
1304 "nothing may be sent for a call the contract refused",
1305 );
1306 }
1307
1308 #[tokio::test]
1309 async fn an_empty_claimed_set_is_permitted_on_every_operation() {
1310 let client = client(
1314 "http://127.0.0.1:8081",
1315 serde_json::json!({"session": {"id": "s-1"}}),
1316 );
1317 let _: Value = client
1318 .call_with_claimed(ops::GET_SESSION, vec![("id", "s-1".to_owned())], &[])
1319 .await
1320 .unwrap();
1321 assert_eq!(
1322 client.transport().seen.borrow()[0],
1323 "http://127.0.0.1:8081/v1/sessions/s-1",
1324 );
1325 }
1326
1327 #[test]
1328 fn every_claim_bearing_operation_is_in_the_vendored_contract() {
1329 for operation in ops::CLAIM_BEARING_OPS {
1333 assert!(
1334 core().unwrap().method(operation).is_ok(),
1335 "ops::CLAIM_BEARING_OPS names {operation:?}, which the vendored contract lacks",
1336 );
1337 }
1338 }
1339}