Initial stabilization #1
4 changed files with 1491 additions and 0 deletions
1
.gitignore
vendored
Normal file
1
.gitignore
vendored
Normal file
|
|
@ -0,0 +1 @@
|
||||||
|
/target
|
||||||
1153
Cargo.lock
generated
Normal file
1153
Cargo.lock
generated
Normal file
File diff suppressed because it is too large
Load diff
13
Cargo.toml
Normal file
13
Cargo.toml
Normal 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
324
src/lib.rs
Normal 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())
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Reference in a new issue