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
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
-- Clear purges separate rows written before a clear from rows written after
-- it by the clear generation, not by time.
--
-- A clear marker used to doom rows with rowid <= max_rowid and
-- created_at <= cleared_at. created_at has one-second resolution and SQLite
-- reuses the highest rowid once that row is deleted, so a row written in the
-- same second as the clear could match both conditions and be purged.
--
-- Every write now stores the generation that is current in its own
-- statement, and a clear stores the generation it bumped to. A row is doomed
-- by a clear marker when its generation is lower than the marker's.
--
-- Rows and markers written before this migration keep generation NULL. A
-- legacy marker still uses the old condition, but only for legacy rows.
--
-- (No semicolons in comments: batches are split on every semicolon.)
ALTER TABLE request_log_payloads ADD COLUMN generation INTEGER;
ALTER TABLE request_log_payload_manifests ADD COLUMN generation INTEGER;
ALTER TABLE request_log_upstream_attempts ADD COLUMN generation INTEGER;
ALTER TABLE request_log_response_links ADD COLUMN generation INTEGER;
ALTER TABLE request_log_payload_purges ADD COLUMN generation INTEGER;
11 changes: 11 additions & 0 deletions crates/core/src/storage/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ pub mod request_log_query;
mod request_logs;
mod storage_space;
pub use request_log_payload_purge::RequestLogPayloadPurgeProgress;
use request_log_payload_purge::REQUEST_LOG_PAYLOAD_GENERATION_TABLES;
pub use storage_space::{
wal_checkpoint_pending, DatabaseSpaceUsage, WalCheckpointOutcome, AUTO_VACUUM_FULL,
AUTO_VACUUM_INCREMENTAL, AUTO_VACUUM_NONE,
Expand Down Expand Up @@ -2381,6 +2382,16 @@ impl Storage {
"142_request_log_payload_purges",
include_str!("../../migrations/142_request_log_payload_purges.sql"),
)?;
self.apply_sql_or_compat_migration(
"143_request_log_payload_purge_generation",
include_str!("../../migrations/143_request_log_payload_purge_generation.sql"),
|storage| {
for table in REQUEST_LOG_PAYLOAD_GENERATION_TABLES {
storage.ensure_column(table, "generation", "INTEGER")?;
}
Ok(())
},
)?;
self.ensure_api_key_rotation_columns()?;
self.ensure_api_key_account_group_filter_column()?;
self.ensure_aggregate_apis_table()?;
Expand Down
151 changes: 108 additions & 43 deletions crates/core/src/storage/request_log_payload_purge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,11 @@
//! Now the clear/prune transaction only makes the old data unreachable and
//! records what has to go:
//! * clear: bumps the generation (in-flight jobs are rejected) and stores, per
//! payload table, the largest rowid and the clear time in
//! `request_log_payload_purges`;
//! payload table, the largest rowid and the new generation in
//! `request_log_payload_purges`. Every payload row stores the generation
//! that was current when it was written, so rows written before the clear
//! are told apart from later rows (even ones that reuse a rowid in the same
//! second) by generation, never by time;
//! * prune: raises `retention_cutoff` (jobs older than it are rejected) and
//! rebases surviving manifests whose parent is about to be deleted.
//!
Expand Down Expand Up @@ -42,6 +45,35 @@ const MANIFESTS_TABLE: &str = "request_log_payload_manifests";
const BLOB_GC_MARKER: &str = "request_log_payload_blobs";
const PURGES_TABLE: &str = "request_log_payload_purges";

/// Tables that carry the `generation` column (migration 143).
pub(super) const REQUEST_LOG_PAYLOAD_GENERATION_TABLES: [&str; 5] = [
"request_log_payloads",
"request_log_payload_manifests",
"request_log_upstream_attempts",
"request_log_response_links",
PURGES_TABLE,
];

/// SQL value for the `generation` column of a payload write: the clear
/// generation current in the same statement, so the value and the write are
/// atomic with respect to a clear on another connection.
pub(super) const CURRENT_PAYLOAD_GENERATION_SQL: &str =
"(SELECT generation FROM request_log_payload_state WHERE id = 1)";

/// Condition that the clear marker given by the `max_rowid`, `cleared_at` and
/// `generation` SQL expressions dooms the row aliased `row`. Markers written since migration 143 carry the
/// generation the clear bumped to and doom rows of any lower generation, plus
/// rows written before the migration (generation NULL). Older markers
/// (generation NULL) keep the time-based condition, restricted to rows
/// written before the migration, so they never match newer rows.
fn cleared_row_sql(row: &str, max_rowid: &str, cleared_at: &str, generation: &str) -> String {
format!(
"{row}.rowid <= {max_rowid}
AND ({row}.generation IS NULL OR {row}.generation < {generation})
AND ({generation} IS NOT NULL OR {row}.created_at <= {cleared_at})"
)
}

const MIN_BATCH_ROWS: i64 = 50;
const MAX_BATCH_ROWS: i64 = 4000;
const INITIAL_BATCH_ROWS: i64 = 1000;
Expand All @@ -63,27 +95,42 @@ pub struct RequestLogPayloadPurgeProgress {
pub pending: bool,
}

/// A pending clear of one payload table, as stored in
/// `request_log_payload_purges`.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct ClearMarker {
max_rowid: i64,
cleared_at: i64,
/// Generation the clear bumped to; `None` for markers written before
/// migration 143.
generation: Option<i64>,
}

#[derive(Debug, Clone, Copy)]
enum Doomed {
/// Rows that existed when a clear committed.
Cleared { max_rowid: i64, cleared_at: i64 },
/// Rows written before a clear committed.
Cleared(ClearMarker),
/// Rows older than the retention cutoff.
Retention { cutoff: i64 },
}

impl Doomed {
fn predicate(self) -> (&'static str, Vec<Value>) {
fn predicate(self, table: &str) -> (String, Vec<Value>) {
match self {
Self::Cleared {
max_rowid,
cleared_at,
} => (
// A rowid reused by a row written after the clear is excluded
// by its later created_at.
"rowid <= ? AND created_at <= ? ORDER BY rowid",
vec![Value::Integer(max_rowid), Value::Integer(cleared_at)],
Self::Cleared(marker) => (
format!(
"{} ORDER BY rowid",
cleared_row_sql(table, "?1", "?2", "?3")
),
vec![
Value::Integer(marker.max_rowid),
Value::Integer(marker.cleared_at),
marker.generation.map_or(Value::Null, Value::Integer),
],
),
Self::Retention { cutoff } => ("created_at < ?", vec![Value::Integer(cutoff)]),
Self::Retention { cutoff } => {
("created_at < ?1".to_string(), vec![Value::Integer(cutoff)])
}
}
}
}
Expand All @@ -100,8 +147,13 @@ pub(super) fn live_payload_row_sql(table: &str, alias: &str) -> String {
AND NOT EXISTS (
SELECT 1 FROM request_log_payload_purges purge
WHERE purge.table_name = '{table}'
AND {alias}.rowid <= purge.max_rowid
AND {alias}.created_at <= purge.cleared_at)"
AND {cleared})",
cleared = cleared_row_sql(
alias,
"purge.max_rowid",
"purge.cleared_at",
"purge.generation"
)
)
}

Expand Down Expand Up @@ -199,6 +251,9 @@ impl Storage {
if !self.has_request_log_payload_purges()? {
return self.delete_request_log_payloads_unbatched();
}
// Rows written before this clear carry a lower generation; rows
// written after it carry this one or higher.
let generation = self.request_log_payload_generation()?;
for table in PURGED_PAYLOAD_TABLES {
if !self.has_table(table)? {
continue;
Expand All @@ -211,13 +266,18 @@ impl Storage {
let Some(max_rowid) = max_rowid else {
continue;
};
// Merging into an older marker: every row it covers predates this
// clear, so the higher generation and rowid bound cover it too.
self.conn.execute(
"INSERT INTO request_log_payload_purges (table_name, max_rowid, cleared_at, cursor_rowid)
VALUES (?1, ?2, ?3, 0)
"INSERT INTO request_log_payload_purges
(table_name, max_rowid, cleared_at, cursor_rowid, generation)
VALUES (?1, ?2, ?3, 0, ?4)
ON CONFLICT(table_name) DO UPDATE SET
max_rowid = MAX(max_rowid, excluded.max_rowid),
cleared_at = MAX(cleared_at, excluded.cleared_at)",
(table, max_rowid, cleared_at),
cleared_at = MAX(cleared_at, excluded.cleared_at),
generation = MAX(COALESCE(generation, excluded.generation),
excluded.generation)",
(table, max_rowid, cleared_at, generation),
)?;
}
self.mark_request_log_payload_blob_gc()
Expand Down Expand Up @@ -328,21 +388,18 @@ impl Storage {
continue;
}
loop {
let Some((max_rowid, cleared_at)) =
self.request_log_payload_purge_marker(table)?
else {
let Some(marker) = self.request_log_payload_purge_marker(table)? else {
break;
};
if out_of_time() {
progress.pending = true;
return Ok(progress);
}
let doomed = Doomed::Cleared {
max_rowid,
cleared_at,
};
let (deleted, finished) =
self.purge_request_log_payload_batch(table, doomed, &mut batch)?;
let (deleted, finished) = self.purge_request_log_payload_batch(
table,
Doomed::Cleared(marker),
&mut batch,
)?;
progress.deleted_rows += deleted;
if finished {
break;
Expand Down Expand Up @@ -393,12 +450,19 @@ impl Storage {
})
}

fn request_log_payload_purge_marker(&self, table: &str) -> Result<Option<(i64, i64)>> {
fn request_log_payload_purge_marker(&self, table: &str) -> Result<Option<ClearMarker>> {
self.conn
.query_row(
"SELECT max_rowid, cleared_at FROM request_log_payload_purges WHERE table_name = ?1",
"SELECT max_rowid, cleared_at, generation
FROM request_log_payload_purges WHERE table_name = ?1",
[table],
|row| Ok((row.get(0)?, row.get(1)?)),
|row| {
Ok(ClearMarker {
max_rowid: row.get(0)?,
cleared_at: row.get(1)?,
generation: row.get(2)?,
})
},
)
.optional()
}
Expand Down Expand Up @@ -428,7 +492,7 @@ impl Storage {
batch: &mut BatchSize,
) -> Result<(usize, bool)> {
let limit = batch.rows(table);
let (predicate, params) = doomed.predicate();
let (predicate, params) = doomed.predicate(table);
let started = Instant::now();
let tx = self.conn.unchecked_transaction()?;
let mut stmt = self.conn.prepare(&format!(
Expand All @@ -455,18 +519,19 @@ impl Storage {
)?;
}
let finished = (rowids.len() as i64) < limit;
if let (
true,
Doomed::Cleared {
max_rowid,
cleared_at,
},
) = (finished, doomed)
{
if let (true, Doomed::Cleared(marker)) = (finished, doomed) {
// Only the marker this batch worked on: a clear that committed
// since then has merged a higher bound into it.
self.conn.execute(
"DELETE FROM request_log_payload_purges
WHERE table_name = ?1 AND max_rowid = ?2 AND cleared_at = ?3",
(table, max_rowid, cleared_at),
WHERE table_name = ?1 AND max_rowid = ?2 AND cleared_at = ?3
AND generation IS ?4",
(
table,
marker.max_rowid,
marker.cleared_at,
marker.generation,
),
)?;
}
tx.commit()?;
Expand Down
Loading
Loading