feat: mem v2 - PR6 (consolidation) (#11374)

This commit is contained in:
jif-oai
2026-02-11 00:02:57 +00:00
committed by GitHub
parent 2c9be54c9a
commit 674799d356
7 changed files with 910 additions and 156 deletions
+109 -97
View File
@@ -194,67 +194,8 @@ WHERE thread_id = ?
}
}
let existing_job = sqlx::query(
let rows_affected = sqlx::query(
r#"
SELECT status, lease_until, retry_at, retry_remaining
FROM jobs
WHERE kind = ? AND job_key = ?
"#,
)
.bind(JOB_KIND_MEMORY_STAGE1)
.bind(thread_id.as_str())
.fetch_optional(&mut *tx)
.await?;
let should_insert = if let Some(existing_job) = existing_job {
let status: String = existing_job.try_get("status")?;
let existing_lease_until: Option<i64> = existing_job.try_get("lease_until")?;
let retry_at: Option<i64> = existing_job.try_get("retry_at")?;
let retry_remaining: i64 = existing_job.try_get("retry_remaining")?;
if retry_remaining <= 0 {
tx.commit().await?;
return Ok(Stage1JobClaimOutcome::SkippedRetryExhausted);
}
if retry_at.is_some_and(|retry_at| retry_at > now) {
tx.commit().await?;
return Ok(Stage1JobClaimOutcome::SkippedRetryBackoff);
}
if status == "running"
&& existing_lease_until.is_some_and(|lease_until| lease_until > now)
{
tx.commit().await?;
return Ok(Stage1JobClaimOutcome::SkippedRunning);
}
false
} else {
true
};
let fresh_running_jobs = sqlx::query(
r#"
SELECT COUNT(*) AS count
FROM jobs
WHERE kind = ?
AND status = 'running'
AND lease_until IS NOT NULL
AND lease_until > ?
"#,
)
.bind(JOB_KIND_MEMORY_STAGE1)
.bind(now)
.fetch_one(&mut *tx)
.await?
.try_get::<i64, _>("count")?;
if fresh_running_jobs >= max_running_jobs {
tx.commit().await?;
return Ok(Stage1JobClaimOutcome::SkippedRunning);
}
if should_insert {
sqlx::query(
r#"
INSERT INTO jobs (
kind,
job_key,
@@ -269,61 +210,96 @@ INSERT INTO jobs (
last_error,
input_watermark,
last_success_watermark
) VALUES (?, ?, 'running', ?, ?, ?, NULL, ?, NULL, ?, NULL, ?, NULL)
"#,
)
.bind(JOB_KIND_MEMORY_STAGE1)
.bind(thread_id.as_str())
.bind(worker_id.as_str())
.bind(ownership_token.as_str())
.bind(now)
.bind(lease_until)
.bind(DEFAULT_RETRY_REMAINING)
.bind(source_updated_at)
.execute(&mut *tx)
.await?;
tx.commit().await?;
return Ok(Stage1JobClaimOutcome::Claimed { ownership_token });
}
let rows_affected = sqlx::query(
r#"
UPDATE jobs
SET
)
SELECT ?, ?, 'running', ?, ?, ?, NULL, ?, NULL, ?, NULL, ?, NULL
WHERE (
SELECT COUNT(*)
FROM jobs
WHERE kind = ?
AND status = 'running'
AND lease_until IS NOT NULL
AND lease_until > ?
) < ?
ON CONFLICT(kind, job_key) DO UPDATE SET
status = 'running',
worker_id = ?,
ownership_token = ?,
started_at = ?,
worker_id = excluded.worker_id,
ownership_token = excluded.ownership_token,
started_at = excluded.started_at,
finished_at = NULL,
lease_until = ?,
lease_until = excluded.lease_until,
retry_at = NULL,
last_error = NULL,
input_watermark = ?
WHERE kind = ? AND job_key = ?
AND (status != 'running' OR lease_until IS NULL OR lease_until <= ?)
AND (retry_at IS NULL OR retry_at <= ?)
AND retry_remaining > 0
input_watermark = excluded.input_watermark
WHERE
(jobs.status != 'running' OR jobs.lease_until IS NULL OR jobs.lease_until <= excluded.started_at)
AND (jobs.retry_at IS NULL OR jobs.retry_at <= excluded.started_at)
AND jobs.retry_remaining > 0
AND (
SELECT COUNT(*)
FROM jobs AS running_jobs
WHERE running_jobs.kind = excluded.kind
AND running_jobs.status = 'running'
AND running_jobs.lease_until IS NOT NULL
AND running_jobs.lease_until > excluded.started_at
AND running_jobs.job_key != excluded.job_key
) < ?
"#,
)
.bind(JOB_KIND_MEMORY_STAGE1)
.bind(thread_id.as_str())
.bind(worker_id.as_str())
.bind(ownership_token.as_str())
.bind(now)
.bind(lease_until)
.bind(DEFAULT_RETRY_REMAINING)
.bind(source_updated_at)
.bind(JOB_KIND_MEMORY_STAGE1)
.bind(thread_id.as_str())
.bind(now)
.bind(now)
.bind(max_running_jobs)
.bind(max_running_jobs)
.execute(&mut *tx)
.await?
.rows_affected();
tx.commit().await?;
if rows_affected == 0 {
Ok(Stage1JobClaimOutcome::SkippedRunning)
} else {
Ok(Stage1JobClaimOutcome::Claimed { ownership_token })
if rows_affected > 0 {
tx.commit().await?;
return Ok(Stage1JobClaimOutcome::Claimed { ownership_token });
}
let existing_job = sqlx::query(
r#"
SELECT status, lease_until, retry_at, retry_remaining
FROM jobs
WHERE kind = ? AND job_key = ?
"#,
)
.bind(JOB_KIND_MEMORY_STAGE1)
.bind(thread_id.as_str())
.fetch_optional(&mut *tx)
.await?;
tx.commit().await?;
if let Some(existing_job) = existing_job {
let status: String = existing_job.try_get("status")?;
let existing_lease_until: Option<i64> = existing_job.try_get("lease_until")?;
let retry_at: Option<i64> = existing_job.try_get("retry_at")?;
let retry_remaining: i64 = existing_job.try_get("retry_remaining")?;
if retry_remaining <= 0 {
return Ok(Stage1JobClaimOutcome::SkippedRetryExhausted);
}
if retry_at.is_some_and(|retry_at| retry_at > now) {
return Ok(Stage1JobClaimOutcome::SkippedRetryBackoff);
}
if status == "running"
&& existing_lease_until.is_some_and(|lease_until| lease_until > now)
{
return Ok(Stage1JobClaimOutcome::SkippedRunning);
}
}
Ok(Stage1JobClaimOutcome::SkippedRunning)
}
pub async fn mark_stage1_job_succeeded(
@@ -627,6 +603,42 @@ WHERE kind = ? AND job_key = ?
Ok(rows_affected > 0)
}
pub async fn mark_global_phase2_job_failed_if_unowned(
&self,
ownership_token: &str,
failure_reason: &str,
retry_delay_seconds: i64,
) -> anyhow::Result<bool> {
let now = Utc::now().timestamp();
let retry_at = now.saturating_add(retry_delay_seconds.max(0));
let rows_affected = sqlx::query(
r#"
UPDATE jobs
SET
status = 'error',
finished_at = ?,
lease_until = NULL,
retry_at = ?,
retry_remaining = retry_remaining - 1,
last_error = ?
WHERE kind = ? AND job_key = ?
AND status = 'running'
AND (ownership_token = ? OR ownership_token IS NULL)
"#,
)
.bind(now)
.bind(retry_at)
.bind(failure_reason)
.bind(JOB_KIND_MEMORY_CONSOLIDATE_GLOBAL)
.bind(MEMORY_CONSOLIDATION_JOB_KEY)
.bind(ownership_token)
.execute(self.pool.as_ref())
.await?
.rows_affected();
Ok(rows_affected > 0)
}
}
async fn enqueue_global_consolidation_with_executor<'e, E>(