pub mod breakdown_trading_data;
pub mod cash_dividend_data;
pub mod daily_stock_prices;
pub mod earnings_calendar;
pub mod financial_statement_details;
pub mod financial_statements;
pub mod futures_prices;
pub mod index_option_prices;
pub mod indicies;
pub mod listed_issue_info;
pub mod morning_session_stock_prices;
pub mod options_prices;
pub mod shared;
pub mod short_sale_by_sector;
pub mod topic_prices;
pub mod trading_by_type_of_investors;
pub mod trading_calendar;
pub mod weekly_margin_trading_outstandings;
use shared::{
auth::{get_id_token_from_api, get_refresh_token_from_api},
responses::error_response::JQuantsErrorResponse,
};
use std::{fmt, sync::Arc};
use tokio::sync::RwLock;
use crate::error::JQuantsError;
use chrono::{DateTime, Local};
use reqwest::{Client, RequestBuilder};
use serde::{de::DeserializeOwned, Serialize};
const BASE_URL: &str = "https://api.jquants.com/v1";
fn build_url(path: &str) -> String {
format!("{}/{}", BASE_URL, path)
}
pub trait JQuantsPlanClient: Clone {
fn new(api_client: JQuantsApiClient) -> Self;
fn new_from_refresh_token(refresh_token: String) -> Self {
let api_client = JQuantsApiClient::new_from_refresh_token(refresh_token);
Self::new(api_client)
}
fn new_from_account(
mailaddress: &str,
password: &str,
) -> impl std::future::Future<Output = Result<Self, JQuantsError>> + Send {
async {
let api_client = JQuantsApiClient::new_from_account(mailaddress, password).await?;
Ok(Self::new(api_client))
}
}
fn get_api_client(&self) -> &JQuantsApiClient;
fn get_current_refresh_token(&self) -> impl std::future::Future<Output = String> + Send {
let api_client = self.get_api_client().clone();
async move {
api_client
.inner
.token_set
.read()
.await
.refresh_token
.clone()
}
}
fn get_refresh_token_from_api(
&self,
mail_address: &str,
password: &str,
) -> impl std::future::Future<Output = Result<String, JQuantsError>> + Send {
let api_client = self.get_api_client().clone();
async move { get_refresh_token_from_api(&api_client.inner.client, mail_address, password).await }
}
fn get_id_token_from_api(
&self,
refresh_token: &str,
) -> impl std::future::Future<Output = Result<String, JQuantsError>> + Send {
let api_client = self.get_api_client().clone();
async move { get_id_token_from_api(&api_client.inner.client, refresh_token).await }
}
fn reset_refresh_token(
&self,
mail_address: &str,
password: &str,
) -> impl std::future::Future<Output = Result<(), JQuantsError>> + Send {
let api_client = self.get_api_client().clone();
async move {
api_client
.inner
.reset_refresh_token(mail_address, password)
.await
}
}
fn reset_id_token(&self) -> impl std::future::Future<Output = Result<(), JQuantsError>> + Send {
let api_client = self.get_api_client().clone();
async move { api_client.inner.reset_id_token().await }
}
fn reauthenticate(
&self,
mail_address: &str,
password: &str,
) -> impl std::future::Future<Output = Result<(), JQuantsError>> + Send {
let api_client = self.get_api_client().clone();
async move { api_client.inner.reset_tokens(mail_address, password).await }
}
}
#[derive(Clone)]
pub struct JQuantsApiClient {
inner: Arc<JQuantsApiClientRef>,
}
impl JQuantsApiClient {
fn new_from_refresh_token(refresh_token: String) -> Self {
Self {
inner: Arc::new(JQuantsApiClientRef::new_from_refresh_token(refresh_token)),
}
}
async fn new_from_account(mailaddress: &str, password: &str) -> Result<Self, JQuantsError> {
let client_ref = JQuantsApiClientRef::new_from_account(mailaddress, password).await?;
Ok(Self {
inner: Arc::new(client_ref),
})
}
}
pub(crate) struct JQuantsApiClientRef {
client: Client,
token_set: Arc<RwLock<TokenSet>>,
}
impl JQuantsApiClientRef {
fn new_from_refresh_token(refresh_token: String) -> Self {
Self {
client: Client::new(),
token_set: Arc::new(RwLock::new(TokenSet {
refresh_token,
id_token: None,
})),
}
}
async fn new_from_account(mailaddress: &str, password: &str) -> Result<Self, JQuantsError> {
let client = Client::new();
let refresh_token = get_refresh_token_from_api(&client, mailaddress, password).await?;
let new_id_token = get_id_token_from_api(&client, &refresh_token).await?;
let id_token_wrapper = IdTokenWrapper::new(new_id_token);
Ok(Self {
client,
token_set: Arc::new(RwLock::new(TokenSet {
refresh_token,
id_token: Some(id_token_wrapper),
})),
})
}
async fn reset_refresh_token(
&self,
mail_address: &str,
password: &str,
) -> Result<(), JQuantsError> {
tracing::debug!("Starting reset a refresh token process.");
match get_refresh_token_from_api(&self.client, mail_address, password).await {
Ok(new_refresh_token) => {
let mut token_set_write = self.token_set.write().await;
token_set_write.refresh_token = new_refresh_token;
tracing::debug!("Refresh token refreshed successfully.");
Ok(())
}
Err(e) => {
tracing::error!("Failed to refresh a refresh token: {:?}", e);
Err(e)
}
}
}
async fn reset_id_token(&self) -> Result<(), JQuantsError> {
tracing::debug!("Starting reset a refresh id process.");
let refresh_token = { self.token_set.read().await.refresh_token.clone() };
match get_id_token_from_api(&self.client, &refresh_token).await {
Ok(new_id_token) => {
let mut token_set_write = self.token_set.write().await;
token_set_write.id_token = Some(IdTokenWrapper::new(new_id_token));
tracing::debug!("ID token refreshed successfully.");
Ok(())
}
Err(e) => {
tracing::error!("Failed to refresh ID token: {:?}", e);
Err(e)
}
}
}
async fn reset_id_token_if_needed(&self) -> Result<(), JQuantsError> {
let needs_refresh = {
let token_set = self.token_set.read().await;
match &token_set.id_token {
Some(token) => !token.is_valid(),
None => true,
}
};
if needs_refresh {
tracing::debug!("ID token is invalid or expired. Attempting to refresh.");
self.reset_id_token().await
} else {
tracing::debug!("ID token is still valid.");
Ok(())
}
}
async fn reset_tokens(&self, mail_address: &str, password: &str) -> Result<(), JQuantsError> {
tracing::debug!("Starting re-authentication process.");
let new_refresh_token = get_refresh_token_from_api(&self.client, mail_address, password)
.await
.map_err(|e| {
tracing::error!("Failed to obtain new refresh token: {:?}", e);
e
})?;
tracing::debug!("Successfully obtained new refresh token.");
let new_id_token = get_id_token_from_api(&self.client, &new_refresh_token)
.await
.map_err(|e| {
tracing::error!("Failed to obtain new ID token: {:?}", e);
e
})?;
tracing::debug!("Successfully obtained new ID token.");
let expires_at = Local::now() + chrono::Duration::hours(24);
let new_id_token_wrapper = Some(IdTokenWrapper {
id_token: new_id_token,
expires_at,
});
{
let mut token_set_write = self.token_set.write().await;
token_set_write.refresh_token = new_refresh_token;
token_set_write.id_token = new_id_token_wrapper;
}
tracing::debug!("Re-authentication process process completed successfully.");
Ok(())
}
async fn get<T: DeserializeOwned + fmt::Debug>(
&self,
path: &str,
params: impl Serialize,
) -> Result<T, JQuantsError> {
let url = format!("{BASE_URL}/{}", path);
let request = self.client.get(&url).query(¶ms);
self.common_send_and_refresh_token_if_needed::<T>(request)
.await
}
async fn common_send_and_refresh_token_if_needed<T: DeserializeOwned + fmt::Debug>(
&self,
request: RequestBuilder,
) -> Result<T, JQuantsError> {
self.reset_id_token_if_needed().await?;
self.common_send(request).await
}
async fn common_send<T: DeserializeOwned + fmt::Debug>(
&self,
request: RequestBuilder,
) -> Result<T, JQuantsError> {
let id_token = {
self.token_set
.read()
.await
.id_token
.as_ref()
.ok_or_else(|| {
tracing::error!("ID token not found.");
JQuantsError::BugError("ID token not found.".to_string())
})?
.id_token
.clone()
};
let request = request.header("Authorization", &format!("Bearer {id_token}"));
if let Some(url) = request
.try_clone()
.and_then(|req| req.build().ok().map(|r| r.url().clone()))
{
tracing::debug!("Sending API request to URL: {url}");
} else {
tracing::debug!("Sending API request.");
}
let response = request.send().await?;
let status = response.status();
let text = response.text().await.unwrap_or_default();
tracing::debug!("Received response with status: {}", status);
if status.is_success() {
match serde_json::from_str::<T>(&text) {
Ok(data) => {
tracing::debug!("Successfully parsed response.");
Ok(data)
}
Err(_) => {
tracing::error!("Failed to parse response");
Err(JQuantsError::InvalidResponseFormat {
status_code: status.as_u16(),
body: text,
})
}
}
} else {
match serde_json::from_str::<JQuantsErrorResponse>(&text) {
Ok(error_response) => match status {
reqwest::StatusCode::UNAUTHORIZED => {
tracing::warn!(
"Received UNAUTHORIZED error. Status code: {}",
status.as_u16()
);
Err(JQuantsError::IdTokenInvalidOrExpired {
body: error_response,
status_code: status.as_u16(),
})
}
_ => {
tracing::error!("API error occurred. Status code: {}", status.as_u16());
Err(JQuantsError::ApiError {
body: error_response,
status_code: status.as_u16(),
})
}
},
Err(_) => {
tracing::error!("Invalid response format. Status code: {}", status.as_u16());
Err(JQuantsError::InvalidResponseFormat {
status_code: status.as_u16(),
body: text,
})
}
}
}
}
}
pub(crate) struct TokenSet {
refresh_token: String,
id_token: Option<IdTokenWrapper>,
}
pub(crate) struct IdTokenWrapper {
id_token: String,
expires_at: DateTime<Local>,
}
impl IdTokenWrapper {
fn new(id_token: String) -> Self {
let expires_at = Local::now() + chrono::Duration::hours(24);
IdTokenWrapper {
id_token,
expires_at,
}
}
fn is_valid(&self) -> bool {
Local::now() < self.expires_at
}
}
impl fmt::Debug for IdTokenWrapper {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let len = self.id_token.len();
let masking_id_token = "*".repeat(len);
f.debug_struct("IdTokenWrapper")
.field("id_token", &masking_id_token)
.field("expires_at", &self.expires_at)
.finish()
}
}