Initial stabilization #1

Merged
dzu merged 10 commits from stabilize-v1 into main 2026-02-20 16:23:43 +01:00
4 changed files with 1491 additions and 0 deletions

1
.gitignore vendored Normal file
View file

@ -0,0 +1 @@
/target

1153
Cargo.lock generated Normal file

File diff suppressed because it is too large Load diff

13
Cargo.toml Normal file
View file

@ -0,0 +1,13 @@
[package]
name = "anode-cache-common"
version = "0.1.0"
edition = "2024"
[dependencies]
anyhow = "1.0.101"
isahc = { version = "1.7.2", features = ["static-ssl", "static-curl"] }
serde_json = "1.0.149"
secrecy = "0.10.3"
thiserror = "2.0.18"
tracing-subscriber = { version = "0.3.22" }
tracing = { version = "0.1.44", features = ["attributes"] }

324
src/lib.rs Normal file
View file

@ -0,0 +1,324 @@
use std::{ffi::OsString, io::Read, path::PathBuf, str::FromStr, string::FromUtf8Error};
use anyhow::bail;
use anyhow::{Context, Result};
use isahc::{
HttpClient, HttpClientBuilder, Request,
auth::{Authentication, Credentials},
config::Configurable,
http::{Method, StatusCode, header},
};
use secrecy::SecretString;
use tracing::{Level, debug, error};
use tracing_subscriber::fmt::format::FmtSpan;
pub fn setup_logger(debug: bool) {
tracing_subscriber::fmt()
.with_ansi(true)
.with_level(true)
.with_max_level(if debug { Level::DEBUG } else { Level::INFO })
.with_span_events(FmtSpan::CLOSE)
.init();
debug!("Logger set up");
}
#[tracing::instrument(level = "debug")]
pub fn paths_parser(input: OsString) -> Result<Vec<PathBuf>, FromUtf8Error> {
let input = String::from_utf8(input.into_encoded_bytes())?;
let mut paths = Vec::new();
for path in input.split('\n').filter(|s| !s.is_empty()) {
paths.push(match PathBuf::from_str(path) {
Ok(ok) => ok,
});
}
debug!(?paths, "Parsed paths");
Ok(paths)
}
#[derive(Debug)]
pub struct CacheKey(String);
#[derive(Debug, thiserror::Error)]
pub enum CacheKeyError {
#[error(transparent)]
Utf8(#[from] FromUtf8Error),
#[error("Cache key is at most 512 characters long. {len} characters long key provided")]
Length { len: usize },
#[error("Cache key cannot contain commas. Comma found at offset {offset}")]
Comma { offset: usize },
}
#[tracing::instrument(level = "debug", fields(return))]
pub fn key_parser(input: OsString) -> Result<CacheKey, CacheKeyError> {
let input = String::from_utf8(input.into_encoded_bytes())?;
let mut len = 0;
for (offset, ch) in input.char_indices() {
if ch == ',' {
return Err(CacheKeyError::Comma { offset });
}
len += 1;
}
if len > 512 {
return Err(CacheKeyError::Length { len });
}
Ok(CacheKey(input))
}
#[derive(Debug, thiserror::Error)]
pub enum ActionsTokenError {
#[error("Could not find action runtime token")]
NotPresent,
#[error("Action runtime token was not UTF-8 -encoded")]
NotUtf8,
}
#[tracing::instrument(level = "debug")]
pub fn actions_token() -> Result<SecretString, ActionsTokenError> {
let os_str = std::env::var_os("ACTIONS_RUNTIME_TOKEN").ok_or(ActionsTokenError::NotPresent)?;
let str =
String::from_utf8(os_str.into_encoded_bytes()).map_err(|_| ActionsTokenError::NotUtf8)?;
Ok(SecretString::new(str.into_boxed_str()))
}
#[derive(Debug)]
pub struct Client {
http_client: HttpClient,
base_url: String,
}
#[derive(Debug, Clone, Copy)]
pub struct CacheId(u64);
impl Client {
#[tracing::instrument(level = "debug", fields(base_url = tracing::field::Empty), skip_all)]
pub fn new(actions_token: impl AsRef<str>, base_url: impl Into<String>) -> Result<Self> {
let http_client = HttpClientBuilder::new()
.authentication(Authentication::basic())
.credentials(Credentials::new("Bearer", actions_token.as_ref()))
.build()
.context("initializing the client")?;
let base_url = base_url.into();
tracing::span::Span::current().record("base_url", &base_url);
Ok(Self {
http_client,
base_url,
})
}
#[tracing::instrument(level = "info", skip(cache))]
pub fn save_cache(&self, key: &CacheKey, cache: Vec<u8>) -> Result<()> {
let cache_len = cache.len();
let cache_id = self.reserve(key, cache_len).context("reserving cache")?;
self.save(cache_id, cache).context("saving the cache")?;
self.commit_cache(cache_id, cache_len)
.context("committing the cache")?;
Ok(())
}
#[tracing::instrument(level = "debug", fields(return))]
fn reserve(&self, key: &CacheKey, cache_size: usize) -> Result<CacheId> {
let uri = format!("{}_apis/artifactcache/caches", self.base_url);
debug!(uri);
let request = Request::builder()
.uri(uri)
.method(Method::POST)
.body(
serde_json::json! {{
"key": &key.0,
"version": "0",
"cacheSize": cache_size,
}}
.to_string(),
)
.context("building the request")?;
let mut response = self
.http_client
.send(request)
.context("sending the request")?;
debug!(?response);
let status = response.status();
let mut body = Vec::with_capacity(100);
response
.body_mut()
.read_to_end(&mut body)
.context("reading response body")?;
if status != StatusCode::OK {
error!(?status, "failed sending a request");
bail!("failed sending a request");
}
let body_value =
serde_json::from_slice::<serde_json::Value>(&body).context("parsing response body")?;
debug!(body = ?body_value);
let cache_id = body_value
.get("cacheId")
.context("response.cacheId")?
.as_u64()
.context("cacheId as u64")?;
Ok(CacheId(cache_id))
}
#[tracing::instrument(level = "debug", skip(archive))]
fn save(&self, id: CacheId, archive: Vec<u8>) -> Result<()> {
let uri = format!("{}_apis/artifactcache/caches/{}", self.base_url, id.0);
debug!(uri);
let request = Request::builder()
.uri(uri)
.method(Method::PATCH)
.header(header::CONTENT_TYPE, "application/octet-stream")
.header(
header::CONTENT_RANGE,
format!("bytes {}-{}/*", 0, archive.len()),
)
.body(archive)
.context("building the request")?;
let mut response = self
.http_client
.send(request)
.context("sending the request")?;
debug!(?response);
let status = response.status();
let mut body = Vec::with_capacity(100);
response
.body_mut()
.read_to_end(&mut body)
.context("reading response body")?;
if status != StatusCode::OK {
error!(?status, body = ?String::from_utf8_lossy(&body), "failed sending a request");
bail!("failed sending a request");
}
Ok(())
}
#[tracing::instrument(level = "debug")]
fn commit_cache(&self, id: CacheId, cache_size: usize) -> Result<()> {
let uri = format!("{}_apis/artifactcache/caches/{}", self.base_url, id.0);
debug!(uri);
let request = Request::builder()
.uri(uri)
.method(Method::POST)
.body(
serde_json::json! {{
"size": cache_size,
}}
.to_string(),
)
.context("building the request")?;
let mut response = self
.http_client
.send(request)
.context("sending the request")?;
debug!(?response);
let status = response.status();
if status != StatusCode::OK {
let mut body = Vec::with_capacity(100);
response
.body_mut()
.read_to_end(&mut body)
.context("reading response body")?;
error!(?status, body = ?String::from_utf8_lossy(&body), "failed sending a request");
bail!("failed sending a request");
}
Ok(())
}
#[tracing::instrument(level = "info")]
pub fn load_cache(&self, key: &CacheKey) -> Result<Option<Box<[u8]>>> {
let Some(cache_location) = self.get_cache_entry(key).context("getting cache entry")? else {
// cache miss
return Ok(None);
};
let cache = self
.download_cache(cache_location)
.context("downloading the cache")?;
Ok(Some(cache))
}
#[tracing::instrument(level = "debug", fields(return))]
fn get_cache_entry(&self, key: &CacheKey) -> Result<Option<String>> {
let uri = format!(
"{}_apis/artifactcache/cache?keys={}&version=0",
self.base_url, key.0
);
debug!(uri);
let request = Request::builder()
.uri(uri)
.method(Method::GET)
.body(())
.context("building the request")?;
let mut response = self
.http_client
.send(request)
.context("sending the request")?;
debug!(?response);
let status = response.status();
if status == StatusCode::NO_CONTENT {
// cache miss
return Ok(None);
}
let mut body = Vec::with_capacity(100);
response
.body_mut()
.read_to_end(&mut body)
.context("reading response body")?;
if status != StatusCode::OK {
error!(?status, body = ?String::from_utf8_lossy(&body), "failed sending a request");
bail!("failed sending a request");
}
let body_value =
serde_json::from_slice::<serde_json::Value>(&body).context("parsing response body")?;
debug!(body = ?body_value);
let location = body_value
.get("archiveLocation")
.context("not cache location field")?
.as_str()
.context("archiveLocation is not a string")?;
Ok(Some(location.to_string()))
}
#[tracing::instrument(level = "debug")]
fn download_cache(&self, location: String) -> Result<Box<[u8]>> {
let request = Request::builder()
.uri(location)
.method(Method::GET)
.body(())
.context("building the request")?;
let mut response = self
.http_client
.send(request)
.context("sending the request")?;
debug!(?response);
let status = response.status();
let mut body = Vec::with_capacity(100);
response
.body_mut()
.read_to_end(&mut body)
.context("reading response body")?;
if status != StatusCode::OK {
error!(?status, body = ?String::from_utf8_lossy(&body), "failed sending a request");
bail!("failed sending a request");
}
Ok(body.into_boxed_slice())
}
}