checkpoint-*: In put_writes clear any previously saved writes for this task_id

- Previously implementations were clearing only tasks if the index matched a previously saved one, which isn't enough to guarantee idempotency
This commit is contained in:
Nuno Campos
2024-08-19 15:40:04 -07:00
parent e15d56c2f2
commit 7aaeedd7ba
6 changed files with 59 additions and 11 deletions
@@ -424,6 +424,16 @@ class SqliteSaver(BaseCheckpointSaver):
task_id (str): Identifier for the task creating the writes.
"""
with self.lock, self.cursor() as cur:
cur.execute(
"DELETE FROM writes WHERE thread_id = ? AND checkpoint_ns = ? AND checkpoint_id = ? AND task_id = ? AND idx >= ?",
(
str(config["configurable"]["thread_id"]),
str(config["configurable"]["checkpoint_ns"]),
str(config["configurable"]["checkpoint_id"]),
task_id,
len(writes),
),
)
cur.executemany(
"INSERT OR REPLACE INTO writes (thread_id, checkpoint_ns, checkpoint_id, task_id, idx, channel, type, value) VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
[