diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/chat/store.py')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/chat/store.py | 39 |
1 files changed, 26 insertions, 13 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/chat/store.py b/packages/meshbay-node/src/meshbay_node/chat/store.py index 1dbcc2b..a5b2d7d 100644 --- a/packages/meshbay-node/src/meshbay_node/chat/store.py +++ b/packages/meshbay-node/src/meshbay_node/chat/store.py @@ -17,17 +17,22 @@ log = logging.getLogger(__name__) _SCHEMA = """ CREATE TABLE IF NOT EXISTS messages ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - sender_id TEXT NOT NULL, - iteration INTEGER NOT NULL, - payload BLOB NOT NULL, - timestamp REAL NOT NULL, - thread_id TEXT DEFAULT NULL + id INTEGER PRIMARY KEY AUTOINCREMENT, + sender_id TEXT NOT NULL, + iteration INTEGER NOT NULL, + payload BLOB NOT NULL, + timestamp REAL NOT NULL, + thread_id TEXT DEFAULT NULL, + sender_name TEXT DEFAULT '' ); CREATE INDEX IF NOT EXISTS idx_messages_ts ON messages(timestamp); CREATE INDEX IF NOT EXISTS idx_messages_thread ON messages(thread_id); """ +_MIGRATE_SENDER_NAME = ( + "ALTER TABLE messages ADD COLUMN sender_name TEXT DEFAULT ''" +) + @dataclass class StoredMessage: @@ -37,6 +42,7 @@ class StoredMessage: payload: bytes timestamp: float thread_id: str | None + sender_name: str = "" class ChatStore: @@ -50,6 +56,10 @@ class ChatStore: self._db_path.parent.mkdir(parents=True, exist_ok=True) self._db = await aiosqlite.connect(str(self._db_path)) await self._db.executescript(_SCHEMA) + try: + await self._db.execute(_MIGRATE_SENDER_NAME) + except Exception: + pass await self._db.commit() async def close(self) -> None: @@ -70,13 +80,14 @@ class ChatStore: iteration: int, payload: bytes, thread_id: str | None = None, + sender_name: str = "", ) -> int: """Store a message. Returns the row id.""" ts = time.time() cursor = await self._db.execute( - "INSERT INTO messages (sender_id, iteration, payload, timestamp, thread_id) " - "VALUES (?, ?, ?, ?, ?)", - (sender_id, iteration, payload, ts, thread_id), + "INSERT INTO messages (sender_id, iteration, payload, timestamp, thread_id, sender_name) " + "VALUES (?, ?, ?, ?, ?, ?)", + (sender_id, iteration, payload, ts, thread_id, sender_name), ) await self._db.commit() return cursor.lastrowid @@ -88,28 +99,30 @@ class ChatStore: ) -> list[StoredMessage]: """Get messages after a timestamp, most recent last.""" cursor = await self._db.execute( - "SELECT id, sender_id, iteration, payload, timestamp, thread_id " + "SELECT id, sender_id, iteration, payload, timestamp, thread_id, sender_name " "FROM messages WHERE timestamp > ? ORDER BY timestamp ASC LIMIT ?", (since, limit), ) rows = await cursor.fetchall() return [ StoredMessage(id=r[0], sender_id=r[1], iteration=r[2], - payload=r[3], timestamp=r[4], thread_id=r[5]) + payload=r[3], timestamp=r[4], thread_id=r[5], + sender_name=r[6] or "") for r in rows ] async def get_thread(self, thread_id: str, limit: int = 100) -> list[StoredMessage]: """Get messages in a thread.""" cursor = await self._db.execute( - "SELECT id, sender_id, iteration, payload, timestamp, thread_id " + "SELECT id, sender_id, iteration, payload, timestamp, thread_id, sender_name " "FROM messages WHERE thread_id = ? ORDER BY timestamp ASC LIMIT ?", (thread_id, limit), ) rows = await cursor.fetchall() return [ StoredMessage(id=r[0], sender_id=r[1], iteration=r[2], - payload=r[3], timestamp=r[4], thread_id=r[5]) + payload=r[3], timestamp=r[4], thread_id=r[5], + sender_name=r[6] or "") for r in rows ] |