Skip to content

fix: library cron fixes - #63

Open
mahatoankitkumar wants to merge 1 commit into
mainfrom
fix/kronos-library-rest-parity
Open

mahatoankitkumar wants to merge 1 commit into
mainfrom
fix/kronos-library-rest-parity

Conversation

@mahatoankitkumar

@mahatoankitkumar mahatoankitkumar commented Aug 11, 2026 •

Copy link
Copy Markdown
Contributor

This pull request refactors the jobs API 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:

  • Refactored the create and cancel functions in crates/api/src/handlers/jobs.rs to delegate job creation and cancellation logic to the shared service::jobs module, 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]
  • Updated the update function to use the shared service::jobs::compute_next_cron for CRON job scheduling, ensuring consistent computation logic.

Idempotency and job creation improvements:

  • Changed the signature of create_immediate and create_delayed in crates/common/src/db/jobs.rs to accept Option<&str> for idempotency_key instead of &str, preventing accidental collisions for keyless jobs and aligning with the partial unique index semantics. [1] [2]

Schema provisioning and reaper job improvements:

  • Made provision_reaper in crates/common/src/db/workspaces.rs public and added a table_prefix parameter 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.
  • Updated workspace creation to call provision_reaper with an empty prefix in the API deployment, matching the schema-per-workspace setup.

Code cleanup and modularization:

  • Removed now-unnecessary local validation and helper functions from the API handler, such as validate_input and the local compute_next_cron, since these are now handled in the service layer.
  • Exposed the service module in crates/common/src/lib.rs for shared workspace-mutation logic between REST and library clients.

Summary by CodeRabbit

  • New Features

    • Added shared job creation and cancellation behavior across REST and library clients.
    • Added support for idempotent job creation, including safe replay handling.
    • Improved delayed and scheduled job validation, timezone-aware cron timing, and input validation.
    • Added endpoint registration options for payload and configuration references.
    • Added configurable workspace cleanup scheduling with repeat-safe provisioning.
  • Bug Fixes

    • Improved consistency of job creation, cancellation, scheduling, and error responses across client interfaces.

@coderabbitai

coderabbitai Bot commented Aug 11, 2026 •

Copy link
Copy Markdown

Review Change Stack

Walkthrough

The 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.

Changes

Shared job service

Layer / File(s) Summary
Service contracts and persistence inputs
crates/common/src/lib.rs, crates/common/src/service/*, crates/common/src/db/jobs.rs
Adds workspace references, transport-neutral service errors, job request and outcome types, and optional idempotency keys.
Job lifecycle service
crates/common/src/service/jobs.rs
Adds validation, idempotency handling, trigger resolution, transactional creation, cancellation, schema validation, and timezone-aware cron computation.
REST service wiring
crates/api/src/handlers/jobs.rs
Delegates job creation, cron computation, and cancellation to the shared service while preserving REST responses and metrics.
Prefix-aware reaper provisioning
crates/common/src/db/workspaces.rs, crates/worker/src/client.rs
Adds configurable reaper schedules, table-prefix support, duplicate-safe provisioning, and conditional cron registration.
Library client and endpoint references
crates/worker/src/client.rs
Delegates library job mutations to the shared service and adds optional payload-spec and config references to endpoint requests and tests.

Estimated code review effort: 4 (Complex) | ~60 minutes

Possibly related issues

Possibly related PRs

  • juspay/kronos#17 — Provides the job APIs and worker client code refactored by this change.
  • juspay/kronos#33 — Previously changed cron handler transaction handling now replaced by shared job services.
  • juspay/kronos#59 — Introduced the shared job-orchestration direction extended by this pull request.

Suggested reviewers: tsdk02, knutties

Poem

A rabbit carries jobs through shared service lanes,
With cron-time carrots and idempotent trains.
Reapers wake once, then sleep by the clock,
REST and library paths share one block.
Endpoints hold references, neat and bright—
The queue hops onward through the night.

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title describes real library CRON changes, but it does not summarize the broader shared service and API refactor.
Docstring Coverage ✅ Passed No functions found in the changed files to evaluate docstring coverage. Skipping docstring coverage check.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch fix/kronos-library-rest-parity

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Actionable comments posted: 2

🧹 Nitpick comments (4)
crates/common/src/service/mod.rs (1)

86-88: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Reuse the Display text instead of repeating the format string.

ServiceError::InternalEndpoint already renders this exact message through thiserror. 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 win

Match only idx_{prefix}jobs_idempotency with SQLSTATE 23505.

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 win

Do not panic on a missing execution in the HTTP path.

execution.expect(...) panics if the service returns Created without 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 to service::jobs::create_job turns 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 win

Validate the reaper CRON expression at configuration time.

with_reaper_schedule stores the string without validation. An invalid expression surfaces later, inside provision_workspace, as a pg_cron scheduling failure during boot. create_job already validates library CRON input with PgCronExpr::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

📥 Commits

Reviewing files that changed from the base of the PR and between 2f6770c and 19e6a49.

📒 Files selected for processing (7)
  • crates/api/src/handlers/jobs.rs
  • crates/common/src/db/jobs.rs
  • crates/common/src/db/workspaces.rs
  • crates/common/src/lib.rs
  • crates/common/src/service/jobs.rs
  • crates/common/src/service/mod.rs
  • crates/worker/src/client.rs

Comment on lines +219 to +239
// `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?;

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🗄️ 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' crates

Repository: 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/migrations

Repository: 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")
PY

Repository: 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:


🏁 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/api

Repository: 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.

Comment on lines +226 to +236
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?;

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🎯 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.

Suggested change
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 branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant