//! HTTP client for a single `twice-server`, including the SSE job stream.
use anyhow::{Context as _, Result, bail};
use futures_util::StreamExt as _;
use twice_core::{
API_KEY_HEADER, API_PREFIX, ApiErrorBody, AppSummary, AppsResponse, Host, ImageRef,
InfrastructureSpec, JobId, JobOutcome, JobView, LogLine, UpdateImageRequest,
};
use crate::profiles::ServerProfile;
/// The two non-destructive lifecycle operations the server exposes.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum Lifecycle {
/// `POST /apps/{host}/start`
Start,
/// `POST /apps/{host}/stop`
Stop,
}
impl Lifecycle {
const fn path_segment(self) -> &'static str {
match self {
Self::Start => "start",
Self::Stop => "stop",
}
}
}
/// A configured connection to one server.
#[derive(Clone, Debug)]
pub(crate) struct ApiClient {
http: reqwest::Client,
base: String,
key: String,
}
impl ApiClient {
/// Builds a client for a saved profile.
///
/// The profile may have been hand-edited on disk, so the fields are
/// trimmed here as well as on save; untrimmed input otherwise reaches
/// `reqwest` as `builder error: invalid port number`.
pub(crate) fn new(profile: &ServerProfile) -> Result<Self> {
let http = reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(20))
.build()
.context("building HTTP client")?;
let base = format!(
"{}{API_PREFIX}",
profile.base_url.trim().trim_end_matches('/')
);
Ok(Self {
http,
base,
key: profile.api_key.trim().to_owned(),
})
}
/// `GET /apps`
pub(crate) async fn apps(&self) -> Result<Vec<AppSummary>> {
let response = self
.http
.get(format!("{}/apps", self.base))
.header(API_KEY_HEADER, &self.key)
.send()
.await
.context("requesting the application list")?;
Ok(decode::<AppsResponse>(response).await?.apps)
}
/// `POST /apps/{host}/{start|stop}`
pub(crate) async fn lifecycle(&self, host: &Host, op: Lifecycle) -> Result<JobView> {
let response = self
.http
.post(format!("{}/apps/{host}/{}", self.base, op.path_segment()))
.header(API_KEY_HEADER, &self.key)
.send()
.await
.with_context(|| format!("requesting {} for {host}", op.path_segment()))?;
decode(response).await
}
/// `POST /apps/{host}/update`
pub(crate) async fn update_image(
&self,
host: &Host,
request: &UpdateImageRequest,
) -> Result<JobView> {
let response = self
.http
.post(format!("{}/apps/{host}/update", self.base))
.header(API_KEY_HEADER, &self.key)
.json(request)
.send()
.await
.with_context(|| format!("requesting an image update for {host}"))?;
decode(response).await
}
/// `POST /infrastructure/apply`
pub(crate) async fn apply_infrastructure(&self, spec: &InfrastructureSpec) -> Result<JobView> {
let response = self
.http
.post(format!("{}/infrastructure/apply", self.base))
.header(API_KEY_HEADER, &self.key)
.json(spec)
.send()
.await
.with_context(|| format!("applying infrastructure for {}", spec.host))?;
decode(response).await
}
/// `GET /jobs/{id}/stream` — calls `emit` for each line and once at the end.
pub(crate) async fn stream_job(
&self,
id: JobId,
emit: &(impl Fn(StreamItem) + Sync),
) -> Result<()> {
let response = self
.http
.get(format!("{}/jobs/{id}/stream", self.base))
.header(API_KEY_HEADER, &self.key)
.send()
.await
.with_context(|| format!("streaming job {id}"))?;
if !response.status().is_success() {
bail!("job stream refused with HTTP {}", response.status());
}
let mut buffer = String::new();
let mut body = response.bytes_stream();
while let Some(chunk) = body.next().await {
let chunk = chunk.context("reading the job stream")?;
buffer.push_str(&String::from_utf8_lossy(&chunk));
while let Some(split) = buffer.find("\n\n") {
let frame: String = buffer.drain(..split).collect();
buffer.drain(..2.min(buffer.len()));
if let Some(item) = parse_frame(&frame) {
let finished = matches!(item, StreamItem::Finished(_));
emit(item);
if finished {
return Ok(());
}
}
}
}
Ok(())
}
}
/// One decoded server-sent event.
#[derive(Clone, Debug)]
pub(crate) enum StreamItem {
/// A line of `once` output.
Line(LogLine),
/// The job's terminal outcome; nothing follows it.
Finished(JobOutcome),
}
/// Parses one `event:`/`data:` SSE frame, ignoring comments and keep-alives.
fn parse_frame(frame: &str) -> Option<StreamItem> {
let mut name = None;
let mut data = String::new();
for line in frame.lines() {
if let Some(value) = line.strip_prefix("event:") {
name = Some(value.trim().to_owned());
} else if let Some(value) = line.strip_prefix("data:") {
data.push_str(value.trim());
}
}
match name.as_deref() {
Some("line") => serde_json::from_str::<LogLine>(&data)
.ok()
.map(StreamItem::Line),
Some("finished") => serde_json::from_str::<JobOutcome>(&data)
.ok()
.map(StreamItem::Finished),
Some(_) | None => None,
}
}
/// Turns a response into `T`, or into the server's typed error message.
async fn decode<T: serde::de::DeserializeOwned>(response: reqwest::Response) -> Result<T> {
let status = response.status();
let body = response.text().await.context("reading the response body")?;
if status.is_success() {
return serde_json::from_str(&body)
.with_context(|| format!("decoding a {status} response"));
}
match serde_json::from_str::<ApiErrorBody>(&body) {
Ok(api_error) => bail!("{} ({})", api_error.message, api_error.error),
Err(_ignored) => bail!("HTTP {status}: {}", body.trim()),
}
}
/// Builds the update request from raw UI text, surfacing parse errors early.
pub(crate) fn update_request(image: &str, auto_update: Option<bool>) -> Result<UpdateImageRequest> {
let image = ImageRef::parse(image.trim()).context("invalid image reference")?;
Ok(UpdateImageRequest { image, auto_update })
}
/// Parses the declarative TOML document shown in the infrastructure editor.
pub(crate) fn infrastructure_spec(document: &str) -> Result<InfrastructureSpec> {
toml::from_str(document).context("invalid infrastructure document")
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_line_frame_decodes_into_a_log_line() {
let frame = "event: line\ndata: {\"seq\":1,\"stream\":\"stdout\",\"text\":\"hi\"}";
match parse_frame(frame) {
Some(StreamItem::Line(line)) => assert_eq!(line.text, "hi"),
other => panic!("expected a line, got {other:?}"),
}
}
#[test]
fn a_finished_frame_decodes_into_an_outcome() {
let frame = "event: finished\ndata: {\"status\":\"succeeded\"}";
match parse_frame(frame) {
Some(StreamItem::Finished(outcome)) => assert_eq!(outcome, JobOutcome::Succeeded),
other => panic!("expected an outcome, got {other:?}"),
}
}
#[test]
fn keep_alive_comments_are_ignored() {
assert!(parse_frame(":\n").is_none());
}
#[test]
fn update_request_rejects_a_flag_shaped_image() {
assert!(update_request("--env=EVIL", None).is_err());
assert!(update_request(" ghcr.io/app:1.2 ", Some(false)).is_ok());
}
#[test]
fn infrastructure_document_declares_source_and_version_separately() {
let document = concat!(
"host = \"writebook.example.com\"\n",
"source = \"ghcr.io/basecamp/writebook\"\n",
"version = \"1.4.0\"\n",
);
let spec = infrastructure_spec(document).unwrap();
assert_eq!(spec.host.as_str(), "writebook.example.com");
assert_eq!(spec.source.as_str(), "ghcr.io/basecamp/writebook");
assert_eq!(spec.version.as_str(), "1.4.0");
}
#[test]
fn infrastructure_document_rejects_a_version_embedded_in_the_source() {
let document = concat!(
"host = \"writebook.example.com\"\n",
"source = \"ghcr.io/basecamp/writebook:latest\"\n",
"version = \"1.4.0\"\n",
);
assert!(infrastructure_spec(document).is_err());
}
#[test]
fn a_hand_edited_profile_with_stray_whitespace_still_builds_a_usable_base_url() {
// Regression: this exact profile shape was written to disk by the UI
// and produced "builder error: invalid port number" on every request.
let profile = ServerProfile {
name: "local ".to_owned(),
base_url: "http://127.0.0.1:7373 ".to_owned(),
api_key: " secret-key ".to_owned(),
};
let client = ApiClient::new(&profile).unwrap();
assert_eq!(client.base, "http://127.0.0.1:7373/api/v1");
assert_eq!(client.key, "secret-key");
}
#[test]
fn a_trailing_slash_does_not_double_up_the_api_prefix() {
let profile = ServerProfile {
name: "local".to_owned(),
base_url: "http://127.0.0.1:7373/".to_owned(),
api_key: "secret".to_owned(),
};
assert_eq!(
ApiClient::new(&profile).unwrap().base,
"http://127.0.0.1:7373/api/v1"
);
}
}