Kode-Animator's picture
remove temporary transport diagnostic
fdb09a5
Raw
History Blame Contribute Delete
19.8 kB
"""
==========================================================================
🧠 Thought Engine Node β€” Sovereign REST Substrate
==========================================================================
Contract: C-THOUGHT-NODE-001 (Phase 1)
Vault: thought-vault (Cloudflare D1)
Stack: FastAPI + httpx + Cloudflare D1 HTTP API
A persistent, queryable, forkable reasoning system with provenance,
pattern memory, and multiple entry surfaces.
==========================================================================
"""
import os
import uuid
import json
from datetime import datetime, timezone
from typing import Optional, List
from dotenv import load_dotenv
load_dotenv()
from fastapi import FastAPI, HTTPException, Query
from fastapi.middleware.cors import CORSMiddleware
import d1_client
from models import (
CreateSessionReq, AddThoughtReq, ForkThoughtReq,
CreateProposalReq, ReviewProposalReq, SearchReq,
SessionInfo, ThoughtInfo, EdgeInfo, TreeNode,
ThoughtClass, ThoughtStatus, EdgeRelation, ReviewAction,
)
# ── Application ────────────────────────────────────────────────
app = FastAPI(
title="Thought Engine Node",
description="Sovereign reasoning substrate β€” C-THOUGHT-NODE-001",
version="0.1.0",
)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
NODE_SEAL = None
def _now() -> str:
return datetime.now(timezone.utc).isoformat()
def _id(prefix: str = "t") -> str:
return f"{prefix}-{uuid.uuid4().hex[:12]}"
# ── Lifecycle ──────────────────────────────────────────────────
@app.on_event("startup")
async def startup():
global NODE_SEAL
seg = uuid.uuid4().hex[:8].upper()
NODE_SEAL = f"⟦ THOUGHT :: PANTHEON-TE :: 🧠-{datetime.now().strftime('%Y%m%d')}-{seg[:4]}-{seg[4:]} :: ACTIVE ⟧"
print(f"\n{'='*60}")
print(f"🧠 Thought Engine Node β€” Phase 1 Sovereign Substrate")
print(f" Seal: {NODE_SEAL}")
print(f" Time: {_now()}")
print(f"{'='*60}\n")
# ── Identity ───────────────────────────────────────────────────
@app.get("/")
async def identity():
count = await d1_client.execute_sql("SELECT COUNT(*) as c FROM sessions")
session_count = count[0]["c"] if count else 0
return {
"node": "Thought Engine Node",
"contract": "C-THOUGHT-NODE-001",
"version": "0.1.0",
"seal": NODE_SEAL,
"sessions": session_count,
"status": "ACTIVE",
"timestamp": _now(),
}
# ══════════════════════════════════════════════════════════════
# SESSIONS
# ══════════════════════════════════════════════════════════════
@app.post("/session")
async def create_session(req: CreateSessionReq):
"""Start a new reasoning session with an initial root thought."""
session_id = _id("ses")
thought_id = _id("th")
title = req.title or req.initial_thought[:80]
# Insert session
await d1_client.execute_sql(
"INSERT INTO sessions (session_id, title, created_by, created_at, status, root_thought_id, active_thought_id) "
"VALUES (?, ?, ?, ?, 'active', ?, ?)",
[session_id, title, req.agent_id, _now(), thought_id, thought_id],
)
# Insert root thought
await d1_client.execute_sql(
"INSERT INTO thought_units (thought_id, session_id, content, thought_class, origin_agent, status, confidence, created_at) "
"VALUES (?, ?, ?, ?, ?, 'open', ?, ?)",
[thought_id, session_id, req.initial_thought, req.thought_class.value, req.agent_id, None, _now()],
)
# Witness event
await d1_client.execute_sql(
"INSERT INTO witness_events (event_id, session_id, thought_id, event_type, actor, detail, created_at) "
"VALUES (?, ?, ?, 'created', ?, 'Session started', ?)",
[_id("ev"), session_id, thought_id, req.agent_id, _now()],
)
return {
"session_id": session_id,
"root_thought_id": thought_id,
"title": title,
"status": "active",
"message": "🧠 Reasoning session created.",
}
@app.get("/session/{session_id}")
async def get_session(session_id: str):
"""Retrieve session metadata."""
rows = await d1_client.execute_sql(
"SELECT * FROM sessions WHERE session_id = ?", [session_id]
)
if not rows:
raise HTTPException(404, "Session not found")
return rows[0]
@app.get("/sessions")
async def list_sessions(
status: Optional[str] = Query(None),
limit: int = Query(20, ge=1, le=100),
):
"""List all sessions, optionally filtered by status."""
if status:
rows = await d1_client.execute_sql(
"SELECT * FROM sessions WHERE status = ? ORDER BY created_at DESC LIMIT ?",
[status, limit],
)
else:
rows = await d1_client.execute_sql(
"SELECT * FROM sessions ORDER BY created_at DESC LIMIT ?", [limit]
)
return {"total": len(rows), "sessions": rows}
# ══════════════════════════════════════════════════════════════
# THOUGHTS
# ══════════════════════════════════════════════════════════════
@app.post("/session/{session_id}/thought")
async def add_thought(session_id: str, req: AddThoughtReq):
"""Add a reasoning step to the session's active chain."""
# Verify session exists
ses = await d1_client.execute_sql(
"SELECT * FROM sessions WHERE session_id = ?", [session_id]
)
if not ses:
raise HTTPException(404, "Session not found")
thought_id = _id("th")
parent_id = req.parent_thought_id or ses[0]["active_thought_id"]
# Insert thought unit
await d1_client.execute_sql(
"INSERT INTO thought_units (thought_id, session_id, content, thought_class, origin_agent, status, confidence, created_at) "
"VALUES (?, ?, ?, ?, ?, 'open', ?, ?)",
[thought_id, session_id, req.content, req.thought_class.value, req.agent_id, req.confidence, _now()],
)
# Insert edge from parent
await d1_client.execute_sql(
"INSERT INTO thought_edges (edge_id, session_id, source_id, target_id, relation, created_at) "
"VALUES (?, ?, ?, ?, 'derives_from', ?)",
[_id("ed"), session_id, parent_id, thought_id, _now()],
)
# Update active pointer
await d1_client.execute_sql(
"UPDATE sessions SET active_thought_id = ?, updated_at = ? WHERE session_id = ?",
[thought_id, _now(), session_id],
)
# Witness
await d1_client.execute_sql(
"INSERT INTO witness_events (event_id, session_id, thought_id, event_type, actor, detail, created_at) "
"VALUES (?, ?, ?, 'created', ?, ?, ?)",
[_id("ev"), session_id, thought_id, req.agent_id, f"Added {req.thought_class.value}", _now()],
)
return {
"thought_id": thought_id,
"parent_id": parent_id,
"thought_class": req.thought_class.value,
"message": f"βœ… Thought added to session.",
}
@app.get("/session/{session_id}/thoughts")
async def list_thoughts(session_id: str):
"""List all thought units in a session."""
rows = await d1_client.execute_sql(
"SELECT * FROM thought_units WHERE session_id = ? ORDER BY created_at ASC",
[session_id],
)
return {"session_id": session_id, "total": len(rows), "thoughts": rows}
# ══════════════════════════════════════════════════════════════
# FORKING (Git-for-Thought Branching)
# ══════════════════════════════════════════════════════════════
@app.post("/session/{session_id}/fork")
async def fork_thought(session_id: str, req: ForkThoughtReq):
"""Fork a thought chain β€” create a branch for alternative exploration."""
# Verify source thought
source = await d1_client.execute_sql(
"SELECT * FROM thought_units WHERE thought_id = ? AND session_id = ?",
[req.source_thought_id, session_id],
)
if not source:
raise HTTPException(404, "Source thought not found in this session")
fork_id = _id("th")
src = source[0]
# Create the forked thought node (copy of source with new id)
await d1_client.execute_sql(
"INSERT INTO thought_units (thought_id, session_id, content, thought_class, origin_agent, status, confidence, created_at, metadata) "
"VALUES (?, ?, ?, ?, ?, 'open', ?, ?, ?)",
[fork_id, session_id, src["content"], src["thought_class"], req.agent_id, src.get("confidence"),
_now(), json.dumps({"forked_from": req.source_thought_id, "branch_label": req.branch_label})],
)
# Edge: forks_from
await d1_client.execute_sql(
"INSERT INTO thought_edges (edge_id, session_id, source_id, target_id, relation, created_at) "
"VALUES (?, ?, ?, ?, 'forks_from', ?)",
[_id("ed"), session_id, req.source_thought_id, fork_id, _now()],
)
# Move active pointer to the fork
await d1_client.execute_sql(
"UPDATE sessions SET active_thought_id = ?, updated_at = ? WHERE session_id = ?",
[fork_id, _now(), session_id],
)
# Witness
await d1_client.execute_sql(
"INSERT INTO witness_events (event_id, session_id, thought_id, event_type, actor, detail, created_at) "
"VALUES (?, ?, ?, 'forked', ?, ?, ?)",
[_id("ev"), session_id, fork_id, req.agent_id, f"Forked from {req.source_thought_id} as '{req.branch_label}'", _now()],
)
return {
"forked_thought_id": fork_id,
"source_thought_id": req.source_thought_id,
"branch_label": req.branch_label,
"message": f"🌿 Forked thought chain: {req.branch_label}",
}
# ══════════════════════════════════════════════════════════════
# PROPOSALS (Git-for-Thought PRs)
# ══════════════════════════════════════════════════════════════
@app.post("/session/{session_id}/proposal")
async def create_proposal(session_id: str, req: CreateProposalReq):
"""Submit a thought proposal (PR) branching from a parent."""
parent = await d1_client.execute_sql(
"SELECT * FROM thought_units WHERE thought_id = ? AND session_id = ?",
[req.parent_thought_id, session_id],
)
if not parent:
raise HTTPException(404, "Parent thought not found")
proposal_id = _id("pr")
# Insert proposal thought
await d1_client.execute_sql(
"INSERT INTO thought_units (thought_id, session_id, content, thought_class, origin_agent, status, confidence, created_at, metadata) "
"VALUES (?, ?, ?, 'proposal', ?, 'proposed', NULL, ?, ?)",
[proposal_id, session_id, req.content, req.agent_id, _now(),
json.dumps({"proposal_note": req.note, "target_branch": req.parent_thought_id})],
)
# Edge
await d1_client.execute_sql(
"INSERT INTO thought_edges (edge_id, session_id, source_id, target_id, relation, created_at) "
"VALUES (?, ?, ?, ?, 'derives_from', ?)",
[_id("ed"), session_id, req.parent_thought_id, proposal_id, _now()],
)
# Review record
await d1_client.execute_sql(
"INSERT INTO thought_reviews (review_id, session_id, thought_id, action, actor, reason, created_at) "
"VALUES (?, ?, ?, 'propose', ?, ?, ?)",
[_id("rv"), session_id, proposal_id, req.agent_id, req.note, _now()],
)
return {
"proposal_id": proposal_id,
"parent_id": req.parent_thought_id,
"status": "proposed",
"message": f"πŸ“ Proposal submitted by {req.agent_id}",
}
@app.get("/session/{session_id}/proposals")
async def list_proposals(session_id: str):
"""List all pending proposals in a session."""
rows = await d1_client.execute_sql(
"SELECT * FROM thought_units WHERE session_id = ? AND status = 'proposed' ORDER BY created_at ASC",
[session_id],
)
return {"session_id": session_id, "total": len(rows), "proposals": rows}
@app.post("/session/{session_id}/proposal/{proposal_id}/review")
async def review_proposal(session_id: str, proposal_id: str, req: ReviewProposalReq):
"""Accept, reject, or supersede a proposal."""
proposal = await d1_client.execute_sql(
"SELECT * FROM thought_units WHERE thought_id = ? AND session_id = ? AND status = 'proposed'",
[proposal_id, session_id],
)
if not proposal:
raise HTTPException(404, "Proposal not found or not in 'proposed' status")
new_status_map = {
ReviewAction.ACCEPT: "accepted",
ReviewAction.REJECT: "rejected",
ReviewAction.SUPERSEDE: "superseded",
}
new_status = new_status_map.get(req.action)
if not new_status:
raise HTTPException(400, "Invalid review action for this endpoint")
# Update thought status
await d1_client.execute_sql(
"UPDATE thought_units SET status = ? WHERE thought_id = ?",
[new_status, proposal_id],
)
# Review record
await d1_client.execute_sql(
"INSERT INTO thought_reviews (review_id, session_id, thought_id, action, actor, reason, created_at) "
"VALUES (?, ?, ?, ?, ?, ?, ?)",
[_id("rv"), session_id, proposal_id, req.action.value, req.actor, req.reason, _now()],
)
# If accepted, move active pointer
if req.action == ReviewAction.ACCEPT:
await d1_client.execute_sql(
"UPDATE sessions SET active_thought_id = ?, updated_at = ? WHERE session_id = ?",
[proposal_id, _now(), session_id],
)
# Witness
await d1_client.execute_sql(
"INSERT INTO witness_events (event_id, session_id, thought_id, event_type, actor, detail, created_at) "
"VALUES (?, ?, ?, 'reviewed', ?, ?, ?)",
[_id("ev"), session_id, proposal_id, req.actor,
f"{req.action.value}: {req.reason or 'No reason given'}", _now()],
)
return {
"proposal_id": proposal_id,
"new_status": new_status,
"action": req.action.value,
"message": f"{'⚑ Merged' if req.action == ReviewAction.ACCEPT else '❌ Rejected' if req.action == ReviewAction.REJECT else 'πŸ”„ Superseded'}: {proposal_id}",
}
# ══════════════════════════════════════════════════════════════
# TREE VISUALIZATION
# ══════════════════════════════════════════════════════════════
@app.get("/session/{session_id}/tree")
async def get_tree(session_id: str):
"""Render the full reasoning tree for a session."""
ses = await d1_client.execute_sql(
"SELECT * FROM sessions WHERE session_id = ?", [session_id]
)
if not ses:
raise HTTPException(404, "Session not found")
thoughts = await d1_client.execute_sql(
"SELECT * FROM thought_units WHERE session_id = ? ORDER BY created_at ASC",
[session_id],
)
edges = await d1_client.execute_sql(
"SELECT * FROM thought_edges WHERE session_id = ?", [session_id]
)
# Build adjacency map (parent -> children)
children_map: dict[str, list[str]] = {}
for edge in edges:
src = edge["source_id"]
tgt = edge["target_id"]
children_map.setdefault(src, []).append(tgt)
thought_map = {t["thought_id"]: t for t in thoughts}
def build_node(tid: str) -> dict:
t = thought_map.get(tid, {})
return {
"thought_id": tid,
"content": t.get("content", ""),
"thought_class": t.get("thought_class", ""),
"status": t.get("status", ""),
"origin_agent": t.get("origin_agent", ""),
"children": [build_node(cid) for cid in children_map.get(tid, [])],
}
root_id = ses[0].get("root_thought_id")
tree = build_node(root_id) if root_id else {}
return {
"session_id": session_id,
"title": ses[0].get("title"),
"total_thoughts": len(thoughts),
"total_edges": len(edges),
"active_thought_id": ses[0].get("active_thought_id"),
"tree": tree,
}
# ══════════════════════════════════════════════════════════════
# SEARCH & EDGES
# ══════════════════════════════════════════════════════════════
@app.post("/session/{session_id}/search")
async def search_thoughts(session_id: str, req: SearchReq):
"""Search thoughts in a session by content pattern."""
rows = await d1_client.execute_sql(
"SELECT * FROM thought_units WHERE session_id = ? AND content LIKE ? ORDER BY created_at ASC",
[session_id, f"%{req.pattern}%"],
)
return {"session_id": session_id, "pattern": req.pattern, "matches": len(rows), "results": rows}
@app.get("/session/{session_id}/edges")
async def list_edges(session_id: str):
"""List all edges (relationships) in a session."""
rows = await d1_client.execute_sql(
"SELECT * FROM thought_edges WHERE session_id = ? ORDER BY created_at ASC",
[session_id],
)
return {"session_id": session_id, "total": len(rows), "edges": rows}
# ══════════════════════════════════════════════════════════════
# WITNESS / PROVENANCE
# ══════════════════════════════════════════════════════════════
@app.get("/session/{session_id}/witness")
async def list_witness_events(session_id: str, limit: int = Query(50, ge=1, le=200)):
"""Get the audit trail for a session."""
rows = await d1_client.execute_sql(
"SELECT * FROM witness_events WHERE session_id = ? ORDER BY created_at DESC LIMIT ?",
[session_id, limit],
)
return {"session_id": session_id, "total": len(rows), "events": rows}
# ══════════════════════════════════════════════════════════════
# HEALTH
# ══════════════════════════════════════════════════════════════
@app.get("/health")
async def health():
return {"status": "ok", "timestamp": _now()}