fix: library cron fixes - #63
mahatoankitkumar wants to merge 1 commit into
Conversation
WalkthroughThe pull request adds a shared, workspace-aware job service for REST and library callers. It centralizes validation, idempotency, trigger handling, transactional creation, cancellation, and cron scheduling. It also adds prefix-aware reaper provisioning and endpoint reference serialization. ChangesShared job service
Estimated code review effort: 4 (Complex) | ~60 minutes Possibly related issues
Possibly related PRs
Suggested reviewers: Poem
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Actionable comments posted: 2
🧹 Nitpick comments (4)
crates/common/src/service/mod.rs (1)
86-88: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueReuse the
Displaytext instead of repeating the format string.
ServiceError::InternalEndpointalready renders this exact message throughthiserror. The duplicated literal can drift if one side changes.♻️ Proposed simplification
- ServiceError::InternalEndpoint(n) => AppError::InvalidRequest(format!( - "Endpoint '{n}' is internal and cannot be used for user-created jobs" - )), + e @ ServiceError::InternalEndpoint(_) => AppError::InvalidRequest(e.to_string()),🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/common/src/service/mod.rs` around lines 86 - 88, Update the ServiceError::InternalEndpoint arm in the AppError conversion to reuse the variant’s existing Display representation when constructing AppError::InvalidRequest, removing the duplicated format string while preserving the current message text.crates/common/src/service/jobs.rs (1)
335-342: 🎯 Functional Correctness | 🔵 Trivial | ⚡ Quick winMatch only
idx_{prefix}jobs_idempotencywith SQLSTATE23505.The index name includes the configured table prefix. Do not use a broad
contains("idempotency")check.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/common/src/service/jobs.rs` around lines 335 - 342, Update duplicate_key_as_conflict to classify an error as ServiceError::Conflict only when the database error has SQLSTATE 23505 and its constraint exactly matches the configured idx_{prefix}jobs_idempotency index; remove the broad constraint-presence match and preserve ServiceError::Db for all other errors.crates/api/src/handlers/jobs.rs (1)
86-125: 🩺 Stability & Availability | 🔵 Trivial | ⚡ Quick winDo not panic on a missing execution in the HTTP path.
execution.expect(...)panics if the service returnsCreatedwithout an execution for an IMMEDIATE or DELAYED job. A panic in the handler aborts the request at the actix worker level instead of returning a controlled error. The invariant lives in another crate, so a future change toservice::jobs::create_jobturns a contract drift into a panic.Map the absent execution to an internal error instead.
🛡️ Proposed change for the IMMEDIATE arm (apply the same pattern to DELAYED)
TriggerType::IMMEDIATE => { - let execution = execution.expect("IMMEDIATE jobs always create an execution"); + let execution = execution.ok_or_else(|| { + AppError::Internal("IMMEDIATE job created without an execution".into()) + })?;🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/api/src/handlers/jobs.rs` around lines 86 - 125, Replace the execution.expect calls in the IMMEDIATE and DELAYED arms of the job response construction with controlled internal-error handling for None, returning the handler’s existing error result instead of panicking. Preserve the current JSON response shape when execution is present and apply the same missing-execution behavior to both TriggerType branches.crates/worker/src/client.rs (1)
210-215: 🩺 Stability & Availability | 🔵 Trivial | ⚡ Quick winValidate the reaper CRON expression at configuration time.
with_reaper_schedulestores the string without validation. An invalid expression surfaces later, insideprovision_workspace, as a pg_cron scheduling failure during boot.create_jobalready validates library CRON input withPgCronExpr::try_from(line 273). Apply the same check here, or make the setter fallible.♻️ Proposed fallible setter
- pub fn with_reaper_schedule(mut self, cron_expression: &str) -> Self { - self.reaper_cron = cron_expression.to_string(); - self - } + pub fn with_reaper_schedule(mut self, cron_expression: &str) -> anyhow::Result<Self> { + let expr = PgCronExpr::try_from(cron_expression.to_string())?; + self.reaper_cron = expr.as_str().to_string(); + Ok(self) + }🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/worker/src/client.rs` around lines 210 - 215, Validate cron_expression in with_reaper_schedule using the existing PgCronExpr::try_from validation used by create_job, or change the setter to a fallible API that returns the validation error; ensure invalid schedules are rejected during client configuration rather than stored for later provisioning.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@crates/common/src/db/workspaces.rs`:
- Around line 219-239: The reaper insertion in the workspace setup can race and
create duplicate active jobs. Add a partial unique index for the active reaper
endpoint scoped to kronos.reaper, then update the INSERT in the existing query
flow to use ON CONFLICT DO NOTHING while retaining the RETURNING-based
registration path for the winning transaction.
In `@crates/common/src/service/jobs.rs`:
- Around line 226-236: Update the TriggerPlan::Delayed arm around
db::jobs::create_delayed to map duplicate-key database errors through the
existing duplicate_key_as_conflict helper, matching the IMMEDIATE arm and
returning a conflict for concurrent idempotency-key races while preserving other
database errors.
---
Nitpick comments:
In `@crates/api/src/handlers/jobs.rs`:
- Around line 86-125: Replace the execution.expect calls in the IMMEDIATE and
DELAYED arms of the job response construction with controlled internal-error
handling for None, returning the handler’s existing error result instead of
panicking. Preserve the current JSON response shape when execution is present
and apply the same missing-execution behavior to both TriggerType branches.
In `@crates/common/src/service/jobs.rs`:
- Around line 335-342: Update duplicate_key_as_conflict to classify an error as
ServiceError::Conflict only when the database error has SQLSTATE 23505 and its
constraint exactly matches the configured idx_{prefix}jobs_idempotency index;
remove the broad constraint-presence match and preserve ServiceError::Db for all
other errors.
In `@crates/common/src/service/mod.rs`:
- Around line 86-88: Update the ServiceError::InternalEndpoint arm in the
AppError conversion to reuse the variant’s existing Display representation when
constructing AppError::InvalidRequest, removing the duplicated format string
while preserving the current message text.
In `@crates/worker/src/client.rs`:
- Around line 210-215: Validate cron_expression in with_reaper_schedule using
the existing PgCronExpr::try_from validation used by create_job, or change the
setter to a fallible API that returns the validation error; ensure invalid
schedules are rejected during client configuration rather than stored for later
provisioning.
🪄 Autofix
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro Plus
Run ID: afffbc49-e121-4699-8886-7ff9aad06e3c
📒 Files selected for processing (7)
crates/api/src/handlers/jobs.rscrates/common/src/db/jobs.rscrates/common/src/db/workspaces.rscrates/common/src/lib.rscrates/common/src/service/jobs.rscrates/common/src/service/mod.rscrates/worker/src/client.rs
| // `WHERE NOT EXISTS` rather than a blind INSERT: a second call must not add | ||
| // a second sweep. Returns no row when one is already installed. | ||
| let existing: Option<(String,)> = sqlx::query_as(&format!( | ||
| "INSERT INTO {tj} ( \ | ||
| endpoint, endpoint_type, trigger_type, \ | ||
| cron_expression, cron_timezone, cron_next_run_at \ | ||
| ) VALUES ($1, $2, $3, $4, 'UTC', now()) \ | ||
| RETURNING job_id", | ||
| ) | ||
| ) SELECT $1, $2, $3, $4, 'UTC', now() \ | ||
| WHERE NOT EXISTS ( \ | ||
| SELECT 1 FROM {tj} WHERE endpoint = $1 AND status = 'ACTIVE' \ | ||
| ) \ | ||
| RETURNING job_id" | ||
| )) | ||
| .bind(REAPER_ENDPOINT_NAME) | ||
| .bind(EndpointType::INTERNAL.to_string()) | ||
| .bind(TriggerType::CRON.as_str()) | ||
| .bind(cron_expression) | ||
| .fetch_one(&mut *tx) | ||
| .fetch_optional(&mut *tx) | ||
| .await?; | ||
|
|
||
| // Provisioning runs in schema-per-workspace mode with unprefixed tables, so | ||
| // the table prefix is empty here. | ||
| register_pg_cron_conn(&mut *tx, "", schema_name, &job_id, cron_expression).await?; | ||
| if let Some((job_id,)) = existing { | ||
| register_pg_cron_conn(&mut tx, table_prefix, schema_name, &job_id, cron_expression).await?; |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🟠 Major | ⚡ Quick win
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Inspect transaction setup and all jobs-table constraints or indexes.
rg -n -C 8 'fn scoped_transaction|async fn scoped_transaction|scoped_transaction\(' crates
rg -n -i -C 8 'CREATE TABLE.*jobs|CREATE.*INDEX.*jobs|UNIQUE.*endpoint|status.*ACTIVE' cratesRepository: juspay/kronos
Length of output: 45606
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '%s\n' '--- scoped transaction implementation ---'
cat -n crates/common/src/db/scoped.rs | sed -n '1,70p'
printf '%s\n' '--- reaper provisioning implementation ---'
cat -n crates/common/src/db/workspaces.rs | sed -n '190,255p'
printf '%s\n' '--- jobs schema and indexes ---'
cat -n crates/common/migrations/workspace_v1.sql | sed -n '45,90p'
printf '%s\n' '--- migration inventory and relevant uniqueness definitions ---'
git ls-files crates/common/migrations
rg -n -i 'unique|create index|isolation|transaction isolation|jobs' crates/common/migrationsRepository: juspay/kronos
Length of output: 8717
🏁 Script executed:
#!/bin/bash
set -euo pipefail
python3 - <<'PY'
from pathlib import Path
import re
scoped = Path("crates/common/src/db/scoped.rs").read_text()
workspaces = Path("crates/common/src/db/workspaces.rs").read_text()
migration = Path("crates/common/migrations/workspace_v1.sql").read_text()
tx_body = re.search(
r"pub async fn scoped_transaction\b.*?(?=\n\})",
scoped,
re.S,
)
assert tx_body, "scoped_transaction was not found"
assert "pool.begin().await?" in tx_body.group(0), "transaction start changed"
reaper = re.search(
r"let existing: Option<\(String,\)>.*?fetch_optional\(&mut \*tx\)",
workspaces,
re.S,
)
assert reaper, "reaper INSERT was not found"
reaper_sql = reaper.group(0)
assert "WHERE NOT EXISTS" in reaper_sql
assert "endpoint = $1 AND status = 'ACTIVE'" in reaper_sql
partial_active_endpoint_uniqueness = re.search(
r"CREATE UNIQUE(?: INDEX| TABLE CONSTRAINT).*?"
r"\bjobs\b.*?"
r"\bendpoint\b.*?"
r"(?:WHERE|status\s*=\s*'ACTIVE')",
migration,
re.I | re.S,
)
assert not partial_active_endpoint_uniqueness, (
"the migration already contains active endpoint uniqueness"
)
print("scoped_transaction uses pool.begin() without an explicit isolation level")
print("reaper provisioning uses an unprotected WHERE NOT EXISTS check")
print("workspace_v1.sql has no partial unique rule for active jobs by endpoint")
print("conclusion: the source does not provide a database uniqueness invariant for this race")
PYRepository: juspay/kronos
Length of output: 299
🌐 Web query:
PostgreSQL documentation READ COMMITTED INSERT SELECT WHERE NOT EXISTS concurrent transactions unique constraint
💡 Result:
Using an INSERT... SELECT... WHERE NOT EXISTS statement is not safe against race conditions in PostgreSQL, even under the default READ COMMITTED isolation level [1][2][3]. In READ COMMITTED mode, each statement sees a snapshot of the database as it existed at the start of that specific statement [4][5]. If two concurrent transactions execute an INSERT... SELECT... WHERE NOT EXISTS for the same key simultaneously, both may perform the SELECT concurrently, both see that the row does not exist, and both proceed to attempt an INSERT [1][2]. This leads to a race condition where both transactions attempt to insert the same value [1][2]. To handle concurrent insertions reliably, you must rely on a UNIQUE constraint or a UNIQUE index on the target column(s) [2][6][7]. PostgreSQL enforces UNIQUE constraints with independent isolation mechanisms that operate correctly regardless of the transaction's isolation level [8][7]. If a concurrent transaction has already inserted the value but has not yet committed, the second transaction will block and wait for the first transaction to complete [8][9]. Once the first transaction commits or rolls back, the second transaction will either detect the conflict and raise a unique violation error (or trigger the ON CONFLICT clause) [1][8]. For atomic, concurrent-safe insertions (often called UPSERT), the recommended approach is to use the INSERT... ON CONFLICT clause [10][11]. This syntax is explicitly designed to handle potential unique constraint violations atomically [10][11]: 1. INSERT... ON CONFLICT DO NOTHING: Silently ignores the insertion if a conflict occurs [10][11]. 2. INSERT... ON CONFLICT DO UPDATE: Performs an update on the existing row if a conflict occurs [10][11]. These clauses are more efficient and reliable than application-level checks or WHERE NOT EXISTS clauses, as they leverage the database's internal locking mechanisms to ensure the desired outcome [10][8][9].
Citations:
- 1: https://www.postgresql.org/message-id/86mwaph2pr.fsf%40jerry.enova.com
- 2: https://stackoverflow.com/questions/54248359/postgres-any-chance-of-duplicates-on-an-insert-select-where-not-exist-in-a-race
- 3: https://stackoverflow.com/questions/77594785/is-concurrent-insert-where-not-exists-in-postgres-safe
- 4: https://www.postgresql.org/docs/current/transaction-iso.html
- 5: https://www.postgresql.org/docs/18/transaction-iso.html
- 6: https://www.exchangetuts.com/index.php/is-concurrent-insert-where-not-exists-in-postgres-safe-1765993802855175
- 7: https://dba.stackexchange.com/questions/293764/transaction-strategy-for-unique-constraint
- 8: https://www.postgresql.org/message-id/CAJJ74VJZV5ramu2zxWfV_2YhpbVZuA7npT%3Deq7Ju272giDCeRQ%40mail.gmail.com
- 9: https://www.postgresql.org/docs/18/index-unique-checks.html
- 10: https://www.postgresql.org/docs/current/sql-insert.html
- 11: https://www.postgresql.org/docs/19/sql-insert.html
🏁 Script executed:
#!/bin/bash
set -euo pipefail
python3 - <<'PY'
from pathlib import Path
import re
scoped = Path("crates/common/src/db/scoped.rs").read_text()
workspaces = Path("crates/common/src/db/workspaces.rs").read_text()
migration = Path("crates/common/migrations/workspace_v1.sql").read_text()
assert "let mut tx = pool.begin().await?;" in scoped
reaper = re.search(
r"let existing: Option<\(String,\)>.*?fetch_optional\(&mut \*tx\)",
workspaces,
re.S,
)
assert reaper and "WHERE NOT EXISTS" in reaper.group(0)
unique_statements = re.findall(
r"CREATE UNIQUE INDEX\b.*?(?=;)",
migration,
re.I | re.S,
)
for statement in unique_statements:
print("UNIQUE:", " ".join(statement.split()))
assert not (
re.search(r"\bon\s+\{p\}jobs\s*\(\s*endpoint\s*\)", statement, re.I)
and re.search(r"\bstatus\s*=\s*'ACTIVE'\b", statement, re.I)
), "found an active endpoint uniqueness index"
print("scoped_transaction uses pool.begin() without an explicit isolation level")
print("reaper provisioning uses WHERE NOT EXISTS")
print("no unique index enforces one ACTIVE reaper job per endpoint")
PY
printf '%s\n' '--- workspace migration application ---'
rg -n -C 10 'workspace_v1\.sql|migrations/workspace|migrate_workspace|workspace migration|CREATE TABLE IF NOT EXISTS' crates/common crates/worker crates/apiRepository: juspay/kronos
Length of output: 12779
Enforce uniqueness for the active reaper job.
Under the default READ COMMITTED isolation, two callers can pass WHERE NOT EXISTS and insert duplicate active reaper jobs. Add a partial unique index scoped to kronos.reaper, then use ON CONFLICT DO NOTHING to make the losing transaction idempotent.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/common/src/db/workspaces.rs` around lines 219 - 239, The reaper
insertion in the workspace setup can race and create duplicate active jobs. Add
a partial unique index for the active reaper endpoint scoped to kronos.reaper,
then update the INSERT in the existing query flow to use ON CONFLICT DO NOTHING
while retaining the RETURNING-based registration path for the winning
transaction.
| TriggerPlan::Delayed { run_at } => { | ||
| let result = db::jobs::create_delayed( | ||
| &mut db, | ||
| req.endpoint, | ||
| &ep.endpoint_type, | ||
| req.idempotency_key, | ||
| req.input, | ||
| run_at, | ||
| max_attempts, | ||
| ) | ||
| .await?; |
There was a problem hiding this comment.
🎯 Functional Correctness | 🟠 Major | ⚡ Quick win
Map duplicate-key failures on the DELAYED path too.
The phase-1 idempotency check runs on a separate connection, before the transaction opens. Two concurrent requests with the same key can both pass it. The IMMEDIATE arm handles the resulting unique-index violation with duplicate_key_as_conflict and returns a conflict. The DELAYED arm does not, so the same race surfaces as ServiceError::Db, which REST renders as a 500 instead of a 409.
🐛 Proposed fix
run_at,
max_attempts,
)
- .await?;
+ .await
+ .map_err(duplicate_key_as_conflict)?;📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| TriggerPlan::Delayed { run_at } => { | |
| let result = db::jobs::create_delayed( | |
| &mut db, | |
| req.endpoint, | |
| &ep.endpoint_type, | |
| req.idempotency_key, | |
| req.input, | |
| run_at, | |
| max_attempts, | |
| ) | |
| .await?; | |
| TriggerPlan::Delayed { run_at } => { | |
| let result = db::jobs::create_delayed( | |
| &mut db, | |
| req.endpoint, | |
| &ep.endpoint_type, | |
| req.idempotency_key, | |
| req.input, | |
| run_at, | |
| max_attempts, | |
| ) | |
| .await | |
| .map_err(duplicate_key_as_conflict)?; |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/common/src/service/jobs.rs` around lines 226 - 236, Update the
TriggerPlan::Delayed arm around db::jobs::create_delayed to map duplicate-key
database errors through the existing duplicate_key_as_conflict helper, matching
the IMMEDIATE arm and returning a conflict for concurrent idempotency-key races
while preserving other database errors.
This pull request refactors the
jobsAPI handler to delegate job creation and cancellation logic to the shared service layer, reducing duplication and centralizing business logic. It also makes several improvements to idempotency handling, schema provisioning, and code clarity. The changes affect both the API and common crates, ensuring that job-related operations are consistent and easier to maintain.API handler refactor and delegation:
createandcancelfunctions incrates/api/src/handlers/jobs.rsto delegate job creation and cancellation logic to the sharedservice::jobsmodule, removing duplicated logic and centralizing validation, guard checks, and transactional operations. This includes moving input validation, endpoint checks, idempotency handling, and transaction boundaries into the service layer. [1] [2]updatefunction to use the sharedservice::jobs::compute_next_cronfor CRON job scheduling, ensuring consistent computation logic.Idempotency and job creation improvements:
create_immediateandcreate_delayedincrates/common/src/db/jobs.rsto acceptOption<&str>foridempotency_keyinstead of&str, preventing accidental collisions for keyless jobs and aligning with the partial unique index semantics. [1] [2]Schema provisioning and reaper job improvements:
provision_reaperincrates/common/src/db/workspaces.rspublic and added atable_prefixparameter to support both library and API deployments, ensuring the reaper job is provisioned into the correct tables. The function now uses parameterized table names and is idempotent, avoiding duplicate reaper jobs.provision_reaperwith an empty prefix in the API deployment, matching the schema-per-workspace setup.Code cleanup and modularization:
validate_inputand the localcompute_next_cron, since these are now handled in the service layer.servicemodule incrates/common/src/lib.rsfor shared workspace-mutation logic between REST and library clients.Summary by CodeRabbit
New Features
Bug Fixes