Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -160,6 +160,18 @@ Releases before v0.27.0 predate the changelog.
limit or a suspension — GitHub does not redeliver, so they are recorded and
wait at claim; only a `deleted` namespace refuses them. Namespaces with no
limit rows behave as before.
- Check-run updates now use a durable, coalescing transactional-outbox
projector and one leased background sender. GitHub API retries, rate-limit
backoff, stale-id reconciliation and restart recovery happen outside
webhook/request paths on both SQLite and PostgreSQL. Each check run carries
an `external_id` of `{run_id}:{job_id}`, and crash-after-POST
reconciliation adopts only a check with that id — never another workflow's
or app's same-named check. The control schema moves to **SQLite v4 /
Postgres v5** in place: an existing control database at the old version is
refused and must be recreated. The Postgres backend must **connect
directly, not through a transaction pooler** (e.g. PgBouncer in transaction
mode): it relies on per-connection `search_path` and on a dedicated
`LISTEN` connection, which transaction pooling breaks.

### Changed

Expand Down
23 changes: 7 additions & 16 deletions crates/preloop-runner-server/src/bootstrap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -511,13 +511,6 @@ pub async fn reap_once(shared: &Arc<SharedState>) {
reason: Some("fork-PR approval window expired".to_owned()),
})
.await;
crate::github::report_check_run_completed(
shared,
*run_id,
&job_id,
ExecutionStatus::Failure,
)
.await;
}
shared
.state
Expand All @@ -541,13 +534,6 @@ pub async fn reap_once(shared: &Arc<SharedState>) {
reason: Some(job.reason.clone()),
})
.await;
crate::github::report_check_run_completed(
shared,
job.run_id,
&job.job_id,
ExecutionStatus::Failure,
)
.await;
}

// Publish each affected run's updated status. A run the sweep just
Expand Down Expand Up @@ -679,8 +665,9 @@ async fn run_history_archiver(shared: Arc<SharedState>) {
}
}

/// Outbox rows are only needed for the live fan-out window (no durable
/// consumer reads them yet): older rows are deleted. Env:
/// Outbox rows are pruned once older than the retention window *and* at or
/// below the slowest durable consumer bookmark (the `check-runs` projector);
/// rows a consumer has not read yet are retained however old they are. Env:
/// `PRELOOP_OUTBOX_RETENTION_SECONDS` (default 3600).
fn outbox_retention() -> Duration {
Duration::from_secs(
Expand Down Expand Up @@ -1694,6 +1681,10 @@ pub async fn serve(config: ServerConfig) -> anyhow::Result<()> {
state: state.clone(),
shutdown: shutdown.clone(),
});
let check_run_sender_shared = shared.clone();
tokio::spawn(async move {
crate::github::run_check_run_sender(check_run_sender_shared).await;
});

// 5s sampler — clone needed state under lock, release, then publish.
let sampler_shared = shared.clone();
Expand Down
68 changes: 68 additions & 0 deletions crates/preloop-runner-server/src/check_run_outbox.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
//! Durable check-run projection and sender support.
//!
//! The projector is deliberately separate from the GitHub client: it only
//! reads committed outbox rows and writes the coalesced desired-state table.

use std::sync::{Arc, Weak};
use std::time::Duration;
use tokio::sync::Notify;

use crate::control::Backend;

const BATCH: usize = 256;
const LEASE: Duration = Duration::from_secs(30);
const POLL: Duration = Duration::from_millis(250);

/// Start the durable projector for either backend. PostgreSQL notifications are
/// converted into the same local dirty signal used by SQLite and by writers on
/// this node; polling remains the recovery path for a lost notification.
pub(crate) fn spawn_consumer(backend: &Arc<Backend>, dirty: Arc<Notify>) {
if let Some(mut notifications) = backend.subscribe_event_notifications() {
let wake = Arc::clone(&dirty);
tokio::spawn(async move {
loop {
match notifications.recv().await {
Ok(_) => wake.notify_one(),
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => wake.notify_one(),
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
}
}
});
}
tokio::spawn(run_consumer(Arc::downgrade(backend), dirty));
}

async fn run_consumer(backend: Weak<Backend>, dirty: Arc<Notify>) {
let owner = uuid::Uuid::new_v4().to_string();
loop {
let Some(backend_ref) = backend.upgrade() else {
return;
};
let mut processed = 0usize;
loop {
match backend_ref
.consume_check_run_outbox(&owner, LEASE, BATCH)
.await
{
Ok(n) => {
processed += n;
if n < BATCH {
break;
}
}
Err(error) => {
tracing::warn!(?error, "check-run outbox projection failed");
break;
}
}
}
drop(backend_ref);
if processed > 0 {
continue;
}
tokio::select! {
_ = dirty.notified() => {},
_ = tokio::time::sleep(POLL) => {},
}
}
}
Loading