Add server client
This commit is contained in:
15
crates/pile-client/Cargo.toml
Normal file
15
crates/pile-client/Cargo.toml
Normal file
@@ -0,0 +1,15 @@
|
||||
[package]
|
||||
name = "pile-client"
|
||||
version = { workspace = true }
|
||||
rust-version = { workspace = true }
|
||||
edition = { workspace = true }
|
||||
|
||||
[lints]
|
||||
workspace = true
|
||||
|
||||
[dependencies]
|
||||
reqwest = { version = "0.12", features = ["json", "stream"] }
|
||||
futures-core = "0.3"
|
||||
serde = { workspace = true }
|
||||
thiserror = { workspace = true }
|
||||
bytes = { workspace = true }
|
||||
230
crates/pile-client/src/lib.rs
Normal file
230
crates/pile-client/src/lib.rs
Normal file
@@ -0,0 +1,230 @@
|
||||
use bytes::Bytes;
|
||||
use futures_core::Stream;
|
||||
use reqwest::{Client, StatusCode, header};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::pin::Pin;
|
||||
use thiserror::Error;
|
||||
|
||||
//
|
||||
// MARK: Error
|
||||
//
|
||||
|
||||
#[derive(Debug, Error)]
|
||||
pub enum ClientError {
|
||||
#[error("invalid bearer token")]
|
||||
InvalidToken,
|
||||
|
||||
#[error("HTTP {status}: {body}")]
|
||||
Http { status: StatusCode, body: String },
|
||||
|
||||
#[error(transparent)]
|
||||
Reqwest(#[from] reqwest::Error),
|
||||
}
|
||||
|
||||
//
|
||||
// MARK: Response types
|
||||
//
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct DatasetInfo {
|
||||
pub name: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
pub struct LookupRequest {
|
||||
pub query: String,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub limit: Option<usize>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct LookupResult {
|
||||
pub score: f32,
|
||||
pub source: String,
|
||||
pub key: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct LookupResponse {
|
||||
pub results: Vec<LookupResult>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct ItemRef {
|
||||
pub source: String,
|
||||
pub key: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct ItemsResponse {
|
||||
pub items: Vec<ItemRef>,
|
||||
pub total: usize,
|
||||
pub offset: usize,
|
||||
pub limit: usize,
|
||||
}
|
||||
|
||||
/// Raw field response: the content-type and body bytes as returned by the server.
|
||||
pub struct FieldResponse {
|
||||
pub content_type: String,
|
||||
pub data: Bytes,
|
||||
}
|
||||
|
||||
//
|
||||
// MARK: PileClient
|
||||
//
|
||||
|
||||
/// A client for a pile server. Use [`PileClient::dataset`] to get a dataset-scoped client.
|
||||
pub struct PileClient {
|
||||
base_url: String,
|
||||
client: Client,
|
||||
}
|
||||
|
||||
impl PileClient {
|
||||
pub fn new(base_url: impl Into<String>, token: Option<&str>) -> Result<Self, ClientError> {
|
||||
let mut headers = header::HeaderMap::new();
|
||||
|
||||
if let Some(token) = token {
|
||||
let value = header::HeaderValue::from_str(&format!("Bearer {token}"))
|
||||
.map_err(|_| ClientError::InvalidToken)?;
|
||||
headers.insert(header::AUTHORIZATION, value);
|
||||
}
|
||||
|
||||
let client = Client::builder()
|
||||
.default_headers(headers)
|
||||
.build()
|
||||
.map_err(ClientError::Reqwest)?;
|
||||
|
||||
Ok(Self {
|
||||
base_url: base_url.into(),
|
||||
client,
|
||||
})
|
||||
}
|
||||
|
||||
/// Returns a client scoped to a specific dataset (i.e. `/{name}/...`).
|
||||
pub fn dataset(&self, name: &str) -> DatasetClient {
|
||||
DatasetClient {
|
||||
base_url: format!("{}/{name}", self.base_url),
|
||||
client: self.client.clone(),
|
||||
}
|
||||
}
|
||||
|
||||
/// `GET /datasets` — list all datasets served by this server.
|
||||
pub async fn list_datasets(&self) -> Result<Vec<DatasetInfo>, ClientError> {
|
||||
let resp = self
|
||||
.client
|
||||
.get(format!("{}/datasets", self.base_url))
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
check_status(resp).await?.json().await.map_err(Into::into)
|
||||
}
|
||||
}
|
||||
|
||||
//
|
||||
// MARK: DatasetClient
|
||||
//
|
||||
|
||||
/// A client scoped to a single dataset on the server.
|
||||
pub struct DatasetClient {
|
||||
base_url: String,
|
||||
client: Client,
|
||||
}
|
||||
|
||||
impl DatasetClient {
|
||||
/// `POST /lookup` — full-text search within this dataset.
|
||||
pub async fn lookup(
|
||||
&self,
|
||||
query: impl Into<String>,
|
||||
limit: Option<usize>,
|
||||
) -> Result<LookupResponse, ClientError> {
|
||||
let body = LookupRequest {
|
||||
query: query.into(),
|
||||
limit,
|
||||
};
|
||||
|
||||
let resp = self
|
||||
.client
|
||||
.post(format!("{}/lookup", self.base_url))
|
||||
.json(&body)
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
check_status(resp).await?.json().await.map_err(Into::into)
|
||||
}
|
||||
|
||||
/// `GET /item` — stream the raw bytes of an item.
|
||||
///
|
||||
/// The returned stream yields chunks as they arrive from the server.
|
||||
pub async fn get_item(
|
||||
&self,
|
||||
source: &str,
|
||||
key: &str,
|
||||
) -> Result<Pin<Box<dyn Stream<Item = Result<Bytes, reqwest::Error>> + Send>>, ClientError> {
|
||||
let resp = self
|
||||
.client
|
||||
.get(format!("{}/item", self.base_url))
|
||||
.query(&[("source", source), ("key", key)])
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
Ok(Box::pin(check_status(resp).await?.bytes_stream()))
|
||||
}
|
||||
|
||||
/// `GET /field` — extract a field from an item by object path (e.g. `$.flac.title`).
|
||||
pub async fn get_field(
|
||||
&self,
|
||||
source: &str,
|
||||
key: &str,
|
||||
path: &str,
|
||||
) -> Result<FieldResponse, ClientError> {
|
||||
let resp = self
|
||||
.client
|
||||
.get(format!("{}/field", self.base_url))
|
||||
.query(&[("source", source), ("key", key), ("path", path)])
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
let resp = check_status(resp).await?;
|
||||
|
||||
let content_type = resp
|
||||
.headers()
|
||||
.get(header::CONTENT_TYPE)
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.unwrap_or("application/octet-stream")
|
||||
.to_owned();
|
||||
|
||||
let data = resp.bytes().await?;
|
||||
|
||||
Ok(FieldResponse { content_type, data })
|
||||
}
|
||||
|
||||
/// `GET /items` — paginate over all items in this dataset, ordered by (source, key).
|
||||
pub async fn list_items(
|
||||
&self,
|
||||
offset: usize,
|
||||
limit: usize,
|
||||
) -> Result<ItemsResponse, ClientError> {
|
||||
let resp = self
|
||||
.client
|
||||
.get(format!("{}/items", self.base_url))
|
||||
.query(&[("offset", offset), ("limit", limit)])
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
check_status(resp).await?.json().await.map_err(Into::into)
|
||||
}
|
||||
}
|
||||
|
||||
//
|
||||
// MARK: helpers
|
||||
//
|
||||
|
||||
async fn check_status(resp: reqwest::Response) -> Result<reqwest::Response, ClientError> {
|
||||
let status = resp.status();
|
||||
if status.is_success() {
|
||||
return Ok(resp);
|
||||
}
|
||||
|
||||
let body = resp.text().await.unwrap_or_default();
|
||||
Err(ClientError::Http { status, body })
|
||||
}
|
||||
Reference in New Issue
Block a user