1use std::{collections::HashMap, result::Result as StdResult, sync::Arc};
29
30use ahash::AHashMap;
31use nautilus_core::{
32 UnixNanos,
33 time::{AtomicTime, get_atomic_clock_realtime},
34};
35use nautilus_model::instruments::InstrumentAny;
36use nautilus_network::{
37 http::{HttpClient, HttpClientError, Method, create_standard_nautilus_headers},
38 retry::{RetryConfig, RetryManager},
39 websocket::proxy::ProxyUrl,
40};
41use rust_decimal::Decimal;
42use serde::{Deserialize, Serialize, de::DeserializeOwned};
43use serde_json::{Value, value::RawValue};
44
45use crate::{
46 common::urls::gamma_api_url,
47 filters::set_market_closed,
48 http::{
49 clob::PolymarketClobPublicClient,
50 error::{Error, Result, decode_response},
51 models::{GammaEvent, GammaMarket, GammaTag, SearchResponse},
52 pagination::{Completion, CursorProtocol, FetchOutcome, Paginator, WindowedCollect},
53 parse::{create_instrument_from_def, enrich_market_fee_schedule, parse_gamma_market},
54 query::{GetGammaEventsParams, GetGammaMarketsParams, GetSearchParams},
55 rate_limits::POLYMARKET_GAMMA_REST_QUOTA,
56 },
57};
58
59const GAMMA_MARKETS_KEYSET_PAGE_LIMIT: u32 = 100;
60const GAMMA_EVENTS_KEYSET_PAGE_LIMIT: u32 = 500;
61
62#[derive(Clone, Copy, Debug, Eq, PartialEq)]
63enum GammaStop {
64 CallerCapped,
65}
66
67#[derive(Debug, Clone)]
72pub struct PolymarketGammaRawHttpClient {
73 client: HttpClient,
74 base_url: String,
75}
76
77impl PolymarketGammaRawHttpClient {
78 pub fn new(base_url: Option<String>, timeout_secs: u64) -> StdResult<Self, HttpClientError> {
84 Self::new_with_proxy(base_url, timeout_secs, None)
85 }
86
87 pub fn new_with_proxy(
93 base_url: Option<String>,
94 timeout_secs: u64,
95 proxy_url: Option<ProxyUrl>,
96 ) -> StdResult<Self, HttpClientError> {
97 Ok(Self {
98 client: HttpClient::builder()
99 .headers(Self::default_headers())
100 .default_quota(*POLYMARKET_GAMMA_REST_QUOTA)
101 .timeout_secs(timeout_secs)
102 .maybe_proxy_url(proxy_url.map(|url| url.expose().to_string()))
103 .build()?,
104 base_url: base_url
105 .unwrap_or_else(|| gamma_api_url().to_string())
106 .trim_end_matches('/')
107 .to_string(),
108 })
109 }
110
111 fn default_headers() -> HashMap<String, String> {
112 let mut headers: HashMap<String, String> =
113 create_standard_nautilus_headers().into_iter().collect();
114 headers.insert("Content-Type".to_string(), "application/json".to_string());
115 headers
116 }
117
118 fn url(&self, path: &str) -> String {
119 format!("{}{path}", self.base_url)
120 }
121
122 async fn send_get<P: Serialize, T: DeserializeOwned>(
123 &self,
124 path: &str,
125 params: Option<&P>,
126 ) -> Result<T> {
127 let url = self.url(path);
128 let response = self
129 .client
130 .request_with_params(Method::GET, url, params, None, None, None, None)
131 .await
132 .map_err(Error::from_http_client)?;
133
134 decode_response(&response)
135 }
136
137 async fn send_get_query_map<T: DeserializeOwned>(
138 &self,
139 path: &str,
140 params: Option<&HashMap<String, Vec<String>>>,
141 ) -> Result<T> {
142 let url = self.url(path);
143 let response = self
144 .client
145 .request(Method::GET, url, params, None, None, None, None)
146 .await
147 .map_err(Error::from_http_client)?;
148
149 decode_response(&response)
150 }
151
152 pub async fn get_gamma_markets(
156 &self,
157 params: GetGammaMarketsParams,
158 ) -> Result<Vec<GammaMarket>> {
159 let query_params = gamma_markets_query_params(params)?;
160 let raw: Box<RawValue> = self
161 .send_get_query_map("/markets", Some(&query_params))
162 .await?;
163 parse_gamma_markets_response(&raw)
164 }
165
166 async fn get_gamma_markets_keyset(
167 &self,
168 mut params: GetGammaMarketsParams,
169 after_cursor: Option<&str>,
170 ) -> Result<GammaMarketsKeysetResponse> {
171 params.validate_keyset().map_err(Error::decode)?;
172 params.offset = None;
173 let mut query_params = gamma_markets_query_params(params)?;
174 if let Some(after_cursor) = after_cursor {
175 query_params.insert("after_cursor".to_string(), vec![after_cursor.to_string()]);
176 }
177 self.send_get_query_map("/markets/keyset", Some(&query_params))
178 .await
179 }
180
181 pub async fn get_gamma_market(&self, market_id: &str) -> Result<GammaMarket> {
183 let path = format!("/markets/{market_id}");
184 self.send_get::<(), _>(&path, None::<&()>).await
185 }
186
187 pub async fn get_gamma_market_by_slug(&self, slug: &str) -> Result<GammaMarket> {
189 let path = format!("/markets/slug/{slug}");
190 self.send_get::<(), _>(&path, None::<&()>).await
191 }
192
193 pub async fn get_gamma_events_by_slug(&self, slug: &str) -> Result<Vec<GammaEvent>> {
195 #[derive(Serialize)]
196 struct EventSlugParams<'a> {
197 slug: &'a str,
198 }
199 let params = EventSlugParams { slug };
200 self.send_get("/events", Some(¶ms)).await
201 }
202
203 pub async fn get_gamma_events(&self, params: GetGammaEventsParams) -> Result<Vec<GammaEvent>> {
205 let query_params = gamma_events_query_params(params)?;
206 self.send_get_query_map("/events", Some(&query_params))
207 .await
208 }
209
210 async fn get_gamma_events_keyset(
211 &self,
212 mut params: GetGammaEventsParams,
213 after_cursor: Option<&str>,
214 ) -> Result<GammaEventsKeysetResponse> {
215 params.validate_keyset().map_err(Error::decode)?;
216 params.offset = None;
217 let mut query_params = gamma_events_query_params(params)?;
218 if let Some(after_cursor) = after_cursor {
219 query_params.insert("after_cursor".to_string(), vec![after_cursor.to_string()]);
220 }
221 self.send_get_query_map("/events/keyset", Some(&query_params))
222 .await
223 }
224
225 pub async fn get_gamma_tags(&self) -> Result<Vec<GammaTag>> {
227 self.send_get::<(), _>("/tags", None::<&()>).await
228 }
229
230 pub async fn get_public_search(&self, params: GetSearchParams) -> Result<SearchResponse> {
232 self.send_get("/public-search", Some(¶ms)).await
233 }
234}
235
236#[derive(Debug, Deserialize)]
237struct GammaMarketsKeysetResponse {
238 markets: Vec<GammaMarket>,
239 next_cursor: Option<String>,
240}
241
242#[derive(Debug, Deserialize)]
243struct GammaEventsKeysetResponse {
244 events: Vec<GammaEvent>,
245 next_cursor: Option<String>,
246}
247
248fn gamma_markets_query_params(
249 params: GetGammaMarketsParams,
250) -> Result<HashMap<String, Vec<String>>> {
251 let mut scalar_params = params;
252 let id = scalar_params.id.take();
253 let slug = scalar_params.slug.take();
254 let clob_token_ids = scalar_params.clob_token_ids.take();
255 let condition_ids = scalar_params.condition_ids.take();
256 let question_ids = scalar_params.question_ids.take();
257 let market_maker_address = scalar_params.market_maker_address.take();
258 let tag_id = scalar_params.tag_id.take();
259 let sports_market_types = scalar_params.sports_market_types.take();
260 let value = serde_json::to_value(&scalar_params).map_err(Error::Serde)?;
261 let fields = value
262 .as_object()
263 .ok_or_else(|| Error::decode("Gamma markets params must encode to an object"))?;
264 let mut params = HashMap::with_capacity(fields.len());
265
266 for (key, value) in fields {
267 if let Some(value) = gamma_query_value(value)? {
268 params.insert(key.clone(), vec![value]);
269 }
270 }
271
272 insert_repeated_param(&mut params, "id", id);
273 insert_repeated_param(&mut params, "slug", slug);
274 insert_repeated_param(&mut params, "clob_token_ids", clob_token_ids);
275 insert_repeated_param(&mut params, "condition_ids", condition_ids);
276 insert_repeated_param(&mut params, "question_ids", question_ids);
277 insert_repeated_param(&mut params, "market_maker_address", market_maker_address);
278 insert_repeated_param(&mut params, "tag_id", tag_id);
279 insert_repeated_param(&mut params, "sports_market_types", sports_market_types);
280
281 Ok(params)
282}
283
284fn gamma_events_query_params(params: GetGammaEventsParams) -> Result<HashMap<String, Vec<String>>> {
285 let mut scalar_params = params;
286 let id = scalar_params.id.take();
287 let slug = scalar_params.slug.take();
288 let tag_id = scalar_params.tag_id.take();
289 let exclude_tag_id = scalar_params.exclude_tag_id.take();
290 let series_id = scalar_params.series_id.take();
291 let game_id = scalar_params.game_id.take();
292 let created_by = scalar_params.created_by.take();
293 let value = serde_json::to_value(&scalar_params).map_err(Error::Serde)?;
294 let fields = value
295 .as_object()
296 .ok_or_else(|| Error::decode("Gamma events params must encode to an object"))?;
297 let mut params = HashMap::with_capacity(fields.len());
298
299 for (key, value) in fields {
300 if let Some(value) = gamma_query_value(value)? {
301 params.insert(key.clone(), vec![value]);
302 }
303 }
304
305 insert_repeated_param(&mut params, "id", id);
306 insert_repeated_param(&mut params, "slug", slug);
307 insert_repeated_param(&mut params, "tag_id", tag_id);
308 insert_repeated_param(&mut params, "exclude_tag_id", exclude_tag_id);
309 insert_repeated_param(&mut params, "series_id", series_id);
310 insert_repeated_param(&mut params, "game_id", game_id);
311 insert_repeated_param(&mut params, "created_by", created_by);
312
313 Ok(params)
314}
315
316fn insert_repeated_param<T: ToString>(
317 params: &mut HashMap<String, Vec<String>>,
318 key: &str,
319 values: Option<Vec<T>>,
320) {
321 let Some(values) = values else {
322 return;
323 };
324
325 params.insert(
326 key.to_string(),
327 values
328 .into_iter()
329 .map(|value| value.to_string().trim().to_string())
330 .collect(),
331 );
332}
333
334fn gamma_query_value(value: &Value) -> Result<Option<String>> {
335 match value {
336 Value::Null => Ok(None),
337 Value::String(value) => Ok(Some(value.clone())),
338 Value::Bool(value) => Ok(Some(value.to_string())),
339 Value::Number(value) => Ok(Some(value.to_string())),
340 other => Err(Error::decode(format!(
341 "Unsupported Gamma query value: {other}"
342 ))),
343 }
344}
345
346fn parse_markets_to_instruments(markets: &[GammaMarket], ts_init: UnixNanos) -> Vec<InstrumentAny> {
347 let (instruments, _transient) = parse_markets_with_transient(markets, ts_init);
348 instruments
349}
350
351pub(crate) fn parse_markets_with_transient(
360 markets: &[GammaMarket],
361 ts_init: UnixNanos,
362) -> (Vec<InstrumentAny>, Vec<String>) {
363 let mut instruments = Vec::new();
364 let mut transient = Vec::new();
365
366 for market in markets {
367 if is_transient_clob_token_ids(&market.clob_token_ids) {
368 transient.push(market.condition_id.clone());
369 continue;
370 }
371
372 match parse_gamma_market(market) {
373 Ok(defs) => {
374 for def in defs {
375 match create_instrument_from_def(&def, ts_init) {
376 Ok(InstrumentAny::BinaryOption(mut binary)) => {
377 set_market_closed(&mut binary, def.closed);
378 instruments.push(InstrumentAny::BinaryOption(binary));
379 }
380 Ok(other) => instruments.push(other),
381 Err(e) => log::warn!("Failed to create instrument: {e}"),
382 }
383 }
384 }
385 Err(e) => log::warn!("Failed to parse gamma market: {e}"),
386 }
387 }
388
389 if !transient.is_empty() {
390 log::debug!(
391 "{} market(s) without usable clob_token_ids deferred as transient (CLOB hydration)",
392 transient.len(),
393 );
394 }
395 (instruments, transient)
396}
397
398fn first_token_id(market: &GammaMarket) -> Option<String> {
400 serde_json::from_str::<Vec<String>>(&market.clob_token_ids)
401 .ok()?
402 .into_iter()
403 .find(|token| !token.is_empty())
404}
405
406fn is_transient_clob_token_ids(raw: &str) -> bool {
410 if raw.is_empty() {
411 return true;
412 }
413
414 match serde_json::from_str::<Vec<String>>(raw) {
415 Ok(ids) => ids.is_empty() || ids.iter().any(|t| t.is_empty()),
416 Err(_) => false,
417 }
418}
419
420pub(crate) fn flatten_event_markets(events: Vec<GammaEvent>) -> Vec<GammaMarket> {
421 events
422 .into_iter()
423 .flat_map(|mut event| {
424 let markets = std::mem::take(&mut event.markets);
425 let event = Arc::new(event);
426
427 markets.into_iter().map(move |mut market| {
428 if market.game_id.is_none() {
429 market.game_id.clone_from(&event.game_id);
430 }
431
432 market.parent_event = Some(event.clone());
433 market
434 })
435 })
436 .collect()
437}
438
439#[derive(Debug, Clone)]
445pub struct PolymarketGammaHttpClient {
446 inner: Arc<PolymarketGammaRawHttpClient>,
447 clock: &'static AtomicTime,
448 retry_manager: Arc<RetryManager<Error>>,
449 clob_client: Option<PolymarketClobPublicClient>,
450 fee_rate_cache: Arc<tokio::sync::Mutex<AHashMap<String, Decimal>>>,
451}
452
453impl PolymarketGammaHttpClient {
454 pub fn new(
460 gamma_base_url: Option<String>,
461 timeout_secs: u64,
462 retry_config: RetryConfig,
463 ) -> StdResult<Self, HttpClientError> {
464 Self::new_with_proxy(gamma_base_url, timeout_secs, retry_config, None)
465 }
466
467 pub fn new_with_proxy(
473 gamma_base_url: Option<String>,
474 timeout_secs: u64,
475 retry_config: RetryConfig,
476 proxy_url: Option<ProxyUrl>,
477 ) -> StdResult<Self, HttpClientError> {
478 Ok(Self {
479 inner: Arc::new(PolymarketGammaRawHttpClient::new_with_proxy(
480 gamma_base_url,
481 timeout_secs,
482 proxy_url,
483 )?),
484 clock: get_atomic_clock_realtime(),
485 retry_manager: Arc::new(RetryManager::new(retry_config)),
486 clob_client: None,
487 fee_rate_cache: Arc::new(tokio::sync::Mutex::new(AHashMap::new())),
488 })
489 }
490
491 pub fn set_clob_client(&mut self, clob_client: PolymarketClobPublicClient) {
493 self.clob_client = Some(clob_client);
494 }
495
496 #[must_use]
498 pub fn clob_client(&self) -> Option<&PolymarketClobPublicClient> {
499 self.clob_client.as_ref()
500 }
501
502 async fn enrich_markets(&self, markets: &mut [GammaMarket]) {
507 for market in markets.iter_mut() {
508 enrich_market_fee_schedule(market);
509 }
510
511 let Some(clob) = self.clob_client.as_ref() else {
512 return;
513 };
514
515 for market in markets.iter_mut() {
516 let needs_fallback = matches!(
517 &market.fee_schedule,
518 Some(schedule) if schedule.rate.is_zero() && !schedule.rebate_rate.is_zero()
519 );
520
521 if !needs_fallback {
522 continue;
523 }
524
525 let Some(token_id) = first_token_id(market) else {
526 continue;
527 };
528
529 if let Some(cached) = self.fee_rate_cache.lock().await.get(&token_id).copied() {
530 if let Some(schedule) = market.fee_schedule.as_mut() {
531 schedule.rate = cached;
532 }
533
534 continue;
535 }
536
537 let rate = match clob.get_fee_rate(&token_id).await {
538 Ok(response) => {
539 let rate = response.to_rate();
540 if rate < Decimal::ZERO {
541 log::warn!("Ignoring negative CLOB fee rate {rate} for token {token_id}");
542 continue;
543 }
544
545 self.fee_rate_cache.lock().await.insert(token_id, rate);
546 rate
547 }
548 Err(e) => {
549 log::warn!(
550 "CLOB fee-rate fallback failed for market {}: {e}",
551 market.id
552 );
553 continue;
554 }
555 };
556
557 if let Some(schedule) = market.fee_schedule.as_mut() {
558 schedule.rate = rate;
559 }
560 }
561 }
562
563 async fn fetch_gamma_markets_paginated(
565 &self,
566 base_params: GetGammaMarketsParams,
567 ) -> anyhow::Result<Vec<GammaMarket>> {
568 let page_size = base_params
569 .limit
570 .unwrap_or(GAMMA_MARKETS_KEYSET_PAGE_LIMIT)
571 .min(GAMMA_MARKETS_KEYSET_PAGE_LIMIT);
572 let protocol = CursorProtocol::<GammaStop>::gamma("Gamma market");
573 let reducer = WindowedCollect::new(
574 base_params.offset.unwrap_or(0) as usize,
575 base_params.max_markets.map(|value| value as usize),
576 GammaStop::CallerCapped,
577 );
578 let paginator = Paginator::new("Gamma market", protocol, reducer);
579 let completed = paginator
580 .run(
581 |position| {
582 let after_cursor = position.map(|cursor| cursor.as_ref().to_string());
583 let params = GetGammaMarketsParams {
584 limit: Some(page_size),
585 offset: None,
586 ..base_params.clone()
587 };
588 async move {
589 let response = self
590 .inner
591 .get_gamma_markets_keyset(params, after_cursor.as_deref())
592 .await?;
593 Ok::<_, anyhow::Error>(FetchOutcome::Page {
594 rows: response.markets,
595 wire: response.next_cursor,
596 })
597 }
598 },
599 anyhow::Error::new,
600 )
601 .await?;
602
603 match completed.completion {
604 Completion::WireExhausted | Completion::Stopped(GammaStop::CallerCapped) => {
605 Ok(completed.output)
606 }
607 }
608 }
609
610 async fn fetch_all_gamma_markets(&self) -> anyhow::Result<Vec<GammaMarket>> {
612 self.fetch_gamma_markets_paginated(GetGammaMarketsParams {
613 active: Some(true),
614 closed: Some(false),
615 ..Default::default()
616 })
617 .await
618 }
619
620 pub async fn request_instruments(&self) -> anyhow::Result<Vec<InstrumentAny>> {
626 let mut markets = self.fetch_all_gamma_markets().await?;
627 self.enrich_markets(&mut markets).await;
628 let ts_init = self.clock.get_time_ns();
629 let instruments = parse_markets_to_instruments(&markets, ts_init);
630 log::debug!("Parsed {} instruments from Gamma API", instruments.len());
631 Ok(instruments)
632 }
633
634 pub async fn request_instruments_by_slugs(
644 &self,
645 slugs: Vec<String>,
646 ) -> anyhow::Result<Vec<InstrumentAny>> {
647 let ts_init = self.clock.get_time_ns();
648
649 let futures = slugs.into_iter().map(|slug| {
650 let inner = Arc::clone(&self.inner);
651 async move {
652 let params = GetGammaMarketsParams {
653 slug: Some(vec![slug.clone()]),
654 ..Default::default()
655 };
656
657 match inner.get_gamma_markets(params).await {
658 Ok(markets) => Some((slug, markets)),
659 Err(e) => {
660 log::warn!("Failed to fetch slug '{slug}': {e}");
661 None
662 }
663 }
664 }
665 });
666
667 let results = futures_util::future::join_all(futures).await;
668
669 let total_slugs = results.len();
670 let succeeded = results.iter().filter(|r| r.is_some()).count();
671 let mut instruments = Vec::new();
672
673 for result in results.into_iter().flatten() {
674 let (slug, mut markets) = result;
675 if markets.is_empty() {
676 log::debug!("No markets found for slug '{slug}'");
677 continue;
678 }
679
680 self.enrich_markets(&mut markets).await;
681 instruments.extend(parse_markets_to_instruments(&markets, ts_init));
682 }
683
684 if succeeded == 0 && total_slugs > 0 {
685 anyhow::bail!("All {total_slugs} slug requests failed");
686 }
687
688 log::debug!("Parsed {} instruments from slug queries", instruments.len());
689 Ok(instruments)
690 }
691
692 pub async fn request_instruments_by_slugs_with_retry(
699 &self,
700 slugs: Vec<String>,
701 ) -> anyhow::Result<Vec<InstrumentAny>> {
702 let inner = Arc::clone(&self.inner);
703 let ts_init = self.clock.get_time_ns();
704
705 let mut markets: Vec<GammaMarket> = self
706 .retry_manager
707 .invocation(
708 "gamma_fetch_by_slugs",
709 || {
710 let inner = Arc::clone(&inner);
711 let slugs = slugs.clone();
712 async move {
713 let futures = slugs.into_iter().map(|slug| {
714 let inner = Arc::clone(&inner);
715 async move {
716 let params = GetGammaMarketsParams {
717 slug: Some(vec![slug.clone()]),
718 ..Default::default()
719 };
720 inner
721 .get_gamma_markets(params)
722 .await
723 .map(|markets| (slug, markets))
724 }
725 });
726
727 let results: Vec<_> = futures_util::future::join_all(futures)
728 .await
729 .into_iter()
730 .collect::<StdResult<Vec<_>, _>>()?;
731
732 let markets: Vec<GammaMarket> = results
733 .into_iter()
734 .flat_map(|(_, markets)| markets)
735 .collect();
736
737 if parse_markets_to_instruments(&markets, ts_init).is_empty() {
738 return Err(Error::transport(
739 "Gamma returned no instruments (indexing lag)",
740 ));
741 }
742
743 Ok(markets)
744 }
745 },
746 |e| e.is_retryable(),
747 |e| Error::transport(e.to_string()),
748 )
749 .execute()
750 .await
751 .map_err(|e| anyhow::anyhow!("{e}"))?;
752
753 self.enrich_markets(&mut markets).await;
754 Ok(parse_markets_to_instruments(&markets, ts_init))
755 }
756
757 pub async fn request_instruments_by_event_slugs(
762 &self,
763 event_slugs: Vec<String>,
764 ) -> anyhow::Result<Vec<InstrumentAny>> {
765 let ts_init = self.clock.get_time_ns();
766
767 let futures = event_slugs.into_iter().map(|slug| {
768 let inner = Arc::clone(&self.inner);
769 async move {
770 match inner.get_gamma_events_by_slug(&slug).await {
771 Ok(events) => Some((slug, events)),
772 Err(e) => {
773 log::warn!("Failed to fetch event slug '{slug}': {e}");
774 None
775 }
776 }
777 }
778 });
779
780 let results = futures_util::future::join_all(futures).await;
781
782 let total = results.len();
783 let succeeded = results.iter().filter(|r| r.is_some()).count();
784 let mut instruments = Vec::new();
785
786 for result in results.into_iter().flatten() {
787 let (slug, events) = result;
788 let mut markets = flatten_event_markets(events);
789 if markets.is_empty() {
790 log::warn!("No markets found in event slug '{slug}'");
791 continue;
792 }
793
794 self.enrich_markets(&mut markets).await;
795 instruments.extend(parse_markets_to_instruments(&markets, ts_init));
796 }
797
798 if succeeded == 0 && total > 0 {
799 anyhow::bail!("All {total} event slug requests failed");
800 }
801
802 log::debug!(
803 "Parsed {} instruments from event slug queries",
804 instruments.len()
805 );
806 Ok(instruments)
807 }
808
809 pub async fn request_instruments_by_params(
811 &self,
812 base_params: GetGammaMarketsParams,
813 ) -> anyhow::Result<Vec<InstrumentAny>> {
814 let mut markets = self.fetch_gamma_markets_paginated(base_params).await?;
815 self.enrich_markets(&mut markets).await;
816 let ts_init = self.clock.get_time_ns();
817 let instruments = parse_markets_to_instruments(&markets, ts_init);
818 log::debug!("Parsed {} instruments from params query", instruments.len());
819 Ok(instruments)
820 }
821
822 pub async fn request_instruments_by_params_with_transient(
828 &self,
829 base_params: GetGammaMarketsParams,
830 ) -> anyhow::Result<(Vec<InstrumentAny>, Vec<String>)> {
831 let mut markets = self.fetch_gamma_markets_paginated(base_params).await?;
832 self.enrich_markets(&mut markets).await;
833 let ts_init = self.clock.get_time_ns();
834 let (instruments, transient) = parse_markets_with_transient(&markets, ts_init);
835 log::debug!(
836 "Parsed {} instruments and {} transient condition_id(s) from params query",
837 instruments.len(),
838 transient.len(),
839 );
840 Ok((instruments, transient))
841 }
842
843 pub async fn request_markets_by_params(
845 &self,
846 base_params: GetGammaMarketsParams,
847 ) -> anyhow::Result<Vec<GammaMarket>> {
848 self.fetch_gamma_markets_paginated(base_params).await
849 }
850
851 pub async fn request_instruments_by_event_query(
860 &self,
861 event_slug: &str,
862 params: GetGammaMarketsParams,
863 ) -> anyhow::Result<Vec<InstrumentAny>> {
864 let events = self.inner.get_gamma_events_by_slug(event_slug).await?;
865 let mut markets = flatten_event_markets(events);
866
867 if markets.is_empty() {
868 log::warn!("No markets found in event slug '{event_slug}'");
869 return Ok(Vec::new());
870 }
871
872 log::debug!("Event '{event_slug}' returned {} markets", markets.len());
873
874 if let Some(ref order_field) = params.order {
876 let ascending = params.ascending.unwrap_or(false);
877 markets.sort_by(|a, b| {
878 let cmp = match order_field.as_str() {
879 "liquidity" => a
880 .liquidity_num
881 .unwrap_or(Decimal::ZERO)
882 .partial_cmp(&b.liquidity_num.unwrap_or(Decimal::ZERO)),
883 "volume" => a
884 .volume_num
885 .unwrap_or(Decimal::ZERO)
886 .partial_cmp(&b.volume_num.unwrap_or(Decimal::ZERO)),
887 "volume24hr" => a
888 .volume_24hr
889 .unwrap_or(Decimal::ZERO)
890 .partial_cmp(&b.volume_24hr.unwrap_or(Decimal::ZERO)),
891 "competitive" => a
892 .competitive
893 .unwrap_or(0.0)
894 .partial_cmp(&b.competitive.unwrap_or(0.0)),
895 "spread" => a
896 .spread
897 .unwrap_or(Decimal::MAX)
898 .partial_cmp(&b.spread.unwrap_or(Decimal::MAX)),
899 "best_bid" => a
900 .best_bid
901 .unwrap_or(Decimal::ZERO)
902 .partial_cmp(&b.best_bid.unwrap_or(Decimal::ZERO)),
903 "one_day_price_change" => a
904 .one_day_price_change
905 .unwrap_or(Decimal::ZERO)
906 .partial_cmp(&b.one_day_price_change.unwrap_or(Decimal::ZERO)),
907 "volume_1wk" => a
908 .volume_1wk
909 .unwrap_or(Decimal::ZERO)
910 .partial_cmp(&b.volume_1wk.unwrap_or(Decimal::ZERO)),
911 _ => None,
912 };
913 let cmp = cmp.unwrap_or(std::cmp::Ordering::Equal);
914 if ascending { cmp } else { cmp.reverse() }
915 });
916 }
917
918 if let Some(cap) = params.max_markets {
920 markets.truncate(cap as usize);
921 }
922
923 self.enrich_markets(&mut markets).await;
924 let ts_init = self.clock.get_time_ns();
925 let instruments = parse_markets_to_instruments(&markets, ts_init);
926 log::debug!(
927 "Parsed {} instruments from event query '{event_slug}'",
928 instruments.len()
929 );
930 Ok(instruments)
931 }
932
933 async fn fetch_gamma_events_paginated(
935 &self,
936 base_params: GetGammaEventsParams,
937 ) -> anyhow::Result<Vec<GammaEvent>> {
938 let page_size = base_params
939 .limit
940 .unwrap_or(GAMMA_EVENTS_KEYSET_PAGE_LIMIT)
941 .min(GAMMA_EVENTS_KEYSET_PAGE_LIMIT);
942 let protocol = CursorProtocol::<GammaStop>::gamma("Gamma event");
943 let reducer = WindowedCollect::new(
944 base_params.offset.unwrap_or(0) as usize,
945 base_params.max_events.map(|value| value as usize),
946 GammaStop::CallerCapped,
947 );
948 let paginator = Paginator::new("Gamma event", protocol, reducer);
949 let completed = paginator
950 .run(
951 |position| {
952 let after_cursor = position.map(|cursor| cursor.as_ref().to_string());
953 let params = GetGammaEventsParams {
954 limit: Some(page_size),
955 offset: None,
956 ..base_params.clone()
957 };
958 async move {
959 let response = self
960 .inner
961 .get_gamma_events_keyset(params, after_cursor.as_deref())
962 .await?;
963 Ok::<_, anyhow::Error>(FetchOutcome::Page {
964 rows: response.events,
965 wire: response.next_cursor,
966 })
967 }
968 },
969 anyhow::Error::new,
970 )
971 .await?;
972
973 match completed.completion {
974 Completion::WireExhausted | Completion::Stopped(GammaStop::CallerCapped) => {
975 Ok(completed.output)
976 }
977 }
978 }
979
980 pub async fn request_instruments_by_event_params(
982 &self,
983 params: GetGammaEventsParams,
984 ) -> anyhow::Result<Vec<InstrumentAny>> {
985 let events = self.fetch_gamma_events_paginated(params).await?;
986 let ts_init = self.clock.get_time_ns();
987 let total_events = events.len();
988 let mut markets = flatten_event_markets(events);
989 let total_markets = markets.len();
990 self.enrich_markets(&mut markets).await;
991 let instruments = parse_markets_to_instruments(&markets, ts_init);
992 log::debug!(
993 "Parsed {} instruments from {total_events} events ({total_markets} markets)",
994 instruments.len(),
995 );
996 Ok(instruments)
997 }
998
999 pub async fn request_events_by_params(
1001 &self,
1002 params: GetGammaEventsParams,
1003 ) -> anyhow::Result<Vec<GammaEvent>> {
1004 self.fetch_gamma_events_paginated(params).await
1005 }
1006
1007 pub async fn request_instruments_by_search(
1009 &self,
1010 params: GetSearchParams,
1011 ) -> anyhow::Result<Vec<InstrumentAny>> {
1012 let response = self.inner.get_public_search(params).await?;
1013 let ts_init = self.clock.get_time_ns();
1014
1015 let mut instruments = Vec::new();
1016
1017 if let Some(markets) = response.markets {
1018 let mut markets = markets;
1019 self.enrich_markets(&mut markets).await;
1020 instruments.extend(parse_markets_to_instruments(&markets, ts_init));
1021 }
1022
1023 if let Some(events) = &response.events {
1024 let mut event_markets = flatten_event_markets(events.clone());
1025 self.enrich_markets(&mut event_markets).await;
1026 instruments.extend(parse_markets_to_instruments(&event_markets, ts_init));
1027 }
1028
1029 log::debug!("Parsed {} instruments from search query", instruments.len());
1030 Ok(instruments)
1031 }
1032
1033 pub async fn request_tags(&self) -> anyhow::Result<Vec<GammaTag>> {
1035 Ok(self.inner.get_gamma_tags().await?)
1036 }
1037
1038 #[must_use]
1040 pub fn inner(&self) -> &Arc<PolymarketGammaRawHttpClient> {
1041 &self.inner
1042 }
1043}
1044
1045fn parse_gamma_markets_response(raw: &RawValue) -> Result<Vec<GammaMarket>> {
1046 #[derive(Deserialize)]
1047 struct MarketsEnvelope {
1048 data: Vec<GammaMarket>,
1049 }
1050
1051 if raw.get().starts_with('[') {
1052 return serde_json::from_str(raw.get()).map_err(Error::Serde);
1053 }
1054 serde_json::from_str::<MarketsEnvelope>(raw.get())
1055 .map(|envelope| envelope.data)
1056 .map_err(Error::Serde)
1057}
1058
1059#[cfg(test)]
1060mod tests {
1061 use std::sync::{
1062 Arc,
1063 atomic::{AtomicUsize, Ordering},
1064 };
1065
1066 use rstest::rstest;
1067 use rust_decimal_macros::dec;
1068
1069 use super::*;
1070
1071 fn load_fee_market(filename: &str) -> GammaMarket {
1072 let path = format!("test_data/{filename}");
1073 let content = std::fs::read_to_string(path).unwrap();
1074 serde_json::from_str(&content).unwrap()
1075 }
1076
1077 async fn fee_rate_test_client() -> (
1078 PolymarketGammaHttpClient,
1079 Arc<AtomicUsize>,
1080 tokio::task::JoinHandle<()>,
1081 ) {
1082 fee_rate_test_client_with(axum::http::StatusCode::OK, r#"{"base_fee":700}"#).await
1083 }
1084
1085 async fn fee_rate_test_client_with(
1086 status: axum::http::StatusCode,
1087 body: &'static str,
1088 ) -> (
1089 PolymarketGammaHttpClient,
1090 Arc<AtomicUsize>,
1091 tokio::task::JoinHandle<()>,
1092 ) {
1093 let calls = Arc::new(AtomicUsize::new(0));
1094 let asserted = Arc::clone(&calls);
1095
1096 let router = axum::Router::new().route(
1097 "/fee-rate",
1098 axum::routing::get(move || {
1099 let calls = Arc::clone(&calls);
1100
1101 async move {
1102 calls.fetch_add(1, Ordering::SeqCst);
1103 (status, body.to_string())
1104 }
1105 }),
1106 );
1107
1108 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1109 let address = listener.local_addr().unwrap();
1110 let server = tokio::spawn(async move { axum::serve(listener, router).await.unwrap() });
1111
1112 let clob = PolymarketClobPublicClient::new(Some(format!("http://{address}")), 5).unwrap();
1113 let mut client = PolymarketGammaHttpClient::new(None, 5, RetryConfig::default()).unwrap();
1114 client.set_clob_client(clob);
1115
1116 (client, asserted, server)
1117 }
1118
1119 fn instrument_fee_schedule(instrument: &InstrumentAny) -> crate::http::models::FeeSchedule {
1120 let InstrumentAny::BinaryOption(binary) = instrument else {
1121 panic!("expected a binary option instrument");
1122 };
1123
1124 let info = binary.info.as_ref().unwrap();
1125 let value = info.get("fee_schedule").unwrap();
1126 serde_json::from_value(value.clone()).unwrap()
1127 }
1128
1129 #[rstest]
1130 fn test_live_instrument_funnel_retains_gamma_metadata() {
1131 let raw = include_str!("../../test_data/gamma_market_metadata.json");
1132 let market: GammaMarket = serde_json::from_str(raw).unwrap();
1133 let expected = raw.trim();
1134 let (instruments, transient) = parse_markets_with_transient(&[market], 1.into());
1135 assert_eq!(instruments.len(), 2);
1136 assert_eq!(transient, Vec::<String>::new());
1137
1138 for instrument in instruments {
1139 let InstrumentAny::BinaryOption(binary) = instrument else {
1140 unreachable!()
1141 };
1142
1143 assert_eq!(binary.event_id.map(|id| id.as_str()), Some("event-456"));
1144 let info = binary.info.unwrap();
1145 assert_eq!(info.get_str("gamma_market"), Some(expected));
1146 assert_eq!(info.get_bool("closed"), Some(false));
1147 }
1148 }
1149
1150 #[rstest]
1151 fn test_event_discovery_retains_parent_metadata() {
1152 let raw = include_str!("../../test_data/gamma_event.json");
1153 let events: Vec<GammaEvent> = serde_json::from_str(raw).unwrap();
1154 let expected: Vec<Value> = serde_json::from_str(raw).unwrap();
1155 let markets = flatten_event_markets(events);
1156 assert_eq!(markets.len(), 2);
1157
1158 for market in &markets {
1159 let parent = market.parent_event.as_ref().unwrap();
1160 let expected_event = expected
1161 .iter()
1162 .find(|event| event["id"] == parent.id)
1163 .unwrap();
1164 let expected_market = expected_event["markets"]
1165 .as_array()
1166 .unwrap()
1167 .iter()
1168 .find(|raw| raw["id"] == market.id)
1169 .unwrap();
1170 assert_eq!(
1171 serde_json::from_str::<Value>(&parent.raw).unwrap(),
1172 *expected_event
1173 );
1174 assert_eq!(
1175 serde_json::from_str::<Value>(&market.raw).unwrap(),
1176 *expected_market
1177 );
1178 let defs = parse_gamma_market(market).unwrap();
1179 for def in defs {
1180 assert_eq!(def.event_id.unwrap().as_str(), parent.id);
1181 assert_eq!(def.gamma_event.as_ref(), Some(&parent.raw));
1182 assert_eq!(def.gamma_market, market.raw);
1183 let instrument = create_instrument_from_def(&def, 1.into()).unwrap();
1184
1185 let InstrumentAny::BinaryOption(binary) = instrument else {
1186 unreachable!()
1187 };
1188
1189 let info = binary.info.unwrap();
1190 assert_eq!(binary.event_id.unwrap().as_str(), parent.id);
1191 assert_eq!(info.get_str("gamma_event"), Some(parent.raw.as_str()));
1192 assert_eq!(info.get_str("gamma_market"), Some(market.raw.as_str()));
1193 }
1194 }
1195 }
1196
1197 #[rstest]
1198 #[case("liquidity")]
1199 #[case("volume")]
1200 #[case("volume24hr")]
1201 #[case("spread")]
1202 #[case("best_bid")]
1203 #[case("one_day_price_change")]
1204 #[case("volume_1wk")]
1205 #[tokio::test]
1206 async fn test_event_sort_preserves_adjacent_decimal_values(#[case] field: &str) {
1207 use nautilus_model::instruments::Instrument;
1208
1209 let mut lower: GammaMarket =
1210 serde_json::from_str(include_str!("../../test_data/gamma_market.json")).unwrap();
1211 lower.clob_token_ids = serde_json::to_string(&["1", "2"]).unwrap();
1212 let mut higher = lower.clone();
1213 higher.clob_token_ids = serde_json::to_string(&["3", "4"]).unwrap();
1214
1215 for (market, value) in [
1216 (&mut lower, dec!(0.1234567890123456789012345678)),
1217 (&mut higher, dec!(0.1234567890123456789012345679)),
1218 ] {
1219 match field {
1220 "liquidity" => market.liquidity_num = Some(value),
1221 "volume" => market.volume_num = Some(value),
1222 "volume24hr" => market.volume_24hr = Some(value),
1223 "spread" => market.spread = Some(value),
1224 "best_bid" => market.best_bid = Some(value),
1225 "one_day_price_change" => market.one_day_price_change = Some(value),
1226 "volume_1wk" => market.volume_1wk = Some(value),
1227 _ => unreachable!(),
1228 }
1229 }
1230 let mut event: GammaEvent =
1231 serde_json::from_str(include_str!("../../test_data/decimal_precision_event.json"))
1232 .unwrap();
1233 event.markets = vec![lower, higher];
1234 let response = serde_json::to_string(&vec![event]).unwrap();
1235 let router = axum::Router::new().route(
1236 "/events",
1237 axum::routing::get(move || {
1238 let response = response.clone();
1239 async move { response }
1240 }),
1241 );
1242 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1243 let address = listener.local_addr().unwrap();
1244 let server = tokio::spawn(async move { axum::serve(listener, router).await.unwrap() });
1245 let client = PolymarketGammaHttpClient::new(
1246 Some(format!("http://{address}")),
1247 5,
1248 RetryConfig::default(),
1249 )
1250 .unwrap();
1251 let instruments = client
1252 .request_instruments_by_event_query(
1253 "precision",
1254 GetGammaMarketsParams {
1255 order: Some(field.into()),
1256 ascending: Some(false),
1257 max_markets: Some(1),
1258 ..Default::default()
1259 },
1260 )
1261 .await
1262 .unwrap();
1263 server.abort();
1264 assert_eq!(instruments.len(), 2);
1265 assert_eq!(instruments[0].raw_symbol().as_str(), "3");
1266 assert_eq!(instruments[1].raw_symbol().as_str(), "4");
1267 }
1268
1269 #[rstest]
1270 #[case(false)]
1271 #[case(true)]
1272 fn test_markets_response_preserves_decimal_precision(#[case] enveloped: bool) {
1273 let market = include_str!("../../test_data/decimal_precision_market.json");
1274 let array = format!("[{market}]");
1275 let raw = if enveloped {
1276 format!("{{\"data\":{array}}}")
1277 } else {
1278 array
1279 };
1280 let markets =
1281 parse_gamma_markets_response(&serde_json::from_str::<Box<RawValue>>(&raw).unwrap())
1282 .unwrap();
1283 assert_eq!(markets.len(), 1);
1284 assert_eq!(
1285 markets[0].best_bid,
1286 Some(dec!(0.1234567890123456789012345678))
1287 );
1288 assert_eq!(markets[0].volume_num, Some(dec!(12345678901.123457)));
1289 assert_eq!(
1290 markets[0].fee_schedule.as_ref().unwrap().rate,
1291 dec!(0.1234567890123456789012345678)
1292 );
1293 }
1294
1295 #[tokio::test]
1296 async fn test_enrich_markets_falls_back_only_for_zero_rate_fee_enabled() {
1297 let (client, asserted, server) = fee_rate_test_client().await;
1298
1299 let mut markets = [
1300 "gamma_market_fee_crypto.json",
1301 "gamma_market_fee_zero_rate.json",
1302 "gamma_market_fee_free.json",
1303 "gamma_market_fee_unclassifiable.json",
1304 ]
1305 .map(load_fee_market);
1306
1307 client.enrich_markets(&mut markets).await;
1308 server.abort();
1309
1310 assert_eq!(asserted.load(Ordering::SeqCst), 1);
1311
1312 let crypto = markets[0].fee_schedule.as_ref().unwrap();
1313 assert_eq!(crypto.rate, dec!(0.07));
1314 assert_eq!(crypto.rebate_rate, dec!(0.20));
1315
1316 let recovered = markets[1].fee_schedule.as_ref().unwrap();
1317 assert_eq!(recovered.rate, dec!(0.07));
1318 assert_eq!(recovered.rebate_rate, dec!(0.20));
1319
1320 assert!(markets[2].fee_schedule.is_none());
1321
1322 let unknown = markets[3].fee_schedule.as_ref().unwrap();
1323 assert_eq!(unknown.rate, Decimal::ZERO);
1324 assert_eq!(unknown.rebate_rate, Decimal::ZERO);
1325 }
1326
1327 #[tokio::test]
1328 async fn test_enrich_markets_caches_fee_rate_by_token() {
1329 let (client, asserted, server) = fee_rate_test_client().await;
1330
1331 let market = load_fee_market("gamma_market_fee_zero_rate.json");
1332 let mut markets = [market.clone(), market];
1333
1334 client.enrich_markets(&mut markets).await;
1335 server.abort();
1336
1337 assert_eq!(asserted.load(Ordering::SeqCst), 1);
1338
1339 for market in &markets {
1340 let schedule = market.fee_schedule.as_ref().unwrap();
1341 assert_eq!(schedule.rate, dec!(0.07));
1342 assert_eq!(schedule.rebate_rate, dec!(0.20));
1343 }
1344 }
1345
1346 #[tokio::test]
1347 async fn test_enrich_markets_keeps_zero_rate_when_fee_rate_fails() {
1348 let (client, asserted, server) = fee_rate_test_client_with(
1349 axum::http::StatusCode::INTERNAL_SERVER_ERROR,
1350 r#"{"error":"unavailable"}"#,
1351 )
1352 .await;
1353
1354 let mut markets = [load_fee_market("gamma_market_fee_zero_rate.json")];
1355
1356 client.enrich_markets(&mut markets).await;
1357 server.abort();
1358
1359 assert_eq!(asserted.load(Ordering::SeqCst), 1);
1360
1361 let schedule = markets[0].fee_schedule.as_ref().unwrap();
1362 assert_eq!(schedule.rate, Decimal::ZERO);
1363 assert_eq!(schedule.rebate_rate, dec!(0.20));
1364 }
1365
1366 #[tokio::test]
1367 async fn test_enrich_markets_ignores_negative_fee_rate() {
1368 let (client, asserted, server) =
1369 fee_rate_test_client_with(axum::http::StatusCode::OK, r#"{"base_fee":-100}"#).await;
1370
1371 let mut markets = [load_fee_market("gamma_market_fee_zero_rate.json")];
1372
1373 client.enrich_markets(&mut markets).await;
1374 server.abort();
1375
1376 assert_eq!(asserted.load(Ordering::SeqCst), 1);
1377
1378 let schedule = markets[0].fee_schedule.as_ref().unwrap();
1379 assert_eq!(schedule.rate, Decimal::ZERO);
1380 assert_eq!(schedule.rebate_rate, dec!(0.20));
1381 }
1382
1383 #[tokio::test]
1384 async fn test_request_instruments_by_slugs_enriches_fee_schedules() {
1385 let market = include_str!("../../test_data/gamma_market_fee_crypto.json");
1386 let response = format!("[{market}]");
1387
1388 let router = axum::Router::new().route(
1389 "/markets",
1390 axum::routing::get(move || {
1391 let response = response.clone();
1392
1393 async move { response }
1394 }),
1395 );
1396
1397 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1398 let address = listener.local_addr().unwrap();
1399 let server = tokio::spawn(async move { axum::serve(listener, router).await.unwrap() });
1400 let client = PolymarketGammaHttpClient::new(
1401 Some(format!("http://{address}")),
1402 5,
1403 RetryConfig::default(),
1404 )
1405 .unwrap();
1406
1407 let instruments = client
1408 .request_instruments_by_slugs(vec!["fee-crypto-1".to_string()])
1409 .await
1410 .unwrap();
1411 server.abort();
1412
1413 assert_eq!(instruments.len(), 2);
1414
1415 for instrument in &instruments {
1416 let schedule = instrument_fee_schedule(instrument);
1417 assert_eq!(schedule.rate, dec!(0.07));
1418 assert_eq!(schedule.rebate_rate, dec!(0.20));
1419 }
1420 }
1421}