pub(crate) mod table;
pub(crate) mod wire;
use std::{io::Read, sync::Arc};
use arrow_array::RecordBatch;
use arrow_schema::SchemaRef;
use serde_json::{Value, json};
use crate::{IndexSpec, InfinoError, Supertable};
const API_KEY_ENV: &str = "INFINO_API_KEY";
pub(crate) struct RemoteCatalog {
agent: ureq::Agent,
base_url: String,
database: String,
api_key: String,
}
impl RemoteCatalog {
pub(crate) fn new(
base_url: String,
database: String,
api_key: Option<String>,
) -> Result<Self, InfinoError> {
let api_key = api_key
.or_else(|| std::env::var(API_KEY_ENV).ok())
.filter(|k| !k.is_empty())
.ok_or_else(|| {
InfinoError::Config(format!(
"a hosted connection needs an API key (pass ConnectOptions::with_api_key or set {API_KEY_ENV})"
))
})?;
Ok(Self {
agent: ureq::agent(),
base_url,
database,
api_key,
})
}
fn url(&self, op: &str) -> String {
format!("{}/v1/{op}/{}", self.base_url, self.database)
}
fn bearer(&self) -> String {
format!("Bearer {}", self.api_key)
}
pub(crate) fn post_json(&self, op: &str, body: Value) -> Result<ureq::Response, InfinoError> {
let request = self
.agent
.post(&self.url(op))
.set("Authorization", &self.bearer())
.set("Accept", wire::ARROW_STREAM_CONTENT_TYPE);
map_send(op, request.send_json(body))
}
pub(crate) fn post_arrow(
&self,
op: &str,
query: &[(&str, &str)],
body: Vec<u8>,
) -> Result<ureq::Response, InfinoError> {
let mut request = self
.agent
.post(&self.url(op))
.set("Authorization", &self.bearer())
.set("Content-Type", wire::ARROW_STREAM_CONTENT_TYPE);
for (key, value) in query {
request = request.query(key, value);
}
map_send(op, request.send_bytes(&body))
}
pub(crate) fn create_database(&self) -> Result<(), InfinoError> {
let url = format!("{}/v1/databases", self.base_url);
let request = self.agent.post(&url).set("Authorization", &self.bearer());
map_send(
"create_database",
request.send_json(json!({ "name": self.database })),
)?;
Ok(())
}
pub(crate) fn create_table(
self: &Arc<Self>,
name: &str,
schema: SchemaRef,
indexes: IndexSpec,
) -> Result<Supertable, InfinoError> {
let body = json!({
"table_name": name,
"schema": wire::schema_to_json(&schema)?,
"indexes": wire::index_spec_to_json(&indexes),
});
self.post_json("create_table", body)?;
Ok(Supertable::from_table(Arc::new(table::RemoteTable::new(
Arc::clone(self),
name.to_string(),
schema,
))))
}
pub(crate) fn open_table(self: &Arc<Self>, name: &str) -> Result<Supertable, InfinoError> {
let response = self.post_json("schema", json!({ "table_name": name }))?;
let value = read_json("schema", response)?;
let fields = value.as_array().ok_or_else(|| {
InfinoError::Backend("schema response was not a JSON array".to_string())
})?;
let schema = Arc::new(wire::json_to_schema(fields)?);
Ok(Supertable::from_table(Arc::new(table::RemoteTable::new(
Arc::clone(self),
name.to_string(),
schema,
))))
}
pub(crate) fn list_tables(&self) -> Result<Vec<String>, InfinoError> {
let response = self.post_json("list_tables", json!({}))?;
let value = read_json("list_tables", response)?;
let names = value.as_array().ok_or_else(|| {
InfinoError::Backend("list_tables response was not a JSON array".to_string())
})?;
names
.iter()
.map(|v| {
v.as_str().map(str::to_owned).ok_or_else(|| {
InfinoError::Backend("list_tables entry was not a string".to_string())
})
})
.collect()
}
pub(crate) fn drop_table(&self, name: &str, purge: bool) -> Result<(), InfinoError> {
self.post_json("drop_table", json!({ "table_name": name, "purge": purge }))?;
Ok(())
}
pub(crate) fn query_sql(&self, sql: &str) -> Result<Vec<RecordBatch>, InfinoError> {
let response = self.post_json("query_sql", json!({ "query": sql }))?;
read_arrow("query_sql", response)
}
}
fn map_send(
op: &str,
result: Result<ureq::Response, ureq::Error>,
) -> Result<ureq::Response, InfinoError> {
match result {
Ok(response) => Ok(response),
Err(ureq::Error::Status(code, response)) => {
let body = response.into_string().unwrap_or_default();
Err(wire::status_to_error(op, code, &body))
}
Err(ureq::Error::Transport(transport)) => Err(InfinoError::Backend(format!(
"{op}: transport error: {transport}"
))),
}
}
pub(crate) fn read_json(op: &str, response: ureq::Response) -> Result<Value, InfinoError> {
response
.into_json::<Value>()
.map_err(|e| InfinoError::Backend(format!("{op}: parsing response: {e}")))
}
pub(crate) fn read_arrow(
op: &str,
response: ureq::Response,
) -> Result<Vec<RecordBatch>, InfinoError> {
let mut buf = Vec::new();
response
.into_reader()
.read_to_end(&mut buf)
.map_err(|e| InfinoError::Backend(format!("{op}: reading response: {e}")))?;
wire::ipc_to_batches(&buf)
}