import argparse
import json
import os
from pathlib import Path
import re
import sqlite3
import threading
import time
from session_context import SessionContext, codex_homes, claude_roots, orca_codex_rollouts
from datetime import datetime
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
ROOT = Path(__file__).parent
HOME = Path.home()
def timestamp(value):
try:
return datetime.fromisoformat(value.replace('Z', '+00:00')).timestamp()
except (ValueError, AttributeError):
return 0
def clean(value, limit=240):
if not isinstance(value, str):
return ''
match = re.search(r'(.*?)', value, re.S)
if match:
value = match[1]
if '## My request:' in value:
value = value.split('## My request:', 1)[1]
return ' '.join(value.split())[:limit]
BODY_LIMIT = 4000
EVENT_TAIL = 40
def body(value, limit=BODY_LIMIT):
"""Message text for the reader: like clean(), but line structure survives so
markdown (tables, lists, code) can render; one extra character marks a cut."""
if not isinstance(value, str):
return ''
match = re.search(r'(.*?)', value, re.S)
if match:
value = match[1]
if '## My request:' in value:
value = value.split('## My request:', 1)[1]
lines = [' '.join(line.split()) if not line.startswith((' ', '\t')) else line.rstrip() for line in value.strip().splitlines()]
text = re.sub(r'\n{3,}', '\n\n', '\n'.join(lines))
return text[:limit + 1]
def records(path, head=False):
try:
with path.open('rb') as f:
size = path.stat().st_size
if not head and size > 262144:
f.seek(size - 262144)
f.readline()
data = f.read(262144)
result = []
for line in data.splitlines():
try:
item = json.loads(line)
if isinstance(item, dict):
result.append(item)
except (ValueError, UnicodeDecodeError):
pass
return result
except OSError:
return []
class Scanner:
def __init__(self, codex=None, claude=None, home=None):
self.home = home or (codex.parent if codex else HOME)
self.codex = codex or self.home / '.codex'
self.context = SessionContext(self.home)
self.claude = claude or self.home / '.claude'
self.cache = {}
def parse(self, path, provider):
try:
stat = path.stat()
except OSError:
return {'events': [], 'updated': 0, 'signal': '', 'cwd': '', 'model': '', 'title': ''}
key = str(path)
version = (stat.st_mtime_ns, stat.st_size)
if key in self.cache and self.cache[key][0] == version:
return self.cache[key][1]
tail = records(path)
head = records(path, head=True)
data = {'events': [], 'updated': 0, 'signal': '', 'cwd': '', 'model': '', 'title': ''}
for row in head + tail:
p = row.get('payload') or {}
if not isinstance(p, dict):
p = {}
data['cwd'] = row.get('cwd') or p.get('cwd') or data['cwd']
msg = row.get('message') or {}
if not isinstance(msg, dict):
msg = {}
data['model'] = row.get('modelId') or msg.get('model') or p.get('model') or data['model']
if provider == 'Pi' and row.get('type') == 'message' and msg.get('role') == 'user' and not data['title']:
content = msg.get('content', [])
data['title'] = clean(' '.join(x.get('text', '') for x in content if isinstance(x, dict) and x.get('type') == 'text'))
if provider == 'Codex' and row.get('type') == 'event_msg' and p.get('type') == 'user_message' and not data['title']:
data['title'] = clean(p.get('message'))
if row.get('type') == 'custom-title':
data['title'] = clean(row.get('customTitle'))
if provider == 'Claude' and row.get('type') == 'user' and not data['title']:
content = msg.get('content')
data['title'] = clean(content if isinstance(content, str) else ' '.join(x.get('text', '') for x in content or [] if isinstance(x, dict) and x.get('type') == 'text'))
for row in tail:
p = row.get('payload') or {}
if not isinstance(p, dict):
p = {}
at = timestamp(row.get('timestamp'))
kind = row.get('type')
label, detail, voice = '', '', False
if provider == 'Codex':
subtype = p.get('type')
if kind == 'event_msg':
if subtype in ('task_started', 'turn_started'):
data['signal'] = 'started'
label = 'Turn started'
elif subtype in ('task_complete', 'turn_complete', 'turn_aborted'):
data['signal'] = 'completed'
label = 'Turn ended'
elif subtype == 'user_message':
label, detail = 'You', body(p.get('message'))
if kind == 'response_item':
if subtype in ('function_call', 'custom_tool_call'):
label, detail = 'Tool', clean(p.get('name'))
elif subtype == 'message' and p.get('role') == 'assistant':
label = 'Assistant'
detail = body('\n'.join(x.get('text', '') for x in p.get('content', []) if isinstance(x, dict)))
if p.get('phase') == 'final':
data['signal'] = 'completed'
if kind == 'realtime_item' and subtype == 'transcript_segment' and p.get('role') == 'user' and isinstance(p.get('text'), str) and p['text'].strip():
label, detail, voice = 'You', body(p['text']), True
else:
msg = row.get('message') or {}
if not isinstance(msg, dict):
msg = {}
if provider == 'Pi' and kind == 'message':
kind = msg.get('role')
content = msg.get('content')
if kind in ('user', 'assistant'):
if isinstance(content, str):
label, detail = ('You' if kind == 'user' else 'Assistant'), body(content)
elif isinstance(content, list):
texts = [x.get('text', '') for x in content if isinstance(x, dict) and x.get('type') == 'text']
calls = [x.get('name', '') for x in content if isinstance(x, dict) and x.get('type') == 'tool_use']
if texts:
label, detail = ('You' if kind == 'user' else 'Assistant'), body('\n'.join(texts))
elif calls:
label, detail = 'Tool', ', '.join(calls)
if kind == 'user':
data['signal'] = 'started'
if msg.get('stop_reason') == 'end_turn' or msg.get('stopReason') == 'stop':
data['signal'] = 'completed'
if kind == 'system' and row.get('subtype') in ('turn_duration', 'stop_hook_summary'):
data['signal'] = 'completed'
label = 'Turn ended'
if at:
data['updated'] = max(data['updated'], at)
if label:
event = {'at': at, 'label': label, 'text': detail}
if voice:
event['voice'] = True
data['events'].append(event)
data['events'] = data['events'][-EVENT_TAIL:]
self.cache[key] = (version, data)
return data
def snapshot(self):
sessions, errors = [], []
for codex in dict.fromkeys([self.codex, *codex_homes(self.home)[1:]]):
databases = sorted(codex.glob('state_*.sqlite'), key=lambda p: int(p.stem.split('_')[-1]))
if databases:
try:
with sqlite3.connect(f'file:{databases[-1]}?mode=ro', uri=True, timeout=1) as db:
db.row_factory = sqlite3.Row
rows = db.execute('select id,title,name,cwd,model,git_branch,git_origin_url,archived,rollout_path,updated_at from threads').fetchall()
for r in rows:
p = self.parse(Path(r['rollout_path']), 'Codex')
sessions.append(dict(id=r['id'], provider='Codex', title=clean(r['name'] or r['title']) or 'Untitled session', cwd=r['cwd'], model=r['model'] or p['model'], branch=r['git_branch'] or '', repo_origin=r['git_origin_url'] or '', archived=bool(r['archived']), updated=p['updated'] or r['updated_at'], signal=p['signal'], events=p['events']))
except (sqlite3.Error, OSError) as exc:
errors.append(f'Codex: {exc}')
known = {s['id'] for s in sessions if s['provider'] == 'Codex'}
for path, meta, archived in orca_codex_rollouts(self.home):
if meta['id'] in known:
continue
p = self.parse(path, 'Codex')
git = meta.get('git') or {}
sessions.append(dict(id=meta['id'], provider='Codex', title=p['title'] or 'Codex conversation', cwd=meta.get('cwd', ''), model=p['model'], branch=git.get('branch', ''), repo_origin=git.get('repository_url', ''), archived=archived, updated=p['updated'], signal=p['signal'], events=p['events'], subagent='subagent' in str(meta.get('source', ''))))
for claude in dict.fromkeys([self.claude / 'projects', *claude_roots(self.home)[1:]]):
for path in claude.rglob('*.jsonl'):
p = self.parse(path, 'Claude')
sessions.append(dict(id=str(path.relative_to(claude)), provider='Claude', title=p['title'] or path.stem, cwd=p['cwd'], model=p['model'], branch='', archived=False, updated=p['updated'] or path.stat().st_mtime, signal=p['signal'], events=p['events'], subagent='subagents' in path.parts))
for path in (self.home / '.pi/agent/sessions').glob('**/*.jsonl'):
header = next((r for r in records(path, head=True) if r.get('type') == 'session'), {})
if not header.get('id'):
continue
p = self.parse(path, 'Pi')
sessions.append(dict(id=header['id'], provider='Pi', title=p['title'] or 'Pi conversation', cwd=header.get('cwd', ''), model=p['model'], branch='', archived=False, updated=p['updated'], signal=p['signal'], events=p['events']))
if not any(s['provider'] == 'Codex' for s in sessions):
errors.append('Codex session history not found')
if not any(root.exists() for root in [self.claude / 'projects', *claude_roots(self.home)[1:]]):
errors.append('Claude Code session directory not found')
unique = {}
for session in sessions:
key = (session['provider'], session['id'])
if key not in unique or session['updated'] > unique[key]['updated']:
unique[key] = session
sessions = list(unique.values())
now = time.time()
for s in sessions:
age = now - s['updated']
s['status'] = 'completed' if s['signal'] == 'completed' else 'recent' if age < 120 else 'quiet'
s['project'] = Path(s['cwd']).name or 'Unknown project'
s.update(self.context.get(s['cwd'], s.get('repo_origin', ''), s['branch']))
if s['provider'] == 'Pi' and s.get('worktree') and s['worktree'] != 'main':
s['display_title'] = s['worktree']
errors.extend(self.context.errors)
sessions.sort(key=lambda s: s['updated'], reverse=True)
return {'sessions': sessions, 'errors': errors, 'at': now}
class Handler(BaseHTTPRequestHandler):
def do_GET(self):
if self.headers.get('Host') not in (f'127.0.0.1:{self.server.server_port}', f'localhost:{self.server.server_port}'):
self.send_error(403)
return
if self.path == '/api/sessions':
with self.server.snapshot_lock:
body = self.server.snapshot_bytes
mime = 'application/json'
elif self.path in ('/', '/app.js', '/style.css'):
name = 'index.html' if self.path == '/' else self.path[1:]
body = (ROOT / name).read_bytes()
mime = {'index.html': 'text/html', 'app.js': 'text/javascript', 'style.css': 'text/css'}[name]
else:
self.send_error(404)
return
self.send_response(200)
self.send_header('Content-Type', mime + '; charset=utf-8')
self.send_header('Cache-Control', 'no-store')
self.send_header('X-Content-Type-Options', 'nosniff')
self.send_header('Content-Security-Policy', "default-src 'self'; script-src 'self'; style-src 'self'; connect-src 'self'; frame-ancestors 'none'")
self.end_headers()
self.wfile.write(body)
def log_message(self, *_):
pass
def main():
parser = argparse.ArgumentParser()
parser.add_argument('--port', type=int, default=8765)
args = parser.parse_args()
scanner = Scanner()
server = ThreadingHTTPServer(('127.0.0.1', args.port), Handler)
server.snapshot_lock = threading.Lock()
server.snapshot_bytes = json.dumps(scanner.snapshot()).encode()
def refresh():
while True:
time.sleep(2)
try:
data = json.dumps(scanner.snapshot()).encode()
with server.snapshot_lock:
server.snapshot_bytes = data
except Exception as exc:
with server.snapshot_lock:
server.snapshot_bytes = json.dumps({'sessions': [], 'errors': [f'Scan failed: {exc}'], 'at': time.time()}).encode()
threading.Thread(target=refresh, daemon=True).start()
print(f'Session Observatory: http://127.0.0.1:{args.port}', flush=True)
server.serve_forever()
if __name__ == '__main__':
main()