Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
73 commits
Select commit Hold shift + click to select a range
ff56be1
fix(scheduler): reject past one-time 'at' cron schedules
Aug 31, 2026
cee7133
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 2, 2026
69f7326
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 2, 2026
887215b
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 2, 2026
fcbee72
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 3, 2026
7b7292c
style(scheduler): apply current formatting
Sep 3, 2026
a1762a2
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 4, 2026
9558a1e
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 4, 2026
9ae4310
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 4, 2026
e3faf53
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 5, 2026
37231ee
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 5, 2026
7af664e
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 5, 2026
735b225
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 6, 2026
172d5c4
Merge commit 'refs/task-a/upstream/main' into agent-tasks/1516
Sep 7, 2026
4ae1088
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 7, 2026
a2d0ee1
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 7, 2026
37f7189
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 7, 2026
433b93a
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 8, 2026
a754b1c
Merge branch 'main' of https://github.com/opensquilla/opensquilla int…
Sep 8, 2026
04c3616
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 8, 2026
d95348f
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 9, 2026
653bb0b
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
e5b2717
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
c92274b
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
b524042
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
3eb4535
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
6061485
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
252d4b5
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
282c067
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
64724f2
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
47e331c
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
429c155
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
194e84f
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
d8a375f
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 10, 2026
90d0adb
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 11, 2026
c0424e7
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 11, 2026
a011723
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 11, 2026
dd0bf44
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 11, 2026
6f2e09f
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 11, 2026
cbe7566
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 11, 2026
bfee5d1
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 11, 2026
c36f218
Merge remote-tracking branch 'upstream/main' into agent-tasks/1516
Sep 12, 2026
f708313
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 13, 2026
6f1f5df
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 14, 2026
df3122d
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 14, 2026
89ae135
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 14, 2026
f79d39d
Merge commit 'refs/task-a/opensquilla/main' into agent-tasks/1516
Sep 14, 2026
90712b5
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 14, 2026
584e951
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 14, 2026
5e33055
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 14, 2026
6b1ee2c
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 14, 2026
44e4145
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
8297705
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
1076d7d
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
bae0185
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
9e21083
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
4a060e8
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
c87e1eb
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
6bf2594
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
9b6a3b2
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
582000d
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
28b6d7d
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 15, 2026
ec50f84
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 16, 2026
4073b8d
Keep cron validation atomic and preserve idempotent retries
Open-Squilla Sep 16, 2026
6a54361
Expose total_changes on the SQLite fallback connection
Open-Squilla Sep 16, 2026
8b74a46
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 16, 2026
45a0317
Record initial generated publication receipts atomically
Open-Squilla Sep 16, 2026
0af192a
Integrate validated mainline fixes for scheduler CI
Open-Squilla Sep 16, 2026
3efb407
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 16, 2026
473e268
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 16, 2026
a3c1257
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 16, 2026
1aa30ed
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 16, 2026
49cda42
Merge remote-tracking branch 'refs/remotes/upstream/main' into agent-…
Sep 16, 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
17 changes: 1 addition & 16 deletions src/opensquilla/gateway/rpc_cron.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,6 @@
DeliveryConfig,
DeliveryMode,
FailureDestination,
JobStatus,
ReplyTargetSnapshot,
ScheduleKind,
SessionTarget,
Expand Down Expand Up @@ -771,21 +770,7 @@ async def _update_cron_job(
patch["tz"] = tz_value if isinstance(tz_value, str) else ""

if "enabled" in params:
# Resolve the enabled toggle but DO NOT early-return: fall through so
# sibling field updates in the same request (text/schedule/…) are
# applied too, and so a DISABLED/FAILED job (not just PAUSED) can be
# revived — those are exactly the states a user runs `--enabled` to fix.
job = await scheduler.get_job(job_id)
if params["enabled"]:
revivable = {
JobStatus.PAUSED.value,
JobStatus.DISABLED.value,
JobStatus.FAILED.value,
}
if job is not None and job.status.value in revivable:
await scheduler.resume_job(job_id)
else:
await scheduler.pause_job(job_id)
patch["enabled"] = bool(params["enabled"])

current_job = await scheduler.get_job(job_id)
if current_job is None:
Expand Down
48 changes: 40 additions & 8 deletions src/opensquilla/scheduler/ops.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,21 @@
)


def _reject_past_at(cron_expr: str, now: datetime) -> None:
"""Reject a one-time ``at`` timestamp that is already in the past.

A past ``at`` fires on the next scheduler tick and the one-shot job is then
deleted, so the payload runs immediately with no future occurrence. That is
almost never what a caller scheduling a one-time reminder intends, so refuse
it at creation time instead of silently running it.
"""
at_dt = parse_iso_at(cron_expr)
if at_dt < now:
raise ValueError(
f"schedule.at is in the past: {cron_expr}; one-time schedules must be in the future"
)


def _validate_structured_schedule(
kind: ScheduleKind | str,
value: str,
Expand Down Expand Up @@ -211,11 +226,7 @@ async def add(
# fall back to ISOLATED instead of failing creation. Headless cron
# callers (no session context) get an isolated run rather than a hard
# error.
if (
session_target == SessionTarget.CURRENT
and not session_key
and not origin_session_key
):
if session_target == SessionTarget.CURRENT and not session_key and not origin_session_key:
session_target = SessionTarget.ISOLATED

origin_session_key = normalize_origin_session_key(session_target, origin_session_key)
Expand Down Expand Up @@ -281,7 +292,13 @@ async def add(
# CRON or EVERY with cron expression: scan forward
job.next_run_at = _next_run(job, now)

return await self._store.create_or_get(job)
# A retry may arrive after the original execution time. Only validate a
# new row, after deduplication and with a fresh clock inside the lock.
validate_new = (
(lambda: _reject_past_at(cron_expr, self._now()))
if kind == ScheduleKind.AT else None
)
return await self._store.create_or_get(job, validate_new=validate_new)

async def update(self, job_id: str, **patch) -> CronJob | None:
"""Apply a partial update to an existing job. Returns None if not found."""
Expand All @@ -302,7 +319,8 @@ async def update(self, job_id: str, **patch) -> CronJob | None:
structured_kind = patch.pop("schedule_kind", None)
structured_value = patch.pop("schedule_value", None)
structured_tz = patch.pop("schedule_tz", None)
if structured_kind is not None and structured_value is not None:
schedule_updated = structured_kind is not None and structured_value is not None
if schedule_updated:
kind, cron_expr = _validate_structured_schedule(structured_kind, structured_value)
if structured_tz is not None:
raw_tz = (structured_tz or "").strip()
Expand All @@ -312,6 +330,7 @@ async def update(self, job_id: str, **patch) -> CronJob | None:
job.schedule_kind = kind
job.cron_expr = cron_expr
if kind == ScheduleKind.AT:
_reject_past_at(cron_expr, now)
job.anchor_at = None
job.next_run_at = datetime.fromisoformat(cron_expr)
elif kind == ScheduleKind.EVERY:
Expand All @@ -326,7 +345,7 @@ async def update(self, job_id: str, **patch) -> CronJob | None:
"pass schedule_kind + schedule_value instead"
)

for field in ("name", "timeout_seconds", "enabled", "origin_session_key"):
for field in ("name", "timeout_seconds", "origin_session_key"):
if field in patch:
setattr(job, field, patch.pop(field))
if "tool_policy" in patch:
Expand Down Expand Up @@ -372,6 +391,19 @@ async def update(self, job_id: str, **patch) -> CronJob | None:
job.origin_session_key,
)

# Persist the enabled toggle with the validated schedule and payload.
# A rejected patch must never resume or pause the existing job.
if "enabled" in patch:
job.enabled = bool(patch.pop("enabled"))
if not job.enabled:
job.status = JobStatus.PAUSED
elif job.status in (JobStatus.PAUSED, JobStatus.DISABLED, JobStatus.FAILED):
job.status = JobStatus.PENDING
job.backoff_until = None
job.consecutive_errors = 0
if not schedule_updated and job.schedule_kind != ScheduleKind.AT:
job.next_run_at = _next_run(job, now)

job.updated_at = now
await self._store.save(job)
return job
Expand Down
12 changes: 9 additions & 3 deletions src/opensquilla/scheduler/persistence.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
import asyncio
import json
import uuid
from collections.abc import AsyncIterator
from collections.abc import AsyncIterator, Callable
from contextlib import asynccontextmanager
from datetime import UTC, datetime

Expand Down Expand Up @@ -620,9 +620,13 @@ async def save(self, job: CronJob) -> None:
await self._execute_save(job)
await self._db().commit()

async def create_or_get(self, job: CronJob) -> CronJob:
"""Atomically create an idempotent job or return the existing row."""
async def create_or_get(
self, job: CronJob, *, validate_new: Callable[[], None] | None = None,
) -> CronJob:
"""Return an existing row, or validate and create under the idempotency lock."""
if not job.idempotency_key:
if validate_new is not None:
validate_new()
await self.save(job)
return job

Expand All @@ -631,6 +635,8 @@ async def create_or_get(self, job: CronJob) -> CronJob:
if existing is not None:
existing.deduplicated = True
return existing
if validate_new is not None:
validate_new()
try:
await self._execute_save(job)
await self._db().commit()
Expand Down
102 changes: 102 additions & 0 deletions tests/test_gateway/test_rpc_cron_update_enabled.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,11 @@

from __future__ import annotations

from datetime import UTC, datetime, timedelta
from pathlib import Path
from unittest.mock import AsyncMock

import pytest

from opensquilla.gateway.rpc import RpcContext
from opensquilla.gateway.rpc_cron import (
Expand Down Expand Up @@ -150,3 +154,101 @@ async def test_update_enabled_false_applies_sibling_fields(tmp_path: Path) -> No
assert payload_text(after.payload, after.session_target) == "new prompt"
finally:
await store.close()


@pytest.mark.parametrize(
("status", "enabled"),
[
(JobStatus.PENDING, False),
(JobStatus.PAUSED, True),
(JobStatus.DISABLED, True),
(JobStatus.FAILED, True),
],
)
@pytest.mark.parametrize(
("patch", "error"),
[
({"schedule": {"kind": "at", "at": "2035-01-01T19:59:59+08:00"}}, "in the past"),
({"tz": "Invalid/Timezone"}, "Unknown timezone"),
],
ids=["past-at", "invalid-timezone"],
)
async def test_update_invalid_patch_does_not_change_enabled_or_persist(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
status: JobStatus,
enabled: bool,
patch: dict,
error: str,
) -> None:
now = datetime(2035, 1, 1, 12, tzinfo=UTC)
store = JobStore(str(tmp_path / "cron.db"))
await store.open()
engine = SchedulerEngine(store, clock=lambda: now)
try:
job = await engine.add_job(
name="original",
schedule_kind=ScheduleKind.AT,
schedule_value=(now + timedelta(hours=1)).isoformat(),
payload=make_agent_turn_payload("original reminder"),
)
job.status = status
job.enabled = not enabled
job.consecutive_errors = 2
job.backoff_until = now + timedelta(minutes=1)
await store.save(job)
before = await store.get(job.id)
save = AsyncMock(wraps=store.save)
monkeypatch.setattr(store, "save", save)

with pytest.raises(ValueError, match=error):
await _handle_cron_update(
{"id": job.id, "enabled": enabled, "name": "changed", **patch},
_ctx(engine),
)

assert await store.get(job.id) == before
save.assert_not_awaited()
finally:
await store.close()


@pytest.mark.parametrize("enabled", [False, True])
async def test_update_enabled_and_schedule_persist_together_once(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch, enabled: bool,
) -> None:
now = datetime(2035, 1, 1, 12, tzinfo=UTC)
store = JobStore(str(tmp_path / "cron.db"))
await store.open()
engine = SchedulerEngine(store, clock=lambda: now)
try:
job = await engine.add_job(
name="original",
enabled=not enabled,
schedule_kind=ScheduleKind.AT,
schedule_value=(now + timedelta(hours=1)).isoformat(),
payload=make_agent_turn_payload("original reminder"),
)
save = AsyncMock(wraps=store.save)
monkeypatch.setattr(store, "save", save)
new_at = (now + timedelta(hours=2)).isoformat()

await _handle_cron_update(
{
"id": job.id,
"enabled": enabled,
"text": "updated reminder",
"schedule": {"kind": "at", "at": new_at},
},
_ctx(engine),
)

after = await store.get(job.id)
assert after is not None
assert after.status == (JobStatus.PENDING if enabled else JobStatus.PAUSED)
assert after.enabled is enabled
assert after.next_run_at == now + timedelta(hours=2)
assert payload_text(after.payload, after.session_target) == "updated reminder"
save.assert_awaited_once()
finally:
await store.close()
Loading
Loading