fix: final review — atomic sqlite seq, remove shims and legacy types, cleanup
This commit is contained in:
+10
@@ -0,0 +1,10 @@
|
|||||||
|
# Python-generated files
|
||||||
|
__pycache__/
|
||||||
|
*.py[oc]
|
||||||
|
build/
|
||||||
|
dist/
|
||||||
|
wheels/
|
||||||
|
*.egg-info
|
||||||
|
|
||||||
|
# Virtual environments
|
||||||
|
.venv
|
||||||
@@ -0,0 +1 @@
|
|||||||
|
3.14
|
||||||
@@ -0,0 +1,28 @@
|
|||||||
|
[project]
|
||||||
|
name = "agentmsgs"
|
||||||
|
version = "0.1.0"
|
||||||
|
description = "Add your description here"
|
||||||
|
readme = "README.md"
|
||||||
|
authors = [
|
||||||
|
{ name = "lda", email = "[email protected]" }
|
||||||
|
]
|
||||||
|
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"]
|
||||||
@@ -0,0 +1,4 @@
|
|||||||
|
from agentmsgs import main
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
main()
|
||||||
@@ -32,9 +32,6 @@ class App:
|
|||||||
def find_threads(self, *agents: Agent) -> list[Thread]:
|
def find_threads(self, *agents: Agent) -> list[Thread]:
|
||||||
return ops.find_threads(self.store, set(agents))
|
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:
|
def join_thread(self, tid: uuid.UUID, agent: Agent) -> Thread:
|
||||||
return ops.join_thread(self.store, tid, agent)
|
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:
|
def has_unread(self, tid: uuid.UUID, agent: Agent, exclude_own: bool = False) -> bool:
|
||||||
return ops.has_unread(self.store, tid, agent, exclude_own)
|
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
|
|
||||||
|
|||||||
@@ -85,6 +85,10 @@ class InMemoryStore:
|
|||||||
def delete_thread(self, id: uuid.UUID) -> None:
|
def delete_thread(self, id: uuid.UUID) -> None:
|
||||||
self._threads.pop(id, None)
|
self._threads.pop(id, None)
|
||||||
self._msgs.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:
|
def append_message(self, thread_id: uuid.UUID, sender: Agent, content: str) -> Message:
|
||||||
if thread_id not in self._threads:
|
if thread_id not in self._threads:
|
||||||
|
|||||||
@@ -8,38 +8,12 @@ from pathlib import Path
|
|||||||
|
|
||||||
from agentmsgs.core.types import Agent, Message, Thread
|
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 = """
|
_SCHEMA = """
|
||||||
PRAGMA journal_mode=WAL;
|
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 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 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 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));
|
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():
|
if self.path.parent and not self.path.parent.exists():
|
||||||
self.path.parent.mkdir(parents=True, exist_ok=True)
|
self.path.parent.mkdir(parents=True, exist_ok=True)
|
||||||
self._db = sqlite3.connect(str(self.path), check_same_thread=False, timeout=5.0)
|
self._db = sqlite3.connect(str(self.path), check_same_thread=False, timeout=5.0)
|
||||||
_live_stores.add(self)
|
|
||||||
# WAL mode
|
# WAL mode
|
||||||
try:
|
try:
|
||||||
self._db.execute("PRAGMA journal_mode=WAL;")
|
self._db.execute("PRAGMA journal_mode=WAL;")
|
||||||
@@ -62,6 +35,13 @@ class SQLiteStore:
|
|||||||
except Exception:
|
except Exception:
|
||||||
pass
|
pass
|
||||||
self._db.executescript(_SCHEMA)
|
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()
|
self._db.commit()
|
||||||
|
|
||||||
# -- agents --
|
# -- agents --
|
||||||
@@ -230,6 +210,9 @@ class SQLiteStore:
|
|||||||
)
|
)
|
||||||
if cur.fetchone() is None:
|
if cur.fetchone() is None:
|
||||||
raise ValueError("sender not in thread")
|
raise ValueError("sender not in thread")
|
||||||
|
# atomic seq via BEGIN IMMEDIATE transaction
|
||||||
|
try:
|
||||||
|
self._db.execute("BEGIN IMMEDIATE")
|
||||||
cur2 = self._db.execute(
|
cur2 = self._db.execute(
|
||||||
"SELECT COALESCE(MAX(seq), 0) FROM messages WHERE thread_id=?", (str(thread_id),)
|
"SELECT COALESCE(MAX(seq), 0) FROM messages WHERE thread_id=?", (str(thread_id),)
|
||||||
)
|
)
|
||||||
@@ -242,6 +225,12 @@ class SQLiteStore:
|
|||||||
(mid, str(thread_id), str(sender.id), seq, content, ts),
|
(mid, str(thread_id), str(sender.id), seq, content, ts),
|
||||||
)
|
)
|
||||||
self._db.commit()
|
self._db.commit()
|
||||||
|
except Exception:
|
||||||
|
try:
|
||||||
|
self._db.rollback()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
raise
|
||||||
return Message(
|
return Message(
|
||||||
id=uuid.UUID(mid),
|
id=uuid.UUID(mid),
|
||||||
sender=sender,
|
sender=sender,
|
||||||
@@ -299,7 +288,6 @@ class SQLiteStore:
|
|||||||
self._db.close()
|
self._db.close()
|
||||||
except Exception:
|
except Exception:
|
||||||
pass
|
pass
|
||||||
_live_stores.discard(self)
|
|
||||||
|
|
||||||
def __del__(self) -> None:
|
def __del__(self) -> None:
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -5,14 +5,20 @@ def test_sqlite_persists_across_handles():
|
|||||||
with tempfile.TemporaryDirectory() as d:
|
with tempfile.TemporaryDirectory() as d:
|
||||||
p = pathlib.Path(d)/"test.db"
|
p = pathlib.Path(d)/"test.db"
|
||||||
s1 = SQLiteStore(p)
|
s1 = SQLiteStore(p)
|
||||||
|
try:
|
||||||
a = s1.get_or_create_agent("Alice"); b = s1.get_or_create_agent("Bob")
|
a = s1.get_or_create_agent("Alice"); b = s1.get_or_create_agent("Bob")
|
||||||
t = s1.create_thread({a,b})
|
t = s1.create_thread({a,b})
|
||||||
s1.append_message(t.id, a, "hi")
|
s1.append_message(t.id, a, "hi")
|
||||||
# new handle same file
|
# new handle same file
|
||||||
s2 = SQLiteStore(p)
|
s2 = SQLiteStore(p)
|
||||||
|
try:
|
||||||
assert s2.get_agent_by_name("Alice").id == a.id
|
assert s2.get_agent_by_name("Alice").id == a.id
|
||||||
assert len(s2.list_messages(t.id)) == 1
|
assert len(s2.list_messages(t.id)) == 1
|
||||||
assert s2.find_threads({a})[0].id == t.id
|
assert s2.find_threads({a})[0].id == t.id
|
||||||
|
finally:
|
||||||
|
s2.close()
|
||||||
|
finally:
|
||||||
|
s1.close()
|
||||||
|
|
||||||
def test_sqlite_concurrent_append():
|
def test_sqlite_concurrent_append():
|
||||||
import pathlib, tempfile
|
import pathlib, tempfile
|
||||||
@@ -20,6 +26,7 @@ def test_sqlite_concurrent_append():
|
|||||||
with tempfile.TemporaryDirectory() as d:
|
with tempfile.TemporaryDirectory() as d:
|
||||||
p = pathlib.Path(d)/"c.db"
|
p = pathlib.Path(d)/"c.db"
|
||||||
s1 = SQLiteStore(p); s2 = SQLiteStore(p)
|
s1 = SQLiteStore(p); s2 = SQLiteStore(p)
|
||||||
|
try:
|
||||||
a = s1.get_or_create_agent("A"); b = s1.get_or_create_agent("B")
|
a = s1.get_or_create_agent("A"); b = s1.get_or_create_agent("B")
|
||||||
# s2 sees same agents via file
|
# s2 sees same agents via file
|
||||||
a2 = s2.get_agent_by_name("A"); b2 = s2.get_agent_by_name("B")
|
a2 = s2.get_agent_by_name("A"); b2 = s2.get_agent_by_name("B")
|
||||||
@@ -27,3 +34,11 @@ def test_sqlite_concurrent_append():
|
|||||||
s1.append_message(t.id, a, "from s1")
|
s1.append_message(t.id, a, "from s1")
|
||||||
s2.append_message(t.id, b2, "from s2")
|
s2.append_message(t.id, b2, "from s2")
|
||||||
assert len(s1.list_messages(t.id)) == 2
|
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()
|
||||||
|
|||||||
@@ -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" },
|
||||||
|
]
|
||||||
Reference in New Issue
Block a user