Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
32 commits
Select commit Hold shift + click to select a range
7565b82
feat(audit): pair every request entry with a response
jeremi Sep 26, 2026
c66be30
fix(casework): pair every audited request with a response
jeremi Sep 26, 2026
e4a9cd8
test(casework): expect a withheld result to pair its request entry
jeremi Sep 26, 2026
f17ed10
fix(scheduling): answer every commitment request entry
jeremi Sep 26, 2026
ca0c54b
fix(render): pair each call's audit entries under a server-drawn id
jeremi Sep 26, 2026
b85ceaf
fix(breg): answer every audited request entry
jeremi Sep 26, 2026
fe92918
fix(hooks): record delivery outcomes only after they commit
jeremi Sep 26, 2026
368ce72
docs: describe the audit request pairing guarantee
jeremi Sep 26, 2026
94e761c
test(breg): expect a faulted mutation to answer its attempt
jeremi Sep 26, 2026
747a32f
test(breg): expect faulted reads and tombstones to answer their attempt
jeremi Sep 26, 2026
8c05a59
docs(breg): enforce BREG-V1-27
jeremi Sep 26, 2026
68157f5
docs: anchor the audit pairing claims in the retention page
jeremi Sep 26, 2026
d7161db
fix(breg): answer maintenance request entries that end early
jeremi Sep 26, 2026
eab4baf
fix(audit): keep request pairing through cancellation
jeremi Sep 26, 2026
474ac05
fix(hooks): resolve an unacknowledged commit before auditing it
jeremi Sep 26, 2026
4c5ed4d
fix(hooks): answer an operator replay whose caller left
jeremi Sep 26, 2026
ba58367
fix(breg): resolve erasure and receipt outcomes before auditing them
jeremi Sep 26, 2026
c398681
fix(audit): claim a request for a response appended in its name
jeremi Sep 26, 2026
91ece62
fix(breg): read back retention and reconcile outcomes after commit er…
jeremi Sep 26, 2026
f2c8221
fix(audit): queue a dropped request's file response outside a runtime
jeremi Sep 26, 2026
f94e9a7
fix(hooks): recognize retry and dead-letter commits after later progress
jeremi Sep 26, 2026
6cad728
fix(audit): answer a request once when a claim races its drop or shut…
jeremi Sep 26, 2026
ebd09b5
fix(breg): hold an Evidence action's attempt until its terminal answe…
jeremi Sep 26, 2026
62d28e2
fix(breg): record a committed erasure before retrying external deletions
jeremi Sep 26, 2026
d8da69b
fix(breg): answer an ingestion transition unfinished when its commit …
jeremi Sep 26, 2026
487a4f8
fix(hooks): answer every delivery transition whose commit fate is unk…
jeremi Sep 26, 2026
476e213
docs(breg): scope audit pairing to a running process in BREG-V1-27
jeremi Sep 26, 2026
fcf54eb
fix(casework): write an operation's whole response tail even when its…
jeremi Sep 26, 2026
2c65e43
fix(casework): read back a commit whose acknowledgment was lost befor…
jeremi Sep 26, 2026
4108d34
fix(casework): name the review request a review note's request entry …
jeremi Sep 26, 2026
d7010a4
fix(scheduling): read back a capacity commit whose acknowledgment was…
jeremi Sep 26, 2026
3dd1be2
docs(site): state what the BReg and Render audit journals record when…
jeremi Sep 26, 2026
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
147 changes: 142 additions & 5 deletions crates/registry-breg/src/action_evidence_maintenance.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,13 +3,22 @@

use std::{path::Path, time::Duration};

use registry_platform_audit::AuditEntry;
use serde_json::{json, Value};
use uuid::Uuid;

use crate::audit::RegistryAudit;
use crate::mutation::{erase_expired_action_evidence, MutationError};
use crate::postgres::{
verify_catalog_identity_for_catalog, verify_migration_role, ConnectionConfig,
ExpectedManagedCatalog, ExpectedRegistryIdentity, RegistryLockKey, SqlIdentifier,
};
use crate::runtime_config::load_runtime_config;

/// Schema of the request and response entries of one expired-Evidence
/// erasure.
pub const EVIDENCE_RETENTION_AUDIT_SCHEMA: &str = "breg-evidence-retention-audit/v1";

/// Package-bound authority for erasing expired protected Evidence material.
/// The migration identity, actual target catalog and registry interlock are
/// checked again in the same transaction that deletes the retained material.
Expand All @@ -22,6 +31,7 @@ pub struct ActionEvidenceRetentionOperatorService {
runtime_role: SqlIdentifier,
lock_timeout: Duration,
statement_timeout: Duration,
audit: RegistryAudit,
}

impl ActionEvidenceRetentionOperatorService {
Expand Down Expand Up @@ -56,6 +66,9 @@ impl ActionEvidenceRetentionOperatorService {
runtime_role: config.database().roles().runtime().clone(),
lock_timeout: config.operational_timeouts().migration_lock,
statement_timeout: config.operational_timeouts().migration_statement,
audit: RegistryAudit::open_companion(&config)
.await
.map_err(|_| MutationError::Unavailable)?,
})
}

Expand All @@ -68,6 +81,7 @@ impl ActionEvidenceRetentionOperatorService {
migration_connection: ConnectionConfig,
migration_role: SqlIdentifier,
runtime_role: SqlIdentifier,
audit: RegistryAudit,
) -> Self {
Self {
expected,
Expand All @@ -78,16 +92,125 @@ impl ActionEvidenceRetentionOperatorService {
runtime_role,
lock_timeout: Duration::from_secs(5),
statement_timeout: Duration::from_secs(10),
audit,
}
}

/// Erase the retained Evidence material whose expiry is before `before`.
///
/// The request entry, naming the threshold, is accepted before the
/// erasure transaction opens, so an audit outage erases nothing. Its
/// response records the erased count once the transaction commits, or
/// the failure when it does not.
pub async fn erase_expired(
&self,
before: chrono::DateTime<chrono::Utc>,
) -> Result<u64, MutationError> {
if before > chrono::Utc::now() {
return Err(MutationError::InvalidRequest);
}
let correlation = Uuid::new_v4().to_string();
let record = |phase: &str, outcome: &str| -> Value {
json!({
"kind": "evidenceRetention",
"phase": phase,
"outcome": outcome,
"packageRevision": self.expected.package_revision,
"actor": "breg:evidence-retention-operator",
"before": before.to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
"correlation": correlation,
})
};
let mut attempt = self
.audit
.begin(
AuditEntry::request(
EVIDENCE_RETENTION_AUDIT_SCHEMA,
correlation.clone(),
record("attempt", "started"),
),
record("terminal", "unfinished"),
)
.await
.map_err(|_| MutationError::Unavailable)?;
let erased = match self.erase_in_transaction(before).await {
Ok(Erasure::Committed(erased)) => Ok(erased),
// A commit that returned an error may still have committed, so
// the outcome recorded for this destructive operation is the one
// the database holds, read on a fresh connection.
Ok(Erasure::Unacknowledged { erased, cutoff }) => {
match self.expired_evidence_remains(cutoff).await {
Some(false) => Ok(erased),
Comment on lines +142 to +143

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Tie erasure recovery to this transaction

When this erasure's commit rolls back after returning an ambiguous error, a concurrent erase-expired invocation can delete the same rows before expired_evidence_remains runs; Some(false) then attributes the second invocation's deletion to the first and records the first request as erased with its speculative count. Fresh evidence beyond the prior read-back fix is that the recovery predicate is not transaction-specific; use the original transaction status or a request-specific durable marker instead.

AGENTS.md reference: products/breg/AGENTS.md:L175-L182

Useful? React with 👍 / 👎.

resolved => {
let outcome = if resolved == Some(true) {
"failed"
} else {
"unfinished"
};
if attempt.respond(record("terminal", outcome)).await.is_err() {
tracing::error!(
"the unacknowledged Evidence retention's response audit entry was not recorded"
);
}
return Err(MutationError::Unavailable);
}
}
}
Err(error) => Err(error),
};
match erased {
Ok(erased) => {
let mut response = record("terminal", "erased");
response["erased"] = json!(erased);
// The erasure committed; a refused entry reports the command
// unavailable, and the writer then refuses every later entry.
attempt
.respond(response)
.await
.map_err(|_| MutationError::Unavailable)?;
Ok(erased)
}
Err(error) => {
if attempt.respond(record("terminal", "failed")).await.is_err() {
tracing::error!(
"the failed Evidence retention's response audit entry was not recorded"
);
}
Err(error)
}
}
}

/// Whether retained Evidence expiring at or before `cutoff` remains, read
/// on a fresh connection after an erasure commit returned an error.
/// `None` when it cannot be read.
async fn expired_evidence_remains(
&self,
cutoff: chrono::DateTime<chrono::Utc>,
) -> Option<bool> {
let pool = self.migration_connection.build_pool().ok()?;
let client = pool.get().await.ok()?;
let row = client
.query_one(
"SELECT EXISTS (SELECT 1 FROM registry_internal.registry_action_evidence_uses
WHERE expires_at <= $1)
OR EXISTS (SELECT 1 FROM registry_internal.registry_request_evidence_uses
WHERE expires_at <= $1)",
&[&cutoff],
)
.await
.ok()?;
row.try_get(0).ok()
}

/// Erase the expired material in one transaction. A commit that returned
/// an error, which does not prove the transaction rolled back, is
/// reported with the count it would have erased and the cutoff it erased
/// through; every earlier error is returned as one.
async fn erase_in_transaction(
&self,
before: chrono::DateTime<chrono::Utc>,
) -> Result<Erasure, MutationError> {
let pool = self
.migration_connection
.build_pool()
Expand Down Expand Up @@ -127,15 +250,29 @@ impl ActionEvidenceRetentionOperatorService {
if !ready {
return Err(MutationError::Unavailable);
}
let erased = erase_expired_action_evidence(&transaction, before).await?;
transaction
.commit()
// The cutoff the deletion applies, fixed by the transaction's start.
let cutoff: chrono::DateTime<chrono::Utc> = transaction
.query_one("SELECT LEAST($1, CURRENT_TIMESTAMP)", &[&before])
.await
.map_err(|_| MutationError::Unavailable)?;
Ok(erased)
.map_err(|_| MutationError::Unavailable)?
.get(0);
let erased = erase_expired_action_evidence(&transaction, before).await?;
if transaction.commit().await.is_err() {
return Ok(Erasure::Unacknowledged { erased, cutoff });
}
Ok(Erasure::Committed(erased))
}
}

/// How an erasure transaction ended once every statement in it succeeded.
enum Erasure {
Committed(u64),
Unacknowledged {
erased: u64,
cutoff: chrono::DateTime<chrono::Utc>,
},
}

/// Erase only material whose declared expiry has passed. Diagnostics contain
/// no connection, selector, assertion or provider values.
pub async fn erase_expired(path: &Path, before: &str) -> Result<u64, MutationError> {
Expand Down
42 changes: 24 additions & 18 deletions crates/registry-breg/src/api/ingestion.rs
Original file line number Diff line number Diff line change
Expand Up @@ -177,15 +177,17 @@ async fn create_run(
.await
{
Ok(run) => ingestion_response(StatusCode::CREATED, json!({ "run": run })),
// A refusal after the request entry was accepted owes the journal
// its response entry, as a refused chunk submission does.
Err(error) => {
// A refusal after the ingestion request entry is already answered in
// the ingestion schema; any other refusal owes the journal its
// single refusal entry here, so an audit outage gates it.
Err(refusal) if refusal.answered => ingestion_problem(refusal.error),
Err(refusal) => {
audited_mutation_refusal(
mutations,
&binding.base,
&surface.context,
None,
ingestion_problem(error),
ingestion_problem(refusal.error),
&correlation,
)
.await
Expand Down Expand Up @@ -395,15 +397,17 @@ async fn cancel_run(
.await
{
Ok(run) => ingestion_response(StatusCode::OK, json!({ "run": run })),
// A refusal after the request entry was accepted owes the journal
// its response entry, as a refused chunk submission does.
Err(error) => {
// A refusal after the ingestion request entry is already answered in
// the ingestion schema; any other refusal owes the journal its
// single refusal entry here, so an audit outage gates it.
Err(refusal) if refusal.answered => ingestion_problem(refusal.error),
Err(refusal) => {
audited_mutation_refusal(
mutations,
&binding.base,
&surface.context,
None,
ingestion_problem(error),
ingestion_problem(refusal.error),
&correlation,
)
.await
Expand Down Expand Up @@ -545,16 +549,17 @@ async fn submit_chunk(
.await
{
Ok(answer) => ingestion_response(StatusCode::OK, answer),
// A submission the run refuses after parsing owes the journal the
// same durable refusal envelope pre-parse failures write, so an audit
// outage gates the refusal instead of passing silently.
Err(error) => {
// A refusal after the ingestion request entry is already answered in
// the ingestion schema; any other refusal owes the journal its
// single refusal entry here, so an audit outage gates it.
Err(refusal) if refusal.answered => ingestion_problem(refusal.error),
Err(refusal) => {
audited_mutation_refusal(
mutations,
&binding.base,
&surface.context,
None,
ingestion_problem(error),
ingestion_problem(refusal.error),
&correlation,
)
.await
Expand Down Expand Up @@ -631,16 +636,17 @@ async fn chunk_receipt(
.await
{
Ok(receipt) => ingestion_response(StatusCode::OK, receipt),
// A recovery the run refuses after its request entry was accepted
// owes the journal the same durable refusal envelope a refused
// cancellation or chunk submission owes.
Err(error) => {
// A refusal after the ingestion request entry is already answered in
// the ingestion schema; any other refusal owes the journal its
// single refusal entry here, so an audit outage gates it.
Err(refusal) if refusal.answered => ingestion_problem(refusal.error),
Err(refusal) => {
audited_mutation_refusal(
mutations,
&binding.base,
&surface.context,
None,
ingestion_problem(error),
ingestion_problem(refusal.error),
&correlation,
)
.await
Expand Down
43 changes: 36 additions & 7 deletions crates/registry-breg/src/attachment_verification_worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,9 @@ impl AttachmentVerificationWorker {
};
transaction.commit().await.map_err(unavailable)?;
drop(client);
self.audit_job(&job, "attempt", "started").await?;
// Held until the terminal entry answers it; a run that ends first,
// including one its time budget cancels, answers it as unfinished.
let _attempt = self.begin_job(&job).await?;

let verdict = match self.content(&job).await {
Ok(bytes) => verifier
Expand Down Expand Up @@ -240,10 +242,42 @@ impl AttachmentVerificationWorker {
Ok(transaction)
}

/// Append the attempt `request` entry of one leased job and return the
/// handle that owes its terminal `response`.
async fn begin_job(
&self,
job: &VerificationJob,
) -> Result<registry_platform_audit::AuditRequest> {
let (reference, record) = self.job_record(job, "attempt", "started")?;
let (_, unfinished) = self.job_record(job, "terminal", "unfinished")?;
self.audit
.begin(
AuditEntry::request(ATTACHMENT_VERIFICATION_AUDIT_SCHEMA, reference, record),
unfinished,
)
.await
.map_err(unavailable)
}

/// Append one verification entry. The attempt is the `request` entry and
/// the terminal outcome is the `response` entry; both are correlated by
/// the keyed verification reference of the leased job.
async fn audit_job(&self, job: &VerificationJob, phase: &str, outcome: &str) -> Result<()> {
let (reference, record) = self.job_record(job, phase, outcome)?;
let entry = if phase == "attempt" {
AuditEntry::request(ATTACHMENT_VERIFICATION_AUDIT_SCHEMA, reference, record)
} else {
AuditEntry::response(ATTACHMENT_VERIFICATION_AUDIT_SCHEMA, reference, record)
};
self.audit.append(entry).await.map_err(unavailable)
}

fn job_record(
&self,
job: &VerificationJob,
phase: &str,
outcome: &str,
) -> Result<(String, serde_json::Value)> {
let hasher = self.audit.profile().key_hasher();
let reference = hasher
.audit_reference_hash(
Expand All @@ -257,12 +291,7 @@ impl AttachmentVerificationWorker {
"packageRevision": self.expected.package_revision,
"actor": "breg:attachment-verifier", "verificationReference": reference,
});
let entry = if phase == "attempt" {
AuditEntry::request(ATTACHMENT_VERIFICATION_AUDIT_SCHEMA, reference, record)
} else {
AuditEntry::response(ATTACHMENT_VERIFICATION_AUDIT_SCHEMA, reference, record)
};
self.audit.append(entry).await.map_err(unavailable)
Ok((reference, record))
}
}

Expand Down
Loading
Loading