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