Skip to content

Commit 4af7e9a

Browse files
authored
Lock on upsert (#95)
* Lock on upsert * Release notes
1 parent 009d42b commit 4af7e9a

5 files changed

Lines changed: 43 additions & 17 deletions

File tree

‎docs/release-notes.md‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,27 @@
11
# Release Notes
22

3+
## v2.2.6
4+
5+
Adds tag filtering for task listings and fixes a race in the task database plugin.
6+
7+
- Task listings can be filtered by tag: the `GET /tasks` endpoint and the `ls`
8+
command of the [task CLI](https://fluid.quantmind.com/reference/task_cli/)
9+
accept a repeatable `tags` option that returns only tasks carrying at least
10+
one of the given tags, and `TaskInfo` now reports each task's tags.
11+
([#94](https://github.com/quantmind/aio-fluid/pull/94))
12+
- The [task database plugin](https://fluid.quantmind.com/reference/task_plugin/#fluid.scheduler.db.TaskDbPlugin)
13+
now serialises its per-run lifecycle writes with a dedicated task-run lock.
14+
`CrudDB.db_upsert` is not atomic — it issues an `UPDATE` and only `INSERT`s
15+
when nothing matched — so when the scheduler wrote the `queued` row and a
16+
consumer wrote the `running` row a few milliseconds later, the consumer's
17+
`UPDATE` could miss the not-yet-committed `INSERT`, fall through to its own
18+
`INSERT` and violate the task-runs primary key. Holding the lock around the
19+
upsert removes the race.
20+
([#95](https://github.com/quantmind/aio-fluid/pull/95))
21+
- [TaskRun.lock](https://fluid.quantmind.com/reference/task_run/#fluid.scheduler.TaskRun.lock)
22+
accepts an optional `name` to acquire a named sub-lock for the task run, and
23+
`timeout` now defaults to `None`. ([#95](https://github.com/quantmind/aio-fluid/pull/95))
24+
325
## v2.2.5
426

527
- The [task decorator](https://fluid.quantmind.com/reference/task/#fluid.scheduler.task)

‎fluid/scheduler/db.py‎

Lines changed: 14 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -125,19 +125,20 @@ async def get_run(self, run_id: str) -> TaskRunHistory:
125125
async def _handle_update(self, task_run: TaskRun) -> None:
126126
if self.skip_db_tag in task_run.task.tags:
127127
return
128-
await self.db.db_upsert(
129-
self.db.tables[self.table_name],
130-
dict(id=task_run.id),
131-
dict(
132-
state=task_run.state,
133-
name=task_run.name,
134-
priority=task_run.priority,
135-
queued=task_run.queued,
136-
start=task_run.start,
137-
end=task_run.end,
138-
params=task_run.params.model_dump(mode="json"),
139-
),
140-
)
128+
async with task_run.lock(timeout=5, name="db_upsert"):
129+
await self.db.db_upsert(
130+
self.db.tables[self.table_name],
131+
dict(id=task_run.id),
132+
dict(
133+
state=task_run.state,
134+
name=task_run.name,
135+
priority=task_run.priority,
136+
queued=task_run.queued,
137+
start=task_run.start,
138+
end=task_run.end,
139+
params=task_run.params.model_dump(mode="json"),
140+
),
141+
)
141142

142143

143144
def task_meta(meta: sa.MetaData, table_name: str = "tasks") -> None:

‎fluid/scheduler/models.py‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -469,9 +469,12 @@ def set_state(
469469
case _:
470470
raise TaskRunError(f"invalid state transition {self.state} -> {state}")
471471

472-
def lock(self, timeout: float | None) -> Lock:
472+
def lock(self, timeout: float | None = None, name: str | None = None) -> Lock:
473473
"""Get a lock for this task run"""
474-
return self.task_manager.broker.lock(f"tasks:{self.name}", timeout=timeout)
474+
lock_name = f"tasks:{self.name}"
475+
if name:
476+
lock_name = f"{lock_name}:{name}"
477+
return self.task_manager.broker.lock(lock_name, timeout=timeout)
475478

476479
async def _execute(self) -> None:
477480
try:

‎pyproject.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
[project]
22
name = "aio-fluid"
3-
version = "2.2.5"
3+
version = "2.2.6"
44
description = "Tools for backend python services"
55
authors = [
66
{ name = "Luca Sbardella", email = "luca@quantmind.com" },

‎uv.lock‎

Lines changed: 1 addition & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

0 commit comments

Comments
 (0)