anvilsign in

collin/anvil

1//! The dial-out half: claim a job, fetch its checkout, report what happened.
2//!
3//! Every request originates here. anvil never connects to a runner, which is
4//! the whole reason this is pull-based — the runner sits on a home network
5//! behind NAT, and the forge sits on a public VPS. An inbound path from the
6//! latter to the former would make the forge's compromise the build host's.
7
8use std::time::Duration;
9
10use anvil_job::{
11 ArtifactSpec,
12 JobResult,
13 JobSpec,
14 RunnerInfo,
15 Stored,
16};
17
18use crate::executor::ArtifactSink;
19
20/// Generous relative to the server's claim poll (55s): a claim that finds work
21/// returns at once, and one that does not still has to come back under this.
22const REQUEST_TIMEOUT: Duration = Duration::from_secs(90);
23
24pub struct Client {
25 http: reqwest::Client,
26 base: String,
27 token: String,
28 info: RunnerInfo,
29}
30
31impl Client {
32 pub fn new(base: &str, token: &str, info: RunnerInfo) -> Result<Self, String> {
33 let http = reqwest::Client::builder()
34 .timeout(REQUEST_TIMEOUT)
35 .build()
36 .map_err(|e| format!("building http client: {e}"))?;
37 Ok(Self {
38 http,
39 base: base.trim_end_matches('/').to_string(),
40 token: token.to_string(),
41 info,
42 })
43 }
44
45 fn url(&self, path: &str) -> String {
46 format!("{}{path}", self.base)
47 }
48
49 /// Token plus identity on every request. The name goes in a header rather
50 /// than only in the claim body so the endpoints with no body of their own
51 /// (the checkout GET, artifact uploads) can be authorized the same way.
52 fn auth(&self, req: reqwest::RequestBuilder) -> reqwest::RequestBuilder {
53 req.header("X-Anvil-Runner-Token", &self.token)
54 .header("X-Anvil-Runner-Name", &self.info.name)
55 }
56
57 /// Long-poll for a job. `Ok(None)` means the poll expired with nothing
58 /// queued, which is the common case and not an error.
59 pub async fn claim(&self) -> Result<Option<JobSpec>, String> {
60 let resp = self
61 .auth(self.http.post(self.url("/-/runner/claim")))
62 .json(&self.info)
63 .send()
64 .await
65 .map_err(|e| format!("claim: {e}"))?;
66 match resp.status() {
67 reqwest::StatusCode::NO_CONTENT => Ok(None),
68 s if s.is_success() => resp
69 .json::<JobSpec>()
70 .await
71 .map(Some)
72 .map_err(|e| format!("decoding job: {e}")),
73 s => Err(format!("claim: {s} {}", body_hint(resp).await)),
74 }
75 }
76
77 pub async fn checkout(&self, run_id: i64) -> Result<Vec<u8>, String> {
78 let resp = self
79 .auth(
80 self.http
81 .get(self.url(&format!("/-/runner/jobs/{run_id}/checkout.tar"))),
82 )
83 .send()
84 .await
85 .map_err(|e| format!("checkout: {e}"))?;
86 if !resp.status().is_success() {
87 return Err(format!(
88 "checkout: {} {}",
89 resp.status(),
90 body_hint(resp).await
91 ));
92 }
93 resp.bytes()
94 .await
95 .map(|b| b.to_vec())
96 .map_err(|e| format!("reading checkout: {e}"))
97 }
98
99 /// Extend the lease. `false` means anvil no longer thinks we hold this run
100 /// — it was requeued out from under us, and the job should be abandoned
101 /// rather than reported.
102 pub async fn heartbeat(&self, run_id: i64) -> bool {
103 let sent = self
104 .auth(
105 self.http
106 .post(self.url(&format!("/-/runner/jobs/{run_id}/heartbeat")))
107 .json(&self.info),
108 )
109 .send()
110 .await;
111 matches!(sent, Ok(r) if r.status().is_success())
112 }
113
114 pub async fn report(&self, run_id: i64, result: &JobResult) -> Result<(), String> {
115 let resp = self
116 .auth(
117 self.http
118 .post(self.url(&format!("/-/runner/jobs/{run_id}/result")))
119 .json(result),
120 )
121 .send()
122 .await
123 .map_err(|e| format!("result: {e}"))?;
124 if resp.status().is_success() {
125 Ok(())
126 } else {
127 Err(format!(
128 "result: {} {}",
129 resp.status(),
130 body_hint(resp).await
131 ))
132 }
133 }
134}
135
136/// Uploads each artifact tar to anvil as it comes out of the container.
137///
138/// The server does the storing, and reports back what it wrote — the runner
139/// never learns anvil's on-disk layout, and `browse` stays the server's call.
140pub struct UploadSink<'a> {
141 pub client: &'a Client,
142 pub run_id: i64,
143}
144
145#[async_trait::async_trait]
146impl ArtifactSink for UploadSink<'_> {
147 async fn put(&mut self, spec: &ArtifactSpec, tar: &[u8]) -> Result<Stored, String> {
148 let url = self.client.url(&format!(
149 "/-/runner/jobs/{}/artifacts/{}",
150 self.run_id, spec.name
151 ));
152 let resp = self
153 .client
154 .auth(
155 self.client
156 .http
157 .post(url)
158 .header("Content-Type", "application/octet-stream")
159 .body(tar.to_vec()),
160 )
161 .send()
162 .await
163 .map_err(|e| format!("upload: {e}"))?;
164 if !resp.status().is_success() {
165 return Err(format!(
166 "upload: {} {}",
167 resp.status(),
168 body_hint(resp).await
169 ));
170 }
171 resp.json::<Stored>()
172 .await
173 .map_err(|e| format!("decoding upload response: {e}"))
174 }
175}
176
177/// A short slice of an error response, for a log line. Bounded because the
178/// body could be an HTML error page.
179async fn body_hint(resp: reqwest::Response) -> String {
180 let text = resp.text().await.unwrap_or_default();
181 let trimmed = text.trim();
182 if trimmed.len() > 200 {
183 format!("{}…", &trimmed[..200])
184 } else {
185 trimmed.to_string()
186 }
187}