mirror of
https://github.com/bitcoinresearchkit/brk.git
synced 2026-07-28 19:28:11 -07:00
bindgen: snap
This commit is contained in:
@@ -11,7 +11,8 @@ use crate::{
|
||||
pub fn generate_imports(output: &mut String) {
|
||||
writeln!(
|
||||
output,
|
||||
r#"use std::sync::Arc;
|
||||
r#"use std::io::Read as _;
|
||||
use std::sync::Arc;
|
||||
use std::ops::{{Bound, RangeBounds}};
|
||||
use serde::de::DeserializeOwned;
|
||||
pub use brk_cohort::*;
|
||||
@@ -59,59 +60,61 @@ impl Default for BrkClientOptions {{
|
||||
}}
|
||||
}}
|
||||
|
||||
/// Base HTTP client for making requests.
|
||||
/// Base HTTP client for making requests. Reuses connections via ureq::Agent.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct BrkClientBase {{
|
||||
agent: ureq::Agent,
|
||||
base_url: String,
|
||||
timeout_secs: u64,
|
||||
}}
|
||||
|
||||
impl BrkClientBase {{
|
||||
/// Create a new client with the given base URL.
|
||||
pub fn new(base_url: impl Into<String>) -> Self {{
|
||||
Self {{
|
||||
base_url: base_url.into(),
|
||||
timeout_secs: 30,
|
||||
}}
|
||||
Self::with_options(BrkClientOptions {{ base_url: base_url.into(), ..Default::default() }})
|
||||
}}
|
||||
|
||||
/// Create a new client with options.
|
||||
pub fn with_options(options: BrkClientOptions) -> Self {{
|
||||
let agent = ureq::Agent::config_builder()
|
||||
.timeout_global(Some(std::time::Duration::from_secs(options.timeout_secs)))
|
||||
.build()
|
||||
.into();
|
||||
Self {{
|
||||
agent,
|
||||
base_url: options.base_url,
|
||||
timeout_secs: options.timeout_secs,
|
||||
}}
|
||||
}}
|
||||
|
||||
fn get(&self, path: &str) -> Result<minreq::Response> {{
|
||||
fn get(&self, path: &str) -> Result<Vec<u8>> {{
|
||||
let base = self.base_url.trim_end_matches('/');
|
||||
let url = format!("{{}}{{}}", base, path);
|
||||
let response = minreq::get(&url)
|
||||
.with_timeout(self.timeout_secs)
|
||||
.send()
|
||||
let mut response = self.agent.get(&url)
|
||||
.call()
|
||||
.map_err(|e| BrkError {{ message: e.to_string() }})?;
|
||||
|
||||
if response.status_code >= 400 {{
|
||||
if response.status().as_u16() >= 400 {{
|
||||
return Err(BrkError {{
|
||||
message: format!("HTTP {{}}", response.status_code),
|
||||
message: format!("HTTP {{}}", response.status().as_u16()),
|
||||
}});
|
||||
}}
|
||||
|
||||
Ok(response)
|
||||
let mut bytes = Vec::new();
|
||||
response.body_mut().as_reader().read_to_end(&mut bytes)
|
||||
.map_err(|e| BrkError {{ message: e.to_string() }})?;
|
||||
Ok(bytes)
|
||||
}}
|
||||
|
||||
/// Make a GET request and deserialize JSON response.
|
||||
pub fn get_json<T: DeserializeOwned>(&self, path: &str) -> Result<T> {{
|
||||
self.get(path)?
|
||||
.json()
|
||||
let bytes = self.get(path)?;
|
||||
serde_json::from_slice(&bytes)
|
||||
.map_err(|e| BrkError {{ message: e.to_string() }})
|
||||
}}
|
||||
|
||||
/// Make a GET request and return raw text response.
|
||||
pub fn get_text(&self, path: &str) -> Result<String> {{
|
||||
self.get(path)?
|
||||
.as_str()
|
||||
.map(|s| s.to_string())
|
||||
let bytes = self.get(path)?;
|
||||
String::from_utf8(bytes)
|
||||
.map_err(|e| BrkError {{ message: e.to_string() }})
|
||||
}}
|
||||
}}
|
||||
|
||||
@@ -13,6 +13,6 @@ exclude = ["examples/"]
|
||||
[dependencies]
|
||||
brk_cohort = { workspace = true }
|
||||
brk_types = { workspace = true }
|
||||
minreq = { workspace = true }
|
||||
ureq = { workspace = true }
|
||||
serde = { workspace = true }
|
||||
serde_json = { workspace = true }
|
||||
|
||||
@@ -7,6 +7,7 @@
|
||||
#![allow(clippy::useless_format)]
|
||||
#![allow(clippy::unnecessary_to_owned)]
|
||||
|
||||
use std::io::Read as _;
|
||||
use std::sync::Arc;
|
||||
use std::ops::{Bound, RangeBounds};
|
||||
use serde::de::DeserializeOwned;
|
||||
@@ -47,59 +48,61 @@ impl Default for BrkClientOptions {
|
||||
}
|
||||
}
|
||||
|
||||
/// Base HTTP client for making requests.
|
||||
/// Base HTTP client for making requests. Reuses connections via ureq::Agent.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct BrkClientBase {
|
||||
agent: ureq::Agent,
|
||||
base_url: String,
|
||||
timeout_secs: u64,
|
||||
}
|
||||
|
||||
impl BrkClientBase {
|
||||
/// Create a new client with the given base URL.
|
||||
pub fn new(base_url: impl Into<String>) -> Self {
|
||||
Self {
|
||||
base_url: base_url.into(),
|
||||
timeout_secs: 30,
|
||||
}
|
||||
Self::with_options(BrkClientOptions { base_url: base_url.into(), ..Default::default() })
|
||||
}
|
||||
|
||||
/// Create a new client with options.
|
||||
pub fn with_options(options: BrkClientOptions) -> Self {
|
||||
let agent = ureq::Agent::config_builder()
|
||||
.timeout_global(Some(std::time::Duration::from_secs(options.timeout_secs)))
|
||||
.build()
|
||||
.into();
|
||||
Self {
|
||||
agent,
|
||||
base_url: options.base_url,
|
||||
timeout_secs: options.timeout_secs,
|
||||
}
|
||||
}
|
||||
|
||||
fn get(&self, path: &str) -> Result<minreq::Response> {
|
||||
fn get(&self, path: &str) -> Result<Vec<u8>> {
|
||||
let base = self.base_url.trim_end_matches('/');
|
||||
let url = format!("{}{}", base, path);
|
||||
let response = minreq::get(&url)
|
||||
.with_timeout(self.timeout_secs)
|
||||
.send()
|
||||
let mut response = self.agent.get(&url)
|
||||
.call()
|
||||
.map_err(|e| BrkError { message: e.to_string() })?;
|
||||
|
||||
if response.status_code >= 400 {
|
||||
if response.status().as_u16() >= 400 {
|
||||
return Err(BrkError {
|
||||
message: format!("HTTP {}", response.status_code),
|
||||
message: format!("HTTP {}", response.status().as_u16()),
|
||||
});
|
||||
}
|
||||
|
||||
Ok(response)
|
||||
let mut bytes = Vec::new();
|
||||
response.body_mut().as_reader().read_to_end(&mut bytes)
|
||||
.map_err(|e| BrkError { message: e.to_string() })?;
|
||||
Ok(bytes)
|
||||
}
|
||||
|
||||
/// Make a GET request and deserialize JSON response.
|
||||
pub fn get_json<T: DeserializeOwned>(&self, path: &str) -> Result<T> {
|
||||
self.get(path)?
|
||||
.json()
|
||||
let bytes = self.get(path)?;
|
||||
serde_json::from_slice(&bytes)
|
||||
.map_err(|e| BrkError { message: e.to_string() })
|
||||
}
|
||||
|
||||
/// Make a GET request and return raw text response.
|
||||
pub fn get_text(&self, path: &str) -> Result<String> {
|
||||
self.get(path)?
|
||||
.as_str()
|
||||
.map(|s| s.to_string())
|
||||
let bytes = self.get(path)?;
|
||||
String::from_utf8(bytes)
|
||||
.map_err(|e| BrkError { message: e.to_string() })
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,7 +13,7 @@ bitcoincore-rpc = ["dep:bitcoincore-rpc"]
|
||||
corepc = ["dep:corepc-client"]
|
||||
fjall = ["dep:fjall"]
|
||||
jiff = ["dep:jiff"]
|
||||
minreq = ["dep:minreq"]
|
||||
ureq = ["dep:ureq"]
|
||||
pco = ["dep:pco"]
|
||||
serde_json = ["dep:serde_json"]
|
||||
tokio = ["dep:tokio"]
|
||||
@@ -25,7 +25,7 @@ bitcoincore-rpc = { workspace = true, optional = true }
|
||||
corepc-client = { workspace = true, optional = true }
|
||||
fjall = { workspace = true, optional = true }
|
||||
jiff = { workspace = true, optional = true }
|
||||
minreq = { workspace = true, optional = true }
|
||||
ureq = { workspace = true, optional = true }
|
||||
pco = { workspace = true, optional = true }
|
||||
serde_json = { workspace = true, optional = true }
|
||||
thiserror = "2.0"
|
||||
|
||||
+16
-26
@@ -35,9 +35,9 @@ pub enum Error {
|
||||
#[error(transparent)]
|
||||
RawDB(#[from] vecdb::RawDBError),
|
||||
|
||||
#[cfg(feature = "minreq")]
|
||||
#[cfg(feature = "ureq")]
|
||||
#[error(transparent)]
|
||||
Minreq(#[from] minreq::Error),
|
||||
Ureq(#[from] ureq::Error),
|
||||
|
||||
#[error(transparent)]
|
||||
SystemTimeError(#[from] time::SystemTimeError),
|
||||
@@ -180,8 +180,8 @@ impl Error {
|
||||
/// Returns false for transient errors worth retrying (timeouts, rate limits, server errors).
|
||||
pub fn is_network_permanently_blocked(&self) -> bool {
|
||||
match self {
|
||||
#[cfg(feature = "minreq")]
|
||||
Error::Minreq(e) => is_minreq_error_permanent(e),
|
||||
#[cfg(feature = "ureq")]
|
||||
Error::Ureq(e) => is_ureq_error_permanent(e),
|
||||
Error::IO(e) => is_io_error_permanent(e),
|
||||
// 403 Forbidden suggests IP/geo blocking; 429 and 5xx are transient
|
||||
Error::HttpStatus { status, .. } => *status == 403,
|
||||
@@ -191,28 +191,18 @@ impl Error {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "minreq")]
|
||||
fn is_minreq_error_permanent(e: &minreq::Error) -> bool {
|
||||
use minreq::Error::*;
|
||||
match e {
|
||||
// DNS resolution failure - likely blocked or misconfigured
|
||||
IoError(io_err) => is_io_error_permanent(io_err),
|
||||
// Check error message for common blocking indicators
|
||||
other => {
|
||||
let msg = format!("{:?}", other);
|
||||
// DNS/connection failures
|
||||
msg.contains("nodename nor servname")
|
||||
|| msg.contains("Name or service not known")
|
||||
|| msg.contains("No such host")
|
||||
|| msg.contains("connection refused")
|
||||
|| msg.contains("Connection refused")
|
||||
// SSL/TLS failures (often due to blocking/MITM)
|
||||
|| msg.contains("certificate")
|
||||
|| msg.contains("SSL")
|
||||
|| msg.contains("TLS")
|
||||
|| msg.contains("handshake")
|
||||
}
|
||||
}
|
||||
#[cfg(feature = "ureq")]
|
||||
fn is_ureq_error_permanent(e: &ureq::Error) -> bool {
|
||||
let msg = format!("{:?}", e);
|
||||
msg.contains("nodename nor servname")
|
||||
|| msg.contains("Name or service not known")
|
||||
|| msg.contains("No such host")
|
||||
|| msg.contains("connection refused")
|
||||
|| msg.contains("Connection refused")
|
||||
|| msg.contains("certificate")
|
||||
|| msg.contains("SSL")
|
||||
|| msg.contains("TLS")
|
||||
|| msg.contains("handshake")
|
||||
}
|
||||
|
||||
fn is_io_error_permanent(e: &std::io::Error) -> bool {
|
||||
|
||||
@@ -9,9 +9,9 @@ repository.workspace = true
|
||||
exclude = ["examples/"]
|
||||
|
||||
[dependencies]
|
||||
brk_error = { workspace = true, features = ["minreq", "serde_json"] }
|
||||
brk_error = { workspace = true, features = ["ureq", "serde_json"] }
|
||||
brk_logger = { workspace = true }
|
||||
brk_types = { workspace = true }
|
||||
tracing = { workspace = true }
|
||||
minreq = { workspace = true }
|
||||
ureq = { workspace = true }
|
||||
serde_json = { workspace = true }
|
||||
|
||||
@@ -9,14 +9,16 @@ use brk_error::{Error, Result};
|
||||
use brk_types::{Date, Height, OHLCCents, Timestamp};
|
||||
use serde_json::Value;
|
||||
use tracing::info;
|
||||
use ureq::Agent;
|
||||
|
||||
use crate::{
|
||||
PriceSource, check_response, default_retry,
|
||||
PriceSource, checked_get, default_retry,
|
||||
ohlc::{compute_ohlc_from_range, date_from_timestamp, ohlc_from_array, timestamp_from_ms},
|
||||
};
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Binance {
|
||||
agent: Agent,
|
||||
path: Option<PathBuf>,
|
||||
_1mn: Option<BTreeMap<Timestamp, OHLCCents>>,
|
||||
_1d: Option<BTreeMap<Date, OHLCCents>>,
|
||||
@@ -24,8 +26,9 @@ pub struct Binance {
|
||||
}
|
||||
|
||||
impl Binance {
|
||||
pub fn init(path: Option<&Path>) -> Self {
|
||||
pub fn new(path: Option<&Path>, agent: Agent) -> Self {
|
||||
Self {
|
||||
agent,
|
||||
path: path.map(|p| p.to_owned()),
|
||||
_1mn: None,
|
||||
_1d: None,
|
||||
@@ -41,7 +44,7 @@ impl Binance {
|
||||
// Try live API data first
|
||||
if self._1mn.as_ref().and_then(|m| m.last_key_value()).is_none_or(|(k, _)| k <= ×tamp)
|
||||
{
|
||||
self._1mn.replace(Self::fetch_1mn()?);
|
||||
self._1mn.replace(self.fetch_1mn()?);
|
||||
}
|
||||
|
||||
let res = compute_ohlc_from_range(
|
||||
@@ -68,11 +71,12 @@ impl Binance {
|
||||
)
|
||||
}
|
||||
|
||||
pub fn fetch_1mn() -> Result<BTreeMap<Timestamp, OHLCCents>> {
|
||||
pub fn fetch_1mn(&self) -> Result<BTreeMap<Timestamp, OHLCCents>> {
|
||||
let agent = &self.agent;
|
||||
default_retry(|_| {
|
||||
let url = Self::url("interval=1m&limit=1000");
|
||||
info!("Fetching {url} ...");
|
||||
let bytes = check_response(minreq::get(&url).with_timeout(30).send()?, &url)?;
|
||||
let bytes = checked_get(agent, &url)?;
|
||||
let json: Value = serde_json::from_slice(&bytes)?;
|
||||
Self::parse_ohlc_array(&json)
|
||||
})
|
||||
@@ -80,7 +84,7 @@ impl Binance {
|
||||
|
||||
pub fn get_from_1d(&mut self, date: &Date) -> Result<OHLCCents> {
|
||||
if self._1d.as_ref().and_then(|m| m.last_key_value()).is_none_or(|(k, _)| k <= date) {
|
||||
self._1d.replace(Self::fetch_1d()?);
|
||||
self._1d.replace(self.fetch_1d()?);
|
||||
}
|
||||
|
||||
self._1d
|
||||
@@ -91,11 +95,12 @@ impl Binance {
|
||||
.ok_or(Error::NotFound("Couldn't find date".into()))
|
||||
}
|
||||
|
||||
pub fn fetch_1d() -> Result<BTreeMap<Date, OHLCCents>> {
|
||||
pub fn fetch_1d(&self) -> Result<BTreeMap<Date, OHLCCents>> {
|
||||
let agent = &self.agent;
|
||||
default_retry(|_| {
|
||||
let url = Self::url("interval=1d");
|
||||
info!("Fetching {url} ...");
|
||||
let bytes = check_response(minreq::get(&url).with_timeout(30).send()?, &url)?;
|
||||
let bytes = checked_get(agent, &url)?;
|
||||
let json: Value = serde_json::from_slice(&bytes)?;
|
||||
Self::parse_date_ohlc_array(&json)
|
||||
})
|
||||
@@ -204,10 +209,8 @@ impl Binance {
|
||||
format!("https://api.binance.com/api/v3/uiKlines?symbol=BTCUSDT&{query}")
|
||||
}
|
||||
|
||||
pub fn ping() -> Result<()> {
|
||||
minreq::get("https://api.binance.com/api/v3/ping")
|
||||
.with_timeout(10)
|
||||
.send()?;
|
||||
pub fn ping(&self) -> Result<()> {
|
||||
self.agent.get("https://api.binance.com/api/v3/ping").call()?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -234,7 +237,7 @@ impl PriceSource for Binance {
|
||||
}
|
||||
|
||||
fn ping(&self) -> Result<()> {
|
||||
Self::ping()
|
||||
self.ping()
|
||||
}
|
||||
|
||||
fn clear(&mut self) {
|
||||
|
||||
@@ -6,16 +6,28 @@ use brk_types::{
|
||||
};
|
||||
use serde_json::Value;
|
||||
use tracing::info;
|
||||
use ureq::Agent;
|
||||
|
||||
use crate::{PriceSource, check_response, default_retry};
|
||||
use crate::{PriceSource, checked_get, default_retry};
|
||||
|
||||
#[derive(Default, Clone)]
|
||||
#[derive(Clone)]
|
||||
#[allow(clippy::upper_case_acronyms)]
|
||||
pub struct BRK {
|
||||
agent: Agent,
|
||||
height_to_ohlc: BTreeMap<Height, Vec<OHLCCents>>,
|
||||
day1_to_ohlc: BTreeMap<Day1, Vec<OHLCCents>>,
|
||||
}
|
||||
|
||||
impl BRK {
|
||||
pub fn new(agent: Agent) -> Self {
|
||||
Self {
|
||||
agent,
|
||||
height_to_ohlc: BTreeMap::new(),
|
||||
day1_to_ohlc: BTreeMap::new(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const API_URL: &str = "https://bitview.space/api/vecs";
|
||||
const CHUNK_SIZE: usize = 10_000;
|
||||
|
||||
@@ -28,7 +40,7 @@ impl BRK {
|
||||
|| ((key + self.height_to_ohlc.get(&key).unwrap().len()) <= height)
|
||||
{
|
||||
self.height_to_ohlc
|
||||
.insert(key, Self::fetch_height_prices(key)?);
|
||||
.insert(key, self.fetch_height_prices(key)?);
|
||||
}
|
||||
|
||||
self.height_to_ohlc
|
||||
@@ -39,7 +51,8 @@ impl BRK {
|
||||
.ok_or(Error::NotFound("Couldn't find height in BRK".into()))
|
||||
}
|
||||
|
||||
fn fetch_height_prices(height: Height) -> Result<Vec<OHLCCents>> {
|
||||
fn fetch_height_prices(&self, height: Height) -> Result<Vec<OHLCCents>> {
|
||||
let agent = &self.agent;
|
||||
default_retry(|_| {
|
||||
let url = format!(
|
||||
"{API_URL}/height-to-price-ohlc?from={}&to={}",
|
||||
@@ -48,7 +61,7 @@ impl BRK {
|
||||
);
|
||||
info!("Fetching {url} ...");
|
||||
|
||||
let bytes = check_response(minreq::get(&url).with_timeout(30).send()?, &url)?;
|
||||
let bytes = checked_get(agent, &url)?;
|
||||
let body: Value = serde_json::from_slice(&bytes)?;
|
||||
|
||||
body.as_array()
|
||||
@@ -68,7 +81,7 @@ impl BRK {
|
||||
if !self.day1_to_ohlc.contains_key(&key)
|
||||
|| ((key + self.day1_to_ohlc.get(&key).unwrap().len()) <= day1)
|
||||
{
|
||||
self.day1_to_ohlc.insert(key, Self::fetch_date_prices(key)?);
|
||||
self.day1_to_ohlc.insert(key, self.fetch_date_prices(key)?);
|
||||
}
|
||||
|
||||
self.day1_to_ohlc
|
||||
@@ -79,7 +92,8 @@ impl BRK {
|
||||
.ok_or(Error::NotFound("Couldn't find date in BRK".into()))
|
||||
}
|
||||
|
||||
fn fetch_date_prices(day1: Day1) -> Result<Vec<OHLCCents>> {
|
||||
fn fetch_date_prices(&self, day1: Day1) -> Result<Vec<OHLCCents>> {
|
||||
let agent = &self.agent;
|
||||
default_retry(|_| {
|
||||
let url = format!(
|
||||
"{API_URL}/day1-to-price-ohlc?from={}&to={}",
|
||||
@@ -88,7 +102,7 @@ impl BRK {
|
||||
);
|
||||
info!("Fetching {url}...");
|
||||
|
||||
let bytes = check_response(minreq::get(&url).with_timeout(30).send()?, &url)?;
|
||||
let bytes = checked_get(agent, &url)?;
|
||||
let body: Value = serde_json::from_slice(&bytes)?;
|
||||
|
||||
body.as_array()
|
||||
@@ -121,8 +135,8 @@ impl BRK {
|
||||
)))
|
||||
}
|
||||
|
||||
pub fn ping() -> Result<()> {
|
||||
minreq::get(API_URL).with_timeout(10).send()?;
|
||||
pub fn ping(&self) -> Result<()> {
|
||||
self.agent.get(API_URL).call()?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -149,7 +163,7 @@ impl PriceSource for BRK {
|
||||
}
|
||||
|
||||
fn ping(&self) -> Result<()> {
|
||||
Self::ping()
|
||||
self.ping()
|
||||
}
|
||||
|
||||
fn clear(&mut self) {
|
||||
|
||||
@@ -4,18 +4,30 @@ use brk_error::{Error, Result};
|
||||
use brk_types::{Date, Height, OHLCCents, Timestamp};
|
||||
use serde_json::Value;
|
||||
use tracing::info;
|
||||
use ureq::Agent;
|
||||
|
||||
use crate::{
|
||||
PriceSource, check_response, default_retry,
|
||||
PriceSource, checked_get, default_retry,
|
||||
ohlc::{compute_ohlc_from_range, date_from_timestamp, ohlc_from_array, timestamp_from_secs},
|
||||
};
|
||||
|
||||
#[derive(Default, Clone)]
|
||||
#[derive(Clone)]
|
||||
pub struct Kraken {
|
||||
agent: Agent,
|
||||
_1mn: Option<BTreeMap<Timestamp, OHLCCents>>,
|
||||
_1d: Option<BTreeMap<Date, OHLCCents>>,
|
||||
}
|
||||
|
||||
impl Kraken {
|
||||
pub fn new(agent: Agent) -> Self {
|
||||
Self {
|
||||
agent,
|
||||
_1mn: None,
|
||||
_1d: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Kraken {
|
||||
fn get_from_1mn(
|
||||
&mut self,
|
||||
@@ -24,7 +36,7 @@ impl Kraken {
|
||||
) -> Result<OHLCCents> {
|
||||
if self._1mn.as_ref().and_then(|m| m.last_key_value()).is_none_or(|(k, _)| k <= ×tamp)
|
||||
{
|
||||
self._1mn.replace(Self::fetch_1mn()?);
|
||||
self._1mn.replace(self.fetch_1mn()?);
|
||||
}
|
||||
compute_ohlc_from_range(
|
||||
self._1mn.as_ref().unwrap(),
|
||||
@@ -34,11 +46,12 @@ impl Kraken {
|
||||
)
|
||||
}
|
||||
|
||||
pub fn fetch_1mn() -> Result<BTreeMap<Timestamp, OHLCCents>> {
|
||||
pub fn fetch_1mn(&self) -> Result<BTreeMap<Timestamp, OHLCCents>> {
|
||||
let agent = &self.agent;
|
||||
default_retry(|_| {
|
||||
let url = Self::url(1);
|
||||
info!("Fetching {url} ...");
|
||||
let bytes = check_response(minreq::get(&url).with_timeout(30).send()?, &url)?;
|
||||
let bytes = checked_get(agent, &url)?;
|
||||
let json: Value = serde_json::from_slice(&bytes)?;
|
||||
Self::parse_ohlc_response(&json)
|
||||
})
|
||||
@@ -46,7 +59,7 @@ impl Kraken {
|
||||
|
||||
fn get_from_1d(&mut self, date: &Date) -> Result<OHLCCents> {
|
||||
if self._1d.as_ref().and_then(|m| m.last_key_value()).is_none_or(|(k, _)| k <= date) {
|
||||
self._1d.replace(Self::fetch_1d()?);
|
||||
self._1d.replace(self.fetch_1d()?);
|
||||
}
|
||||
self._1d
|
||||
.as_ref()
|
||||
@@ -56,11 +69,12 @@ impl Kraken {
|
||||
.ok_or(Error::NotFound("Couldn't find date".into()))
|
||||
}
|
||||
|
||||
pub fn fetch_1d() -> Result<BTreeMap<Date, OHLCCents>> {
|
||||
pub fn fetch_1d(&self) -> Result<BTreeMap<Date, OHLCCents>> {
|
||||
let agent = &self.agent;
|
||||
default_retry(|_| {
|
||||
let url = Self::url(1440);
|
||||
info!("Fetching {url} ...");
|
||||
let bytes = check_response(minreq::get(&url).with_timeout(30).send()?, &url)?;
|
||||
let bytes = checked_get(agent, &url)?;
|
||||
let json: Value = serde_json::from_slice(&bytes)?;
|
||||
Self::parse_date_ohlc_response(&json)
|
||||
})
|
||||
@@ -95,10 +109,8 @@ impl Kraken {
|
||||
format!("https://api.kraken.com/0/public/OHLC?pair=XBTUSD&interval={interval}")
|
||||
}
|
||||
|
||||
pub fn ping() -> Result<()> {
|
||||
minreq::get("https://api.kraken.com/0/public/Time")
|
||||
.with_timeout(10)
|
||||
.send()?;
|
||||
pub fn ping(&self) -> Result<()> {
|
||||
self.agent.get("https://api.kraken.com/0/public/Time").call()?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -125,7 +137,7 @@ impl PriceSource for Kraken {
|
||||
}
|
||||
|
||||
fn ping(&self) -> Result<()> {
|
||||
Self::ping()
|
||||
self.ping()
|
||||
}
|
||||
|
||||
fn clear(&mut self) {
|
||||
|
||||
@@ -1,10 +1,12 @@
|
||||
#![doc = include_str!("../README.md")]
|
||||
|
||||
use std::io::Read as _;
|
||||
use std::{path::Path, thread::sleep, time::Duration};
|
||||
|
||||
use brk_error::{Error, Result};
|
||||
use brk_types::{Date, Height, OHLCCents, Timestamp};
|
||||
use tracing::{info, warn};
|
||||
use ureq::Agent;
|
||||
|
||||
mod binance;
|
||||
mod brk;
|
||||
@@ -22,21 +24,32 @@ pub use source::{PriceSource, TrackedSource};
|
||||
|
||||
const MAX_RETRIES: usize = 12 * 60; // 12 hours of retrying
|
||||
|
||||
/// Check HTTP response status and return bytes or error
|
||||
pub fn check_response(response: minreq::Response, url: &str) -> Result<Vec<u8>> {
|
||||
let status = response.status_code as u16;
|
||||
if (200..300).contains(&status) {
|
||||
Ok(response.into_bytes())
|
||||
} else {
|
||||
Err(Error::HttpStatus {
|
||||
/// Create a shared HTTP agent with connection pooling and default timeout.
|
||||
pub fn new_agent(timeout_secs: u64) -> Agent {
|
||||
Agent::config_builder()
|
||||
.timeout_global(Some(Duration::from_secs(timeout_secs)))
|
||||
.build()
|
||||
.into()
|
||||
}
|
||||
|
||||
/// Perform a GET request and check the response status.
|
||||
pub fn checked_get(agent: &Agent, url: &str) -> Result<Vec<u8>> {
|
||||
let mut response = agent.get(url).call()?;
|
||||
let status = response.status().as_u16();
|
||||
if status >= 400 {
|
||||
return Err(Error::HttpStatus {
|
||||
status,
|
||||
url: url.to_string(),
|
||||
})
|
||||
});
|
||||
}
|
||||
let mut bytes = Vec::new();
|
||||
response.body_mut().as_reader().read_to_end(&mut bytes)?;
|
||||
Ok(bytes)
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Fetcher {
|
||||
pub agent: Agent,
|
||||
pub binance: TrackedSource<Binance>,
|
||||
pub kraken: TrackedSource<Kraken>,
|
||||
pub brk: TrackedSource<BRK>,
|
||||
@@ -48,10 +61,12 @@ impl Fetcher {
|
||||
}
|
||||
|
||||
pub fn new(hars_path: Option<&Path>) -> Result<Self> {
|
||||
let agent = new_agent(30);
|
||||
Ok(Self {
|
||||
binance: TrackedSource::new(Binance::init(hars_path)),
|
||||
kraken: TrackedSource::new(Kraken::default()),
|
||||
brk: TrackedSource::new(BRK::default()),
|
||||
binance: TrackedSource::new(Binance::new(hars_path, agent.clone())),
|
||||
kraken: TrackedSource::new(Kraken::new(agent.clone())),
|
||||
brk: TrackedSource::new(BRK::new(agent.clone())),
|
||||
agent,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user