From 8adb9984f026375182020730db106ebfeba827c0 Mon Sep 17 00:00:00 2001 From: lda Date: Tue, 1 Sep 2026 02:10:35 +0700 Subject: [PATCH] =?UTF-8?q?fix:=20final=20review=20=E2=80=94=20atomic=20sq?= =?UTF-8?q?lite=20seq,=20remove=20shims=20and=20legacy=20types,=20cleanup?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 10 +++++ .python-version | 1 + README.md | 0 pyproject.toml | 28 ++++++++++++ src/agentmsgs/__main__.py | 4 ++ src/agentmsgs/core/app.py | 12 ------ src/agentmsgs/stores/memory.py | 4 ++ src/agentmsgs/stores/sqlite.py | 70 +++++++++++++----------------- tests/__init__.py | 0 tests/test_store_sqlite.py | 45 ++++++++++++------- uv.lock | 79 ++++++++++++++++++++++++++++++++++ 11 files changed, 185 insertions(+), 68 deletions(-) create mode 100644 .gitignore create mode 100644 .python-version create mode 100644 README.md create mode 100644 pyproject.toml create mode 100644 src/agentmsgs/__main__.py create mode 100644 tests/__init__.py create mode 100644 uv.lock diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..505a3b1 --- /dev/null +++ b/.gitignore @@ -0,0 +1,10 @@ +# Python-generated files +__pycache__/ +*.py[oc] +build/ +dist/ +wheels/ +*.egg-info + +# Virtual environments +.venv diff --git a/.python-version b/.python-version new file mode 100644 index 0000000..6324d40 --- /dev/null +++ b/.python-version @@ -0,0 +1 @@ +3.14 diff --git a/README.md b/README.md new file mode 100644 index 0000000..e69de29 diff --git a/pyproject.toml b/pyproject.toml new file mode 100644 index 0000000..0b66922 --- /dev/null +++ b/pyproject.toml @@ -0,0 +1,28 @@ +[project] +name = "agentmsgs" +version = "0.1.0" +description = "Add your description here" +readme = "README.md" +authors = [ + { name = "lda", email = "lda@ldlda.com" } +] +requires-python = ">=3.14" +dependencies = [] + +[tool.uv] +package = true + +[project.scripts] +agentmsgs = "agentmsgs:main" + +[build-system] +requires = ["uv_build>=0.12.3,<0.13.0"] +build-backend = "uv_build" + +[dependency-groups] +dev = [ + "pytest>=9.1.1", +] + +[tool.pytest.ini_options] +pythonpath = ["src"] diff --git a/src/agentmsgs/__main__.py b/src/agentmsgs/__main__.py new file mode 100644 index 0000000..e155bdc --- /dev/null +++ b/src/agentmsgs/__main__.py @@ -0,0 +1,4 @@ +from agentmsgs import main + +if __name__ == "__main__": + main() diff --git a/src/agentmsgs/core/app.py b/src/agentmsgs/core/app.py index f769eef..ab5ed3b 100644 --- a/src/agentmsgs/core/app.py +++ b/src/agentmsgs/core/app.py @@ -32,9 +32,6 @@ class App: def find_threads(self, *agents: Agent) -> list[Thread]: return ops.find_threads(self.store, set(agents)) - def find_thread(self, *agents: Agent) -> list[Thread]: - return self.find_threads(*agents) - def join_thread(self, tid: uuid.UUID, agent: Agent) -> Thread: return ops.join_thread(self.store, tid, agent) @@ -52,12 +49,3 @@ class App: def has_unread(self, tid: uuid.UUID, agent: Agent, exclude_own: bool = False) -> bool: return ops.has_unread(self.store, tid, agent, exclude_own) - - # compat shims for old tests - def add_thread(self, *a: Agent) -> Thread: - return self.create_thread(*a) - - def add_thread_2(self, thread: Thread): - if hasattr(self.store, "_threads"): - self.store._threads[thread.id] = thread # type: ignore[attr-defined] - return None diff --git a/src/agentmsgs/stores/memory.py b/src/agentmsgs/stores/memory.py index 918c434..547b42c 100644 --- a/src/agentmsgs/stores/memory.py +++ b/src/agentmsgs/stores/memory.py @@ -85,6 +85,10 @@ class InMemoryStore: def delete_thread(self, id: uuid.UUID) -> None: self._threads.pop(id, None) self._msgs.pop(id, None) + # also delete cursors for that thread + for k in list(self._cursors.keys()): + if k[0] == id: + self._cursors.pop(k, None) def append_message(self, thread_id: uuid.UUID, sender: Agent, content: str) -> Message: if thread_id not in self._threads: diff --git a/src/agentmsgs/stores/sqlite.py b/src/agentmsgs/stores/sqlite.py index 9205314..22b086f 100644 --- a/src/agentmsgs/stores/sqlite.py +++ b/src/agentmsgs/stores/sqlite.py @@ -8,38 +8,12 @@ from pathlib import Path from agentmsgs.core.types import Agent, Message, Thread -# track live stores for Windows temp-file cleanup (PermissionError if db still open) -import tempfile as _tempfile - -_live_stores: set["SQLiteStore"] = set() -_orig_temp_cleanup = _tempfile.TemporaryDirectory.cleanup - -def _patched_temp_cleanup(self): # type: ignore[no-untyped-def] - # close any SQLiteStore whose path lives inside this temp dir - try: - tpath = str(self.name) - except Exception: - tpath = "" - for s in list(_live_stores): - try: - sp = str(getattr(s, "path", "")) - if sp.startswith(tpath) if tpath else False: - try: - s.close() - except Exception: - pass - except Exception: - pass - return _orig_temp_cleanup(self) - -_tempfile.TemporaryDirectory.cleanup = _patched_temp_cleanup # type: ignore[method-assign,assignment] - _SCHEMA = """ PRAGMA journal_mode=WAL; CREATE TABLE IF NOT EXISTS agents (id TEXT PRIMARY KEY, name TEXT UNIQUE NOT NULL, created_at TEXT); CREATE TABLE IF NOT EXISTS threads (id TEXT PRIMARY KEY, created_at TEXT NOT NULL); CREATE TABLE IF NOT EXISTS thread_participants (thread_id TEXT NOT NULL, agent_id TEXT NOT NULL, PRIMARY KEY(thread_id, agent_id)); -CREATE TABLE IF NOT EXISTS messages (id TEXT PRIMARY KEY, thread_id TEXT NOT NULL, sender_id TEXT NOT NULL, seq INTEGER NOT NULL, content TEXT NOT NULL, ts TEXT NOT NULL); +CREATE TABLE IF NOT EXISTS messages (id TEXT PRIMARY KEY, thread_id TEXT NOT NULL, sender_id TEXT NOT NULL, seq INTEGER NOT NULL, content TEXT NOT NULL, ts TEXT NOT NULL, UNIQUE(thread_id, seq)); CREATE TABLE IF NOT EXISTS cursors (thread_id TEXT NOT NULL, agent_id TEXT NOT NULL, seq INTEGER NOT NULL, PRIMARY KEY(thread_id, agent_id)); """ @@ -51,7 +25,6 @@ class SQLiteStore: if self.path.parent and not self.path.parent.exists(): self.path.parent.mkdir(parents=True, exist_ok=True) self._db = sqlite3.connect(str(self.path), check_same_thread=False, timeout=5.0) - _live_stores.add(self) # WAL mode try: self._db.execute("PRAGMA journal_mode=WAL;") @@ -62,6 +35,13 @@ class SQLiteStore: except Exception: pass self._db.executescript(_SCHEMA) + # ensure unique index for DBs created before constraint was added + try: + self._db.execute( + "CREATE UNIQUE INDEX IF NOT EXISTS idx_messages_thread_seq ON messages(thread_id, seq)" + ) + except Exception: + pass self._db.commit() # -- agents -- @@ -230,18 +210,27 @@ class SQLiteStore: ) if cur.fetchone() is None: raise ValueError("sender not in thread") - cur2 = self._db.execute( - "SELECT COALESCE(MAX(seq), 0) FROM messages WHERE thread_id=?", (str(thread_id),) - ) - max_seq = cur2.fetchone()[0] - seq = int(max_seq) + 1 - mid = str(uuid.uuid4()) - ts = datetime.datetime.now(datetime.timezone.utc).isoformat() - self._db.execute( - "INSERT INTO messages (id, thread_id, sender_id, seq, content, ts) VALUES (?, ?, ?, ?, ?, ?)", - (mid, str(thread_id), str(sender.id), seq, content, ts), - ) - self._db.commit() + # atomic seq via BEGIN IMMEDIATE transaction + try: + self._db.execute("BEGIN IMMEDIATE") + cur2 = self._db.execute( + "SELECT COALESCE(MAX(seq), 0) FROM messages WHERE thread_id=?", (str(thread_id),) + ) + max_seq = cur2.fetchone()[0] + seq = int(max_seq) + 1 + mid = str(uuid.uuid4()) + ts = datetime.datetime.now(datetime.timezone.utc).isoformat() + self._db.execute( + "INSERT INTO messages (id, thread_id, sender_id, seq, content, ts) VALUES (?, ?, ?, ?, ?, ?)", + (mid, str(thread_id), str(sender.id), seq, content, ts), + ) + self._db.commit() + except Exception: + try: + self._db.rollback() + except Exception: + pass + raise return Message( id=uuid.UUID(mid), sender=sender, @@ -299,7 +288,6 @@ class SQLiteStore: self._db.close() except Exception: pass - _live_stores.discard(self) def __del__(self) -> None: try: diff --git a/tests/__init__.py b/tests/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/test_store_sqlite.py b/tests/test_store_sqlite.py index f29e80c..f4b9ebf 100644 --- a/tests/test_store_sqlite.py +++ b/tests/test_store_sqlite.py @@ -5,14 +5,20 @@ def test_sqlite_persists_across_handles(): with tempfile.TemporaryDirectory() as d: p = pathlib.Path(d)/"test.db" s1 = SQLiteStore(p) - a = s1.get_or_create_agent("Alice"); b = s1.get_or_create_agent("Bob") - t = s1.create_thread({a,b}) - s1.append_message(t.id, a, "hi") - # new handle same file - s2 = SQLiteStore(p) - assert s2.get_agent_by_name("Alice").id == a.id - assert len(s2.list_messages(t.id)) == 1 - assert s2.find_threads({a})[0].id == t.id + try: + a = s1.get_or_create_agent("Alice"); b = s1.get_or_create_agent("Bob") + t = s1.create_thread({a,b}) + s1.append_message(t.id, a, "hi") + # new handle same file + s2 = SQLiteStore(p) + try: + assert s2.get_agent_by_name("Alice").id == a.id + assert len(s2.list_messages(t.id)) == 1 + assert s2.find_threads({a})[0].id == t.id + finally: + s2.close() + finally: + s1.close() def test_sqlite_concurrent_append(): import pathlib, tempfile @@ -20,10 +26,19 @@ def test_sqlite_concurrent_append(): with tempfile.TemporaryDirectory() as d: p = pathlib.Path(d)/"c.db" s1 = SQLiteStore(p); s2 = SQLiteStore(p) - a = s1.get_or_create_agent("A"); b = s1.get_or_create_agent("B") - # s2 sees same agents via file - a2 = s2.get_agent_by_name("A"); b2 = s2.get_agent_by_name("B") - t = s1.create_thread({a,b}) - s1.append_message(t.id, a, "from s1") - s2.append_message(t.id, b2, "from s2") - assert len(s1.list_messages(t.id)) == 2 + try: + a = s1.get_or_create_agent("A"); b = s1.get_or_create_agent("B") + # s2 sees same agents via file + a2 = s2.get_agent_by_name("A"); b2 = s2.get_agent_by_name("B") + t = s1.create_thread({a,b}) + s1.append_message(t.id, a, "from s1") + s2.append_message(t.id, b2, "from s2") + assert len(s1.list_messages(t.id)) == 2 + # verify seq uniqueness and monotonicity + msgs = s1.list_messages(t.id) + seqs = [m.seq for m in msgs] + assert seqs == [1, 2] + assert len(set(seqs)) == 2 + finally: + s2.close() + s1.close() diff --git a/uv.lock b/uv.lock new file mode 100644 index 0000000..d30f384 --- /dev/null +++ b/uv.lock @@ -0,0 +1,79 @@ +version = 1 +revision = 3 +requires-python = ">=3.14" + +[[package]] +name = "agentmsgs" +version = "0.1.0" +source = { editable = "." } + +[package.dev-dependencies] +dev = [ + { name = "pytest" }, +] + +[package.metadata] + +[package.metadata.requires-dev] +dev = [{ name = "pytest", specifier = ">=9.1.1" }] + +[[package]] +name = "colorama" +version = "0.4.6" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/d8/53/6f443c9a4a8358a93a6792e2acffb9d9d5cb0a5cfd8802644b7b1c9a02e4/colorama-0.4.6.tar.gz", hash = "sha256:08695f5cb7ed6e0531a20572697297273c47b8cae5a63ffc6d6ed5c201be6e44", size = 27697, upload-time = "2022-10-25T02:36:22.414Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/d1/d6/3965ed04c63042e047cb6a3e6ed1a63a35087b6a609aa3a15ed8ac56c221/colorama-0.4.6-py2.py3-none-any.whl", hash = "sha256:4f1d9991f5acc0ca119f9d443620b77f9d6b33703e51011c16baf57afb285fc6", size = 25335, upload-time = "2022-10-25T02:36:20.889Z" }, +] + +[[package]] +name = "iniconfig" +version = "2.3.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/72/34/14ca021ce8e5dfedc35312d08ba8bf51fdd999c576889fc2c24cb97f4f10/iniconfig-2.3.0.tar.gz", hash = "sha256:c76315c77db068650d49c5b56314774a7804df16fee4402c1f19d6d15d8c4730", size = 20503, upload-time = "2025-10-18T21:55:43.219Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/cb/b1/3846dd7f199d53cb17f49cba7e651e9ce294d8497c8c150530ed11865bb8/iniconfig-2.3.0-py3-none-any.whl", hash = "sha256:f631c04d2c48c52b84d0d0549c99ff3859c98df65b3101406327ecc7d53fbf12", size = 7484, upload-time = "2025-10-18T21:55:41.639Z" }, +] + +[[package]] +name = "packaging" +version = "26.3" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/7d/fa/3944b40b07da9ce895c0e6303a5ab7d53da063554f534556b134a54d6093/packaging-26.3.tar.gz", hash = "sha256:94edc256424af38762eb31306eed28beb9f0efc50a8837492c9d6fd6004aed79", size = 313412, upload-time = "2026-08-04T18:15:28.737Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/63/34/ba1c580383c9eada3711951fef0795c80b829a078d72188184bcab9dd527/packaging-26.3-py3-none-any.whl", hash = "sha256:d7193f7c8e4e93f444fde0262bf90af30e16fa0ad0ad44cb553c87339b23cd1c", size = 129956, upload-time = "2026-08-04T18:15:27.159Z" }, +] + +[[package]] +name = "pluggy" +version = "1.6.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/f9/e2/3e91f31a7d2b083fe6ef3fa267035b518369d9511ffab804f839851d2779/pluggy-1.6.0.tar.gz", hash = "sha256:7dcc130b76258d33b90f61b658791dede3486c3e6bfb003ee5c9bfb396dd22f3", size = 69412, upload-time = "2025-05-15T12:30:07.975Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/54/20/4d324d65cc6d9205fabedc306948156824eb9f0ee1633355a8f7ec5c66bf/pluggy-1.6.0-py3-none-any.whl", hash = "sha256:e920276dd6813095e9377c0bc5566d94c932c33b27a3e3945d8389c374dd4746", size = 20538, upload-time = "2025-05-15T12:30:06.134Z" }, +] + +[[package]] +name = "pygments" +version = "2.21.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/49/2e/ced460408999b33da6b31b0021b0f37d329e202d4169aeb164493778f25b/pygments-2.21.0.tar.gz", hash = "sha256:610ca751c9bc2492b38eb9a38a7fbc93edbbb2d7182edaf34e66ae493dee5c8c", size = 5005329, upload-time = "2026-08-17T08:02:48.824Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/71/46/17f022dd3e953bf20a04a028a21ec746d942f8d2af30fa0f124fa0e6a684/pygments-2.21.0-py3-none-any.whl", hash = "sha256:2363c69b61c4a97c838da3b130dcd6468f4848992b21a82f2a63ec34377137d9", size = 1250147, upload-time = "2026-08-17T08:02:44.912Z" }, +] + +[[package]] +name = "pytest" +version = "9.1.1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "colorama", marker = "sys_platform == 'win32'" }, + { name = "iniconfig" }, + { name = "packaging" }, + { name = "pluggy" }, + { name = "pygments" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/e4/47/b9efed96c114afcfa3c9d3fe98a76a1d14c74a9e266d397cf6eb64be5e01/pytest-9.1.1.tar.gz", hash = "sha256:1088fbde8f2b49d95a549a195707afa7a76a3ce9bcadc26b6d71f0ffda5fe313", size = 1636369, upload-time = "2026-06-19T10:58:32.857Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/24/25/1de2678b631f5a49215c6c96fff41ba892b0a34df68d6d80292b1b48aa7f/pytest-9.1.1-py3-none-any.whl", hash = "sha256:37a86b45efb9a47a61a36449063e8e18d0cab3161329fc099eb21783169c4f0c", size = 386536, upload-time = "2026-06-19T10:58:31.347Z" }, +]