From 1802ffd40b2fc02e699dca86e45742a045df230a Mon Sep 17 00:00:00 2001 From: xiaowei Date: Mon, 25 May 2026 01:46:34 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20Phase=201.1=20=E6=95=B0=E6=8D=AE?= =?UTF-8?q?=E6=A8=A1=E5=9E=8B=20+=20=E6=8C=81=E4=B9=85=E5=8C=96=E9=98=9F?= =?UTF-8?q?=E5=88=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Episode 模型 (src/models/episode.py) - Distilled 模型 (src/models/distilled.py) - SQLite 持久化队列 (src/distill/queue.py) - 模型测试 (tests/test_models.py) - 队列测试 (tests/test_queue.py) --- .opencode/opencode.db-shm | Bin 0 -> 32768 bytes .opencode/opencode.db-wal | Bin 0 -> 78312 bytes src/distill/queue.py | 89 ++++++++++++++++++++++++++++++++++++++ src/models/__init__.py | 4 ++ src/models/distilled.py | 47 ++++++++++++++++++++ src/models/episode.py | 43 ++++++++++++++++++ tests/test_models.py | 48 ++++++++++++++++++++ tests/test_queue.py | 67 ++++++++++++++++++++++++++++ 8 files changed, 298 insertions(+) create mode 100644 .opencode/opencode.db-shm create mode 100644 .opencode/opencode.db-wal create mode 100644 src/distill/queue.py create mode 100644 src/models/__init__.py create mode 100644 src/models/distilled.py create mode 100644 src/models/episode.py create mode 100644 tests/test_models.py create mode 100644 tests/test_queue.py diff --git a/.opencode/opencode.db-shm b/.opencode/opencode.db-shm new file mode 100644 index 0000000000000000000000000000000000000000..a8e275b911947958fbd34bc005ec01d65fd5a82f GIT binary patch literal 32768 zcmeI)J4!=A6b9hQ!}lvHNK&Lp?*hce6^MnUm5Xo*?m*BkpnDKpj@z(^=M;g9osi~o zz8`+L6YibNdL_VHzT|1o>}dVih%^ZoVS zx$mlz?(g{=|J{bs(VypWthb`u(VggSbT7IeJ%}Dg^Ss^fV<`j(5FkK+009C72oNAZ zfB*pk1PBlyK!5-N0t5&UAV7cs0RjXF5FkK+009C72oNAZfB*pk1PBlyK!5-N0t5&U zAV7cs0Rp`c7_@Oa?ZrtQlE9_~hA~;?Y}1?Ck3cU3Mln+;rxFQ)dIC91O9<2y$SHC{ zpr%01PZR<*1#)tz5U44TGjoMNO@W*iEd*)`rAW&9dQugL9 gPJjRb0t5&UAV7cs0RjXF5FkK+009C72-Fk!1!;jK{r~^~ literal 0 HcmV?d00001 diff --git a/.opencode/opencode.db-wal b/.opencode/opencode.db-wal new file mode 100644 index 0000000000000000000000000000000000000000..ed64eefc874ed907f8ea6883352a2669bf9e8551 GIT binary patch literal 78312 zcmeI5Yit}>7037CefpJ@z*MqaT`0vn@~)q8qQ-TV&3f!?W3Qd{t`e(_R%`FrnRIt% zy)$lHJ9oNVhg{}%a`WHMcBdA;vii;w zbH9jw_mi~9U0;53>HgpLfAxd+=FUzjhD?^UN-1xUVcOc`ah)WD%6jOJ*KA{NX8RU? zK);#3y5nZM9#`L{*iLIa9k01MyXXft5C8!X009sH0T2KI5C8!X009u#4}n>aU+C?1 z&l~xLqO7kJ=~Hb^mseJ0b+Pf#n@CIXj6^c=Q&SRYkd3JMk{qQi#q|YuK#0ZM8<&~R zWlhuNTwx*ih+NSXO>I2(IJ7p_l1MldQVJxQ%1Dz^n#`t?)A95?IU~)Jcs4VWOws1k zQYsUra`qW=b@Hj13`u3Drp#(Z&*jVIq9WJS)mIz3kuQ}=COIw5W#ZGbWL%nvXQwhG zQpg*!p_JrEkE&hn5yfywq{z+Z{ObFHE5@5||7GQibLKOQH>VnavC;K+%gqbEj2s9EO!{d@U5f?xh# z`nEj%=I^Zc0$Z-mt=rmgTs#PX00@8p2!H?xfB*=900@8p2!Mb^;1h0dEYvuUfmS)z z5&ZdI=Y((cUI^~8j=+3vz_&KK^YPXZbiC|3?Codk2-@V1FSkAudb)L`bs}_M@WbFS z|7g(Xd)*gk>7pXoK;Tv;u&x{qbj^1MH4cRW%`BbBHa?+ekw716L_(X;Zk4sY$o%}SDB3xVelr))4g|oBc=Bz0D z5_B>rWx|ARj%y)f4aYGu7(RDeN=t;!RvIH(vCzkjPzs}ADK);X9J(=E6;+oj2EV!+ zwAE^C+nP+xN$Jc@v$fZ0t14gKX6r;}pes)Odau*3+?br@szNTmi+xQR*f!a2+R=Z! zrRhIAGR)Ia-{?kV$(sIkkKPj{)+Dfc=a~lCPE$Gp!f3a9J)o!sc}*#-@sQN3F{%{! zsoHHHB~#7uNnrr#2%O^`F2RkWhqth}7B7`n596y&A+YSAD)4_@f8beA-mjZDfqUaZJ8 zcNt_py{~#YLSiwJRh2bjlDB5&DJm%`#iF9ii<(-{+4a`qqB=$JhIK3{s-j5OYy(hye@XXP#LMm!kKb+0^0o#`I#=GpuiBc|ZP)I0r=Wd44bHl| zT3yd>C`R)eh24W$)yGY`8Mm_QBtN}ZI)du&7kK}VpL^=`qkntb zT1Rlr)p_j}HXD})0w4eaAOHd&00JNY0w4eaAOHd&utx$%=rRI#<1zyCcMh^O0Q~n0 zeC?UH|8ea{FTKGQ3Rr)_V2}DC1_B@e0w4eaAOHd&00JNY0w4eaAaENJFqdLf=L`Jd zl~2!J`Ps$~!)77g3z!tJfdB}A00@8p2!H?xfB*=900@A7~EFsE@z(S~gY5CQ)yWe~2DZCdr$S#L%AOHd&00JNY z0w4eaAOHd&00JOTB@i|R9JYY~2!H?xfB*=900@8p2!H?xfIyW%tGUPD z3%vE6fydr@@|zd%dk3q0umb`h00JNY0w4eaAOHd&00JNY0tcBun`s<>FR(Jd^z*O3 z|C@fy7dXf+hixDL0w4eaAOHd&00JNY0w4eaAW$V>P7P@1?*(4p`s+`>cj)SCcrQ@p zgB=h60T2KI5C8!X009sH0T2KI5ID#L+@^8-y}+UW9C`kQ3yCM^&Q2+YOqR4tDQ}Qr z*CCh3<2p$Q74^^`w#23XzB9yH&b!b${ve*144JV`|+%iUnt7Dtm}%V>ioMU zk(S~aiDcrZrX<2kNhD0@rW8mEv`gJx|U^^HEk&&R1mB$l2mKTh6L86r(6x z?NT!tlFCj^SrsLE>in{tThvxngCtWKX;MlPE_Jo6(V$M#B2`qB4GQE2v~6qTP^cS(5}j&dS6|h#MU6_Ct@^l0 ztK~x8kPA5)jg4TUgv4Sbt14^6P$&~4Un&ztB_*X;RCIY!QwutC%3?*{NeA`i?mD7G zUoDmLl}B@&yzN(UT{#@+n(q#7EE^SNd0DPlms|BemuH^MQJ4jpQkCsuorhRDIca)) zZ;iw8i3}AUm!>2do---JwS`YflgU&#J3DU1Z8c)&qzt3rxE31fjLd-2neIy?)8A87UVUhu*O4Jug0iS;HPT0ZAhlZrSminnhMu@fkbLC zB3_+rWBuu7p7p0(DgSQefLMRZOLyx}YxvR))?rLgr@$tK>B}G;Ai;g~T&*z_*NlVGelzBsPIEQw<6EwL~32Cm@ zGenevXu8wxN)qw8M0{Kdi~a3>A<^x2v6#_UigZMf%dZ+5JFW%LAa`Qh`tjC4*QgNO z=y9gYl0xsRb)nCWc^o)9y++Y(S;7RIT^uXgUw%9g_6wrmUaY2vbsgoranU;Awv?y| znCPrah?w%?Lm|H~BDk+K=&FS<6Ata5u_jihw^k@#75xE05Zs%CEWMa6_TlG1Sz_52 zbp26EO|3IY>`??Yrdc_?p&E2>)2!+dxnd5UnpNE@S#=XK!@`qICC_Fh=5jhPuttq^ zpcf@J1Tid%LBB92?36&vMY$pN6%+%ziw-i;!1{!*=}pPLhWbR~MLaZ|xyEbjGILwH zaT#aE@zVVzUB}CwUN-W+8zXOPsO}gjW8_^Oy^H-V*Si(*aLf(zJs*8MWPfOH*5em? zd)@O^gjx46U0zw0)kWF9rSU^A&l2lN&nX*G^CdY-MaA_6cR+~6+#8pf&Sg#0JQ zXQzGRvB#mcv6hZ|LLoZR=Fc>B50WIF&CDcIwE47@%0#K0eFpPslFRdlL$g}ZbNO<) zsL0K$4eP-ulbn|3GV$rz>hom8e6q3zl57YmhC?F7XFlgw-wS*+dGgmE{ouli`LOTp zeBRahZs)7?3mXW400@8p2!H?xfB*=900@8p2!O!?~`u3%`mNewm0OMj~C2p;f21n^nsI literal 0 HcmV?d00001 diff --git a/src/distill/queue.py b/src/distill/queue.py new file mode 100644 index 0000000..18c4315 --- /dev/null +++ b/src/distill/queue.py @@ -0,0 +1,89 @@ +import sqlite3 +import json +import threading +from datetime import datetime +from pathlib import Path +from typing import Optional +import sys +sys.path.insert(0, str(Path(__file__).parent.parent)) +from models.episode import Episode + +class PersistenceQueue: + def __init__(self, db_path: str = 'zhiyi.db'): + self.db_path = db_path + self.lock = threading.Lock() + self.conn = sqlite3.connect(db_path, check_same_thread=False) + self.conn.row_factory = sqlite3.Row + self._init_db() + + def _init_db(self): + self.conn.execute(''' + CREATE TABLE IF NOT EXISTS episode_queue ( + id TEXT PRIMARY KEY, + data TEXT NOT NULL, + enqueued_at TEXT NOT NULL, + dequeued_at TEXT, + status TEXT DEFAULT 'pending' + ) + ''') + self.conn.execute('CREATE INDEX IF NOT EXISTS idx_status ON episode_queue(status)') + self.conn.commit() + + def enqueue(self, episode: Episode) -> bool: + with self.lock: + try: + data = json.dumps(episode.to_dict()) + self.conn.execute( + 'INSERT INTO episode_queue (id, data, enqueued_at, status) VALUES (?, ?, ?, ?)', + (episode.id, data, datetime.now().isoformat(), 'pending') + ) + self.conn.commit() + return True + except sqlite3.IntegrityError: + return False + + def dequeue(self) -> Optional[dict]: + with self.lock: + row = self.conn.execute( + 'SELECT * FROM episode_queue WHERE status = ? ORDER BY enqueued_at ASC LIMIT 1', + ('pending',) + ).fetchone() + if row is None: + return None + self.conn.execute( + 'UPDATE episode_queue SET status = ?, dequeued_at = ? WHERE id = ?', + ('dequeued', datetime.now().isoformat(), row['id']) + ) + self.conn.commit() + return json.loads(row['data']) + + def peek(self) -> Optional[dict]: + with self.lock: + row = self.conn.execute( + 'SELECT data FROM episode_queue WHERE status = ? ORDER BY enqueued_at ASC LIMIT 1', + ('pending',) + ).fetchone() + return json.loads(row['data']) if row else None + + def size(self) -> int: + with self.lock: + row = self.conn.execute( + 'SELECT COUNT(*) as cnt FROM episode_queue WHERE status = ?', + ('pending',) + ).fetchone() + return row['cnt'] if row else 0 + + def is_empty(self) -> bool: + return self.size() == 0 + + def requeue(self, episode: Episode) -> bool: + with self.lock: + self.conn.execute( + 'UPDATE episode_queue SET status = ?, dequeued_at = ? WHERE id = ?', + ('pending', None, episode.id) + ) + self.conn.commit() + return True + + def close(self): + self.conn.close() \ No newline at end of file diff --git a/src/models/__init__.py b/src/models/__init__.py new file mode 100644 index 0000000..4da3e42 --- /dev/null +++ b/src/models/__init__.py @@ -0,0 +1,4 @@ +from .episode import Episode +from .distilled import Distilled + +__all__ = ["Episode", "Distilled"] \ No newline at end of file diff --git a/src/models/distilled.py b/src/models/distilled.py new file mode 100644 index 0000000..a8e14ff --- /dev/null +++ b/src/models/distilled.py @@ -0,0 +1,47 @@ +from dataclasses import dataclass, field +from datetime import datetime + +TYPE_DECISION = 'decision' +TYPE_REQUEST = 'request' +TYPE_FACT = 'fact' +TYPE_PATTERN = 'pattern' + +STATUS_PENDING = 'pending' +STATUS_VALIDATED = 'validated' +STATUS_DEPRECATED = 'deprecated' + +@dataclass +class Distilled: + id: str + episode_id: str + type: str + summary: str + entities: list[str] = field(default_factory=list) + facts: list[str] = field(default_factory=list) + confidence: float = 0.5 + status: str = 'pending' + importance: int = 0 + created_at: datetime = field(default_factory=datetime.now) + updated_at: datetime = field(default_factory=datetime.now) + + def to_dict(self) -> dict: + return { + 'id': self.id, + 'episode_id': self.episode_id, + 'type': self.type, + 'summary': self.summary, + 'entities': self.entities, + 'facts': self.facts, + 'confidence': self.confidence, + 'status': self.status, + 'importance': self.importance, + 'created_at': self.created_at.isoformat(), + 'updated_at': self.updated_at.isoformat(), + } + + @classmethod + def from_dict(cls, d: dict): + d = d.copy() + d['created_at'] = datetime.fromisoformat(d['created_at']) + d['updated_at'] = datetime.fromisoformat(d['updated_at']) + return cls(**d) \ No newline at end of file diff --git a/src/models/episode.py b/src/models/episode.py new file mode 100644 index 0000000..2e6f54a --- /dev/null +++ b/src/models/episode.py @@ -0,0 +1,43 @@ +from dataclasses import dataclass, field +from datetime import datetime +from typing import Optional +import uuid + +@dataclass +class Episode: + id: str + timestamp: datetime + content: str + entities: list[str] = field(default_factory=list) + facts: list[str] = field(default_factory=list) + metadata: dict = field(default_factory=dict) + source: str = 'hermes' + + def to_dict(self) -> dict: + return { + 'id': self.id, + 'timestamp': self.timestamp.isoformat(), + 'content': self.content, + 'entities': self.entities, + 'facts': self.facts, + 'metadata': self.metadata, + 'source': self.source, + } + + @classmethod + def from_dict(cls, d: dict): + d = d.copy() + d['timestamp'] = datetime.fromisoformat(d['timestamp']) + return cls(**d) + + @classmethod + def create(cls, content: str, source: str = 'hermes', entities: list = None, facts: list = None, metadata: dict = None): + return cls( + id=str(uuid.uuid4()), + timestamp=datetime.now(), + content=content, + entities=entities or [], + facts=facts or [], + metadata=metadata or {}, + source=source, + ) \ No newline at end of file diff --git a/tests/test_models.py b/tests/test_models.py new file mode 100644 index 0000000..65ab79b --- /dev/null +++ b/tests/test_models.py @@ -0,0 +1,48 @@ +import sys +from pathlib import Path +sys.path.insert(0, str(Path(__file__).parent.parent / 'src')) +from datetime import datetime +from models.episode import Episode +from models.distilled import Distilled, TYPE_FACT, STATUS_PENDING, STATUS_VALIDATED + +def test_episode_to_dict_from_dict(): + ep = Episode.create(content='test content', source='test', entities=['entity1'], facts=['fact1']) + d = ep.to_dict() + assert d['content'] == 'test content' + assert d['source'] == 'test' + ep2 = Episode.from_dict(d) + assert ep2.content == ep.content + assert ep2.id == ep.id + +def test_episode_create(): + ep = Episode.create(content='hello') + assert ep.id is not None + assert ep.content == 'hello' + assert ep.source == 'hermes' + assert isinstance(ep.timestamp, datetime) + +def test_distilled_to_dict_from_dict(): + ep = Episode.create(content='test') + d = Distilled( + id='d1', episode_id=ep.id, type=TYPE_FACT, + summary='a summary', confidence=0.8, status=STATUS_PENDING + ) + dd = d.to_dict() + assert dd['type'] == TYPE_FACT + assert dd['confidence'] == 0.8 + d2 = Distilled.from_dict(dd) + assert d2.id == d.id + assert d2.type == d.type + +def test_distilled_defaults(): + d = Distilled(id='d1', episode_id='e1', type=TYPE_FACT, summary='s') + assert d.status == STATUS_PENDING + assert d.confidence == 0.5 + assert d.importance == 0 + +if __name__ == '__main__': + test_episode_to_dict_from_dict() + test_episode_create() + test_distilled_to_dict_from_dict() + test_distilled_defaults() + print('All model tests passed!') \ No newline at end of file diff --git a/tests/test_queue.py b/tests/test_queue.py new file mode 100644 index 0000000..4dfd2f2 --- /dev/null +++ b/tests/test_queue.py @@ -0,0 +1,67 @@ +import sys +import os +import threading +from pathlib import Path +sys.path.insert(0, str(Path(__file__).parent.parent / 'src')) +from models.episode import Episode +from distill.queue import PersistenceQueue + +def test_queue_basic(): + q = PersistenceQueue(':memory:') + ep = Episode.create(content='test1') + assert q.enqueue(ep) == True + assert q.size() == 1 + assert q.is_empty() == False + dequeued = q.dequeue() + assert dequeued is not None + assert dequeued['content'] == 'test1' + assert q.size() == 0 + assert q.is_empty() == True + q.close() + +def test_queue_requeue(): + q = PersistenceQueue(':memory:') + ep = Episode.create(content='test2') + q.enqueue(ep) + deq = q.dequeue() + d = Episode.from_dict(deq) + assert q.requeue(d) == True + assert q.size() == 1 + q.close() + +def test_queue_peek(): + q = PersistenceQueue(':memory:') + ep1 = Episode.create(content='first') + ep2 = Episode.create(content='second') + q.enqueue(ep1) + q.enqueue(ep2) + peeked = q.peek() + assert peeked['content'] == 'first' + assert q.size() == 2 # peek does not remove + q.close() + +def test_queue_concurrent(): + q = PersistenceQueue(':memory:') + errors = [] + def enqueue_many(n): + try: + for i in range(n): + ep = Episode.create(content=f'thread-{i}') + q.enqueue(ep) + except Exception as e: + errors.append(str(e)) + threads = [threading.Thread(target=enqueue_many, args=(20,)) for _ in range(5)] + for t in threads: + t.start() + for t in threads: + t.join() + assert len(errors) == 0, errors + assert q.size() == 100 + q.close() + +if __name__ == '__main__': + test_queue_basic() + test_queue_requeue() + test_queue_peek() + test_queue_concurrent() + print('All queue tests passed!') \ No newline at end of file