-
Notifications
You must be signed in to change notification settings - Fork 2
feat(archwiz): novel work — Linear sync + path normalization + error logging #13
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
b10c62f
375b88e
a8ae313
78c30d0
40feea7
5abfd00
6642c41
4bc507c
0fe2e26
7812177
567d35f
7302e3c
3da2b3c
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Large diffs are not rendered by default.
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,120 @@ | ||
| import json | ||
| import re | ||
| import hashlib | ||
| from pathlib import Path | ||
| from typing import List, Dict, Optional, Any, Tuple | ||
| from datetime import datetime, timezone | ||
| from dataclasses import dataclass, field | ||
| from archwiz.config import ARCHWIZ_ROOT | ||
|
|
||
| @dataclass | ||
| class Pointer: | ||
| """Lightweight reference: (session_id, msg_idx, blk_idx) -> content_hash.""" | ||
| session_id: str | ||
| message_index: int | ||
| block_index: int | ||
| content_hash: str | ||
| start_line: int = 0 | ||
| end_line: int = 0 | ||
|
|
||
| def to_key(self) -> str: | ||
| return f"{self.session_id}:{self.message_index}:{self.block_index}" | ||
|
|
||
| def to_dict(self) -> Dict: | ||
| return { | ||
| "session_id": self.session_id, | ||
| "message_index": self.message_index, | ||
| "block_index": self.block_index, | ||
| "content_hash": self.content_hash, | ||
| "start_line": self.start_line, | ||
| "end_line": self.end_line, | ||
| } | ||
|
|
||
| @dataclass | ||
| class CodeBlock: | ||
| """Represents an extracted code block.""" | ||
| content: str | ||
| language: str | ||
| session_id: str | ||
| message_index: int | ||
| block_index: int | ||
| provider: str | ||
| timestamp: datetime = field(default_factory=lambda: datetime.now(timezone.utc)) | ||
| metadata: Dict = field(default_factory=dict) | ||
| content_hash: str = "" | ||
|
|
||
| def __post_init__(self): | ||
| if not self.content_hash and self.content: | ||
| self.content_hash = hashlib.sha256(self.content.encode()).hexdigest()[:16] | ||
|
|
||
| class CodexIndex: | ||
| """ | ||
| Content-addressed hierarchical index for code blocks. | ||
| Salvaged and normalized from PR #6 (TER-9). | ||
| """ | ||
| CODE_BLOCK_PATTERN = re.compile(r'```(\w+)?\n(.*?)```', re.DOTALL) | ||
|
|
||
| def __init__(self, provider: str = "global"): | ||
| self.provider = provider | ||
| self.base_dir = ARCHWIZ_ROOT / ".archwiz" / "codex" / provider | ||
| self.base_dir.mkdir(parents=True, exist_ok=True) | ||
| self.blobs_dir = self.base_dir / "blobs" | ||
| self.blobs_dir.mkdir(exist_ok=True) | ||
|
|
||
| self.index_file = self.base_dir / "index.json" | ||
| self.pointers: List[Pointer] = [] | ||
| self.blobs: Dict[str, str] = {} # hash -> path | ||
|
|
||
| self._load() | ||
|
|
||
| def _load(self): | ||
| if self.index_file.exists(): | ||
| try: | ||
| data = json.loads(self.index_file.read_text()) | ||
| for p_data in data.get("pointers", []): | ||
| self.pointers.append(Pointer(**p_data)) | ||
| except Exception: | ||
| pass | ||
|
|
||
| for blob_file in self.blobs_dir.glob("*.blob"): | ||
| self.blobs[blob_file.stem] = str(blob_file) | ||
|
|
||
| def _save(self): | ||
| data = { | ||
| "version": 1, | ||
| "provider": self.provider, | ||
| "updated_at": datetime.now(timezone.utc).isoformat(), | ||
| "pointers": [p.to_dict() for p in self.pointers] | ||
| } | ||
| self.index_file.write_text(json.dumps(data, indent=2)) | ||
|
|
||
| def harvest(self, session_id: str, messages: List[Dict[str, Any]]) -> int: | ||
| """Extract and index code blocks from messages.""" | ||
| count = 0 | ||
| for msg_idx, msg in enumerate(messages): | ||
| content = msg.get("content", "") | ||
| for blk_idx, match in enumerate(self.CODE_BLOCK_PATTERN.finditer(content)): | ||
| lang = (match.group(1) or "text").lower() | ||
| code = match.group(2) | ||
|
Comment on lines
+94
to
+98
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 📝 Info: Harvest assumes message content is a string
Was this helpful? React with 👍 or 👎 to provide feedback. |
||
| ch = hashlib.sha256(code.encode()).hexdigest()[:16] | ||
|
|
||
| # Store blob | ||
| blob_path = self.blobs_dir / f"{ch}.blob" | ||
| if not blob_path.exists(): | ||
| blob_path.write_text(code) | ||
| self.blobs[ch] = str(blob_path) | ||
|
|
||
| # Record pointer | ||
| p = Pointer(session_id, msg_idx, blk_idx, ch) | ||
| self.pointers.append(p) | ||
| count += 1 | ||
|
Comment on lines
+107
to
+110
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🟡 Code index accumulates duplicate entries each time a session is re-scanned Extracted code references are always added to the list ( No dedup on (session_id, message_index, block_index)
Prompt for agentsWas this helpful? React with 👍 or 👎 to provide feedback. |
||
|
|
||
| if count > 0: | ||
| self._save() | ||
| return count | ||
|
|
||
| def get_code(self, content_hash: str) -> Optional[str]: | ||
| path = self.blobs.get(content_hash) | ||
| if path and Path(path).exists(): | ||
| return Path(path).read_text() | ||
| return None | ||
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| @@ -1,46 +1,102 @@ | ||||||||||||||||||||||||||||
| #!/usr/bin/env python3 | ||||||||||||||||||||||||||||
| """Additive post‑ingestion dispatch. Called by core.py after every session fetch. | ||||||||||||||||||||||||||||
| Non‑blocking: all heavy operations run via background scripts.""" | ||||||||||||||||||||||||||||
| """ | ||||||||||||||||||||||||||||
| Enhanced event-sourced dispatch pipeline for ArchWiz. | ||||||||||||||||||||||||||||
| Decouples session ingestion from downstream updates (SSOT, Codex, Linear, etc.) | ||||||||||||||||||||||||||||
| """ | ||||||||||||||||||||||||||||
| import json | ||||||||||||||||||||||||||||
| import logging | ||||||||||||||||||||||||||||
| import sys | ||||||||||||||||||||||||||||
| import time | ||||||||||||||||||||||||||||
| from pathlib import Path | ||||||||||||||||||||||||||||
| from typing import List, Dict, Any, Callable | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| HOME = Path.home() | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| def update_all(session_id: str, account: str = "primary"): | ||||||||||||||||||||||||||||
| """Run downstream updates for a session. Lightweight — no subprocess calls.""" | ||||||||||||||||||||||||||||
| store = HOME / '.deepcli' / 'session_store' / account / f'{session_id}.json' | ||||||||||||||||||||||||||||
| # Fallback to flat store if not in account subdir | ||||||||||||||||||||||||||||
| if not store.exists(): | ||||||||||||||||||||||||||||
| store = HOME / '.deepcli' / 'session_store' / f'{session_id}.json' | ||||||||||||||||||||||||||||
| if not store.exists(): | ||||||||||||||||||||||||||||
| return | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| try: | ||||||||||||||||||||||||||||
| with open(store) as f: | ||||||||||||||||||||||||||||
| msgs = json.load(f) | ||||||||||||||||||||||||||||
| except Exception: | ||||||||||||||||||||||||||||
| return | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| # 1. Write session.json into export dir | ||||||||||||||||||||||||||||
| export_dir = HOME / 'synthegration_exports' / account / session_id | ||||||||||||||||||||||||||||
| export_dir.mkdir(parents=True, exist_ok=True) | ||||||||||||||||||||||||||||
| (export_dir / 'session.json').write_text(json.dumps(msgs, indent=2)) | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| # 2. Lexicon harvest (lightweight: single session) | ||||||||||||||||||||||||||||
| try: | ||||||||||||||||||||||||||||
| sys.path.insert(0, str(HOME / 'archwiz')) | ||||||||||||||||||||||||||||
| from lexicon_harvest import harvest_session | ||||||||||||||||||||||||||||
| harvest_session(session_id) | ||||||||||||||||||||||||||||
| except Exception: | ||||||||||||||||||||||||||||
| pass | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| # 3. Codex index (incremental add — does NOT rebuild full index) | ||||||||||||||||||||||||||||
| try: | ||||||||||||||||||||||||||||
| sys.path.insert(0, str(HOME / 'cli-synthegration')) | ||||||||||||||||||||||||||||
| from synthegration_index import CodexIndex | ||||||||||||||||||||||||||||
| codex = CodexIndex(HOME / 'cli-synthegration' / 'codex') | ||||||||||||||||||||||||||||
| codex.index_conversation(session_id, session_id[:8], msgs) | ||||||||||||||||||||||||||||
| codex._save() | ||||||||||||||||||||||||||||
| except Exception: | ||||||||||||||||||||||||||||
| pass | ||||||||||||||||||||||||||||
| # Add root to path | ||||||||||||||||||||||||||||
| sys.path.insert(0, str(Path(__file__).resolve().parents[1])) | ||||||||||||||||||||||||||||
| from archwiz.config import ARCHWIZ_ROOT, LOG_DIR, SSOT_DIR | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| # Setup logging | ||||||||||||||||||||||||||||
| logging.basicConfig( | ||||||||||||||||||||||||||||
| filename=LOG_DIR / "dispatch.log", | ||||||||||||||||||||||||||||
| level=logging.INFO, | ||||||||||||||||||||||||||||
| format="%(asctime)s [%(levelname)s] %(message)s" | ||||||||||||||||||||||||||||
| ) | ||||||||||||||||||||||||||||
| logger = logging.getLogger("dispatch") | ||||||||||||||||||||||||||||
|
Comment on lines
+18
to
+23
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🔍 Importing the dispatch pipeline reconfigures the whole application's logging
Was this helpful? React with 👍 or 👎 to provide feedback. |
||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| class DispatchPipeline: | ||||||||||||||||||||||||||||
| def __init__(self): | ||||||||||||||||||||||||||||
| self.dispatchers: List[Callable[[str, List[Dict[str, Any]]], None]] = [] | ||||||||||||||||||||||||||||
| self._register_default_dispatchers() | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| def register(self, func: Callable[[str, List[Dict[str, Any]]], None]): | ||||||||||||||||||||||||||||
| self.dispatchers.append(func) | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| def _register_default_dispatchers(self): | ||||||||||||||||||||||||||||
| self.register(self.dispatch_ssot) | ||||||||||||||||||||||||||||
| self.register(self.dispatch_codex) | ||||||||||||||||||||||||||||
| self.register(self.dispatch_linear_hint) | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| def dispatch_ssot(self, session_id: str, messages: List[Dict[str, Any]]): | ||||||||||||||||||||||||||||
| """Sync to Session SSOT.""" | ||||||||||||||||||||||||||||
| try: | ||||||||||||||||||||||||||||
| from archwiz.session_ssot import SessionSSOT | ||||||||||||||||||||||||||||
| ssot = SessionSSOT() | ||||||||||||||||||||||||||||
| ssot.sync_session(session_id, messages) | ||||||||||||||||||||||||||||
| logger.info(f"SSOT sync successful for {session_id}") | ||||||||||||||||||||||||||||
| except Exception as e: | ||||||||||||||||||||||||||||
| logger.error(f"SSOT dispatch failed for {session_id}: {e}") | ||||||||||||||||||||||||||||
|
Comment on lines
+40
to
+46
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🔴 Session data is never written to the new canonical store The canonical session store is asked to save data with no session identity supplied ( Constructor signature and missing method mismatch
Suggested change
Was this helpful? React with 👍 or 👎 to provide feedback. |
||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| def dispatch_codex(self, session_id: str, messages: List[Dict[str, Any]]): | ||||||||||||||||||||||||||||
| """Harvest code blocks into Codex.""" | ||||||||||||||||||||||||||||
| try: | ||||||||||||||||||||||||||||
| from archwiz.codex import CodexIndex | ||||||||||||||||||||||||||||
| # We use a global index for the pipeline | ||||||||||||||||||||||||||||
| codex = CodexIndex(provider="pipeline") | ||||||||||||||||||||||||||||
| count = codex.harvest(session_id, messages) | ||||||||||||||||||||||||||||
| logger.info(f"Codex harvested {count} blocks from {session_id}") | ||||||||||||||||||||||||||||
| except Exception as e: | ||||||||||||||||||||||||||||
| logger.error(f"Codex dispatch failed for {session_id}: {e}") | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| def dispatch_linear_hint(self, session_id: str, messages: List[Dict[str, Any]]): | ||||||||||||||||||||||||||||
| """Hint that a session might need Linear sync if it contains task updates.""" | ||||||||||||||||||||||||||||
| # Simple heuristic: look for "done" or "task" in the last message | ||||||||||||||||||||||||||||
| if not messages: | ||||||||||||||||||||||||||||
| return | ||||||||||||||||||||||||||||
| last_msg = messages[-1].get("content", "").lower() | ||||||||||||||||||||||||||||
| if any(kw in last_msg for kw in ["done", "fixed", "implemented", "task"]): | ||||||||||||||||||||||||||||
| logger.info(f"Session {session_id} marked for Linear sync review") | ||||||||||||||||||||||||||||
| # In a full implementation, we might trigger linear_sync.py here | ||||||||||||||||||||||||||||
| # For now, we just log the hint. | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| def run(self, session_id: str, messages: List[Dict[str, Any]]): | ||||||||||||||||||||||||||||
| """Execute all registered dispatchers.""" | ||||||||||||||||||||||||||||
| start_time = time.time() | ||||||||||||||||||||||||||||
| logger.info(f"Starting dispatch for session {session_id} ({len(messages)} messages)") | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| for dispatcher in self.dispatchers: | ||||||||||||||||||||||||||||
| try: | ||||||||||||||||||||||||||||
| dispatcher(session_id, messages) | ||||||||||||||||||||||||||||
| except Exception as e: | ||||||||||||||||||||||||||||
| logger.error(f"Dispatcher {dispatcher.__name__} crashed: {e}") | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| duration = time.time() - start_time | ||||||||||||||||||||||||||||
| logger.info(f"Dispatch completed for {session_id} in {duration:.2f}s") | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| def trigger_dispatch(session_id: str, messages: List[Dict[str, Any]]): | ||||||||||||||||||||||||||||
| """Entry point for core.py and other ingestors.""" | ||||||||||||||||||||||||||||
| pipeline = DispatchPipeline() | ||||||||||||||||||||||||||||
| pipeline.run(session_id, messages) | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| if __name__ == "__main__": | ||||||||||||||||||||||||||||
| if len(sys.argv) < 2: | ||||||||||||||||||||||||||||
| print("Usage: dispatch_pipeline.py <session_id> [path_to_json]") | ||||||||||||||||||||||||||||
| sys.exit(1) | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| sid = sys.argv[1] | ||||||||||||||||||||||||||||
| msgs = [] | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| if len(sys.argv) > 2: | ||||||||||||||||||||||||||||
| p = Path(sys.argv[2]) | ||||||||||||||||||||||||||||
| if p.exists(): | ||||||||||||||||||||||||||||
| msgs = json.loads(p.read_text()) | ||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||
| trigger_dispatch(sid, msgs) | ||||||||||||||||||||||||||||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
📝 Info: Pointer loading trusts arbitrary keys from the index file
Pointer(**p_data)will raise TypeError on unexpected/missing keys from a hand-edited or older-schema index.json; the surroundingexcept Exception: passswallows it but leavesself.pointerspartially populated with whatever was appended before the failure, and the next_save()then silently truncates the index to that partial state. Filteringp_datato known fields would avoid silent data loss.Was this helpful? React with 👍 or 👎 to provide feedback.