Skip to content
1,127 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R21//! The services GitHub's Actions toolkit calls from inside a job, so that
2//! actions built on `@actions/cache` and `@actions/artifact` (such as
3//! `actions/setup-node` with `cache:`, or `Swatinem/rust-cache`) work on
4//! g1t unchanged. A job is told where they are in its variables
5//! (`ACTIONS_RUNTIME_TOKEN`, `ACTIONS_RESULTS_URL`, `ACTIONS_CACHE_URL`,
6//! `ACTIONS_CACHE_SERVICE_V2`; see `runtime_variables`).
7//!
8//! - Twirp, at `/twirp/github.actions.results.api.v1.CacheService/…` and
9//! `…ArtifactService/…`: the cache's newer protocol (`CreateCacheEntry`,
10//! `FinalizeCacheEntryUpload`, `GetCacheEntryDownloadURL`) and the
11//! artifacts' (`CreateArtifact`, `FinalizeArtifact`, `ListArtifacts`,
12//! `GetSignedArtifactURL`, `DeleteArtifact`). JSON, the toolkit's field
Merge remote-tracking branch 'origin/main' into workspace-chat13//! names; the cache's methods also in protobuf (`application/protobuf`),
14//! which other clients of the protocol send (sccache, through OpenDAL).
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R215//! - The cache's older protocol, at `{ACTIONS_CACHE_URL}_apis/artifactcache/…`,
16//! which the toolkit's client uses whenever the server it runs against is
Merge remote-tracking branch 'origin/main' into workspace-chat17//! not github.com: on g1t, that is the one it uses, and sccache's too.
18//! Its entries are sent in 32 MB chunks, or in one chunk of any size.
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R219//! - Blobs, at `/actions/toolkit/blobs/{token}`: the signed links those
Merge remote-tracking branch 'origin/main' into workspace-chat20//! hand out. Downloads are a GET, of the whole blob or of one byte range
21//! (`Range`), as the toolkit's client fetches large entries in segments.
22//! Uploads speak the part of Azure Blob Storage's protocol the toolkit's
23//! client uses (Put Blob, Put Block, Put Block List), mapped onto an R2
24//! multipart upload: a block's id ends in its index, which is its part's
25//! number. An upload link carries a query, as an Azure SAS link does,
26//! which clients that sign their requests with it need.
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R227//!
28//! Every call carries the job's runtime token; the actions service checks
29//! it and keeps the entries (cache.rs, artifacts.rs, runtime.rs there).
30
31use base64::Engine;
32use base64::engine::general_purpose::{STANDARD, URL_SAFE_NO_PAD};
33use g1t_contracts::actions::{
34 ARTIFACT_MAX_BYTES, Artifact, ArtifactBlob, ArtifactCommitArgs, ArtifactReservation, ArtifactReserveArgs, BlobArgs, BlobGrant, BlobPart,
35 BlobSignArgs, CACHE_MAX_ENTRY_BYTES, CACHE_PART_BYTES, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation,
36 CacheReserveArgs, CacheUploadArgs, JobArtifactsArgs,
37};
38use g1t_contracts::{Failure, FailureCode, Outcome};
39use serde_json::{Map, Value, json};
40use worker::{Bucket, Env, Request, Response, Result, UploadedPart};
41
42use crate::artifacts::blob_url;
43use crate::operations::Services;
44
45/// The Twirp services, by name.
46pub const CACHE_SERVICE: &str = "github.actions.results.api.v1.CacheService";
47pub const ARTIFACT_SERVICE: &str = "github.actions.results.api.v1.ArtifactService";
48/// Where the cache's older protocol is: `ACTIONS_CACHE_URL`.
49pub const CACHE_PATH: &str = "/actions/toolkit/";
50/// The largest single block or blob a request may carry.
51const MAX_BLOCK_BYTES: u64 = 256 * 1024 * 1024;
52
53/// The bearer token of a request.
54pub fn bearer(request: &Request) -> String {
55 request
56 .headers()
57 .get("authorization")
58 .ok()
59 .flatten()
60 .and_then(|h| h.split_once(' ').map(|(_, t)| t.trim().to_owned()))
61 .unwrap_or_default()
62}
63
64fn claims(token: &str) -> Option<Value> {
65 serde_json::from_slice(&URL_SAFE_NO_PAD.decode(token.split('.').nth(1)?).ok()?).ok()
66}
67
68/// The job a runtime token names, unchecked: the actions service checks it.
69pub fn runtime_job(token: &str) -> Option<String> {
70 claims(token)?["job"].as_str().map(str::to_owned)
71}
72
73/// The variables a job gets for the toolkit: its runtime token and where
74/// the services are, and where to ask for an OIDC token when it may.
75pub fn runtime_variables(api: &str, token: &str, id_token: bool) -> Map<String, Value> {
76 let mut vars = Map::new();
77 let mut set = |k: &str, v: String| {
78 vars.insert(k.to_owned(), Value::String(v));
79 };
80 set("ACTIONS_RUNTIME_TOKEN", token.to_owned());
81 set("ACTIONS_RESULTS_URL", format!("{api}/"));
82 set("ACTIONS_CACHE_URL", format!("{api}{CACHE_PATH}"));
83 set("ACTIONS_CACHE_SERVICE_V2", "True".to_owned());
84 if id_token {
85 set("ACTIONS_ID_TOKEN_REQUEST_URL", format!("{}/token?api-version=2.0", crate::oidc::issuer(api)));
86 set("ACTIONS_ID_TOKEN_REQUEST_TOKEN", token.to_owned());
87 }
88 vars
89}
90
91// ── Twirp ───────────────────────────────────────────────────────────────────
92
Merge remote-tracking branch 'origin/main' into workspace-chat93/// Twirp's binary encoding.
94const PROTOBUF: &str = "application/protobuf";
95
96/// Whether a request's `Content-Type` is Twirp's protobuf encoding.
97fn is_protobuf(content_type: &str) -> bool {
98 let kind = content_type.split(';').next().unwrap_or_default().trim().to_ascii_lowercase();
99 kind == PROTOBUF || kind == "application/x-protobuf"
100}
101
102/// The cache service's messages in protobuf, read into and written from the
103/// JSON the handlers use (`results/api/v1/cache.proto`, field numbers as
104/// there). Only what the cache's three methods carry: strings, a repeated
105/// string, an int64 and a bool. `metadata` (field 1 of each request) is
106/// skipped: the runtime token says whose cache it is.
107mod proto {
108 use serde_json::{Map, Value};
109
110 #[derive(Clone, Copy)]
111 enum Kind {
112 Text,
113 Texts,
114 Int,
115 Bool,
116 }
117
118 /// A message's fields: number, JSON name, kind.
119 type Fields = &'static [(u64, &'static str, Kind)];
120
121 fn request_fields(method: &str) -> Option<Fields> {
122 Some(match method {
123 "CreateCacheEntry" => &[(2, "key", Kind::Text), (3, "version", Kind::Text)],
124 "FinalizeCacheEntryUpload" => &[(2, "key", Kind::Text), (3, "size_bytes", Kind::Int), (4, "version", Kind::Text)],
125 "GetCacheEntryDownloadURL" => &[(2, "key", Kind::Text), (3, "restore_keys", Kind::Texts), (4, "version", Kind::Text)],
126 _ => return None,
127 })
128 }
129
130 fn response_fields(method: &str) -> Fields {
131 match method {
132 "CreateCacheEntry" => &[(1, "ok", Kind::Bool), (2, "signed_upload_url", Kind::Text), (3, "message", Kind::Text)],
133 "FinalizeCacheEntryUpload" => &[(1, "ok", Kind::Bool), (2, "entry_id", Kind::Int), (3, "message", Kind::Text)],
134 "GetCacheEntryDownloadURL" => &[(1, "ok", Kind::Bool), (2, "signed_download_url", Kind::Text), (3, "matched_key", Kind::Text)],
135 _ => &[],
136 }
137 }
138
139 fn varint(bytes: &[u8], at: &mut usize) -> Option<u64> {
140 let mut value = 0u64;
141 for shift in (0..64).step_by(7) {
142 let byte = *bytes.get(*at)?;
143 *at += 1;
144 value |= u64::from(byte & 0x7f) << shift;
145 if byte & 0x80 == 0 {
146 return Some(value);
147 }
148 }
149 None
150 }
151
152 fn put_varint(out: &mut Vec<u8>, mut value: u64) {
153 while value >= 0x80 {
154 out.push((value as u8 & 0x7f) | 0x80);
155 value >>= 7;
156 }
157 out.push(value as u8);
158 }
159
160 /// A request of `method` as JSON, or None when it is not one.
161 pub fn request(method: &str, bytes: &[u8]) -> Option<Value> {
162 let fields = request_fields(method)?;
163 let mut out = Map::new();
164 let mut at = 0;
165 while at < bytes.len() {
166 let tag = varint(bytes, &mut at)?;
167 let (number, wire) = (tag >> 3, tag & 7);
168 let known = fields.iter().find(|(n, _, _)| *n == number);
169 match wire {
170 0 => {
171 let value = varint(bytes, &mut at)?;
172 if let Some((_, name, Kind::Int)) = known {
173 // An int64 is sent as its two's complement.
174 out.insert((*name).to_owned(), Value::String((value as i64).to_string()));
175 }
176 }
177 2 => {
178 let length = usize::try_from(varint(bytes, &mut at)?).ok()?;
179 let end = at.checked_add(length).filter(|end| *end <= bytes.len())?;
180 let raw = &bytes[at..end];
181 at = end;
182 match known {
183 Some((_, name, Kind::Text)) => {
184 out.insert((*name).to_owned(), Value::String(String::from_utf8(raw.to_vec()).ok()?));
185 }
186 Some((_, name, Kind::Texts)) => {
187 let text = Value::String(String::from_utf8(raw.to_vec()).ok()?);
188 match out.entry((*name).to_owned()).or_insert_with(|| Value::Array(Vec::new())) {
189 Value::Array(list) => list.push(text),
190 _ => return None,
191 }
192 }
193 _ => {}
194 }
195 }
196 1 => at = at.checked_add(8).filter(|end| *end <= bytes.len())?,
197 5 => at = at.checked_add(4).filter(|end| *end <= bytes.len())?,
198 _ => return None,
199 }
200 }
201 Some(Value::Object(out))
202 }
203
204 /// A response of `method` from its JSON. Defaults are left out, as
205 /// proto3 does.
206 pub fn response(method: &str, value: &Value) -> Vec<u8> {
207 let mut out = Vec::new();
208 for (number, name, kind) in response_fields(method) {
209 let field = &value[*name];
210 match kind {
211 Kind::Bool if field.as_bool() == Some(true) => {
212 put_varint(&mut out, number << 3);
213 put_varint(&mut out, 1);
214 }
215 Kind::Int => {
216 let n = field.as_i64().or_else(|| field.as_str().and_then(|s| s.parse().ok())).unwrap_or(0);
217 if n != 0 {
218 put_varint(&mut out, number << 3);
219 put_varint(&mut out, n as u64);
220 }
221 }
222 Kind::Text => {
223 let text = field.as_str().unwrap_or_default();
224 if !text.is_empty() {
225 put_varint(&mut out, (number << 3) | 2);
226 put_varint(&mut out, text.len() as u64);
227 out.extend_from_slice(text.as_bytes());
228 }
229 }
230 _ => {}
231 }
232 }
233 out
234 }
235}
236
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2237/// A Twirp error: its code and message, at the status Twirp gives it.
238fn twirp_error(code: &str, message: &str) -> Result<Response> {
239 let status = match code {
240 "unauthenticated" => 401,
241 "permission_denied" => 403,
242 "not_found" => 404,
243 "already_exists" => 409,
Merge remote-tracking branch 'origin/main' into workspace-chat244 "invalid_argument" | "malformed" => 400,
245 "bad_route" => 404,
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2246 "failed_precondition" => 412,
247 "resource_exhausted" => 429,
248 _ => 500,
249 };
250 Ok(Response::from_json(&json!({ "code": code, "msg": message }))?.with_status(status))
251}
252
253fn twirp_failure(failure: &Failure) -> Result<Response> {
254 let code = match failure.code {
255 FailureCode::Unauthenticated => "unauthenticated",
256 FailureCode::Forbidden => "permission_denied",
257 FailureCode::NotFound => "not_found",
258 FailureCode::Conflict => "already_exists",
259 FailureCode::Invalid => "invalid_argument",
260 _ => "failed_precondition",
261 };
262 twirp_error(code, &failure.message)
263}
264
265/// A field in the toolkit's spelling (`snake_case`), or its JSON name.
266fn field<'a>(body: &'a Value, name: &str) -> &'a Value {
267 if !body[name].is_null() {
268 return &body[name];
269 }
270 let camel: String = name.split('_').enumerate().map(|(i, p)| if i == 0 { p.to_owned() } else { p[..1].to_uppercase() + &p[1..] }).collect();
271 &body[camel]
272}
273
274fn text(body: &Value, name: &str) -> String {
275 match field(body, name) {
276 Value::String(s) => s.clone(),
277 Value::Number(n) => n.to_string(),
278 // A wrapper written as an object, `{ "value": … }`.
279 Value::Object(o) => o.get("value").map(|v| v.as_str().map_or_else(|| v.to_string(), str::to_owned)).unwrap_or_default(),
280 _ => String::new(),
281 }
282}
283
284fn number(body: &Value, name: &str) -> Option<u64> {
285 text(body, name).trim().parse().ok()
286}
287
288/// The run and job a runtime token names, which a request's backend ids
289/// must match.
290fn backend_ids(token: &str) -> (String, String) {
291 let c = claims(token).unwrap_or_default();
292 (c["run"].as_str().unwrap_or_default().to_owned(), c["job"].as_str().unwrap_or_default().to_owned())
293}
294
295/// What the toolkit's artifact client lists.
296fn listed(artifact: &Artifact) -> Value {
297 json!({
298 "workflow_run_backend_id": artifact.run_id,
299 "workflow_job_run_backend_id": artifact.job_id,
300 "database_id": artifact.id.to_string(),
301 "name": artifact.name,
302 "size": artifact.size.to_string(),
303 "created_at": artifact.created_at,
304 "digest": artifact.digest,
305 })
306}
307
Merge remote-tracking branch 'origin/main' into workspace-chat308/// The Azure Storage version g1t's blob links answer as.
309const AZURE_VERSION: &str = "2024-11-04";
310
311/// An upload link: the blob's, with a query as an Azure SAS link has one.
312/// A client that treats it as a container, a blob and a SAS token (OpenDAL,
313/// which sccache uses) refuses a link without one; the token in the path
314/// is what g1t checks.
315pub fn upload_url(api: &str, blob: &str) -> String {
316 format!("{}?sv={AZURE_VERSION}", blob_url(api, blob))
317}
318
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2319/// Starts an R2 upload for an entry the service reserved, and the signed
320/// link the toolkit sends it to.
321async fn start_upload(bucket: &Bucket, services: &Services, job: &str, token: &str, kind: &str, id: &str, object: &str) -> Result<Outcome<String>> {
322 let upload = bucket.create_multipart_upload(object).execute().await?;
323 let upload = upload.upload_id().await;
324 let signed: Outcome<String> = g1t_kit::call(
325 &services.actions,
326 "blob_sign",
327 &BlobSignArgs { job: job.to_owned(), token: token.to_owned(), kind: kind.to_owned(), id: id.to_owned(), upload },
328 )
329 .await?;
330 Ok(match signed {
Merge remote-tracking branch 'origin/main' into workspace-chat331 Outcome::Ok(blob) => Outcome::Ok(upload_url(&services.addresses.api, &blob)),
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2332 Outcome::Fail(refused) => Outcome::Fail(refused),
333 })
334}
335
336/// `POST /twirp/{service}/{method}`.
337pub async fn twirp(mut request: Request, env: &Env, services: &Services, service: &str, method: &str) -> Result<Response> {
338 let token = bearer(&request);
339 let Some(job) = runtime_job(&token) else {
340 return twirp_error("unauthenticated", "Send the job's ACTIONS_RUNTIME_TOKEN as a bearer token.");
341 };
Merge remote-tracking branch 'origin/main' into workspace-chat342 // Twirp clients send JSON or protobuf, and are answered in kind.
343 let binary = is_protobuf(&request.headers().get("content-type")?.unwrap_or_default());
344 let body: Value = if binary {
345 match proto::request(method, &request.bytes().await.unwrap_or_default()) {
346 Some(body) => body,
347 None => return twirp_error("malformed", "That is not a protobuf message this method takes."),
348 }
349 } else {
350 request.json().await.unwrap_or(Value::Null)
351 };
352 let answer = |value: Value| -> Result<Response> {
353 if binary {
354 let mut response = Response::from_bytes(proto::response(method, &value))?;
355 response.headers_mut().set("content-type", PROTOBUF)?;
356 Ok(response)
357 } else {
358 Response::from_json(&value)
359 }
360 };
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2361 let bucket = env.bucket("ACTIONS_CACHE")?;
362 let actions = &services.actions;
363 let (run, own_job) = backend_ids(&token);
364 // An artifact call names its run, and for an upload its job: the
365 // token's.
366 if service == ARTIFACT_SERVICE {
367 let asked_run = text(&body, "workflow_run_backend_id");
368 if !asked_run.is_empty() && asked_run != run {
369 return twirp_error("permission_denied", "The runtime token is for another run.");
370 }
371 let asked_job = text(&body, "workflow_job_run_backend_id");
372 if matches!(method, "CreateArtifact" | "FinalizeArtifact") && !asked_job.is_empty() && asked_job != own_job {
373 return twirp_error("permission_denied", "The runtime token is for another job.");
374 }
375 }
376 match (service, method) {
377 (CACHE_SERVICE, "GetCacheEntryDownloadURL") => {
378 let restore: Vec<String> = field(&body, "restore_keys").as_array().map(|k| k.iter().filter_map(|v| v.as_str().map(str::to_owned)).collect()).unwrap_or_default();
379 let args = CacheLookupArgs { job, token, key: text(&body, "key"), restore, version: Some(text(&body, "version")) };
380 let found: Outcome<Option<CacheHit>> = g1t_kit::call(actions, "cache_lookup", &args).await?;
381 match found {
382 Outcome::Ok(Some(CacheHit { key, blob: Some(blob), .. })) => {
Merge remote-tracking branch 'origin/main' into workspace-chat383 answer(json!({ "ok": true, "signed_download_url": blob_url(&services.addresses.api, &blob), "matched_key": key }))
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2384 }
Merge remote-tracking branch 'origin/main' into workspace-chat385 Outcome::Ok(_) => answer(json!({ "ok": false, "signed_download_url": "", "matched_key": "" })),
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2386 Outcome::Fail(refused) => twirp_failure(&refused),
387 }
388 }
389 (CACHE_SERVICE, "CreateCacheEntry") => {
390 let args = CacheReserveArgs { job: job.clone(), token: token.clone(), key: text(&body, "key"), size: 0, version: Some(text(&body, "version")) };
391 let reserved: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_reserve", &args).await?;
392 let reserved = match reserved {
393 Outcome::Ok(reserved) => reserved,
Merge remote-tracking branch 'origin/main' into workspace-chat394 // A key already saved, or being saved by another job, to a
395 // protobuf client (OpenDAL's) is Twirp's `already_exists`
396 // (409), which it takes as "someone else has it": sccache
397 // then still writes. An `ok: false` would read as a broken
398 // cache, and sccache would only read from it.
399 Outcome::Fail(refused) if binary && refused.code == FailureCode::Conflict => return twirp_failure(&refused),
400 // The toolkit's client logs this ("another job may be
401 // creating this cache") and goes on.
402 Outcome::Fail(refused) => return answer(json!({ "ok": false, "signed_upload_url": "", "message": refused.message })),
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2403 };
404 match start_upload(&bucket, services, &job, &token, "cache", &reserved.id, &reserved.object).await? {
Merge remote-tracking branch 'origin/main' into workspace-chat405 Outcome::Ok(url) => answer(json!({ "ok": true, "signed_upload_url": url })),
406 Outcome::Fail(refused) => answer(json!({ "ok": false, "signed_upload_url": "", "message": refused.message })),
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2407 }
408 }
409 (CACHE_SERVICE, "FinalizeCacheEntryUpload") => {
410 let args = CacheUploadArgs { job: job.clone(), token: token.clone(), number: None, key: Some(text(&body, "key")), version: Some(text(&body, "version")) };
411 let pending: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_upload", &args).await?;
412 let pending = match pending {
413 Outcome::Ok(pending) => pending,
Merge remote-tracking branch 'origin/main' into workspace-chat414 Outcome::Fail(refused) => return answer(json!({ "ok": false, "entry_id": "0", "message": refused.message })),
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2415 };
416 let Some(object) = bucket.head(&pending.object).await? else {
Merge remote-tracking branch 'origin/main' into workspace-chat417 return answer(json!({ "ok": false, "entry_id": "0", "message": "Nothing was uploaded for that entry." }));
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2418 };
419 match commit_cache(&bucket, services, &job, &token, &pending.id, object.size()).await? {
Merge remote-tracking branch 'origin/main' into workspace-chat420 Outcome::Ok(()) => answer(json!({ "ok": true, "entry_id": pending.number.to_string() })),
421 Outcome::Fail(refused) => answer(json!({ "ok": false, "entry_id": "0", "message": refused.message })),
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2422 }
423 }
424 (ARTIFACT_SERVICE, "CreateArtifact") => {
425 let expires_at = Some(text(&body, "expires_at")).filter(|e| !e.is_empty());
426 let args = ArtifactReserveArgs { job: job.clone(), token: token.clone(), name: text(&body, "name"), expires_at, ..ArtifactReserveArgs::default() };
427 let reserved: Outcome<ArtifactReservation> = g1t_kit::call(actions, "artifact_reserve", &args).await?;
428 let reserved = match reserved {
429 Outcome::Ok(reserved) => reserved,
430 Outcome::Fail(refused) => return twirp_failure(&refused),
431 };
432 match start_upload(&bucket, services, &job, &token, "artifact", &reserved.id.to_string(), &reserved.object).await? {
433 Outcome::Ok(url) => Response::from_json(&json!({ "ok": true, "signed_upload_url": url })),
434 Outcome::Fail(refused) => twirp_failure(&refused),
435 }
436 }
437 (ARTIFACT_SERVICE, "FinalizeArtifact") => {
438 let digest = Some(text(&body, "hash")).filter(|h| !h.is_empty());
439 // The size it was measured at as it was stored, not the one it
440 // says.
441 let args = ArtifactCommitArgs { job, token, id: None, name: Some(text(&body, "name")), size: 0, digest };
442 let done: Outcome<Artifact> = g1t_kit::call(actions, "artifact_commit", &args).await?;
443 match done {
444 Outcome::Ok(artifact) => Response::from_json(&json!({ "ok": true, "artifact_id": artifact.id.to_string() })),
445 Outcome::Fail(refused) => twirp_failure(&refused),
446 }
447 }
448 (ARTIFACT_SERVICE, "ListArtifacts") => {
449 let name = Some(text(&body, "name_filter")).filter(|n| !n.is_empty());
450 let id = number(&body, "id_filter");
451 let args = JobArtifactsArgs { job, token, run_id: None, name, id };
452 let found: Outcome<Vec<Artifact>> = g1t_kit::call(actions, "job_artifacts", &args).await?;
453 match found {
454 Outcome::Ok(found) => Response::from_json(&json!({ "artifacts": found.iter().map(listed).collect::<Vec<_>>() })),
455 Outcome::Fail(refused) => twirp_failure(&refused),
456 }
457 }
458 (ARTIFACT_SERVICE, "GetSignedArtifactURL") => {
459 let args = JobArtifactsArgs { job, token, run_id: None, name: Some(text(&body, "name")), id: None };
460 let found: Outcome<ArtifactBlob> = g1t_kit::call(actions, "job_artifact", &args).await?;
461 match found {
462 Outcome::Ok(found) if !found.blob.is_empty() => Response::from_json(&json!({ "signed_url": blob_url(&services.addresses.api, &found.blob) })),
463 Outcome::Ok(_) => twirp_error("failed_precondition", "Download links are not set up on this installation."),
464 Outcome::Fail(refused) => twirp_failure(&refused),
465 }
466 }
467 (ARTIFACT_SERVICE, "DeleteArtifact") => {
468 let args = JobArtifactsArgs { job, token, run_id: None, name: Some(text(&body, "name")), id: None };
469 let done: Outcome<Artifact> = g1t_kit::call(actions, "job_delete_artifact", &args).await?;
470 match done {
471 Outcome::Ok(artifact) => Response::from_json(&json!({ "ok": true, "artifact_id": artifact.id.to_string() })),
472 Outcome::Fail(refused) => twirp_failure(&refused),
473 }
474 }
475 _ => twirp_error("bad_route", &format!("No method {method} on {service}.")),
476 }
477}
478
479/// Marks an uploaded cache entry ready, and deletes what that evicted.
480async fn commit_cache(bucket: &Bucket, services: &Services, job: &str, token: &str, id: &str, size: u64) -> Result<Outcome<()>> {
481 let committed: Outcome<CacheCommitted> = g1t_kit::call(
482 &services.actions,
483 "cache_commit",
484 &CacheCommitArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned(), size },
485 )
486 .await?;
487 Ok(match committed {
488 Outcome::Ok(committed) => {
489 if !committed.evicted.is_empty() {
490 bucket.delete_multiple(committed.evicted.iter().map(String::as_str).collect()).await?;
491 }
492 Outcome::Ok(())
493 }
494 Outcome::Fail(refused) => Outcome::Fail(refused),
495 })
496}
497
498// ── The cache's older protocol ──────────────────────────────────────────────
499
500fn plain_error(status: u16, message: &str) -> Result<Response> {
501 Ok(Response::from_json(&json!({ "message": message, "error": { "message": message } }))?.with_status(status))
502}
503
504fn query(request: &Request, name: &str) -> Option<String> {
505 request.url().ok()?.query_pairs().find(|(k, _)| k == name).map(|(_, v)| v.into_owned())
506}
507
508/// The part a chunk of the older protocol is, from its `Content-Range`:
Merge remote-tracking branch 'origin/main' into workspace-chat509/// chunks are `CACHE_PART_BYTES` apart, as the toolkit sends them. A first
510/// chunk may be larger, up to `MAX_BLOCK_BYTES`: a client that sends an
511/// entry in one request (sccache does) sends a single chunk from 0.
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2512pub fn chunk_part(range: &str) -> Option<(u16, u64)> {
513 let range = range.trim().strip_prefix("bytes ")?;
514 let (span, _) = range.split_once('/')?;
515 let (start, end) = span.split_once('-')?;
516 let (start, end): (u64, u64) = (start.trim().parse().ok()?, end.trim().parse().ok()?);
Merge remote-tracking branch 'origin/main' into workspace-chat517 if end < start {
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2518 return None;
519 }
Merge remote-tracking branch 'origin/main' into workspace-chat520 let length = end - start + 1;
521 if start == 0 && length <= MAX_BLOCK_BYTES {
522 return Some((1, length));
523 }
524 if start % CACHE_PART_BYTES != 0 || length > CACHE_PART_BYTES {
525 return None;
526 }
527 Some(((start / CACHE_PART_BYTES + 1) as u16, length))
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2528}
529
Merge remote-tracking branch 'origin/main' into workspace-chat530/// The bytes a download's `Range` header asks for, out of `size`: first and
531/// last, inclusive. None to send the whole blob (no header, or one this
532/// does not read, such as several ranges); `Some(None)` when the range is
533/// past the end (416).
534pub fn byte_range(header: &str, size: u64) -> Option<Option<(u64, u64)>> {
535 let spec = header.trim().strip_prefix("bytes=")?.trim();
536 if spec.contains(',') {
537 return None;
538 }
539 let (first, last) = spec.split_once('-')?;
540 let (first, last) = (first.trim(), last.trim());
541 let range = if first.is_empty() {
542 // The last `n` bytes.
543 let n: u64 = last.parse().ok()?;
544 if n == 0 || size == 0 {
545 return Some(None);
546 }
547 (size.saturating_sub(n), size - 1)
548 } else {
549 let first: u64 = first.parse().ok()?;
550 let last: u64 = if last.is_empty() { u64::MAX } else { last.parse().ok()? };
551 if last < first {
552 return None;
553 }
554 if first >= size {
555 return Some(None);
556 }
557 (first, last.min(size - 1))
558 };
559 Some(Some(range))
560}
561
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2562/// `{ACTIONS_CACHE_URL}_apis/artifactcache/…`. `rest` is the path after it.
563pub async fn cache_v1(mut request: Request, env: &Env, services: &Services, method: &str, rest: &str) -> Result<Response> {
564 let token = bearer(&request);
565 let Some(job) = runtime_job(&token) else {
566 return plain_error(401, "Send the job's ACTIONS_RUNTIME_TOKEN as a bearer token.");
567 };
568 let bucket = env.bucket("ACTIONS_CACHE")?;
569 let actions = &services.actions;
570 let parts: Vec<&str> = rest.split('/').filter(|p| !p.is_empty()).collect();
571 match (method, parts.as_slice()) {
572 ("GET", ["cache"]) => {
573 let keys: Vec<String> = query(&request, "keys").unwrap_or_default().split(',').map(|k| k.trim().to_owned()).filter(|k| !k.is_empty()).collect();
574 let Some((key, restore)) = keys.split_first() else {
575 return plain_error(400, "Give keys.");
576 };
577 let version = query(&request, "version").unwrap_or_default();
578 let args = CacheLookupArgs { job, token, key: key.clone(), restore: restore.to_vec(), version: Some(version.clone()) };
579 let found: Outcome<Option<CacheHit>> = g1t_kit::call(actions, "cache_lookup", &args).await?;
580 match found {
581 Outcome::Ok(Some(CacheHit { key, blob: Some(blob), created_at, .. })) => Response::from_json(&json!({
582 "cacheKey": key,
583 "cacheVersion": version,
584 "scope": "",
585 "creationTime": created_at,
586 "archiveLocation": blob_url(&services.addresses.api, &blob),
587 })),
588 Outcome::Ok(_) => Ok(Response::empty()?.with_status(204)),
589 Outcome::Fail(refused) => plain_error(refused.code.http_status(), &refused.message),
590 }
591 }
592 ("POST", ["caches"]) => {
593 let body: Value = request.json().await.unwrap_or(Value::Null);
594 let size = body["cacheSize"].as_u64().unwrap_or(0);
595 let key = body["key"].as_str().unwrap_or_default().to_owned();
596 let version = body["version"].as_str().unwrap_or_default().to_owned();
597 let args = CacheReserveArgs { job: job.clone(), token: token.clone(), key, size, version: Some(version) };
598 let reserved: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_reserve", &args).await?;
599 let reserved = match reserved {
600 Outcome::Ok(reserved) => reserved,
601 Outcome::Fail(refused) => return plain_error(refused.code.http_status(), &refused.message),
602 };
603 match start_upload(&bucket, services, &job, &token, "cache", &reserved.id, &reserved.object).await? {
604 Outcome::Ok(_) => Ok(Response::from_json(&json!({ "cacheId": reserved.number }))?.with_status(201)),
605 Outcome::Fail(refused) => plain_error(refused.code.http_status(), &refused.message),
606 }
607 }
608 (_, ["caches", number]) => {
609 let Ok(number) = number.parse::<u64>() else {
610 return plain_error(404, "No such cache entry.");
611 };
612 let args = CacheUploadArgs { job: job.clone(), token: token.clone(), number: Some(number), key: None, version: None };
613 let pending: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_upload", &args).await?;
614 let pending = match pending {
615 Outcome::Ok(pending) => pending,
616 Outcome::Fail(refused) => return plain_error(refused.code.http_status(), &refused.message),
617 };
618 let (Some(blob), Some(upload_id)) = (pending.blob.clone(), pending.upload.clone()) else {
619 return plain_error(409, "That entry's upload was not started.");
620 };
621 match method {
622 "PATCH" => {
623 let range = request.headers().get("content-range")?.unwrap_or_default();
624 let Some((part, length)) = chunk_part(&range) else {
625 return plain_error(400, &format!("Send the entry in chunks of {} MB, each with its Content-Range.", CACHE_PART_BYTES / 1_048_576));
626 };
627 let Some(body) = request.inner().body() else { return plain_error(400, "The chunk is empty.") };
628 let upload = bucket.resume_multipart_upload(&pending.object, &upload_id)?;
629 let uploaded = upload.upload_part(part, body).await?;
630 let recorded = BlobArgs { blob, part: u32::from(part), etag: uploaded.etag(), size: length };
631 let _: Outcome<bool> = g1t_kit::call(actions, "blob_part", &recorded).await?;
632 Ok(Response::empty()?.with_status(204))
633 }
634 "POST" => {
635 let parts: Outcome<Vec<BlobPart>> = g1t_kit::call(actions, "blob_parts", &BlobArgs { blob: blob.clone(), ..BlobArgs::default() }).await?;
636 let parts = match parts {
637 Outcome::Ok(parts) => parts,
638 Outcome::Fail(refused) => return plain_error(refused.code.http_status(), &refused.message),
639 };
640 let size = match finish(&bucket, &pending.object, &upload_id, &parts).await {
641 Ok(size) => size,
642 Err(problem) => return plain_error(400, &format!("The entry could not be completed: {problem}")),
643 };
644 let _: Outcome<bool> = g1t_kit::call(actions, "blob_done", &BlobArgs { blob, size, ..BlobArgs::default() }).await?;
645 match commit_cache(&bucket, services, &job, &token, &pending.id, size).await? {
646 Outcome::Ok(()) => Ok(Response::empty()?.with_status(204)),
647 Outcome::Fail(refused) => plain_error(refused.code.http_status(), &refused.message),
648 }
649 }
650 _ => plain_error(405, "PATCH a chunk, or POST to commit."),
651 }
652 }
653 _ => plain_error(404, "No such endpoint."),
654 }
655}
656
657/// Completes an R2 upload from its recorded parts, in order; the size it
658/// came to. An upload with no parts is an empty object.
659async fn finish(bucket: &Bucket, object: &str, upload_id: &str, parts: &[BlobPart]) -> std::result::Result<u64, String> {
660 let upload = bucket.resume_multipart_upload(object, upload_id).map_err(|e| e.to_string())?;
661 if parts.is_empty() {
662 let _ = upload.abort().await;
663 bucket.put(object, Vec::<u8>::new()).execute().await.map_err(|e| e.to_string())?;
664 return Ok(0);
665 }
666 let done = upload
667 .complete(parts.iter().map(|p| UploadedPart::new(p.part as u16, p.etag.clone())))
668 .await
669 .map_err(|e| e.to_string())?;
670 Ok(done.size())
671}
672
673// ── Blobs ───────────────────────────────────────────────────────────────────
674
675/// A block's index, from its id: the toolkit's client (Azure's SDK) makes
676/// block ids as base64 of a prefix and the index padded with zeros.
677pub fn block_index(id: &str) -> Option<u32> {
678 let decoded = STANDARD.decode(id.trim()).ok()?;
679 let text = String::from_utf8(decoded).ok()?;
680 let digits: String = text.chars().rev().take_while(char::is_ascii_digit).collect::<Vec<_>>().into_iter().rev().collect();
681 if digits.is_empty() {
682 return None;
683 }
684 digits.parse().ok()
685}
686
687/// The block ids of a Put Block List body, in order.
688pub fn block_list(xml: &str) -> Vec<String> {
689 let mut ids = Vec::new();
690 let mut rest = xml;
691 while let Some(open) = rest.find('<') {
692 rest = &rest[open + 1..];
693 let Some(close) = rest.find('>') else { break };
694 let tag = &rest[..close];
695 rest = &rest[close + 1..];
696 if matches!(tag, "Latest" | "Committed" | "Uncommitted") {
697 let Some(end) = rest.find("</") else { break };
698 ids.push(rest[..end].trim().to_owned());
699 rest = &rest[end..];
700 }
701 }
702 ids
703}
704
705/// The parts a block list names, as recorded: each block must have been
706/// sent, as the part its index says, and they must run from the first.
707pub fn parts_for(ids: &[String], recorded: &[BlobPart]) -> std::result::Result<Vec<BlobPart>, String> {
708 let mut out = Vec::with_capacity(ids.len());
709 for (position, id) in ids.iter().enumerate() {
710 let index = block_index(id).ok_or_else(|| format!("The block id {id} does not end in its index."))?;
711 let part = index + 1;
712 if part as usize != position + 1 {
713 return Err("The blocks must be listed in the order they were numbered.".to_owned());
714 }
715 let found = recorded.iter().find(|p| p.part == part).ok_or_else(|| format!("Block {id} was never sent."))?;
716 out.push(found.clone());
717 }
718 Ok(out)
719}
720
721fn azure(status: u16) -> Result<Response> {
722 let mut response = Response::empty()?.with_status(status);
723 let headers = response.headers_mut();
724 headers.set("x-ms-request-id", &g1t_contracts::new_id("req", g1t_kit::now_ms()))?;
725 headers.set("x-ms-version", "2024-11-04")?;
726 headers.set("x-ms-request-server-encrypted", "true")?;
727 Ok(response)
728}
729
730fn azure_error(status: u16, code: &str, message: &str) -> Result<Response> {
731 let body = format!("<?xml version=\"1.0\" encoding=\"utf-8\"?><Error><Code>{code}</Code><Message>{message}</Message></Error>");
732 let mut response = Response::ok(body)?.with_status(status);
733 response.headers_mut().set("content-type", "application/xml")?;
734 response.headers_mut().set("x-ms-error-code", code)?;
735 Ok(response)
736}
737
738/// `/actions/toolkit/blobs/{token}`: GET or HEAD a download, PUT an upload.
739pub async fn blob(mut request: Request, env: &Env, services: &Services, method: &str, token: &str) -> Result<Response> {
740 let opened: Outcome<BlobGrant> = g1t_kit::call(&services.actions, "blob_open", &BlobArgs { blob: token.to_owned(), ..BlobArgs::default() }).await?;
741 let grant = match opened {
742 Outcome::Ok(grant) => grant,
743 Outcome::Fail(refused) => {
744 let status = if refused.code == FailureCode::NotFound { 404 } else { 403 };
745 return azure_error(status, if status == 404 { "BlobNotFound" } else { "AuthenticationFailed" }, &refused.message);
746 }
747 };
748 let bucket = env.bucket("ACTIONS_CACHE")?;
749 match (method, grant.upload.as_deref()) {
750 ("GET" | "HEAD", None) => {
751 let headers = |response: &mut Response, size: u64| -> Result<()> {
752 let headers = response.headers_mut();
753 headers.set("content-length", &size.to_string())?;
754 headers.set("content-type", grant.content_type.as_deref().unwrap_or("application/octet-stream"))?;
755 headers.set("x-ms-blob-type", "BlockBlob")?;
Merge remote-tracking branch 'origin/main' into workspace-chat756 headers.set("accept-ranges", "bytes")?;
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2757 if let Some(name) = &grant.filename {
758 headers.set("content-disposition", &format!("attachment; filename=\"{}\"", name.replace('"', "")))?;
759 }
760 Ok(())
761 };
762 if method == "HEAD" {
763 let Some(object) = bucket.head(&grant.object).await? else { return azure_error(404, "BlobNotFound", "It is gone.") };
764 let mut response = Response::empty()?;
765 headers(&mut response, object.size())?;
766 return Ok(response);
767 }
Merge remote-tracking branch 'origin/main' into workspace-chat768 // One byte range (`Range`, or Azure's `x-ms-range`): the
769 // toolkit's client fetches a large entry in segments, side by
770 // side, and writes each where its range says.
771 let asked = match request.headers().get("x-ms-range")? {
772 Some(range) => Some(range),
773 None => request.headers().get("range")?,
774 };
775 if let Some(asked) = asked.filter(|r| !r.trim().is_empty()) {
776 let Some(object) = bucket.head(&grant.object).await? else { return azure_error(404, "BlobNotFound", "It is gone.") };
777 let size = object.size();
778 match byte_range(&asked, size) {
779 Some(Some((first, last))) => {
780 let length = last - first + 1;
781 let Some(object) = bucket.get(&grant.object).range(worker::Range::OffsetWithLength { offset: first, length }).execute().await? else {
782 return azure_error(404, "BlobNotFound", "It is gone.");
783 };
784 let Some(body) = object.body() else { return azure_error(404, "BlobNotFound", "It is gone.") };
785 let mut response = Response::from_body(body.response_body()?)?.with_status(206);
786 headers(&mut response, length)?;
787 response.headers_mut().set("content-range", &format!("bytes {first}-{last}/{size}"))?;
788 return Ok(response);
789 }
790 Some(None) => {
791 let mut response = azure_error(416, "InvalidRange", "The range is past the end of the blob.")?;
792 response.headers_mut().set("content-range", &format!("bytes */{size}"))?;
793 return Ok(response);
794 }
795 // Not a range this reads: the whole blob.
796 None => {}
797 }
798 }
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2799 let Some(object) = bucket.get(&grant.object).execute().await? else { return azure_error(404, "BlobNotFound", "It is gone.") };
800 let size = object.size();
801 let Some(body) = object.body() else { return azure_error(404, "BlobNotFound", "It is gone.") };
802 let mut response = Response::from_body(body.response_body()?)?;
803 headers(&mut response, size)?;
804 Ok(response)
805 }
806 ("PUT", Some(upload_id)) => {
807 let comp = query(&request, "comp").unwrap_or_default();
808 let limit = if grant.kind == "cache" { CACHE_MAX_ENTRY_BYTES } else { ARTIFACT_MAX_BYTES };
809 match comp.as_str() {
810 // Put Block, or Put Blob: one part.
811 "block" | "" => {
812 let part = if comp == "block" {
813 match query(&request, "blockid").as_deref().and_then(block_index) {
814 Some(index) if index < 10_000 => index + 1,
815 _ => return azure_error(400, "InvalidQueryParameterValue", "A block id ends in its index, from 0."),
816 }
817 } else {
818 1
819 };
820 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
821 if length > MAX_BLOCK_BYTES || length > limit {
822 return azure_error(413, "RequestBodyTooLarge", "That block is larger than g1t takes at once.");
823 }
824 let upload = bucket.resume_multipart_upload(&grant.object, upload_id)?;
825 let uploaded = match request.inner().body() {
826 Some(body) if length > 0 => upload.upload_part(part as u16, body).await?,
827 _ => upload.upload_part(part as u16, Vec::<u8>::new()).await?,
828 };
829 let recorded = BlobArgs { blob: token.to_owned(), part, etag: uploaded.etag(), size: length };
830 let _: Outcome<bool> = g1t_kit::call(&services.actions, "blob_part", &recorded).await?;
831 if comp == "block" {
832 return azure(201);
833 }
834 // Put Blob is the whole thing: finish it now.
835 complete_blob(&bucket, services, token, &grant, upload_id, &[], limit, true).await
836 }
837 "blocklist" => {
838 let xml = request.text().await.unwrap_or_default();
839 let ids = block_list(&xml);
840 complete_blob(&bucket, services, token, &grant, upload_id, &ids, limit, false).await
841 }
842 _ => azure_error(400, "InvalidQueryParameterValue", "g1t takes Put Blob, Put Block and Put Block List."),
843 }
844 }
845 _ => azure_error(405, "UnsupportedHttpVerb", "That link is not for this."),
846 }
847}
848
849/// Finishes an upload from its parts: the blocks a list names, or the one
850/// part of a Put Blob.
851#[allow(clippy::too_many_arguments)]
852async fn complete_blob(
853 bucket: &Bucket,
854 services: &Services,
855 token: &str,
856 grant: &BlobGrant,
857 upload_id: &str,
858 ids: &[String],
859 limit: u64,
860 whole: bool,
861) -> Result<Response> {
862 let recorded: Outcome<Vec<BlobPart>> = g1t_kit::call(&services.actions, "blob_parts", &BlobArgs { blob: token.to_owned(), ..BlobArgs::default() }).await?;
863 let recorded = match recorded {
864 Outcome::Ok(recorded) => recorded,
865 Outcome::Fail(refused) => return azure_error(403, "AuthenticationFailed", &refused.message),
866 };
867 let parts = if whole {
868 recorded.into_iter().filter(|p| p.part == 1).collect()
869 } else {
870 match parts_for(ids, &recorded) {
871 Ok(parts) => parts,
872 Err(problem) => return azure_error(400, "InvalidBlockList", &problem),
873 }
874 };
875 let size = match finish(bucket, &grant.object, upload_id, &parts).await {
876 Ok(size) => size,
877 Err(problem) => return azure_error(400, "InvalidBlockList", &format!("The upload could not be completed: {problem}")),
878 };
879 if size > limit {
880 bucket.delete(&grant.object).await?;
881 return azure_error(413, "RequestBodyTooLarge", &format!("It is {} MB, more than g1t keeps ({} MB).", size / 1_048_576, limit / 1_048_576));
882 }
883 let _: Outcome<bool> = g1t_kit::call(&services.actions, "blob_done", &BlobArgs { blob: token.to_owned(), size, ..BlobArgs::default() }).await?;
884 azure(201)
885}
886
887#[cfg(test)]
888mod tests {
889 use super::*;
890
891 /// What Azure's SDK sends as block ids: base64 of a 36-character uuid
892 /// prefix and the index padded to 48 characters in all.
893 fn azure_block_id(index: u32) -> String {
894 let prefix = "4a2f0d2e-8a44-4f1b-9d55-6f1a2b3c4d5e";
895 let padded = format!("{prefix}{index:0>width$}", width = 48 - prefix.len());
896 STANDARD.encode(padded)
897 }
898
899 #[test]
900 fn block_ids_give_their_index() {
901 assert_eq!(block_index(&azure_block_id(0)), Some(0));
902 assert_eq!(block_index(&azure_block_id(17)), Some(17));
903 assert_eq!(block_index(&STANDARD.encode("no-digits")), None);
904 assert_eq!(block_index("not base64!"), None);
905 }
906
907 #[test]
908 fn a_block_list_is_read_in_order_and_matched_to_parts() {
909 // As the SDK's commitBlockList sends it.
910 let xml = format!(
911 "<?xml version=\"1.0\" encoding=\"UTF-8\" standalone=\"yes\"?><BlockList><Latest>{}</Latest><Latest>{}</Latest></BlockList>",
912 azure_block_id(0),
913 azure_block_id(1)
914 );
915 let ids = block_list(&xml);
916 assert_eq!(ids, [azure_block_id(0), azure_block_id(1)]);
917 // Sent out of order, as they are concurrently.
918 let recorded = vec![
919 BlobPart { part: 2, etag: "b".into(), size: 3 },
920 BlobPart { part: 1, etag: "a".into(), size: 8 },
921 ];
922 let parts = parts_for(&ids, &recorded).unwrap();
923 assert_eq!(parts.iter().map(|p| p.etag.as_str()).collect::<Vec<_>>(), ["a", "b"]);
924 // A block never sent, or listed out of order.
925 assert!(parts_for(&[azure_block_id(0), azure_block_id(2)], &recorded).is_err());
926 assert!(parts_for(&[azure_block_id(1), azure_block_id(0)], &recorded).is_err());
927 assert!(block_list("<BlockList></BlockList>").is_empty());
928 }
929
930 #[test]
931 fn older_protocol_chunks_are_parts() {
932 let mb32 = CACHE_PART_BYTES;
933 assert_eq!(chunk_part(&format!("bytes 0-{}/*", mb32 - 1)), Some((1, mb32)));
934 assert_eq!(chunk_part(&format!("bytes {}-{}/*", mb32 * 2, mb32 * 2 + 99)), Some((3, 100)));
Merge remote-tracking branch 'origin/main' into workspace-chat935 // A whole entry in one chunk, as sccache (OpenDAL) sends it: its
936 // check file is 13 bytes, a compiled crate can be well over 32 MB.
937 assert_eq!(chunk_part("bytes 0-12/*"), Some((1, 13)));
938 assert_eq!(chunk_part(&format!("bytes 0-{}/*", mb32)), Some((1, mb32 + 1)));
939 assert_eq!(chunk_part(&format!("bytes 0-{}/*", MAX_BLOCK_BYTES - 1)), Some((1, MAX_BLOCK_BYTES)));
940 // Not on a chunk's boundary, too long, backwards, or not a range.
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2941 assert_eq!(chunk_part("bytes 5-10/*"), None);
Merge remote-tracking branch 'origin/main' into workspace-chat942 assert_eq!(chunk_part(&format!("bytes {mb32}-{}/*", mb32 * 2)), None);
943 assert_eq!(chunk_part(&format!("bytes 0-{}/*", MAX_BLOCK_BYTES)), None);
944 assert_eq!(chunk_part("bytes 10-5/*"), None);
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2945 assert_eq!(chunk_part("0-10"), None);
946 }
947
Merge remote-tracking branch 'origin/main' into workspace-chat948 #[test]
949 fn downloads_read_one_byte_range() {
950 // OpenDAL's stat: the first byte.
951 assert_eq!(byte_range("bytes=0-0", 100), Some(Some((0, 0))));
952 // The toolkit's segments, the last one cut at the end.
953 assert_eq!(byte_range("bytes=0-49", 100), Some(Some((0, 49))));
954 assert_eq!(byte_range("bytes=50-999", 100), Some(Some((50, 99))));
955 assert_eq!(byte_range("bytes=90-", 100), Some(Some((90, 99))));
956 assert_eq!(byte_range("bytes=-10", 100), Some(Some((90, 99))));
957 assert_eq!(byte_range("bytes=-1000", 100), Some(Some((0, 99))));
958 // Past the end: 416.
959 assert_eq!(byte_range("bytes=100-200", 100), Some(None));
960 assert_eq!(byte_range("bytes=0-0", 0), Some(None));
961 assert_eq!(byte_range("bytes=-0", 100), Some(None));
962 // Not read: the whole blob.
963 assert_eq!(byte_range("bytes=0-1,5-6", 100), None);
964 assert_eq!(byte_range("bytes=9-3", 100), None);
965 assert_eq!(byte_range("items=0-1", 100), None);
966 assert_eq!(byte_range("bytes=a-b", 100), None);
967 }
968
969 #[test]
970 fn upload_links_carry_a_query_as_sas_links_do() {
971 let url = upload_url("https://api.g1t.sh", "tok.sig");
972 assert_eq!(url, "https://api.g1t.sh/actions/toolkit/blobs/tok.sig?sv=2024-11-04");
973 // How OpenDAL reads a signed upload link: a container, a blob in
974 // it, and a SAS query, all of which must be there.
975 let rest = url.strip_prefix("https://api.g1t.sh/").unwrap();
976 let (path, query) = rest.split_once('?').unwrap();
977 let (container, blob) = path.split_once('/').unwrap();
978 assert_eq!((container, blob, query), ("actions", "toolkit/blobs/tok.sig", "sv=2024-11-04"));
979 }
980
981 /// A protobuf length-delimited field, as prost writes it.
982 fn pb_text(number: u8, text: &str) -> Vec<u8> {
983 let mut out = vec![(number << 3) | 2, text.len() as u8];
984 out.extend_from_slice(text.as_bytes());
985 out
986 }
987
988 /// The requests sccache 0.18 sends (OpenDAL 0.58's `ghac` service, with
989 /// prost): fields in number order, defaults left out, no metadata.
990 #[test]
991 fn twirp_reads_sccaches_protobuf_requests() {
992 assert!(is_protobuf("application/protobuf"));
993 assert!(is_protobuf("Application/Protobuf; charset=utf-8"));
994 assert!(!is_protobuf("application/json"));
995 assert!(!is_protobuf(""));
996
997 let key = "sccache/f/c/b/fcb0a1d2e3";
998 let version = "sccache-v0.18.0";
999 let create = [pb_text(2, key), pb_text(3, version)].concat();
1000 let read = proto::request("CreateCacheEntry", &create).unwrap();
1001 assert_eq!((text(&read, "key"), text(&read, "version")), (key.to_owned(), version.to_owned()));
1002
1003 // size_bytes is field 3, a varint: 300 is 0xac 0x02.
1004 let finalize = [pb_text(2, key), vec![0x18, 0xac, 0x02], pb_text(4, version)].concat();
1005 let read = proto::request("FinalizeCacheEntryUpload", &finalize).unwrap();
1006 assert_eq!(number(&read, "size_bytes"), Some(300));
1007 assert_eq!(text(&read, "version"), version);
1008
1009 let lookup = [pb_text(2, key), pb_text(4, version)].concat();
1010 let read = proto::request("GetCacheEntryDownloadURL", &lookup).unwrap();
1011 assert_eq!(text(&read, "key"), key);
1012 assert!(field(&read, "restore_keys").is_null());
1013 // The toolkit's own lookup, with metadata (skipped) and restore keys.
1014 let metadata = vec![0x0a, 0x02, 0x08, 0x07];
1015 let with_restore = [metadata, pb_text(2, "k"), pb_text(3, "k-"), pb_text(3, "x-"), pb_text(4, "v")].concat();
1016 let read = proto::request("GetCacheEntryDownloadURL", &with_restore).unwrap();
1017 assert_eq!(field(&read, "restore_keys"), &json!(["k-", "x-"]));
1018 assert_eq!(text(&read, "version"), "v");
1019
1020 // The same three, as prost 0.14 encodes them with OpenDAL's
1021 // generated types (the `ghac` crate, 0.3.0), byte for byte.
1022 let recorded = |hex: &str| -> Vec<u8> { (0..hex.len()).step_by(2).map(|i| u8::from_str_radix(&hex[i..i + 2], 16).unwrap()).collect() };
1023 let prefix = "1218736363616368652f662f632f622f66636230613164326533";
1024 let suffix = "0f736363616368652d76302e31382e30";
1025 assert_eq!(recorded(&format!("{prefix}1a{suffix}")), create);
1026 assert_eq!(recorded(&format!("{prefix}18ac0222{suffix}")), finalize);
1027 assert_eq!(recorded(&format!("{prefix}22{suffix}")), lookup);
1028
1029 // Cut short, or not a cache method.
1030 assert!(proto::request("CreateCacheEntry", &create[..create.len() - 1]).is_none());
1031 assert!(proto::request("CreateCacheEntry", &[0x12, 0xff]).is_none());
1032 assert!(proto::request("CreateArtifact", &create).is_none());
1033 assert_eq!(proto::request("CreateCacheEntry", &[]), Some(json!({})));
1034 }
1035
1036 #[test]
1037 fn twirp_answers_in_protobuf_as_prost_reads_it() {
1038 let url = "https://api.g1t.sh/actions/toolkit/blobs/t?sv=2024-11-04";
1039 let created = proto::response("CreateCacheEntry", &json!({ "ok": true, "signed_upload_url": url }));
1040 assert_eq!(created, [vec![0x08, 0x01], pb_text(2, url)].concat());
1041 // Refused: ok false is the default, so only the message is sent.
1042 let refused = proto::response("CreateCacheEntry", &json!({ "ok": false, "signed_upload_url": "", "message": "no" }));
1043 assert_eq!(refused, pb_text(3, "no"));
1044 // entry_id is an int64, given as a string in JSON.
1045 let finalized = proto::response("FinalizeCacheEntryUpload", &json!({ "ok": true, "entry_id": "300" }));
1046 assert_eq!(finalized, vec![0x08, 0x01, 0x10, 0xac, 0x02]);
1047 let found = proto::response("GetCacheEntryDownloadURL", &json!({ "ok": true, "signed_download_url": "u", "matched_key": "k" }));
1048 assert_eq!(found, [vec![0x08, 0x01], pb_text(2, "u"), pb_text(3, "k")].concat());
1049 // A miss is an empty message: ok false.
1050 assert!(proto::response("GetCacheEntryDownloadURL", &json!({ "ok": false, "signed_download_url": "", "matched_key": "" })).is_empty());
1051 }
1052
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R21053 /// The toolkit's requests, as `@actions/cache` 4 and `@actions/artifact`
1054 /// 2 send them (protobuf-ts, proto field names, no defaults).
1055 #[test]
1056 fn twirp_requests_read_in_the_toolkits_spelling() {
1057 let create_cache = json!({ "key": "node-cache-Linux-x64-npm-abc", "version": "a7f2c1e0" });
1058 assert_eq!(text(&create_cache, "key"), "node-cache-Linux-x64-npm-abc");
1059 let lookup = json!({ "key": "k", "restore_keys": ["k-", "x-"], "version": "v" });
1060 assert_eq!(field(&lookup, "restore_keys").as_array().unwrap().len(), 2);
1061 let finalize = json!({ "key": "k", "size_bytes": "1048576", "version": "v" });
1062 assert_eq!(number(&finalize, "size_bytes"), Some(1_048_576));
1063 let create_artifact = json!({
1064 "workflow_run_backend_id": "run_1",
1065 "workflow_job_run_backend_id": "job_1",
1066 "name": "dist",
1067 "expires_at": "2026-10-13T00:00:00Z",
1068 "version": 4
1069 });
1070 assert_eq!(text(&create_artifact, "workflow_job_run_backend_id"), "job_1");
1071 let list = json!({ "workflow_run_backend_id": "run_1", "workflow_job_run_backend_id": "job_1", "id_filter": "42", "name_filter": "dist" });
1072 assert_eq!(number(&list, "id_filter"), Some(42));
1073 assert_eq!(text(&list, "name_filter"), "dist");
1074 // A client writing JSON names instead reads the same.
1075 let camel = json!({ "workflowRunBackendId": "run_1", "sizeBytes": 3 });
1076 assert_eq!(text(&camel, "workflow_run_backend_id"), "run_1");
1077 assert_eq!(number(&camel, "size_bytes"), Some(3));
1078 }
1079
1080 #[test]
1081 fn the_runtime_token_names_its_run_and_job() {
1082 let payload = URL_SAFE_NO_PAD.encode(json!({ "job": "job_1", "run": "run_1" }).to_string());
1083 let token = format!("h.{payload}.s");
1084 assert_eq!(runtime_job(&token).as_deref(), Some("job_1"));
1085 assert_eq!(backend_ids(&token), ("run_1".to_owned(), "job_1".to_owned()));
1086 assert_eq!(runtime_job("deadbeef"), None);
1087 }
1088
1089 #[test]
1090 fn a_job_is_told_where_the_toolkit_s_services_are() {
1091 let vars = runtime_variables("https://api.g1t.sh", "tok", false);
1092 // Twirp paths are resolved against the root of ACTIONS_RESULTS_URL.
1093 assert_eq!(vars["ACTIONS_RESULTS_URL"], "https://api.g1t.sh/");
1094 // The older protocol appends `_apis/artifactcache/…`.
1095 assert_eq!(vars["ACTIONS_CACHE_URL"], "https://api.g1t.sh/actions/toolkit/");
1096 assert_eq!(vars["ACTIONS_CACHE_SERVICE_V2"], "True");
1097 assert!(vars.get("ACTIONS_ID_TOKEN_REQUEST_URL").is_none());
1098 let vars = runtime_variables("https://api.g1t.sh", "tok", true);
1099 // core.getIDToken appends `&audience=…`.
1100 assert_eq!(vars["ACTIONS_ID_TOKEN_REQUEST_URL"], "https://api.g1t.sh/actions/oidc/token?api-version=2.0");
1101 assert_eq!(vars["ACTIONS_ID_TOKEN_REQUEST_TOKEN"], "tok");
1102 }
1103
1104 #[test]
1105 fn artifacts_are_listed_as_the_toolkit_reads_them() {
1106 let artifact = Artifact {
1107 id: 7,
1108 name: "dist".into(),
1109 size: 10,
1110 digest: None,
1111 format: "zip".into(),
1112 run_id: "run_1".into(),
1113 job_id: "job_1".into(),
1114 repo_id: "repo_1".into(),
1115 expired: false,
1116 created_at: "2026-10-08T12:00:00.000Z".into(),
1117 updated_at: "2026-10-08T12:00:00.000Z".into(),
1118 expires_at: "2026-10-22T12:00:00.000Z".into(),
1119 head_branch: None,
1120 head_sha: None,
1121 };
1122 let shown = listed(&artifact);
1123 assert_eq!(shown["database_id"], "7");
1124 assert_eq!(shown["size"], "10");
1125 assert_eq!(shown["workflow_job_run_backend_id"], "job_1");
1126 }
1127}

This file's history is long; its oldest lines are credited to the oldest commit read.