From d33807b5eaa9bf2c3b20c3bb92e77b386ba8a296 Mon Sep 17 00:00:00 2001 From: vbarda Date: Wed, 7 Aug 2024 17:28:24 -0400 Subject: [PATCH] checkpoint-postgres: remove unset channel values from blobs --- .../langgraph/checkpoint/postgres/base.py | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/libs/checkpoint-postgres/langgraph/checkpoint/postgres/base.py b/libs/checkpoint-postgres/langgraph/checkpoint/postgres/base.py index 7bbce9641..97b13ed30 100644 --- a/libs/checkpoint-postgres/langgraph/checkpoint/postgres/base.py +++ b/libs/checkpoint-postgres/langgraph/checkpoint/postgres/base.py @@ -40,7 +40,7 @@ MIGRATIONS = [ channel TEXT NOT NULL, version TEXT NOT NULL, type TEXT NOT NULL, - blob BYTEA NOT NULL, + blob BYTEA, PRIMARY KEY (thread_id, checkpoint_ns, channel, version) );""", """CREATE TABLE IF NOT EXISTS checkpoint_writes ( @@ -54,6 +54,7 @@ MIGRATIONS = [ blob BYTEA NOT NULL, PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id, task_id, idx) );""", + "ALTER TABLE checkpoint_blobs ALTER COLUMN blob DROP not null;", ] SELECT_SQL = """ @@ -140,6 +141,7 @@ class BasePostgresSaver(BaseCheckpointSaver): return { k.decode(): self.serde.loads_typed((t.decode(), v)) for k, t, v in blob_values + if t.decode() != "empty" } def _dump_blobs( @@ -162,10 +164,13 @@ class BasePostgresSaver(BaseCheckpointSaver): checkpoint_ns, k, ver, - *self.serde.dumps_typed(values[k]), + *( + self.serde.dumps_typed(values[k]) + if k in values + else ("empty", None) + ), ) for k, ver in versions.items() - if k in values ] def _load_writes(