From ff56be1e641173e16d1d3b8cfb0cad23f205652c Mon Sep 17 00:00:00 2001 From: "liumingyao.marvin" Date: Mon, 31 Aug 2026 23:00:59 +0800 Subject: [PATCH 1/5] fix(scheduler): reject past one-time 'at' cron schedules MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A one-time 'at' schedule with a timestamp already in the past is accepted, fires on the next scheduler tick, and the one-shot job is then deleted — so the payload runs immediately with no future occurrence and disappears. That is almost never what a caller scheduling a one-time reminder intends. Validate the 'at' timestamp against the current time in SchedulerOps.add and SchedulerOps.update (the two structured-schedule entry points) and raise a field-named ValueError before persisting, matching the existing timezone and interval validation style. Recurring cron/every schedules are unaffected. Fixes #1516 --- src/opensquilla/scheduler/ops.py | 18 ++++++++ .../test_ops_strict_schedule.py | 45 +++++++++++++++++++ 2 files changed, 63 insertions(+) diff --git a/src/opensquilla/scheduler/ops.py b/src/opensquilla/scheduler/ops.py index 3a55b5da6f..86e2816006 100644 --- a/src/opensquilla/scheduler/ops.py +++ b/src/opensquilla/scheduler/ops.py @@ -32,6 +32,22 @@ ) +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, @@ -270,6 +286,7 @@ async def add( ) if kind == ScheduleKind.AT: + _reject_past_at(cron_expr, now) job.delete_after_run = True job.next_run_at = datetime.fromisoformat(cron_expr) elif kind == ScheduleKind.EVERY and cron_expr.isdigit(): @@ -312,6 +329,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: diff --git a/tests/test_scheduler/test_ops_strict_schedule.py b/tests/test_scheduler/test_ops_strict_schedule.py index 4ac451e23e..04aa36c989 100644 --- a/tests/test_scheduler/test_ops_strict_schedule.py +++ b/tests/test_scheduler/test_ops_strict_schedule.py @@ -326,6 +326,51 @@ async def test_ops_add_at_rejects_naive_iso(tmp_path: Path) -> None: await store.close() +async def test_ops_add_at_rejects_past_timestamp(tmp_path: Path) -> None: + """A one-time ``at`` in the past would fire immediately then delete itself; + reject it at creation time instead (issue #1516).""" + store, ops = await _open_ops(tmp_path) + try: + past = (datetime.now(UTC) - timedelta(hours=1)).isoformat() + with pytest.raises(ValueError, match="in the past"): + await ops.add( + name="stale", + handler_key="agent_run", + payload=make_agent_turn_payload("ping"), + session_target=SessionTarget.ISOLATED, + schedule_kind=ScheduleKind.AT, + schedule_value=past, + ) + # Nothing should have been persisted. + assert await store.list_active() == [] + finally: + await store.close() + + +async def test_ops_update_at_rejects_past_timestamp(tmp_path: Path) -> None: + """Repointing a job to a past one-time ``at`` is rejected as well.""" + store, ops = await _open_ops(tmp_path) + try: + future = (datetime.now(UTC) + timedelta(hours=1)).isoformat() + job = await ops.add( + name="once", + handler_key="agent_run", + payload=make_agent_turn_payload("ping"), + session_target=SessionTarget.ISOLATED, + schedule_kind=ScheduleKind.AT, + schedule_value=future, + ) + past = (datetime.now(UTC) - timedelta(hours=1)).isoformat() + with pytest.raises(ValueError, match="in the past"): + await ops.update( + job.id, + schedule_kind=ScheduleKind.AT, + schedule_value=past, + ) + finally: + await store.close() + + async def test_ops_add_every_rejects_zero_seconds(tmp_path: Path) -> None: store, ops = await _open_ops(tmp_path) try: From 7b7292c24ec06e99d5b690b8053aa7529ea5027a Mon Sep 17 00:00:00 2001 From: "liumingyao.marvin" Date: Thu, 3 Sep 2026 14:45:02 +0800 Subject: [PATCH 2/5] style(scheduler): apply current formatting --- src/opensquilla/scheduler/ops.py | 9 ++------- tests/test_scheduler/test_ops_strict_schedule.py | 4 +--- 2 files changed, 3 insertions(+), 10 deletions(-) diff --git a/src/opensquilla/scheduler/ops.py b/src/opensquilla/scheduler/ops.py index 86e2816006..602540e76c 100644 --- a/src/opensquilla/scheduler/ops.py +++ b/src/opensquilla/scheduler/ops.py @@ -43,8 +43,7 @@ def _reject_past_at(cron_expr: str, now: datetime) -> None: 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" + f"schedule.at is in the past: {cron_expr}; one-time schedules must be in the future" ) @@ -227,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) diff --git a/tests/test_scheduler/test_ops_strict_schedule.py b/tests/test_scheduler/test_ops_strict_schedule.py index 04aa36c989..5484c3d734 100644 --- a/tests/test_scheduler/test_ops_strict_schedule.py +++ b/tests/test_scheduler/test_ops_strict_schedule.py @@ -132,9 +132,7 @@ async def test_cron_creator_authority_survives_persistence_without_widening_owne creator_is_owner or expected_persisted_host_execute ) assert bool(envelope.metadata.get("cron_trusted_owner")) is creator_is_owner - assert bool(envelope.metadata.get("cron_trusted_host")) is ( - expected_persisted_host_execute - ) + assert bool(envelope.metadata.get("cron_trusted_host")) is (expected_persisted_host_execute) assert bool(envelope.metadata.get(PRINCIPAL_HOST_EXECUTE_METADATA_KEY)) is ( expected_persisted_host_execute ) From 4073b8d1b85ca77a1d0c94675eec37b83002c29a Mon Sep 17 00:00:00 2001 From: Open-Squilla <275096992+Open-Squilla@users.noreply.github.com> Date: Wed, 16 Sep 2026 17:47:57 +0800 Subject: [PATCH 3/5] Keep cron validation atomic and preserve idempotent retries --- src/opensquilla/gateway/rpc_cron.py | 17 +-- src/opensquilla/scheduler/ops.py | 27 ++++- src/opensquilla/scheduler/persistence.py | 12 ++- .../test_rpc_cron_update_enabled.py | 102 ++++++++++++++++++ .../test_ops_strict_schedule.py | 102 ++++++++++++++++++ 5 files changed, 237 insertions(+), 23 deletions(-) diff --git a/src/opensquilla/gateway/rpc_cron.py b/src/opensquilla/gateway/rpc_cron.py index c9b83057f0..9164e69a77 100644 --- a/src/opensquilla/gateway/rpc_cron.py +++ b/src/opensquilla/gateway/rpc_cron.py @@ -61,7 +61,6 @@ DeliveryConfig, DeliveryMode, FailureDestination, - JobStatus, ReplyTargetSnapshot, ScheduleKind, SessionTarget, @@ -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: diff --git a/src/opensquilla/scheduler/ops.py b/src/opensquilla/scheduler/ops.py index 602540e76c..be093aa7f4 100644 --- a/src/opensquilla/scheduler/ops.py +++ b/src/opensquilla/scheduler/ops.py @@ -281,7 +281,6 @@ async def add( ) if kind == ScheduleKind.AT: - _reject_past_at(cron_expr, now) job.delete_after_run = True job.next_run_at = datetime.fromisoformat(cron_expr) elif kind == ScheduleKind.EVERY and cron_expr.isdigit(): @@ -293,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.""" @@ -314,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() @@ -339,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: @@ -385,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 diff --git a/src/opensquilla/scheduler/persistence.py b/src/opensquilla/scheduler/persistence.py index 81bb4cdc36..ca714433ad 100644 --- a/src/opensquilla/scheduler/persistence.py +++ b/src/opensquilla/scheduler/persistence.py @@ -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 @@ -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 @@ -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() diff --git a/tests/test_gateway/test_rpc_cron_update_enabled.py b/tests/test_gateway/test_rpc_cron_update_enabled.py index d5885f9650..fc2f6864d6 100644 --- a/tests/test_gateway/test_rpc_cron_update_enabled.py +++ b/tests/test_gateway/test_rpc_cron_update_enabled.py @@ -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 ( @@ -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() diff --git a/tests/test_scheduler/test_ops_strict_schedule.py b/tests/test_scheduler/test_ops_strict_schedule.py index 5484c3d734..2fb8504307 100644 --- a/tests/test_scheduler/test_ops_strict_schedule.py +++ b/tests/test_scheduler/test_ops_strict_schedule.py @@ -358,6 +358,7 @@ async def test_ops_update_at_rejects_past_timestamp(tmp_path: Path) -> None: schedule_kind=ScheduleKind.AT, schedule_value=future, ) + before = await store.get(job.id) past = (datetime.now(UTC) - timedelta(hours=1)).isoformat() with pytest.raises(ValueError, match="in the past"): await ops.update( @@ -365,6 +366,107 @@ async def test_ops_update_at_rejects_past_timestamp(tmp_path: Path) -> None: schedule_kind=ScheduleKind.AT, schedule_value=past, ) + assert await store.get(job.id) == before + finally: + await store.close() + + +async def test_ops_at_idempotent_retries_after_due_time_return_existing_job(tmp_path: Path) -> None: + store, _ = await _open_ops(tmp_path) + now = datetime(2035, 1, 1, 12, tzinfo=UTC) + ops = SchedulerOps(store, clock=lambda: now) + at = (now + timedelta(seconds=1)).isoformat() + try: + async def add_once(key: str = "same-request"): + return await ops.add( + name="once", + payload=make_agent_turn_payload("synthetic reminder"), + schedule_kind=ScheduleKind.AT, + schedule_value=at, + idempotency_key=key, + ) + + original = await add_once() + now += timedelta(seconds=2) + retries = await asyncio.gather(*(add_once() for _ in range(8))) + + assert all(job.id == original.id and job.deduplicated for job in retries) + with pytest.raises(ValueError, match="in the past"): + await add_once("different-request") + assert len(await store.list_active()) == 1 + finally: + await store.close() + + +async def test_ops_at_checks_time_after_waiting_for_creation_lock(tmp_path: Path) -> None: + store, _ = await _open_ops(tmp_path) + now = datetime(2035, 1, 1, 12, tzinfo=UTC) + started = asyncio.Event() + + def clock() -> datetime: + started.set() + return now + + ops = SchedulerOps(store, clock=clock) + try: + async with store._idempotent_create_lock: + pending = asyncio.create_task(ops.add( + name="once", + payload=make_agent_turn_payload("synthetic reminder"), + schedule_kind=ScheduleKind.AT, + schedule_value=(now + timedelta(seconds=1)).isoformat(), + idempotency_key="new-request", + )) + await started.wait() + now += timedelta(seconds=2) + with pytest.raises(ValueError, match="in the past"): + await pending + assert await store.list_active() == [] + finally: + await store.close() + + +@pytest.mark.parametrize( + "at", ["2035-01-01T12:00:00Z", "2035-01-01T20:00:00+08:00"], +) +async def test_ops_at_accepts_the_current_instant_with_either_offset( + tmp_path: Path, at: str, +) -> None: + store, _ = await _open_ops(tmp_path) + now = datetime(2035, 1, 1, 12, tzinfo=UTC) + ops = SchedulerOps(store, clock=lambda: now) + try: + job = await ops.add( + name="once", + payload=make_agent_turn_payload("synthetic reminder"), + schedule_kind=ScheduleKind.AT, + schedule_value=at, + ) + assert job.next_run_at == now + finally: + await store.close() + + +async def test_ops_overdue_at_metadata_update_preserves_schedule(tmp_path: Path) -> None: + store, _ = await _open_ops(tmp_path) + now = datetime(2035, 1, 1, 12, tzinfo=UTC) + ops = SchedulerOps(store, clock=lambda: now) + try: + job = await ops.add( + name="once", + payload=make_agent_turn_payload("synthetic reminder"), + schedule_kind=ScheduleKind.AT, + schedule_value="2035-01-01T21:00:00+08:00", + ) + assert job.next_run_at == now + timedelta(hours=1) + now += timedelta(hours=2) + + changed = await ops.update(job.id, name="renamed", enabled=False) + + assert changed is not None + assert changed.name == "renamed" + assert changed.next_run_at == job.next_run_at + assert changed.enabled is False finally: await store.close() From 6a543612ba36fae91088b9f3352813a261e14b15 Mon Sep 17 00:00:00 2001 From: Open-Squilla <275096992+Open-Squilla@users.noreply.github.com> Date: Wed, 16 Sep 2026 18:10:13 +0800 Subject: [PATCH 4/5] Expose total_changes on the SQLite fallback connection --- src/opensquilla/compat/aiosqlite.py | 7 ++++ tests/test_compat/test_aiosqlite.py | 51 +++++++++++++++++++++++++++++ 2 files changed, 58 insertions(+) diff --git a/src/opensquilla/compat/aiosqlite.py b/src/opensquilla/compat/aiosqlite.py index 8008e2cd51..5a77775826 100644 --- a/src/opensquilla/compat/aiosqlite.py +++ b/src/opensquilla/compat/aiosqlite.py @@ -80,6 +80,9 @@ def row_factory(self, value: Any) -> None: ... @property def in_transaction(self) -> bool: ... + @property + def total_changes(self) -> int: ... + def execute(self, sql: str, params: Iterable[Any] = ()) -> CursorContext: ... def executemany( @@ -211,6 +214,10 @@ def row_factory(self, value: Any) -> None: def in_transaction(self) -> bool: return self._conn.in_transaction + @property + def total_changes(self) -> int: + return self._conn.total_changes + async def _execute(self, sql: str, params: Iterable[Any] = ()) -> _AsyncCursor: async with self._locked: cursor = await asyncio.to_thread(self._conn.execute, sql, tuple(params)) diff --git a/tests/test_compat/test_aiosqlite.py b/tests/test_compat/test_aiosqlite.py index 175bfb661a..53fa1cc827 100644 --- a/tests/test_compat/test_aiosqlite.py +++ b/tests/test_compat/test_aiosqlite.py @@ -50,3 +50,54 @@ async def test_sqlite3_fallback_execute_supports_await_and_async_context( await conn.close() monkeypatch.delenv("OPENSQUILLA_FORCE_SQLITE3_BACKEND", raising=False) importlib.reload(aiosqlite) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("backend", ["native", "sqlite3"]) +async def test_total_changes_tracks_writes_and_survives_rollback( + monkeypatch: pytest.MonkeyPatch, backend: str, +) -> None: + monkeypatch.setattr(aiosqlite, "_FORCE_SQLITE3_FALLBACK", backend == "sqlite3") + monkeypatch.setattr(aiosqlite, "_prefer_native", backend == "native") + conn = await aiosqlite.connect(":memory:") + try: + if backend == "native": + assert isinstance(conn, aiosqlite._native_aiosqlite.Connection) + else: + assert isinstance(conn, aiosqlite._AsyncConnection) + + assert conn.total_changes == 0 + await conn.execute("CREATE TABLE items (id INTEGER PRIMARY KEY, name TEXT NOT NULL)") + async with conn.execute("SELECT COUNT(*) FROM items") as cursor: + assert (await cursor.fetchone())[0] == 0 + assert conn.total_changes == 0 + + await conn.execute("INSERT INTO items VALUES (1, 'alpha')") + assert conn.total_changes == 1 + await conn.executemany( + "INSERT INTO items VALUES (?, ?)", [(2, "beta"), (3, "gamma")], + ) + await conn.commit() + assert conn.total_changes == 3 + + await conn.execute("UPDATE items SET name = 'changed' WHERE id = 1") + assert conn.total_changes == 4 + await conn.rollback() + async with conn.execute("SELECT name FROM items WHERE id = 1") as cursor: + assert (await cursor.fetchone())[0] == "alpha" + # SQLite counts writes even when their transaction is rolled back. + assert conn.total_changes == 4 + + await conn.execute("DELETE FROM items WHERE id = 2") + assert conn.total_changes == 5 + await conn.rollback() + async with conn.execute("SELECT COUNT(*), total_changes() FROM items") as cursor: + row = await cursor.fetchone() + assert tuple(row) == (3, 5) + assert conn.total_changes == 5 + + with pytest.raises(AttributeError): + conn.total_changes = 0 + assert conn.total_changes == 5 + finally: + await conn.close() From 45a0317f1141d7b9f8a78bcaec5de2de77105d90 Mon Sep 17 00:00:00 2001 From: Open-Squilla <275096992+Open-Squilla@users.noreply.github.com> Date: Wed, 16 Sep 2026 18:41:06 +0800 Subject: [PATCH 5/5] Record initial generated publication receipts atomically --- .../artifact_session/repository.py | 19 ++++++- src/opensquilla/artifact_session/service.py | 5 ++ .../gateway/generated_artifact_adoption.py | 1 + .../gateway/workbench_resource_runtime.py | 2 + .../test_generated_source_continuity.py | 51 +++++++++++++++++++ 5 files changed, 77 insertions(+), 1 deletion(-) diff --git a/src/opensquilla/artifact_session/repository.py b/src/opensquilla/artifact_session/repository.py index e40726ef31..d672614973 100644 --- a/src/opensquilla/artifact_session/repository.py +++ b/src/opensquilla/artifact_session/repository.py @@ -3,6 +3,7 @@ from __future__ import annotations import asyncio +import hashlib import json import secrets import time @@ -568,6 +569,7 @@ async def _append_audit( lease_id: str | None = None, payload: dict[str, Any] | None = None, created_at: int | None = None, + event_id: str | None = None, ) -> None: await conn.execute( """ @@ -578,7 +580,7 @@ async def _append_audit( ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( - self._id_factory("audit"), + event_id if event_id is not None else self._id_factory("audit"), document_id, event_type, actor.kind.value, @@ -827,6 +829,7 @@ async def adopt_generated_deliverable( deliverable: ArtifactBlobRef, actor: Actor, working_source: dict[str, str] | None = None, + publication_id: str = "", ) -> tuple[CommitResult, DocumentSourceBinding, bool]: """Atomically adopt and bind one public generated deliverable. @@ -980,6 +983,20 @@ async def adopt_generated_deliverable( conn, commit=commit, session_key=session_key, session_id=session_id, working_source=working_source, ) + if publication_id and not rows: + # The initial snapshot is already this occurrence's version. + # Record its receipt before another publisher can advance the + # head, so delayed adoption cannot reapply the old snapshot. + identity = f"{commit.document.document_id}\0{publication_id}".encode() + await self._append_audit( + conn, + document_id=commit.document.document_id, + event_type="document.source_published", + event_id=f"working-publish:{hashlib.sha256(identity).hexdigest()}", + actor=actor, + revision_id=commit.revision.revision_id, + payload={"artifact_id": deliverable.artifact_id}, + ) return commit, binding, True async def retire_legacy_html_state(self) -> None: diff --git a/src/opensquilla/artifact_session/service.py b/src/opensquilla/artifact_session/service.py index 7ab6222f06..d2945c37c0 100644 --- a/src/opensquilla/artifact_session/service.py +++ b/src/opensquilla/artifact_session/service.py @@ -265,6 +265,7 @@ async def adopt_generated_deliverable( deliverable: ArtifactBlobRef, actor: Actor, working_source: dict[str, str] | None = None, + publication_id: str = "", ) -> tuple[CommitResult, DocumentSourceBinding, bool]: """Adopt one immutable generated deliverable as a stable Document. @@ -282,6 +283,10 @@ async def adopt_generated_deliverable( deliverable=_blob(deliverable), actor=_actor(actor), working_source=working_source, + publication_id=( + _bounded_text(publication_id, "publication_id", max_bytes=512) + if publication_id else "" + ), ) async def reserve_document_import_attempt( diff --git a/src/opensquilla/gateway/generated_artifact_adoption.py b/src/opensquilla/gateway/generated_artifact_adoption.py index 6cc44a5e6f..fc019ba2d1 100644 --- a/src/opensquilla/gateway/generated_artifact_adoption.py +++ b/src/opensquilla/gateway/generated_artifact_adoption.py @@ -319,6 +319,7 @@ async def __call__(self, event: ArtifactEvent) -> None: session_id=self.session_id, ref=ref, working_source=source, + publication_id=event.publication_id, ) if adopted is None: return diff --git a/src/opensquilla/gateway/workbench_resource_runtime.py b/src/opensquilla/gateway/workbench_resource_runtime.py index c798e3e26e..4d5b2a56a6 100644 --- a/src/opensquilla/gateway/workbench_resource_runtime.py +++ b/src/opensquilla/gateway/workbench_resource_runtime.py @@ -557,6 +557,7 @@ async def adopt_generated_deliverable_if_editable( ref: ArtifactRef, actor: Actor | None = None, working_source: dict[str, str] | None = None, + publication_id: str = "", ) -> tuple[Document, Revision, DocumentSourceBinding, bool] | None: """Associate a generated HTML deliverable and its complete resource bundle with a document. @@ -615,6 +616,7 @@ async def adopt_generated_deliverable_if_editable( ), actor=actor or Actor(kind=ActorKind.SYSTEM, actor_id="generated-deliverable"), working_source=working_source, + publication_id=publication_id, ) return commit.document, commit.revision, binding, created diff --git a/tests/test_gateway/test_generated_source_continuity.py b/tests/test_gateway/test_generated_source_continuity.py index 93cbc2b46f..699eccc26a 100644 --- a/tests/test_gateway/test_generated_source_continuity.py +++ b/tests/test_gateway/test_generated_source_continuity.py @@ -472,6 +472,57 @@ async def test_concurrent_first_source_publications_create_one_document( assert (await cursor.fetchone())[0] == 1 +@pytest.mark.parametrize("changed", [False, True]) +async def test_initial_source_receipt_survives_a_concurrent_later_publication( + source_environment, monkeypatch: pytest.MonkeyPatch, changed: bool, +): + import opensquilla.gateway.generated_artifact_adoption as adoption + + env = source_environment + adopters, events = [], [] + for index in range(2): + context, adopter = _turn(env.service, env.store, env.workspace, env.media, env.preview) + if index and changed: + env.style.write_text("h1{color:crimson}") + token = current_tool_context.set(context) + try: + await publish_artifact(path="index.html") + finally: + current_tool_context.reset(token) + adopters.append(adopter) + events.append(ArtifactEvent(**context.published_artifacts[0])) + + created, release = asyncio.Event(), asyncio.Event() + original = adoption.adopt_generated_deliverable_if_editable + + async def pause_after_creation(**kwargs): + result = await original(**kwargs) + if result is not None and result[3]: + created.set() + await release.wait() + return result + + monkeypatch.setattr(adoption, "adopt_generated_deliverable_if_editable", pause_after_creation) + first = asyncio.create_task(adopters[0](events[0])) + try: + await asyncio.wait_for(created.wait(), timeout=5) + await adopters[1](events[1]) + finally: + release.set() + await first + + binding = await _binding(env) + revisions = await env.service.list_revisions(binding.document_id) + assert len(revisions) == (2 if changed else 1) + head = await env.service.get_document_head(binding.document_id) + assert head.revision.artifact_id == events[1 if changed else 0].id + # Replaying either completed occurrence must not revert the newer version. + await adopters[0](events[0]) + await adopters[1](events[1]) + assert await env.service.get_document_head(binding.document_id) == head + assert await env.service.list_revisions(binding.document_id) == revisions + + async def test_source_identity_is_scoped_to_session_and_workspace(source_environment): env = source_environment await _publish(env.context, env.adopter)