Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
74 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
a1d4213
Integrate validated one-shot schedule fixes
Open-Squilla Sep 16, 2026
87f5dac
Merge branch 'main' of https://github.com/TokenRhythm/opensquilla int…
Open-Squilla Sep 16, 2026
05cdba6
Wait for rendered frames before preview pointer input
Open-Squilla 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
25 changes: 19 additions & 6 deletions desktop/electron/scripts/live-html-journey-business.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -64,11 +64,22 @@ async function semanticControl(spec) {
const element = candidates[0] || pool.find(item => name(item).includes(wanted))
if (!element) throw new Error('SEMANTIC_CONTROL_MISSING')
element.scrollIntoView({ block: 'center', inline: 'nearest', behavior: 'instant' })
let rect = element.getBoundingClientRect()
const nextFrameRect = () => new Promise((resolve, reject) => {
let frame
const unavailable = () => reject(new Error('SEMANTIC_CONTROL_UNAVAILABLE'))
const remaining = spec.deadline - Date.now()
if (remaining <= 0) { unavailable(); return }
const timer = setTimeout(() => { cancelAnimationFrame(frame); unavailable() }, remaining)
frame = requestAnimationFrame(() => {
clearTimeout(timer)
if (Date.now() >= spec.deadline) unavailable()
else resolve(element.getBoundingClientRect())
})
})
let rect = await nextFrameRect()
let stable = false
for (let attempt = 0; attempt < 8; attempt++) {
await new Promise(resolve => setTimeout(resolve, 16))
const next = element.getBoundingClientRect()
const next = await nextFrameRect()
stable = ['x', 'y', 'width', 'height'].every(key => Math.abs(next[key] - rect[key]) <= 0.5)
rect = next
if (stable) break
Expand Down Expand Up @@ -101,12 +112,14 @@ export function createBusinessDriver(app, getPreviewId, capture) {
const result = await app.evaluate(async ({ webContents }, request) => {
const contents = webContents.fromId(request.id)
if (!contents || contents.isDestroyed()) throw new Error('SELECTED_PREVIEW_DESTROYED')
let point
let point, lastSemanticError
const until = Date.now() + 4000
while (!point) {
try { point = await contents.executeJavaScript(`(${request.resolver})(${JSON.stringify({ ...request.spec, kind: request.kind })})`, true) }
try { point = await contents.executeJavaScript(`(${request.resolver})(${JSON.stringify({ ...request.spec, kind: request.kind, deadline: until })})`, true) }
catch (error) {
if (!/SEMANTIC_CONTROL_(?:MISSING|UNAVAILABLE|COVERED)\b/.test(String(error)) || Date.now() >= until) throw error
if (!/SEMANTIC_CONTROL_(?:MISSING|UNAVAILABLE|COVERED)\b/.test(String(error))) throw error
if (Date.now() >= until) throw lastSemanticError || error
lastSemanticError = error
await new Promise(resolve => setTimeout(resolve, 100))
}
}
Expand Down
99 changes: 99 additions & 0 deletions desktop/electron/scripts/test-live-html-business.mjs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import assert from 'node:assert/strict'
import { readFile } from 'node:fs/promises'
import { after, before, test } from 'node:test'
import { runInNewContext } from 'node:vm'
import { chromium } from 'playwright'
import { createBusinessDriver, verifyBusinessCase } from './live-html-journey-business.mjs'
import { createSupplementalClient } from './live-html-supplemental-client.mjs'
Expand Down Expand Up @@ -49,6 +50,104 @@ test('moving modal close waits for actionability and receives a real pointer cli
} finally { await f.close() }
})

function controlledFrameDriver({ renderFrames = true, covered = false } = {}) {
let now = 0, frame = 0, nextId = 0, pumping = false
const tasks = new Map(), inputs = [], deadlines = [], cancelledFrames = []
const pump = () => {
if (pumping || ![...tasks.values()].some(task => Number.isFinite(task.at))) return
pumping = true
setImmediate(() => {
pumping = false
const next = [...tasks.entries()].filter(([, task]) => Number.isFinite(task.at)).sort((a, b) => a[1].at - b[1].at)[0]
if (!next) return
const [id, task] = next
tasks.delete(id)
now = task.at
if (task.frame) frame++
task.callback(now)
pump()
})
}
const schedule = (callback, at, isFrame = false) => {
const id = ++nextId
tasks.set(id, { callback, at, frame: isFrame })
pump()
return id
}
const clock = {
Date: { now: () => now },
setTimeout(callback, milliseconds) { deadlines.push(now + milliseconds); return schedule(callback, now + milliseconds) },
clearTimeout(id) { tasks.delete(id) },
// At 30 fps, a 16 ms timer can run without a new rendering frame.
requestAnimationFrame(callback) { return schedule(callback, renderFrames ? now + 1000 / 30 : Infinity, true) },
cancelAnimationFrame(id) { cancelledFrames.push(id); tasks.delete(id) },
}
class Control {
tagName = 'BUTTON'
innerText = 'Close'
isConnected = true
getBoundingClientRect() {
const x = Math.max(100, 300 - frame * 100)
return { x, y: 100, width: 40, height: 20, left: x, right: x + 40, top: 100, bottom: 120 }
}
getAttribute() { return null }
hasAttribute() { return false }
scrollIntoView() {}
matches() { return false }
contains(other) { return other === this }
}
const button = new Control()
const renderer = {
...clock, Element: Control,
document: { getElementById: () => null, querySelectorAll: () => [button], elementFromPoint: () => covered ? null : button },
getComputedStyle: () => ({ display: 'block', visibility: 'visible', opacity: '1' }),
innerWidth: 800, innerHeight: 700,
}
const contents = {
isDestroyed: () => false,
focus() {},
executeJavaScript: expression => runInNewContext(expression, renderer),
debugger: { isAttached: () => true, async sendCommand(command, args) { inputs.push({ command, ...args }) } },
}
const app = {
evaluate: (fn, request) => runInNewContext(`(${fn.toString()})(environment, request)`, {
...clock, request, environment: { webContents: { fromId: () => contents } },
}),
}
return {
driver: createBusinessDriver(app, () => 1, async () => {}),
inputs, deadlines, cancelledFrames,
elapsed: () => now,
pending: () => tasks.size,
}
}

test('timer samples within one rendering frame do not make moving controls actionable', async () => {
const f = controlledFrameDriver()
const point = await f.driver.click('Close')
assert.equal(point.x, 120, 'click must use the position after movement stops between frames')
assert.equal(f.inputs.length, 3)
assert.ok(f.inputs.every(input => input.command === 'Input.dispatchMouseEvent' && input.x === 120))
assert.ok(f.deadlines.every(deadline => deadline === 4000), 'frame waits must share the action deadline')
assert.equal(f.pending(), 0, 'successful frame samples must clear their watchdogs')
})

test('a renderer that stops producing frames fails within the action budget without input', async () => {
const f = controlledFrameDriver({ renderFrames: false })
await assert.rejects(f.driver.click('Close'), /SEMANTIC_CONTROL_UNAVAILABLE/)
assert.equal(f.elapsed(), 4000)
assert.deepEqual(f.inputs, [])
assert.equal(f.cancelledFrames.length, 1)
assert.equal(f.pending(), 0, 'timed-out frame callbacks must be cancelled')
})

test('the action deadline preserves the last confirmed covered-control failure', async () => {
const f = controlledFrameDriver({ covered: true })
await assert.rejects(f.driver.click('Close'), /SEMANTIC_CONTROL_COVERED/)
assert.deepEqual(f.inputs, [])
assert.equal(f.pending(), 0)
})

for (const kind of ['disabled', 'covered']) {
test(`${kind} control remains a failure without forced input`, async () => {
const f = await fixture(`
Expand Down
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
Loading
Loading