use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use std::process::Stdio;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use anyhow::{Context, Result, anyhow, bail};
use async_trait::async_trait;
use base64::Engine;
use base64::engine::general_purpose::STANDARD;
use tokio::io::AsyncReadExt;
use tokio::process::Command;
use tokio::sync::{OnceCell, Semaphore};
use url::Url;
use verify_trust::{MergeFacts, MergeOutcome, RangeCommit};
use vgi_forge::Resource;
use vgi_forge_github::checks::check_sha;
use vgi_forge_github::{CheckConclusion, CheckTrigger, CheckTriggerKind};
use zeroize::Zeroizing;
use crate::bridge::Bridge;
use crate::config::{BridgeConfig, CheckConfig};
use crate::jobs::Ctx;
use crate::store::{LastCheck, NamespaceRecord, NamespaceState, RepoRecord, Table, repo_key};
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct CommitLine {
pub sha: String,
pub passes: bool,
pub verdict: String,
}
impl CommitLine {
pub fn new(sha: impl Into<String>, passes: bool, verdict: impl Into<String>) -> Self {
CommitLine {
sha: sha.into(),
passes,
verdict: verdict.into(),
}
}
}
#[async_trait]
pub trait CommitVerifier: Send + Sync {
async fn verify(
&self,
commits: &[RangeCommit],
resource: &str,
fallback: &str,
) -> Result<Vec<CommitLine>>;
async fn verify_with_facts(
&self,
commits: &[RangeCommit],
_facts: &MergeFacts,
resource: &str,
fallback: &str,
) -> Result<Vec<CommitLine>> {
self.verify(commits, resource, fallback).await
}
fn exempts_platform_merges(&self, _host: &str) -> bool {
false
}
async fn commit_sign_granted(&self, _did: &str, _resource: &str) -> Result<Option<bool>> {
Ok(None)
}
}
pub struct VerifyTrustVerifier {
registry_did: String,
vtc_did: String,
max_signers: usize,
keyrings: BTreeMap<String, PathBuf>,
tdk: OnceCell<Arc<affinidi_tdk::TDK>>,
ttl: std::time::Duration,
registry_route: tokio::sync::Mutex<Option<(trql_client::TransportChoice, std::time::Instant)>>,
re_resolved: tokio::sync::Mutex<std::collections::HashMap<String, std::time::Instant>>,
registry_override: Option<String>,
transport: verify_trust::TransportSelector,
channels: Vec<Arc<dyn verify_trust::RegistryChannel>>,
}
impl VerifyTrustVerifier {
pub fn new(cfg: &BridgeConfig) -> Self {
VerifyTrustVerifier {
registry_did: cfg.trust_registry_did.clone(),
vtc_did: cfg.vtc_did.clone(),
max_signers: cfg.checks.max_signers,
keyrings: cfg
.github
.iter()
.filter_map(|g| g.platform_keyring_file.clone().map(|k| (g.host.clone(), k)))
.collect(),
tdk: OnceCell::new(),
ttl: std::time::Duration::from_secs(cfg.did_cache_ttl_secs),
registry_route: tokio::sync::Mutex::new(None),
re_resolved: Default::default(),
registry_override: None,
transport: bridge_transport(cfg.verify_trust.transport),
channels: Vec::new(),
}
}
pub fn with_registry_url(mut self, url: impl Into<String>) -> Self {
self.registry_override = Some(url.into());
self
}
#[doc(hidden)]
pub fn with_route(self, route: trql_client::TransportChoice) -> Self {
*self.registry_route.try_lock().expect("not yet shared") =
Some((route, std::time::Instant::now()));
self
}
pub fn with_channel(mut self, channel: Arc<dyn verify_trust::RegistryChannel>) -> Self {
self.channels.push(channel);
self
}
fn supported(&self) -> Vec<trql_client::TransportKind> {
let mut kinds: Vec<_> = self.channels.iter().map(|c| c.kind()).collect();
kinds.push(trql_client::TransportKind::Https);
kinds
}
async fn registry(&self, tdk: &affinidi_tdk::TDK) -> Result<verify_trust::Registry> {
if let Some(u) = &self.registry_override {
return verify_trust::Registry::https(u, &self.registry_did);
}
let mut cached = self.registry_route.lock().await;
let route = match cached.as_ref() {
Some((r, at)) if at.elapsed() < self.ttl => r.clone(),
_ => {
let r = verify_trust::registry::discover_registry_route(
tdk,
&self.registry_did,
self.transport,
&self.supported(),
)
.await?;
tracing::info!(binding = %r.kind, "querying the Trust Registry");
*cached = Some((r.clone(), std::time::Instant::now()));
r
}
};
if route.kind == trql_client::TransportKind::Https {
return verify_trust::Registry::https(&route.endpoint, &self.registry_did);
}
match self.channels.iter().find(|c| c.kind() == route.kind) {
Some(channel) => Ok(verify_trust::Registry::over_channel(
Arc::clone(channel),
&self.registry_did,
)),
None => bail!("the bridge cannot query the registry over {}", route.kind),
}
}
async fn resolver(&self) -> Result<&Arc<affinidi_tdk::TDK>> {
let ttl = u32::try_from(self.ttl.as_secs()).unwrap_or(u32::MAX);
self.tdk
.get_or_try_init(|| async {
verify_trust::build_resolver_with_cache_ttl(false, ttl)
.await
.map(Arc::new)
})
.await
}
async fn forget_endpoint(&self) {
*self.registry_route.lock().await = None;
}
}
fn bridge_transport(t: vgi_forge::VerifyTransport) -> verify_trust::TransportSelector {
match t {
vgi_forge::VerifyTransport::Didcomm => verify_trust::TransportSelector::Didcomm,
vgi_forge::VerifyTransport::Https => verify_trust::TransportSelector::Https,
vgi_forge::VerifyTransport::Tsp => verify_trust::TransportSelector::Tsp,
_ => verify_trust::TransportSelector::Auto,
}
}
#[async_trait]
impl CommitVerifier for VerifyTrustVerifier {
async fn verify(
&self,
range: &[RangeCommit],
resource: &str,
fallback: &str,
) -> Result<Vec<CommitLine>> {
self.verify_with_facts(range, &MergeFacts::default(), resource, fallback)
.await
}
fn exempts_platform_merges(&self, host: &str) -> bool {
self.keyrings.contains_key(host)
}
async fn verify_with_facts(
&self,
range: &[RangeCommit],
facts: &MergeFacts,
resource: &str,
fallback: &str,
) -> Result<Vec<CommitLine>> {
let claimed = verify_trust::claimed_signer_dids(range, self.max_signers)?;
let tdk = self.resolver().await?;
let lines = self
.verify_once(tdk, range, facts, resource, fallback, &claimed)
.await?;
let stale: Vec<String> = lines
.1
.iter()
.filter_map(|s| match s {
verify_trust::CommitStatus::UnknownKey { did, .. }
| verify_trust::CommitStatus::UnresolvedSigner { did, .. } => Some(did.clone()),
_ => None,
})
.collect();
if stale.is_empty() {
return Ok(lines.0);
}
let stale: Vec<String> = {
let mut last = self.re_resolved.lock().await;
let now = std::time::Instant::now();
last.retain(|_, t| now.duration_since(*t) < crate::wire::RE_RESOLVE_EVERY);
stale
.into_iter()
.filter(|d| last.insert(d.clone(), now).is_none())
.collect()
};
if stale.is_empty() {
return Ok(lines.0);
}
for did in &stale {
let _ = tdk.did_resolver().remove(did).await;
}
tracing::info!(signers = ?stale, "re-resolving signer DIDs after an unknown key");
Ok(self
.verify_once(tdk, range, facts, resource, fallback, &claimed)
.await?
.0)
}
async fn commit_sign_granted(&self, did: &str, resource: &str) -> Result<Option<bool>> {
use trql_client::TrqpQuery;
let tdk = self.resolver().await?;
let registry = self.registry(tdk).await?;
let query = TrqpQuery::new(did, &self.vtc_did, "git.commit.sign", resource);
match registry.client().authorization(query).await {
Ok(r) => Ok(Some(r.authorized)),
Err(trql_client::TrqlError::Rejected { .. }) => Ok(Some(false)),
Err(e) => Err(e.into()),
}
}
}
impl VerifyTrustVerifier {
async fn verify_once(
&self,
tdk: &affinidi_tdk::TDK,
range: &[RangeCommit],
facts: &MergeFacts,
resource: &str,
fallback: &str,
claimed: &[String],
) -> Result<(Vec<CommitLine>, Vec<verify_trust::CommitStatus>)> {
let registry = self.registry(tdk).await?;
let signers = verify_trust::resolve_signer_keys(tdk, claimed).await?;
let host = resource.split('/').next().unwrap_or_default();
let exempt = match self.keyrings.get(host) {
Some(p) => Some(verify_trust::pgp_exempt::ExemptKeyring::load(p)?),
None => None,
};
let args = verify_trust::VerifyTrustArgs {
repo_dir: PathBuf::new(),
range: String::new(),
max_signers: self.max_signers,
registry_url: None,
transport: self.transport,
registry_did: self.registry_did.clone(),
vtc_did: self.vtc_did.clone(),
action: "git.commit.sign".into(),
resource: resource.into(),
fallback_resource: Some(fallback.into()),
exempt_keyring: None,
resolve_agent_names: false,
json: false,
};
let report = match verify_trust::verify_prepared_with_facts(
&args,
range,
&signers,
exempt.as_ref(),
®istry,
facts,
)
.await
{
Ok(r) => r,
Err(e) => {
self.forget_endpoint().await;
return Err(e);
}
};
if report.commits.iter().any(|c| {
matches!(
c.status,
verify_trust::CommitStatus::RegistryUnavailable { .. }
)
}) {
self.forget_endpoint().await;
}
let statuses = report.commits.iter().map(|c| c.status.clone()).collect();
let lines = report
.commits
.into_iter()
.map(|c| {
let verdict = serde_json::to_value(&c.status)
.ok()
.and_then(|v| v.get("status").and_then(|s| s.as_str()).map(str::to_string))
.unwrap_or_else(|| format!("{:?}", c.status));
CommitLine::new(c.sha, c.status.passes(), verdict)
})
.collect();
Ok((lines, statuses))
}
}
pub const MAX_PLATFORM_MERGES: usize = 8;
pub fn boundary_of(commits: &[RangeCommit]) -> BTreeSet<String> {
let listed: BTreeSet<&str> = commits.iter().map(|c| c.sha.as_str()).collect();
commits
.iter()
.flat_map(|c| verify_trust::commit_parents(&c.raw))
.filter(|parent| !listed.contains(parent.as_str()))
.collect()
}
pub fn platform_merge_candidates(commits: &[RangeCommit]) -> Vec<(String, String, String)> {
commits
.iter()
.filter(|c| {
matches!(
vgi_core::split_signed_commit(&c.raw),
Ok(Some((_, pem))) if pem.starts_with("-----BEGIN PGP SIGNATURE-----")
)
})
.filter_map(|c| match verify_trust::commit_parents(&c.raw).as_slice() {
[ours, theirs] => Some((c.sha.clone(), ours.clone(), theirs.clone())),
_ => None,
})
.collect()
}
pub async fn merge_facts<F, Fut>(
fetcher: &GitFetcher,
dir: &Path,
token: Option<&str>,
commits: &[RangeCommit],
recompute: bool,
merge_base: F,
) -> MergeFacts
where
F: Fn(String, String) -> Fut,
Fut: std::future::Future<Output = Result<String>>,
{
let mut facts = MergeFacts {
boundary: boundary_of(commits),
merges: BTreeMap::new(),
};
if !recompute {
return facts;
}
for (n, (merge, ours, theirs)) in platform_merge_candidates(commits).into_iter().enumerate() {
let outcome = if n >= MAX_PLATFORM_MERGES {
MergeOutcome::Unavailable(format!(
"more than {MAX_PLATFORM_MERGES} platform-signed merges in one check"
))
} else {
match merge_base(ours.clone(), theirs.clone()).await {
Ok(base) => {
fetcher
.recompute_merge(dir, token, &ours, &theirs, &base)
.await
}
Err(e) => MergeOutcome::Unavailable(format!("no merge base from the forge: {e:#}")),
}
};
facts.merges.insert(merge, outcome);
}
facts
}
const MAX_MISSING_OBJECTS: usize = 2_000;
#[derive(Debug, Clone)]
pub struct GitFetcher {
git: PathBuf,
timeout: Duration,
max_bytes: u64,
allow_file: bool,
remote_override: Option<Url>,
}
#[derive(Debug)]
#[non_exhaustive]
pub struct Fetched {
pub commits: Vec<RangeCommit>,
pub dir: tempfile::TempDir,
}
impl GitFetcher {
pub fn new(cfg: &CheckConfig) -> Self {
GitFetcher {
git: cfg.git.clone(),
timeout: Duration::from_secs(cfg.fetch_timeout_secs),
max_bytes: cfg.max_fetch_bytes,
allow_file: false,
remote_override: None,
}
}
#[doc(hidden)]
pub fn with_local_remote(mut self, remote: Url) -> Self {
self.allow_file = true;
self.remote_override = Some(remote);
self
}
fn command(&self, dir: &Path, args: &[&str], token: Option<&str>, lazy: bool) -> Command {
let mut cmd = Command::new(&self.git);
cmd.arg("-C").arg(dir);
for c in [
"core.hooksPath=/dev/null",
"http.followRedirects=false",
"protocol.allow=never",
"core.fsmonitor=false",
"submodule.recurse=false",
"fetch.recurseSubmodules=false",
"transfer.fsckObjects=true",
"protocol.https.allow=always",
] {
cmd.arg("-c").arg(c);
}
if self.allow_file {
cmd.arg("-c").arg("protocol.file.allow=always");
}
cmd.args(args)
.env_clear()
.env("PATH", std::env::var_os("PATH").unwrap_or_default())
.env("HOME", dir)
.env("GIT_CONFIG_NOSYSTEM", "1")
.env("GIT_CONFIG_GLOBAL", "/dev/null")
.env("GIT_TERMINAL_PROMPT", "0")
.env("GIT_ASKPASS", "/bin/false")
.env("GIT_PROTOCOL_FROM_USER", "0")
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.kill_on_drop(true);
if !lazy {
cmd.env("GIT_NO_LAZY_FETCH", "1");
}
if let Some(t) = token {
let header = Zeroizing::new(format!(
"Authorization: Basic {}",
STANDARD.encode(format!("x-access-token:{t}"))
));
cmd.env("GIT_CONFIG_COUNT", "1")
.env("GIT_CONFIG_KEY_0", "http.extraHeader")
.env("GIT_CONFIG_VALUE_0", header.as_str());
}
cmd
}
async fn git(
&self,
dir: &Path,
args: &[&str],
token: Option<&str>,
watch: bool,
) -> Result<Vec<u8>> {
let (status, out, err) = self.run(dir, args, token, watch, false).await?;
if !status.success() {
let err = String::from_utf8_lossy(&err);
let line = err
.lines()
.find(|l| l.contains("rejected]"))
.or_else(|| err.lines().last())
.unwrap_or("failed");
bail!("git {}: {}", args[0], line.trim());
}
Ok(out)
}
async fn run(
&self,
dir: &Path,
args: &[&str],
token: Option<&str>,
watch: bool,
lazy: bool,
) -> Result<(std::process::ExitStatus, Vec<u8>, Vec<u8>)> {
let mut child = self
.command(dir, args, token, lazy)
.spawn()
.context("running git")?;
let mut stdout = child.stdout.take().expect("piped");
let mut stderr = child.stderr.take().expect("piped");
let out_task = tokio::spawn(async move {
let mut b = Vec::new();
let _ = stdout.read_to_end(&mut b).await;
b
});
let err_task = tokio::spawn(async move {
let mut b = Vec::new();
let _ = stderr.read_to_end(&mut b).await;
b
});
let deadline = tokio::time::Instant::now() + self.timeout;
let mut tick = tokio::time::interval(Duration::from_millis(100));
let status = loop {
tokio::select! {
s = child.wait() => break s.context("waiting for git")?,
_ = tick.tick() => {
if tokio::time::Instant::now() > deadline {
let _ = child.kill().await;
bail!("git {} timed out", args[0]);
}
if watch && dir_size(dir) > self.max_bytes {
let _ = child.kill().await;
bail!(
"the fetch passed {} bytes; refusing a change this large",
self.max_bytes
);
}
}
}
};
let out = out_task.await.unwrap_or_default();
let err = err_task.await.unwrap_or_default();
if watch && dir_size(dir) > self.max_bytes {
bail!(
"the fetch passed {} bytes; refusing a change this large",
self.max_bytes
);
}
Ok((status, out, err))
}
pub(crate) async fn complete_for_push(
&self,
dir: &Path,
token: Option<&str>,
new_head: &str,
old_head: &str,
) -> Result<()> {
check_sha(new_head).map_err(|e| anyhow!("{e}"))?;
check_sha(old_head).map_err(|e| anyhow!("{e}"))?;
let not_old = format!("^{old_head}");
let listed = self
.git(
dir,
&[
"rev-list",
"--objects",
"--missing=print",
new_head,
¬_old,
],
None,
false,
)
.await?;
let missing: Vec<String> = String::from_utf8_lossy(&listed)
.lines()
.filter_map(|l| l.strip_prefix('?'))
.map(|l| l.trim().to_string())
.collect();
if missing.is_empty() {
return Ok(());
}
if missing.len() > MAX_MISSING_OBJECTS {
bail!(
"the change touches {} files; more than the bridge re-signs ({MAX_MISSING_OBJECTS})",
missing.len()
);
}
for m in &missing {
check_sha(m).map_err(|e| anyhow!("{e}"))?;
}
let mut args = vec![
"fetch",
"--quiet",
"--no-tags",
"--no-write-fetch-head",
"--no-recurse-submodules",
"--filter=blob:none",
"origin",
];
args.extend(missing.iter().map(String::as_str));
self.git(dir, &args, token, true).await?;
Ok(())
}
pub(crate) async fn git_with_input(
&self,
dir: &Path,
args: &[&str],
input: &[u8],
) -> Result<Vec<u8>> {
use tokio::io::AsyncWriteExt;
let mut cmd = self.command(dir, args, None, false);
cmd.stdin(Stdio::piped());
let mut child = cmd.spawn().context("running git")?;
let mut stdin = child.stdin.take().expect("piped");
let input = input.to_vec();
let writer = tokio::spawn(async move {
let _ = stdin.write_all(&input).await;
drop(stdin);
});
let out = tokio::time::timeout(self.timeout, child.wait_with_output())
.await
.map_err(|_| anyhow!("git {} timed out", args[0]))?
.context("waiting for git")?;
let _ = writer.await;
if !out.status.success() {
let err = String::from_utf8_lossy(&out.stderr);
bail!(
"git {}: {}",
args[0],
err.lines().last().unwrap_or("failed")
);
}
Ok(out.stdout)
}
pub(crate) async fn branch_head(
&self,
dir: &Path,
token: Option<&str>,
branch: &str,
) -> Result<Option<String>> {
crate::resign::check_branch_name(branch)?;
let refname = format!("refs/heads/{branch}");
let out = self
.git(dir, &["ls-remote", "origin", &refname], token, false)
.await?;
let out = String::from_utf8_lossy(&out);
for line in out.lines() {
if let Some((sha, name)) = line.split_once('\t')
&& name.trim() == refname
{
check_sha(sha).map_err(|e| anyhow!("{e}"))?;
return Ok(Some(sha.to_string()));
}
}
Ok(None)
}
pub(crate) async fn push(
&self,
dir: &Path,
token: Option<&str>,
new_head: &str,
branch: &str,
lease: &str,
) -> Result<()> {
check_sha(new_head).map_err(|e| anyhow!("{e}"))?;
check_sha(lease).map_err(|e| anyhow!("{e}"))?;
crate::resign::check_branch_name(branch)?;
let lease = format!("--force-with-lease=refs/heads/{branch}:{lease}");
let refspec = format!("{new_head}:refs/heads/{branch}");
self.git(
dir,
&[
"push",
"--quiet",
"--no-verify",
"--no-thin",
"--no-recurse-submodules",
&lease,
"origin",
&refspec,
],
token,
false,
)
.await
.map(|_| ())
}
pub async fn fetch(
&self,
remote: &Url,
token: Option<&str>,
head: &str,
commits: &[String],
) -> Result<Fetched> {
self.fetch_with(remote, token, head, commits, &["tree:0", "blob:none"], 0)
.await
}
pub(crate) async fn fetch_for_rewrite(
&self,
remote: &Url,
token: Option<&str>,
head: &str,
commits: &[String],
) -> Result<Fetched> {
self.fetch_with(remote, token, head, commits, &["blob:none"], 1)
.await
}
async fn fetch_with(
&self,
remote: &Url,
token: Option<&str>,
head: &str,
commits: &[String],
filters: &[&str],
extra_depth: usize,
) -> Result<Fetched> {
check_sha(head).map_err(|e| anyhow!("{e}"))?;
for c in commits {
check_sha(c).map_err(|e| anyhow!("{e}"))?;
}
let remote = self.remote_override.as_ref().unwrap_or(remote);
match remote.scheme() {
"https" => {}
"file" if self.allow_file => {}
s => bail!("refusing to fetch over `{s}`"),
}
let dir = tempfile::tempdir()?;
let d = dir.path();
self.git(d, &["init", "--quiet", "--bare"], None, false)
.await?;
for (k, v) in [
("core.repositoryformatversion", "1"),
("extensions.partialClone", "origin"),
("remote.origin.url", remote.as_str()),
("remote.origin.promisor", "true"),
] {
self.git(d, &["config", k, v], None, false).await?;
}
let depth = format!("--depth={}", commits.len().max(1) + extra_depth);
let mut last_err = None;
for &filter in filters {
self.git(
d,
&["config", "remote.origin.partialclonefilter", filter],
None,
false,
)
.await?;
let f = format!("--filter={filter}");
match self
.git(
d,
&[
"fetch",
"--quiet",
"--no-tags",
"--no-write-fetch-head",
"--no-recurse-submodules",
&f,
&depth,
"origin",
head,
],
token,
true,
)
.await
{
Ok(_) => {
last_err = None;
break;
}
Err(e) if e.to_string().contains("bytes") => return Err(e),
Err(e) => last_err = Some(e),
}
}
if let Some(e) = last_err {
return Err(e);
}
let mut out = Vec::with_capacity(commits.len());
for sha in commits {
let raw = self
.git(d, &["cat-file", "commit", sha], None, false)
.await
.with_context(|| format!("commit {sha} did not arrive"))?;
out.push(RangeCommit {
sha: sha.clone(),
raw,
});
}
Ok(Fetched { commits: out, dir })
}
pub async fn recompute_merge(
&self,
dir: &Path,
token: Option<&str>,
ours: &str,
theirs: &str,
merge_base: &str,
) -> MergeOutcome {
match self
.try_recompute_merge(dir, token, ours, theirs, merge_base)
.await
{
Ok(outcome) => outcome,
Err(e) => MergeOutcome::Unavailable(format!("{e:#}")),
}
}
async fn try_recompute_merge(
&self,
dir: &Path,
token: Option<&str>,
ours: &str,
theirs: &str,
merge_base: &str,
) -> Result<MergeOutcome> {
for sha in [ours, theirs, merge_base] {
check_sha(sha).map_err(|e| anyhow!("{e}"))?;
}
self.git(
dir,
&["config", "remote.origin.partialclonefilter", "blob:none"],
None,
false,
)
.await?;
self.git(
dir,
&[
"fetch",
"--quiet",
"--no-tags",
"--no-write-fetch-head",
"--no-recurse-submodules",
"--filter=blob:none",
"--depth=1",
"origin",
ours,
theirs,
merge_base,
],
token,
true,
)
.await
.context("fetching the merge's parents")?;
let args = verify_trust::merge_tree_args(ours, theirs, Some(merge_base));
let args: Vec<&str> = args.iter().map(String::as_str).collect();
let (status, stdout, stderr) = self.run(dir, &args, token, true, true).await?;
Ok(verify_trust::merge_outcome(status.code(), &stdout, &stderr))
}
}
fn dir_size(dir: &Path) -> u64 {
let mut total = 0;
let mut stack = vec![dir.to_path_buf()];
while let Some(p) = stack.pop() {
let Ok(rd) = std::fs::read_dir(&p) else {
continue;
};
for e in rd.flatten() {
match e.metadata() {
Ok(m) if m.is_dir() => stack.push(e.path()),
Ok(m) => total += m.len(),
Err(_) => {}
}
}
}
total
}
pub struct CheckRunner {
pub(crate) verifier: Arc<dyn CommitVerifier>,
pub(crate) fetcher: GitFetcher,
max_commits: usize,
permits: Semaphore,
in_flight: Mutex<BTreeSet<(String, String, String)>>,
reports: Mutex<BTreeMap<String, ReportSlot>>,
}
#[derive(Debug, Default)]
struct ReportSlot {
last: Option<Instant>,
pending: bool,
}
const REPORT_INTERVAL: Duration = Duration::from_secs(60);
impl CheckRunner {
pub fn new(cfg: &CheckConfig, verifier: Arc<dyn CommitVerifier>, fetcher: GitFetcher) -> Self {
CheckRunner {
verifier,
fetcher,
max_commits: cfg.max_commits,
permits: Semaphore::new(cfg.concurrency.max(1)),
in_flight: Mutex::new(BTreeSet::new()),
reports: Mutex::new(BTreeMap::new()),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum CheckOutcome {
Posted(CheckConclusion),
Skipped(&'static str),
}
pub(crate) fn spawn(
bridge: &Arc<Bridge>,
triggers: Vec<CheckTrigger>,
on_done: impl FnOnce(&Bridge) + Send + 'static,
) {
let bridge = Arc::clone(bridge);
tokio::spawn(async move {
let mut all_ok = true;
for t in &triggers {
match run(&bridge, t).await {
Ok(o) => {
tracing::info!(repo = %t.repo, head = %t.head_sha, outcome = ?o, "check");
if matches!(o, CheckOutcome::Posted(_)) {
report_check(&bridge, t);
}
}
Err(e) => {
all_ok = false;
tracing::warn!(repo = %t.repo, head = %t.head_sha, error = %e, "check not posted");
}
}
}
if all_ok {
on_done(&bridge);
}
});
}
fn check_namespace(bridge: &Bridge, repo: &Resource) -> Option<Ctx> {
let ns = bridge
.store
.list::<NamespaceRecord>(Table::Namespaces)
.ok()?
.into_iter()
.map(|(_, n)| n)
.find(|n| n.state == NamespaceState::Bound && n.resource.contains(repo))?;
let ctx = Ctx::load(bridge, &ns.id).ok()?;
ctx.adapter
.forge()
.capabilities(&ctx.namespace)
.bridge_posted_check
.then_some(ctx)
}
struct InFlight<'a> {
set: &'a Mutex<BTreeSet<(String, String, String)>>,
key: (String, String, String),
}
impl Drop for InFlight<'_> {
fn drop(&mut self) {
self.set.lock().expect("lock").remove(&self.key);
}
}
pub async fn run(bridge: &Bridge, trigger: &CheckTrigger) -> Result<CheckOutcome> {
let Some(ctx) = check_namespace(bridge, &trigger.repo) else {
return Ok(CheckOutcome::Skipped("not a bridge-check namespace"));
};
let g = ctx
.adapter
.github()
.context("bridge checks are GitHub's")?
.clone();
let (base_ref, base_sha, pr_info) = match trigger.kind {
CheckTriggerKind::PullRequest { number } | CheckTriggerKind::Rerequested { number } => {
let pr = g.pull_request(&trigger.repo, number).await?;
if !pr.open {
return Ok(CheckOutcome::Skipped("the pull request is closed"));
}
if pr.head_sha != trigger.head_sha {
return Ok(CheckOutcome::Skipped("the pull request has moved on"));
}
(pr.base_ref.clone(), pr.base_sha.clone(), Some(pr))
}
_ => (trigger.base_ref.clone(), trigger.base_sha.clone(), None),
};
check_sha(&base_sha).map_err(|e| anyhow!("{e}"))?;
let protected = g.default_branch(&trigger.repo).await?;
if base_ref != protected {
return Ok(CheckOutcome::Skipped("the base branch is not protected"));
}
let key = (
trigger.repo.to_string(),
trigger.head_sha.clone(),
base_ref.clone(),
);
if !bridge
.checks
.in_flight
.lock()
.expect("lock")
.insert(key.clone())
{
return Ok(CheckOutcome::Skipped("the same check is already running"));
}
let _in_flight = InFlight {
set: &bridge.checks.in_flight,
key,
};
let check_name = bridge
.adapters
.vgi_for(&ctx.ns.resource)
.map(|v| v.required_check)
.unwrap_or_else(|| vgi_forge::DEFAULT_REQUIRED_CHECK.into());
let _permit = bridge.checks.permits.acquire().await?;
let external = format!("{base_ref}@{}", trigger.head_sha);
let id = g
.start_check_run(&trigger.repo, &trigger.head_sha, &check_name, &external)
.await?;
let (conclusion, title, mut summary) =
match verify(bridge, &g, &ctx, trigger, &base_ref, &base_sha).await {
Ok(v) => v,
Err(e) => (
CheckConclusion::Failure,
"The commits could not be verified".to_string(),
format!(
"The bridge could not complete the check, so it fails closed.\n\n`{}`",
e.to_string().replace('`', "'")
),
),
};
if conclusion == CheckConclusion::Failure
&& let Some(pr) = &pr_info
&& let Some(hint) = crate::resign::check_hint(bridge, &ctx, trigger, pr, &protected)
{
summary.push_str(&hint);
}
g.finish_check_run(&trigger.repo, id, conclusion, &title, &summary)
.await?;
let conclusion_word = match conclusion {
CheckConclusion::Success => "success",
_ => "failure",
};
let _ = bridge.store.update::<RepoRecord, _>(
Table::Repos,
&repo_key(trigger.repo.host(), trigger.repo_id),
|r| {
Ok((
r.map(|mut r| {
r.last_check = Some(LastCheck::new(
trigger.head_sha.clone(),
conclusion_word,
crate::bridge::now(),
));
r
}),
(),
))
},
);
Ok(CheckOutcome::Posted(conclusion))
}
fn report_check(bridge: &Arc<Bridge>, t: &CheckTrigger) {
let key = repo_key(t.repo.host(), t.repo_id);
if !matches!(
bridge.store.get::<RepoRecord>(Table::Repos, &key),
Ok(Some(_))
) {
return;
}
let delay = {
let mut slots = bridge.checks.reports.lock().expect("lock");
let slot = slots.entry(key.clone()).or_default();
if slot.pending {
return;
}
slot.pending = true;
slot.last
.map(|l| (l + REPORT_INTERVAL).saturating_duration_since(Instant::now()))
.unwrap_or_default()
};
let weak = Arc::downgrade(bridge);
tokio::spawn(async move {
tokio::time::sleep(delay).await;
let Some(bridge) = weak.upgrade() else {
return;
};
{
let mut slots = bridge.checks.reports.lock().expect("lock");
let slot = slots.entry(key.clone()).or_default();
slot.pending = false;
slot.last = Some(Instant::now());
}
let Ok(Some(rec)) = bridge.store.get::<RepoRecord>(Table::Repos, &key) else {
return;
};
let Ok(ctx) = Ctx::load(&bridge, &rec.namespace) else {
return;
};
let lock = bridge.ns_lock(&rec.namespace);
let _g = lock.lock().await;
if let Err(e) = crate::jobs::inspect_repo(&bridge, &ctx, &rec.resource, true, None).await {
tracing::warn!(repo = %rec.resource, error = %e, "could not report a posted check");
}
});
}
async fn verify(
bridge: &Bridge,
g: &vgi_forge_github::GitHubForge,
ctx: &Ctx,
trigger: &CheckTrigger,
base_ref: &str,
base_sha: &str,
) -> Result<(CheckConclusion, String, String)> {
let cmp = g
.compare_commits(&trigger.repo, base_sha, &trigger.head_sha)
.await?;
let max = bridge.checks.max_commits;
if cmp.total as usize > max {
bail!(
"{} commits are more than this bridge checks at once ({max}); split the change",
cmp.total
);
}
if (cmp.commits.len() as u64) < cmp.total {
bail!(
"GitHub listed {} of the {} commits; the range cannot be checked from a partial list",
cmp.commits.len(),
cmp.total
);
}
if cmp.commits.is_empty() {
return Ok(if trigger.head_sha == base_sha {
(
CheckConclusion::Success,
format!("The head is the tip of `{base_ref}`; nothing to verify"),
"Every commit here is already part of the protected branch.".into(),
)
} else {
(
CheckConclusion::Failure,
format!("The head is already part of `{base_ref}`"),
"There is nothing to merge: the head commit is already contained in the \
protected branch."
.into(),
)
});
}
let token = g.contents_read_token(&trigger.repo).await?;
let remote = g.clone_url(&trigger.repo)?;
let fetched = bridge
.checks
.fetcher
.fetch(
&remote,
Some(token.expose()),
&trigger.head_sha,
&cmp.commits,
)
.await?;
let facts = merge_facts(
&bridge.checks.fetcher,
fetched.dir.path(),
Some(token.expose()),
&fetched.commits,
bridge
.checks
.verifier
.exempts_platform_merges(trigger.repo.host()),
|ours, theirs| async move { Ok(g.merge_base(&trigger.repo, &ours, &theirs).await?) },
)
.await;
drop(token);
let resource = trigger.repo.as_str();
let fallback = ctx.ns.resource.as_str();
let lines = bridge
.checks
.verifier
.verify_with_facts(&fetched.commits, &facts, resource, fallback)
.await?;
Ok(summarise(&lines, base_ref))
}
fn summarise(lines: &[CommitLine], base_ref: &str) -> (CheckConclusion, String, String) {
let failed = lines.iter().filter(|l| !l.passes).count();
let conclusion = if failed == 0 && !lines.is_empty() {
CheckConclusion::Success
} else {
CheckConclusion::Failure
};
let title = if lines.is_empty() {
"No commits were verified".to_string()
} else if failed == 0 {
format!("All {} commits are signed by trusted DIDs", lines.len())
} else {
format!("{failed} of {} commits are not trusted", lines.len())
};
let mut summary = format!(
"Checked by the community's bridge against its Trust Registry, for merging into \
`{base_ref}`.\n\n| Commit | Verdict |\n|---|---|\n"
);
for l in lines {
summary.push_str(&format!(
"| `{}` | {} {} |\n",
&l.sha[..l.sha.len().min(12)],
if l.passes { "✅" } else { "❌" },
l.verdict.replace('|', "/")
));
}
(conclusion, title, summary)
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn the_registry_is_queried_over_tsp_then_didcomm_then_https() {
use crate::registry_channel::{BridgeRegistryChannel, RegistryReplies};
use crate::transport::{Via, memory::ChannelLink};
use trql_client::TransportKind;
let cfg = crate::config::BridgeConfig::parse(crate::config::tests::EXAMPLE).unwrap();
let replies = Arc::new(RegistryReplies::new(cfg.trust_registry_did.clone()));
let (link, _sent) = ChannelLink::new();
let link: Arc<dyn crate::transport::VtcLink> = Arc::new(link);
let channel = |via| {
Arc::new(BridgeRegistryChannel::new(
Arc::clone(&link),
"did:key:z6MkBridge",
Arc::clone(&replies),
via,
))
};
let v = VerifyTrustVerifier::new(&cfg)
.with_channel(channel(Via::Tsp))
.with_channel(channel(Via::Didcomm));
let ours = v.supported();
assert_eq!(
ours,
[
TransportKind::Tsp,
TransportKind::Didcomm,
TransportKind::Https
]
);
let caps = |services: serde_json::Value| {
trql_client::ServiceCapabilities::from_document(&serde_json::json!({
"id": cfg.trust_registry_did, "service": services
}))
};
let both = caps(serde_json::json!([
{ "id": "#didcomm", "type": "DIDCommMessaging", "serviceEndpoint": "did:web:m.example" },
{ "id": "#tsp", "type": "TSPTransport", "serviceEndpoint": "did:web:m.example" },
{ "id": "#rest", "type": "TRQPRest", "serviceEndpoint": "https://r.example" }
]));
let pick = |caps: &trql_client::ServiceCapabilities, t| {
verify_trust::registry::choose_route(caps, bridge_transport(t), &ours)
};
assert_eq!(
pick(&both, vgi_forge::VerifyTransport::Auto).unwrap().kind,
TransportKind::Tsp
);
assert_eq!(
pick(&both, vgi_forge::VerifyTransport::Tsp).unwrap().kind,
TransportKind::Tsp
);
assert_eq!(
pick(&both, vgi_forge::VerifyTransport::Didcomm)
.unwrap()
.kind,
TransportKind::Didcomm
);
let didcomm_only = caps(serde_json::json!([
{ "id": "#didcomm", "type": "DIDCommMessaging", "serviceEndpoint": "did:web:m.example" }
]));
assert!(pick(&didcomm_only, vgi_forge::VerifyTransport::Tsp).is_err());
assert_eq!(
pick(&didcomm_only, vgi_forge::VerifyTransport::Auto)
.unwrap()
.kind,
TransportKind::Didcomm
);
}
fn commit(sha: &str, parents: &[&str], armor: Option<&str>) -> RangeCommit {
let mut raw = String::from("tree 4b825dc642cb6eb9a060e54bf8d69288fbee4904\n");
for p in parents {
raw.push_str(&format!("parent {p}\n"));
}
raw.push_str("committer A <a@example.com> 1700000000 +0000\n");
if let Some(armor) = armor {
raw.push_str(&format!(
"gpgsig -----BEGIN {armor} SIGNATURE-----\n AAAA\n -----END {armor} SIGNATURE-----\n"
));
}
raw.push_str("\nmsg\n");
RangeCommit {
sha: sha.to_string(),
raw: raw.into_bytes(),
}
}
fn sha(n: u8) -> String {
format!("{n:040x}")
}
#[test]
fn the_boundary_is_the_parents_outside_the_listing() {
let (base, a, b, m) = (sha(1), sha(2), sha(3), sha(4));
let commits = [
commit(&a, &[&base], None),
commit(&m, &[&a, &b], Some("PGP")),
];
assert_eq!(boundary_of(&commits), [base, b].into());
}
#[test]
fn only_pgp_signed_two_parent_merges_are_candidates() {
let (x, y, z) = (sha(1), sha(2), sha(3));
let commits = [
commit(&sha(10), &[&x, &y], Some("PGP")),
commit(&sha(11), &[&x, &y], Some("SSH")),
commit(&sha(12), &[&x, &y], None),
commit(&sha(13), &[&x], Some("PGP")),
commit(&sha(14), &[&x, &y, &z], Some("PGP")),
];
assert_eq!(platform_merge_candidates(&commits), vec![(sha(10), x, y)]);
}
#[tokio::test]
async fn merges_past_the_cap_or_without_a_base_fail_closed() {
let commits: Vec<RangeCommit> = (0..=MAX_PLATFORM_MERGES as u8)
.map(|n| commit(&sha(100 + n), &[&sha(1), &sha(2)], Some("PGP")))
.collect();
let fetcher = GitFetcher::new(&CheckConfig::default());
let dir = tempfile::tempdir().unwrap();
let facts = merge_facts(&fetcher, dir.path(), None, &commits, true, |_, _| async {
Err(anyhow!("no route"))
})
.await;
assert_eq!(facts.merges.len(), commits.len());
for (n, c) in commits.iter().enumerate() {
let MergeOutcome::Unavailable(why) = &facts.merges[&c.sha] else {
panic!("{:?}", facts.merges[&c.sha]);
};
if n < MAX_PLATFORM_MERGES {
assert!(why.contains("no merge base"), "{why}");
} else {
assert!(why.contains("more than"), "{why}");
}
}
let facts = merge_facts(&fetcher, dir.path(), None, &commits, false, |_, _| async {
Err(anyhow!("must not be asked"))
})
.await;
assert!(facts.merges.is_empty());
}
#[test]
fn the_summary_fails_on_any_untrusted_commit_and_on_none() {
let (c, t, s) = summarise(
&[
CommitLine::new("a".repeat(40), true, "trusted"),
CommitLine::new("b".repeat(40), false, "unauthorized"),
],
"main",
);
assert_eq!(c, CheckConclusion::Failure);
assert_eq!(t, "1 of 2 commits are not trusted");
assert!(s.contains("unauthorized") && s.contains("`main`"));
let (c, _, _) = summarise(&[CommitLine::new("a".repeat(40), true, "trusted")], "main");
assert_eq!(c, CheckConclusion::Success);
let (c, _, _) = summarise(&[], "main");
assert_eq!(
c,
CheckConclusion::Failure,
"a verifier that saw nothing passes nothing"
);
}
#[tokio::test]
async fn the_fetcher_refuses_anything_but_https_and_commit_ids() {
let f = GitFetcher::new(&CheckConfig::default());
let head = "a".repeat(40);
for bad in [
"file:///tmp/x",
"ext::sh -c touch% /tmp/pwned",
"http://example.org/r",
] {
let Ok(u) = Url::parse(bad) else { continue };
assert!(f.fetch(&u, None, &head, &[]).await.is_err(), "{bad}");
}
let u = Url::parse("https://example.org/r.git").unwrap();
assert!(f.fetch(&u, None, "--upload-pack=x", &[]).await.is_err());
}
}