use std::sync::Arc;
use http::{HeaderMap, StatusCode};
use serde_json::json;
use vgi_forge::{ForgeEventKind, HookDecision, InstallationChange, Resource};
use crate::bridge::{Bridge, now};
use crate::jobs::{self, Ctx};
use crate::mapping::Report;
use crate::store::{BranchLedger, NamespaceRecord, NamespaceState, RepoRecord, Table, repo_key};
const NAMESPACE_REPOS: [&str; 1] = [".vgi"];
fn webhook_adapter(
bridge: &Bridge,
host: &str,
owner: Option<&str>,
) -> Option<crate::registry::Adapter> {
match owner {
Some(o) => bridge.adapters.get_github(host, o),
None => bridge.adapters.get(host),
}
}
#[derive(Debug, Clone)]
pub(crate) struct Scope {
host: String,
owner: Option<String>,
}
impl Scope {
pub(crate) fn covers(&self, r: &Resource) -> bool {
r.host() == self.host
&& self
.owner
.as_ref()
.is_none_or(|o| r.owner().eq_ignore_ascii_case(o))
}
pub(crate) fn admits_repo(
&self,
bridge: &Bridge,
repo: &Resource,
id: u64,
need_record: bool,
) -> bool {
if !self.covers(repo) {
return false;
}
match repo_record(bridge, &self.host, id) {
Some(rec) => {
self.covers(&rec.resource)
&& namespace_record(bridge, &rec.namespace)
.is_some_and(|n| self.covers(&n.resource))
}
None => !need_record,
}
}
fn admits(&self, kind: &ForgeEventKind) -> bool {
match kind {
ForgeEventKind::RepoCreated { repo, .. }
| ForgeEventKind::RepoDeleted { repo, .. }
| ForgeEventKind::RepoArchived { repo, .. }
| ForgeEventKind::RepoVisibilityChanged { repo, .. }
| ForgeEventKind::CollaboratorChanged { repo, .. }
| ForgeEventKind::PullRequestOpened { repo, .. } => self.covers(repo),
ForgeEventKind::RepoRenamed { from, to, .. } => self.covers(from) && self.covers(to),
ForgeEventKind::RepoTransferred {
from_namespace, to, ..
} => self.covers(to) || from_namespace.as_ref().is_some_and(|f| self.covers(f)),
ForgeEventKind::OrgMembershipChanged { namespace, .. }
| ForgeEventKind::TeamMembershipChanged { namespace, .. }
| ForgeEventKind::InstallationChanged { namespace, .. } => self.covers(namespace),
ForgeEventKind::ProtectionChanged {
repo, namespace, ..
} => self.covers(namespace) && repo.as_ref().is_none_or(|r| self.covers(r)),
_ => false,
}
}
}
fn repo_record(bridge: &Bridge, host: &str, id: u64) -> Option<RepoRecord> {
bridge
.store
.get::<RepoRecord>(Table::Repos, &repo_key(host, id))
.ok()
.flatten()
}
fn namespace_record(bridge: &Bridge, id: &str) -> Option<NamespaceRecord> {
bridge
.store
.get::<NamespaceRecord>(Table::Namespaces, id)
.ok()
.flatten()
}
pub(crate) async fn on_webhook(
bridge: &Arc<Bridge>,
host: &str,
owner: Option<&str>,
headers: &HeaderMap,
body: &[u8],
) -> StatusCode {
let Some(adapter) = webhook_adapter(bridge, host, owner) else {
return StatusCode::NOT_FOUND;
};
let scope = Scope {
host: host.to_string(),
owner: owner.map(str::to_ascii_lowercase),
};
#[cfg(feature = "forge-github")]
if let Some(g) = adapter.github() {
match g.parse_push(headers, body) {
Ok(Some(push)) => {
if !scope.admits_repo(bridge, &push.repo, push.repo_id, true) {
tracing::warn!(%host, owner = ?scope.owner, repo = %push.repo, id = push.repo_id,
"dropping a push for a repository outside the App's organisation");
return StatusCode::NO_CONTENT;
}
return crate::resign::on_push(bridge, host, push);
}
Ok(None) => {}
Err(e) => {
tracing::warn!(%host, error = %e, "refused a webhook");
return StatusCode::UNAUTHORIZED;
}
}
if let Ok(Some(ev)) = adapter.forge().parse_event(headers, body)
&& matches!(ev.kind, ForgeEventKind::PullRequestOpened { .. })
{
pull_request_opened(bridge, &scope, ev.delivery_id.as_deref(), ev.kind);
}
match g.parse_check_trigger(headers, body) {
Ok(mut triggers) if !triggers.is_empty() => {
let before = triggers.len();
triggers.retain(|t| scope.admits_repo(bridge, &t.repo, t.repo_id, false));
if triggers.len() != before {
tracing::warn!(%host, owner = ?scope.owner,
"dropping check triggers for repositories outside the App's organisation");
}
if triggers.is_empty() {
return StatusCode::NO_CONTENT;
}
let key = triggers[0]
.delivery_id
.as_deref()
.map(|d| format!("{host}#{d}"));
if let Some(k) = &key
&& already_seen(bridge, k)
{
return StatusCode::OK;
}
crate::resign::on_pull_request_triggers(bridge, &triggers);
crate::checks::spawn(bridge, triggers, move |b| {
if let Some(k) = key {
let _ = b.store.put(Table::Deliveries, &k, &now());
}
});
return StatusCode::ACCEPTED;
}
Ok(_) => {}
Err(e) => {
tracing::warn!(%host, error = %e, "refused a webhook");
return StatusCode::UNAUTHORIZED;
}
}
}
let event = match adapter.forge().parse_event(headers, body) {
Ok(Some(e)) => e,
Ok(None) => return StatusCode::NO_CONTENT,
Err(e) => {
tracing::warn!(%host, error = %e, "refused a webhook");
return StatusCode::UNAUTHORIZED;
}
};
if !scope.admits(&event.kind) {
tracing::warn!(%host, owner = ?scope.owner, event = ?event.kind,
"dropping an event that names a namespace or repository outside the App's organisation");
return StatusCode::NO_CONTENT;
}
if let ForgeEventKind::PullRequestOpened { .. } = &event.kind {
pull_request_opened(bridge, &scope, event.delivery_id.as_deref(), event.kind);
return StatusCode::ACCEPTED;
}
if seen(bridge, host, event.delivery_id.as_deref()) {
return StatusCode::OK;
}
let event = match adapter.hooks().on_event(&event) {
HookDecision::Modify(e) => e,
HookDecision::Abort(why) => {
tracing::info!(%why, "an adapter hook dropped an event");
return StatusCode::NO_CONTENT;
}
_ => event,
};
let me = Arc::clone(bridge);
tokio::spawn(async move { handle(&me, &scope, event.kind).await });
StatusCode::ACCEPTED
}
fn already_seen(bridge: &Bridge, key: &str) -> bool {
matches!(bridge.store.get::<i64>(Table::Deliveries, key), Ok(Some(_)))
}
fn seen(bridge: &Bridge, host: &str, delivery: Option<&str>) -> bool {
let Some(d) = delivery else { return false };
!bridge
.store
.put_new(Table::Deliveries, &format!("{host}#{d}"), &now())
.unwrap_or(true)
}
fn namespace_in(bridge: &Bridge, scope: &Scope, r: &Resource) -> Option<NamespaceRecord> {
if !scope.covers(r) {
return None;
}
bridge
.store
.list::<NamespaceRecord>(Table::Namespaces)
.ok()?
.into_iter()
.map(|(_, n)| n)
.find(|n| {
n.state == NamespaceState::Bound && n.resource.contains(r) && scope.covers(&n.resource)
})
}
fn repo_in(bridge: &Bridge, ns: &NamespaceRecord, id: u64) -> Option<RepoRecord> {
repo_record(bridge, ns.resource.host(), id).filter(|r| r.namespace == ns.id)
}
fn claimed_elsewhere(bridge: &Bridge, ns: &NamespaceRecord, id: u64) -> bool {
let host = ns.resource.host();
if repo_record(bridge, host, id).is_some_and(|r| r.namespace != ns.id) {
return true;
}
bridge
.store
.list::<NamespaceRecord>(Table::Namespaces)
.unwrap_or_default()
.into_iter()
.any(|(_, n)| n.id != ns.id && n.resource.host() == host && n.managed.contains(&id))
}
fn detach_from(bridge: &Bridge, ns: &NamespaceRecord, id: u64) -> bool {
if claimed_elsewhere(bridge, ns, id) {
tracing::warn!(namespace = %ns.id, forge_id = id,
"not detaching a repository another namespace governs");
return false;
}
detach(bridge, ns.resource.host(), id);
true
}
async fn confirm_transfer_in(bridge: &Bridge, to_ns: &NamespaceRecord, forge_id: u64) {
let Some(rec) = repo_record(bridge, to_ns.resource.host(), forge_id) else {
return;
};
if rec.namespace == to_ns.id {
return;
}
let Some(old_ns) = namespace_record(bridge, &rec.namespace) else {
return;
};
#[cfg(feature = "forge-github")]
{
let Some(g) = bridge
.adapters
.for_resource(&to_ns.resource)
.and_then(|a| a.github().cloned())
else {
return;
};
match g.repository_by_id(&to_ns.resource, forge_id).await {
Ok(Some(now)) if to_ns.resource.contains(&now) && !old_ns.resource.contains(&now) => {
tracing::info!(from = %rec.resource, to = %now,
"GitHub confirms a repository moved to another organisation; detaching it there");
transferred_out(bridge, &old_ns, &rec.resource, &now, forge_id).await;
}
Ok(other) => tracing::warn!(forge_id, now = ?other,
"a transfer in that GitHub does not confirm; the other organisation's record is left alone"),
Err(e) => tracing::warn!(forge_id, error = %e,
"could not confirm a transfer with GitHub; the other organisation's record is left alone"),
}
}
#[cfg(not(feature = "forge-github"))]
let _ = old_ns;
}
pub(crate) fn detach(bridge: &Bridge, host: &str, forge_id: u64) {
let _ = bridge.store.delete(Table::Repos, &repo_key(host, forge_id));
let namespaces = bridge
.store
.list::<NamespaceRecord>(Table::Namespaces)
.unwrap_or_default();
for (id, ns) in namespaces {
if ns.resource.host() != host || !ns.managed.contains(&forge_id) {
continue;
}
let managed = bridge
.store
.update::<NamespaceRecord, _>(Table::Namespaces, &id, |n| {
let Some(mut n) = n else {
return Ok((None, None));
};
n.managed.remove(&forge_id);
let m = n.managed.clone();
Ok((Some(n), Some(m)))
});
#[cfg(feature = "forge-github")]
if let (Ok(Some(m)), Some(g)) = (
managed,
bridge
.adapters
.for_resource(&ns.resource)
.and_then(|a| a.github().cloned()),
) {
g.set_managed_repositories(&ns.resource, m);
}
#[cfg(not(feature = "forge-github"))]
let _ = managed;
}
let prefix = format!("{host}#{forge_id}#");
if let Ok(ledgers) = bridge.store.list::<BranchLedger>(Table::Branches) {
for (key, _) in ledgers.into_iter().filter(|(k, _)| k.starts_with(&prefix)) {
let _ = bridge.store.delete(Table::Branches, &key);
}
}
}
pub(crate) fn detach_reused_name(
bridge: &Bridge,
ns_id: &str,
resource: &Resource,
forge_id: u64,
) -> bool {
let stale: Vec<RepoRecord> = bridge
.store
.list::<RepoRecord>(Table::Repos)
.unwrap_or_default()
.into_iter()
.map(|(_, r)| r)
.filter(|r| r.namespace == ns_id && r.resource == *resource && r.forge_id != forge_id)
.collect();
for r in &stale {
tracing::warn!(
%resource, old = r.forge_id, new = forge_id,
"a new repository took a governed name; the old one is no longer managed"
);
detach(bridge, resource.host(), r.forge_id);
}
!stale.is_empty()
}
pub(crate) async fn report_unmanaged(
bridge: &Bridge,
ns: &NamespaceRecord,
resource: &Resource,
forge_id: u64,
) {
if NAMESPACE_REPOS.contains(&resource.repo_name().unwrap_or_default())
|| repo_in(bridge, ns, forge_id).is_some()
{
return;
}
detach_reused_name(bridge, &ns.id, resource, forge_id);
let ev = json!({ "type": "repoCreatedUnmanaged", "forgeId": forge_id.to_string(), "resource": resource.as_str() });
report(bridge, &ns.id, ev).await;
}
pub(crate) async fn transferred_out(
bridge: &Bridge,
from_ns: &NamespaceRecord,
from: &Resource,
to: &Resource,
forge_id: u64,
) {
if !detach_from(bridge, from_ns, forge_id) {
return;
}
if NAMESPACE_REPOS.contains(&from.repo_name().unwrap_or_default()) {
return;
}
let ev = json!({ "type": "repoTransferred", "forgeId": forge_id.to_string(), "from": from.as_str(), "to": to.as_str() });
report(bridge, &from_ns.id, ev).await;
let same_owner = Scope {
host: from_ns.resource.host().to_string(),
owner: Some(from_ns.resource.owner().to_ascii_lowercase()),
};
if let Some(to_ns) = namespace_in(bridge, &same_owner, to)
&& to_ns.id != from_ns.id
{
report_unmanaged(bridge, &to_ns, to, forge_id).await;
}
}
async fn inspect_in(bridge: &Bridge, ns: &NamespaceRecord, repo: &Resource) {
let Ok(ctx) = Ctx::load(bridge, &ns.id) else {
return;
};
let lock = bridge.ns_lock(&ns.id);
let _g = lock.lock().await;
if let Err(e) = jobs::inspect_repo(bridge, &ctx, repo, true, None).await {
tracing::warn!(%repo, error = %e, "could not inspect after a webhook");
}
}
pub(crate) fn pull_request_opened(
bridge: &Arc<Bridge>,
scope: &Scope,
delivery_id: Option<&str>,
kind: ForgeEventKind,
) -> bool {
let ForgeEventKind::PullRequestOpened {
repo,
forge_id,
number,
reopened,
author,
actor,
draft,
from_fork,
} = kind
else {
return false;
};
if !bridge.cfg.event_version.reports_pull_requests() {
tracing::debug!(%repo, number,
"not reporting a pull request: event_version is below 0.4 (no pull-request gate)");
return false;
}
if !scope.admits_repo(bridge, &repo, forge_id, true) {
tracing::debug!(%repo, forge_id, "not reporting a pull request on an unmanaged repository");
return false;
}
let Some(ns) = namespace_in(bridge, scope, &repo) else {
return false;
};
if NAMESPACE_REPOS.contains(&repo.repo_name().unwrap_or_default())
|| !ns.managed.contains(&forge_id)
|| repo_in(bridge, &ns, forge_id).is_none()
{
return false;
}
if let Some(d) = delivery_id {
let key = format!("{}#{d}#pullRequestOpened", repo.host());
if !bridge
.store
.put_new(Table::Deliveries, &key, &now())
.unwrap_or(false)
{
return false;
}
}
let host = repo.host();
let mut ev = json!({
"type": "pullRequestOpened",
"forgeId": forge_id.to_string(),
"resource": repo.as_str(),
"number": number,
"action": if reopened { "reopened" } else { "opened" },
"author": crate::mapping::wire_account(host, &author),
"actor": crate::mapping::wire_account(host, &actor),
});
if let Some(d) = draft {
ev["draft"] = json!(d);
}
if let Some(f) = from_fork {
ev["fromFork"] = json!(f);
}
let me = Arc::clone(bridge);
tokio::spawn(async move { report(&me, &ns.id, ev).await });
true
}
async fn report(bridge: &Bridge, ns: &str, ev: serde_json::Value) {
if let Err(e) = bridge.send_event(ns, ev, None).await {
tracing::error!(error = %e, "could not report an event");
}
}
async fn handle(bridge: &Arc<Bridge>, scope: &Scope, kind: ForgeEventKind) {
match kind {
ForgeEventKind::RepoCreated { repo, forge_id } => {
let Some(ns) = namespace_in(bridge, scope, &repo) else {
return;
};
let name = repo.repo_name().unwrap_or_default();
if NAMESPACE_REPOS.contains(&name) || repo_in(bridge, &ns, forge_id).is_some() {
return;
}
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
report_unmanaged(bridge, &ns, &repo, forge_id).await;
}
ForgeEventKind::RepoDeleted { repo, forge_id } => {
let Some(ns) = namespace_in(bridge, scope, &repo) else {
return;
};
if NAMESPACE_REPOS.contains(&repo.repo_name().unwrap_or_default()) {
return;
}
if !detach_from(bridge, &ns, forge_id) {
return;
}
let ev = json!({ "type": "repoDeleted", "forgeId": forge_id.to_string(), "resource": repo.as_str() });
report(bridge, &ns.id, ev).await;
}
ForgeEventKind::RepoRenamed { forge_id, from, to } => {
let Some(ns) = namespace_in(bridge, scope, &to) else {
return;
};
if !ns.resource.contains(&from) {
tracing::warn!(%from, %to, "ignoring a rename that crosses namespaces");
return;
}
if claimed_elsewhere(bridge, &ns, forge_id) {
tracing::warn!(%from, %to, "ignoring a rename of a repository another namespace governs");
return;
}
detach_reused_name(bridge, &ns.id, &to, forge_id);
let ns_id = ns.id.clone();
let _ = bridge.store.update::<RepoRecord, _>(
Table::Repos,
&repo_key(to.host(), forge_id),
|r| {
Ok((
r.map(|mut r| {
if r.namespace == ns_id {
r.resource = to.clone();
}
r
}),
(),
))
},
);
let ev = json!({ "type": "repoRenamed", "forgeId": forge_id.to_string(), "from": from.as_str(), "to": to.as_str() });
report(bridge, &ns.id, ev).await;
}
ForgeEventKind::RepoTransferred {
forge_id,
from_namespace,
to,
} => {
let rec = repo_record(bridge, to.host(), forge_id).filter(|r| {
scope.covers(&r.resource)
&& namespace_record(bridge, &r.namespace)
.is_some_and(|n| scope.covers(&n.resource))
});
let from = rec.as_ref().map(|r| r.resource.clone()).or_else(|| {
from_namespace
.as_ref()
.filter(|n| scope.covers(n))
.and_then(|n| n.join(to.repo_name().unwrap_or_default()).ok())
});
let from_ns = from.as_ref().and_then(|f| namespace_in(bridge, scope, f));
match (from_ns, from) {
(Some(from_ns), Some(from)) if from_ns.resource.contains(&from) => {
transferred_out(bridge, &from_ns, &from, &to, forge_id).await;
}
_ => {
if let Some(t) = namespace_in(bridge, scope, &to) {
confirm_transfer_in(bridge, &t, forge_id).await;
report_unmanaged(bridge, &t, &to, forge_id).await;
}
}
}
}
ForgeEventKind::CollaboratorChanged {
repo,
forge_id,
account,
..
} => {
let Some(ns) = namespace_in(bridge, scope, &repo) else {
return;
};
let Ok(ctx) = Ctx::load(bridge, &ns.id) else {
return;
};
let role = match ctx.adapter.forge().inspect(&repo).await {
Ok(state) => state
.collaborators
.iter()
.find(|c| c.account.id == account.id)
.map(|c| c.role.to_string()),
Err(e) => {
tracing::warn!(%repo, error = %e, "could not read roles after a webhook");
return;
}
};
let mut ev = json!({
"type": "roleChanged",
"forgeId": forge_id.to_string(),
"resource": repo.as_str(),
"account": crate::mapping::wire_account(repo.host(), &account),
});
if let Some(r) = role {
ev["role"] = json!(r);
}
report(bridge, &ns.id, ev).await;
if repo_in(bridge, &ns, forge_id).is_some() {
inspect_in(bridge, &ns, &repo).await;
}
}
ForgeEventKind::ProtectionChanged {
repo: Some(repo), ..
} => {
if let Some(ns) = namespace_in(bridge, scope, &repo) {
inspect_in(bridge, &ns, &repo).await;
}
}
ForgeEventKind::ProtectionChanged {
repo: None,
namespace,
..
} => {
if let Some(ns) = namespace_in(bridge, scope, &namespace)
&& let Ok(ctx) = Ctx::load(bridge, &ns.id)
{
let lock = bridge.ns_lock(&ns.id);
let _g = lock.lock().await;
let mut r = Report::default();
jobs::sweep(bridge, &ctx, false, &mut r).await;
}
}
ForgeEventKind::InstallationChanged {
namespace,
installation_id,
change,
} => {
let Some(ns) = namespace_in(bridge, scope, &namespace) else {
return;
};
let recorded = ns
.binding
.as_ref()
.and_then(|b| b.namespace.installation_id);
if recorded != Some(installation_id) {
tracing::warn!(namespace = %ns.id, installation_id, ?recorded,
"ignoring an installation event for another installation");
return;
}
match change {
InstallationChange::Deleted | InstallationChange::Suspended => {
report(bridge, &ns.id, json!({ "type": "installationRemoved" })).await;
}
_ => bridge.probe_bridge_checks(&ns.id).await,
}
}
other => tracing::debug!(event = ?other, "no event to report for this delivery"),
}
}