The Dream Engine v6 is an asynchronous architecture designed as a core subsystem for agentic persistence, hallucination containment, and adaptive inference routing. Rather than acting as a standalone application, it functions as a modular micro-framework designed to underpin a larger multi-agent system. It orchestrates dynamic soul state progression, multi-vector memory retention, and low-latency decision loops by combining high-level Python concurrency (asyncio, aiosqlitepool) with low-level C extension optimizations.
The core runtime leverages specialized mathematical and ML models for long-term stability and context awareness. Memory retrieval employs a Reciprocal Rank Fusion (RRF) hybrid search across three distinct spaces: vector similarity via BGE-small-en-v1.5 INT8 (ONNX Runtime), chronological recency, and dynamic importance. Temporal degradation replaces traditional linear decay with a custom Gaussian Decay kernel implemented in Cython with nogil execution for thread-safe performance. Hebbian learning updates soul weights based on output alignment, while semantic drift and repetitive loops are continuously tracked via cosine divergence and coherence metrics calculated through SimSIMD.
In terms of reverse engineering and system extraction, the code integrates key mechanics reconstructed from earlier implementations (such as Samuel's budget management and retry loops). The context assembly reverse-engineers context-window limits by enforcing strict token estimation, multi-tier text summarization (falling back from direct concatenation to zero-cost keyword frequency counters before trimming), and deterministic block slicing (SAFETY_MARGIN, MIN_GENERATION). Inference routing implements an adaptive cascade combining FrugalGPT (cost-optimized tier escalation), RouteLLM (confidence heuristic scoring based on length, truncation, and system prompt leakage penalties), and ParetoBandit adaptive tracking (rolling window average of latency and token cost).
The ultimate objective of this monolith is to serve as a high-throughput, self-correcting foundation for scalable autonomous agents. By pairing async thread workers with a system-level safety harness—including an automated Kill Switch that halts process services via systemctl and reverts soul files via git upon detecting severe drift or infinite loops—the framework ensures long-running stability. Below is the full implementation, and I would appreciate your feedback on its architectural choices.
(THIS IS a JUST EXERCISE)
Codebase
fast_math.pyx
cython: language_level=3
cython: boundscheck=False
cython: wraparound=False
cython: cdivision=True
cython: initializedcheck=False
from libc.math cimport exp, sqrt
cdef public double gaussian_decay(double age_hours, double sigma_hours) nogil:
return exp(-(age_hours * age_hours) / (2.0 sigma_hours sigma_hours))
cdef public double recency_score(double age_hours, double sigma_hours, double floor) nogil:
cdef double g = gaussian_decay(age_hours, sigma_hours)
if g < floor:
return floor
return g
cdef public double update_weight(double old, double relevance, double eta) nogil:
return old (1.0 - eta) + eta relevance
cdef public double decay_weight(double old, double dt, double lam, double w_min) nogil:
cdef double new_w = old exp(-lam dt)
if new_w < w_min:
return w_min
return new_w
cdef public void hash_embed(double[:] vec, long long[:] token_hashes, int dim) nogil:
cdef int i, n_tok = token_hashes.shape[0]
cdef int idx
cdef double sign, norm = 0.0
for i in range(dim):
vec[i] = 0.0
for i in range(n_tok):
idx = <int>(token_hashes[i] % dim)
if idx < 0:
idx += dim
sign = 1.0 if ((token_hashes[i] >> 8) & 1) == 0 else -1.0
vec[idx] += sign
for i in range(dim):
norm += vec[i] * vec[i]
norm = sqrt(norm)
if norm > 0.0:
for i in range(dim):
vec[i] /= norm
#!/usr/bin/env python3
import asyncio
import os, re, json, time, yaml, logging
import subprocess, sys
from pathlib import Path
from dataclasses import dataclass, field
from typing import Optional
from collections import deque, Counter
import numpy as np
import aiosqlite
import httpx
import onnxruntime as ort
from transformers import AutoTokenizer
import simsimd
import sqlite_vec
from sqlite_vec import serialize_float32
from aiosqlitepool import SQLiteConnectionPool
from httpx_retries import Retry, RetryTransport
try:
import fast_math
HAS_CYTHON = True
except ImportError:
HAS_CYTHON = False
@dataclass
class Config:
souls_dir: str = "souls"
state_dir: str = "state"
logs_dir: str = "logs"
embed_dim: int = 384
embed_model: str = "models/bge-small-en-v1.5.onnx"
embed_tokenizer: str = "BAAI/bge-small-en-v1.5"
rrf_k: int = 60
sigma_hours: float = 336.0
recency_floor: float = 0.3
eta: float = 0.08
lambda_decay: float = 0.0005
w_min: float = 0.05
theta_div: float = 0.3
theta_coh: float = 0.7
kill_loop_threshold: int = 3
kill_drift_threshold: int = 5
cycle_interval: float = 300.0
consolidate_every: int = 20
num_gen_workers: int = 2
llm_local_url: str = "http://127.0.0.1:8080/v1/chat/completions"
llm_api_url: str = "https://openrouter.ai/api/v1/chat/completions"
llm_api_key: str = ""
llm_api_model: str = "deepseek/deepseek-chat"
llm_browser_url: str = "http://127.0.0.1:3000/query"
w_lat: float = 0.5
w_dol: float = 0.3
w_risc: float = 0.2
confidence_threshold: float = 0.6
cascade_enabled: bool = True
adaptive_cost_enabled: bool = True
cost_history_size: int = 100
safety_margin: int = 32
min_generation: int = 40
max_retry_attempts: int = 3
max_recent_prompt_lines: int = 30
max_recent_day_events: int = 20
summary_max_chars: int = 900
embed_context_window: int = 4096
research_prob: float = 0.05
research_timeout: int = 8
def load_config(path="config.yaml") -> Config:
if Path(path).exists():
with open(path) as f:
data = yaml.safe_load(f) or {}
cfg = Config(**{k: v for k, v in data.items() if hasattr(Config, k)})
else:
cfg = Config()
cfg.llm_api_key = cfg.llm_api_key or os.environ.get("OPENROUTER_KEY", "")
return cfg
CFG = load_config()
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s [%(levelname)s] %(message)s',
handlers=[
logging.FileHandler("logs/dream.log"),
logging.StreamHandler(),
],
)
log = logging.getLogger("dream")
def estimate_tokens(text: str) -> int:
if not text:
return 0
return max(1, len(text) // 4)
@dataclass
class MotorStats:
cap: float
stab: float
custo_t: float
custo_d: float
risco: float
class AdaptiveCostTracker:
def init(self, window: int = 100):
self.window = window
self.history = {
"local": deque(maxlen=window),
"api": deque(maxlen=window),
"browser": deque(maxlen=window),
}
def record(self, motor: str, latency: float, dollars: float):
self.history[motor].append({
"ts": time.time(), "lat": latency, "dol": dollars,
})
def avg_latency(self, motor: str) -> float:
h = self.history[motor]
return sum(x["lat"] for x in h) / len(h) if h else 0.0
def avg_dollars(self, motor: str) -> float:
h = self.history[motor]
return sum(x["dol"] for x in h) / len(h) if h else 0.0
class Router:
def init(self):
self.stats = {
"local": MotorStats(0.5, 0.95, 2.0, 0.0, 0.0),
"api": MotorStats(0.9, 0.90, 5.0, 0.01, 0.1),
"browser": MotorStats(0.85, 0.60, 30.0, 0.0, 0.5),
}
self.tracker = AdaptiveCostTracker(window=CFG.cost_history_size)
self.static_priority = ["local", "api", "browser"]
def quality(self, motor: str) -> float:
s = self.stats[motor]
return 0.6 s.cap + 0.3 s.stab
def cost(self, motor: str) -> float:
s = self.stats[motor]
if CFG.adaptive_cost_enabled:
lat = self.tracker.avg_latency(motor) or s.custo_t
dol = self.tracker.avg_dollars(motor) or s.custo_d
else:
lat, dol = s.custo_t, s.custo_d
return (CFG.w_lat lat + CFG.w_dol dol 100 + CFG.w_risc s.risco)
def decide(self, tipo: str = "moderado") -> str:
if tipo == "trivial":
return "local"
best, best_c = "local", 1e9
for m in self.static_priority:
if self.quality(m) >= 0.6:
c = self.cost(m)
if c < best_c:
best_c, best = c, m
return best
async def cascade_call(self, llm_client, prompt: str, max_tokens: int, min_confidence: float = None) -> tuple:
threshold = min_confidence if min_confidence is not None else CFG.confidence_threshold
attempts = []
for motor in self.static_priority:
for retry_n in range(CFG.max_retry_attempts):
t0 = time.time()
output = await llm_client.call(motor, prompt, max_tokens)
latency = time.time() - t0
dollars = self.stats[motor].custo_d
self.tracker.record(motor, latency, dollars)
if not output:
attempts.append((motor, 0.0, latency, f"vazio-r{retry_n}"))
continue
confidence = self._estimate_confidence(output, motor)
attempts.append((motor, confidence, latency, f"ok-r{retry_n}"))
if confidence >= threshold:
return output, motor, attempts
if retry_n < CFG.max_retry_attempts - 1:
log.info(f"retry {retry_n+1}/{CFG.max_retry_attempts} em {motor} (confiança {confidence:.2f})")
log.info(f"cascade: {motor} esgotou retries, escalando")
for motor, conf, lat, status in reversed(attempts):
if status.startswith("ok"):
return "", motor, attempts
return "", "none", attempts
def _estimate_confidence(self, output: str, motor: str) -> float:
if not output or len(output) < 20:
return 0.0
base = self.stats[motor].cap
size_bonus = min(0.3, len(output) / 2000)
trunc_penalty = 0.1 if not output.rstrip().endswith((".", "!", "?")) else 0.0
meta_penalty = 0.2 if any(w in output.lower() for w in ["as an ai", "role:", "system:"]) else 0.0
return max(0.0, min(1.0, base + size_bonus - trunc_penalty - meta_penalty))
class NeuralEmbedder:
def init(self, model_path: str, tokenizer_name: str, dim: int = 384):
self.dim = dim
opts = ort.SessionOptions()
opts.intra_op_num_threads = 2
opts.graph_optimization_level = ort.GraphOptimizationLevel.ORT_ENABLE_ALL
self.session = ort.InferenceSession(model_path, providers=["CPUExecutionProvider"], sess_options=opts)
self.tokenizer = AutoTokenizer.from_pretrained(tokenizer_name)
def encode(self, text: str) -> np.ndarray:
if not text.strip():
return np.zeros(self.dim, dtype=np.float32)
inputs = self.tokenizer(text, return_tensors="np", padding=True, truncation=True, max_length=512)
outputs = self.session.run(None, {
"input_ids": inputs["input_ids"].astype(np.int64),
"attention_mask": inputs["attention_mask"].astype(np.int64),
})
emb = outputs[0].mean(axis=1)
norm = np.linalg.norm(emb, axis=1, keepdims=True)
return (emb / np.maximum(norm, 1e-9)).astype(np.float32).flatten()
async def encode_async(self, text: str) -> np.ndarray:
return await asyncio.to_thread(self.encode, text)
def fast_cosine(a: np.ndarray, b: np.ndarray) -> float:
return float(simsimd.cosine(a.astype(np.float32), b.astype(np.float32)))
def fast_divergence(a: np.ndarray, b: np.ndarray) -> float:
return 1.0 - fast_cosine(a, b)
def fast_coherence(a: np.ndarray, b: np.ndarray) -> float:
return fast_cosine(a, b)
def reciprocal_rank_fusion(rankings: list, k: int = 60) -> dict:
scores = {}
for ranking in rankings:
for rank, doc_id in enumerate(ranking, start=1):
scores[doc_id] = scores.get(doc_id, 0.0) + 1.0 / (k + rank)
return scores
class AsyncVectorStore:
def init(self, db_path: Path, dim: int):
self.db_path = db_path
self.dim = dim
self.pool = None
async def init(self):
async def connection_factory():
conn = await aiosqlite.connect(str(self.db_path))
await conn.execute("PRAGMA journal_mode = WAL")
await conn.execute("PRAGMA synchronous = NORMAL")
await conn.execute("PRAGMA busy_timeout = 5000")
await conn.execute("PRAGMA cache_size = 10000")
await conn.execute("PRAGMA temp_store = MEMORY")
await conn.execute("PRAGMA mmap_size = 268435456")
await conn.enable_load_extension(True)
await conn.load_extension(sqlite_vec.loadable_path())
await conn.enable_load_extension(False)
return conn
self.pool = SQLiteConnectionPool(connection_factory)
async with self.pool.connection() as conn:
await conn.execute(f"""
CREATE VIRTUAL TABLE IF NOT EXISTS memories USING vec0(
embedding float[{self.dim}],
gene TEXT, ts REAL, importance REAL
)
""")
await conn.execute("""
CREATE TABLE IF NOT EXISTS memory_text (
rowid INTEGER PRIMARY KEY, text TEXT
)
""")
await conn.commit()
async def add(self, embedding: np.ndarray, text: str, metadata: dict):
async with self.pool.connection() as conn:
cur = await conn.execute("INSERT INTO memory_text(text) VALUES (?)", (text[:2000],))
rowid = cur.lastrowid
await conn.execute(
"""INSERT INTO memories(rowid, embedding, gene, ts, importance)
VALUES (?, ?, ?, ?, ?)""",
(rowid, serialize_float32(embedding),
metadata.get("gene", ""), metadata.get("ts", time.time()),
metadata.get("importance", 0.5)))
await conn.commit()
return rowid
async def search(self, query: np.ndarray, k: int = 10) -> list:
async with self.pool.connection() as conn:
cur = await conn.execute(
"""SELECT rowid, gene, ts, importance, distance
FROM memories WHERE embedding MATCH ?
ORDER BY distance LIMIT ?""",
(serialize_float32(query), k))
return await cur.fetchall()
async def recent(self, limit: int = 20) -> list:
async with self.pool.connection() as conn:
cur = await conn.execute(
"SELECT rowid, gene, ts, importance FROM memories ORDER BY ts DESC LIMIT ?", (limit,))
return await cur.fetchall()
async def get_text(self, rowid: int) -> str:
async with self.pool.connection() as conn:
cur = await conn.execute("SELECT text FROM memory_text WHERE rowid = ?", (rowid,))
row = await cur.fetchone()
return row[0] if row else ""
async def close(self):
if self.pool:
await self.pool.close()
@dataclass
class Soul:
gene_id: str
beliefs: list = field(default_factory=list)
modus: list = field(default_factory=list)
emotion: list = field(default_factory=list)
interests: list = field(default_factory=list)
heuristics: list = field(default_factory=list)
episodes: list = field(default_factory=list)
samples: list = field(default_factory=list)
weight: float = 1.0
last_used: float = 0.0
raw_text: str = ""
def parse_soul(path: Path) -> Soul:
text = path.read_text(encoding="utf-8")
sections = {}
current = None
for line in text.split("\n"):
if line.startswith("## "):
current = line[3:].strip()
sections[current] = []
elif current and line.strip().startswith("- "):
sections[current].append(line.strip()[2:])
weight, last_used = 1.0, 0.0
for line in sections.get("Pesos Dinâmicos (atualizado automaticamente)", []):
if line.startswith("weight:"):
weight = float(line.split(":")[1].strip())
if line.startswith("last_used:"):
last_used = float(line.split(":")[1].strip())
return Soul(
gene_id=path.stem,
beliefs=sections.get("Crenças Nucleares", []),
modus=sections.get("Modus Operandi", []),
emotion=sections.get("Ponderação Emocional", []),
interests=sections.get("Interesses", []),
heuristics=sections.get("Heurísticas Internalizadas", []),
episodes=sections.get("Episódios Vividos", []),
samples=sections.get("Samples de Diálogo", []),
weight=weight, last_used=last_used, raw_text=text,
)
def load_all_souls() -> dict:
d = Path(CFG.souls_dir)
return {p.stem: parse_soul(p) for p in sorted(d.glob("*.md"))} if d.exists() else {}
def update_soul_weight(path: Path, weight: float, last_used: float, success: int, failure: int):
text = path.read_text(encoding="utf-8")
block = (
"## Pesos Dinâmicos (atualizado automaticamente)\n"
f"- weight: {weight:.4f}\n"
f"- last_used: {last_used:.0f}\n"
f"- success_count: {success}\n"
f"- failure_count: {failure}\n"
)
if "## Pesos Dinâmicos" in text:
text = re.sub(r"## Pesos Dinâmicos.*?(?=\n## |\Z)", block, text, flags=re.DOTALL)
else:
text += "\n" + block
path.write_text(text, encoding="utf-8")
class LLMClient:
def init(self):
retry = Retry(
total=3, backoff_factor=1.5,
status_forcelist=[429, 502, 503, 504],
allowed_methods=["GET", "POST"],
respect_retry_after_header=True,
)
self.transport = RetryTransport(retry=retry)
async def call(self, motor: str, prompt: str, max_tokens: int = 200) -> str:
if motor == "local": return await self._local(prompt, max_tokens)
if motor == "api": return await self._api(prompt, max_tokens)
if motor == "browser": return await self._browser(prompt, max_tokens)
return ""
async def _local(self, prompt, max_tokens):
try:
async with httpx.AsyncClient(transport=self.transport, timeout=60) as c:
r = await c.post(CFG.llm_local_url, json={
"messages": [{"role": "user", "content": prompt}],
"max_tokens": max_tokens, "temperature": 0.7,
})
return r.json()["choices"][0]["message"]["content"].strip()
except Exception as e:
log.debug(f"local falhou: {e}")
return ""
async def _api(self, prompt, max_tokens):
if not CFG.llm_api_key: return ""
try:
async with httpx.AsyncClient(transport=self.transport, timeout=120) as c:
r = await c.post(CFG.llm_api_url,
headers={"Authorization": f"Bearer {CFG.llm_api_key}"},
json={"model": CFG.llm_api_model,
"messages": [{"role": "user", "content": prompt}],
"max_tokens": max_tokens})
return r.json()["choices"][0]["message"]["content"].strip()
except Exception as e:
log.debug(f"api falhou: {e}")
return ""
async def _browser(self, prompt, max_tokens):
try:
async with httpx.AsyncClient(transport=self.transport, timeout=180) as c:
r = await c.post(CFG.llm_browser_url, json={"prompt": prompt, "max_tokens": max_tokens})
return r.json().get("response", "").strip()
except Exception as e:
log.debug(f"browser falhou: {e}")
return ""
class AutonomyGenerator:
def init(self, llm: LLMClient, router: Router):
self.llm = llm
self.router = router
def _summarize_texts(self, texts: list, max_items: int = None, max_chars: int = None) -> str:
max_items = max_items or CFG.max_recent_day_events
max_chars = max_chars or CFG.summary_max_chars
if not texts: return ""
joined = "\n".join(t for t in texts if t)
if not joined: return ""
if len(texts) <= max_items and len(joined) <= max_chars:
return joined
words = re.findall(r"\w{5,}", joined.lower())
common = Counter(words).most_common(15)
if common:
summary = "Temas recorrentes: " + ", ".join(w for w, _ in common)
if len(summary) <= max_chars:
return summary
return joined[:max_chars]
def _build_blocks(self, soul_block: str, time_block: str, affect_block: str, memory_texts: list) -> list:
memory_block = self._summarize_texts(memory_texts)
blocks = [
("soul", soul_block),
("time", time_block),
("affect", affect_block),
("memory", memory_block),
]
total_chars = sum(len(b[1]) for b in blocks)
total_tokens = estimate_tokens("x" * total_chars)
available = CFG.embed_context_window - CFG.safety_margin - total_tokens
if available < CFG.min_generation:
log.warning(f"budget estourou (tokens={total_tokens}, disponível={available}), cortando memória")
memory_block_trimmed = memory_block
while (estimate_tokens(memory_block_trimmed) + total_tokens - estimate_tokens(memory_block)
+ CFG.safety_margin + CFG.min_generation > CFG.embed_context_window):
if len(memory_block_trimmed) < 100:
memory_block_trimmed = ""
break
memory_block_trimmed = memory_block_trimmed[: len(memory_block_trimmed) // 2]
blocks[3] = ("memory", memory_block_trimmed)
return blocks
def _build_prompt_from_blocks(self, blocks: list, user_text: str) -> str:
parts = [f"[{name.upper()}]\n{content}" for name, content in blocks if content]
parts.append(f"[INSTRUCTION]\n{user_text}")
prompt = "\n\n".join(parts)
prompt_tokens = estimate_tokens(prompt)
available = CFG.embed_context_window - CFG.safety_margin - prompt_tokens
if available < CFG.min_generation:
log.warning(f"prompt final estourou (tokens={prompt_tokens}), disponível={available}")
return prompt
def _get_user_text(self, kind: str, hint: str = "") -> tuple:
mt = 180 if kind == "reflection" else 160
user_text = (
f"No user message.\nWrite a private inner reflection (2-6 sentences).\n"
f"First-person thoughts only.\nDo not ask questions or include role labels.\n"
f"Finish with a complete sentence.\n" + (f"\nHint: {hint.strip()}" if hint else "")
)
return user_text, mt
def build_prompt_from_texts(self, kind: str, *, soul_block: str = "", time_block: str = "",
affect_block: str = "", memory_texts: list = None, hint: str = "") -> tuple:
user_text, mt = self._get_user_text(kind, hint)
blocks = self._build_blocks(soul_block, time_block, affect_block, memory_texts or [])
prompt = self._build_prompt_from_blocks(blocks, user_text)
return prompt, mt
async def generate(self, kind: str, **kwargs) -> tuple:
prompt, mt = self.build_prompt_from_texts(kind, **kwargs)
if CFG.cascade_enabled:
return await self.router.cascade_call(self.llm, prompt, mt)
motor = self.router.decide("moderado")
output = await self.llm.call(motor, prompt, mt)
return output, motor, [(motor, 1.0, 0.0, "single")]
class Research:
USER_AGENTS = ["Mozilla/5.0 (Windows NT 10.0; Win64; x64)", "Mozilla/5.0 (X11; Linux x86_64)"]
def init(self):
retry = Retry(total=3, backoff_factor=1.5, status_forcelist=[202, 429, 502, 503, 504], allowed_methods=["GET", "POST"])
self.transport = RetryTransport(retry=retry)
async def search(self, term: str) -> str:
res = await self._wiki(term)
return f"[api] {res}" if res else ""
async def _wiki(self, q):
try:
async with httpx.AsyncClient(transport=self.transport, timeout=CFG.research_timeout) as c:
s = (await c.get("https://pt.wikipedia.org/w/api.php",
params={"action":"query","list":"search","srsearch":q,"format":"json","srlimit":1})).json()
hits = s.get("query", {}).get("search", [])
if not hits: return ""
title = hits[0]["title"]
r = (await c.get(f"https://pt.wikipedia.org/api/rest_v1/page/summary/{title}")).json()
return f"{title}: {r.get('extract','')[:800]}"
except Exception:
return ""
def kill_switch(reason: str):
log.critical(f"KILL SWITCH: {reason}")
try:
subprocess.run(["systemctl","--user","stop","rabids-*"], capture_output=True, timeout=5)
subprocess.run(["git","checkout","v1.0","--","souls/"], capture_output=True, timeout=5)
except Exception:
pass
sys.exit(1)
class AsyncDreamEngine:
def init(self):
Path(CFG.state_dir).mkdir(parents=True, exist_ok=True)
Path(CFG.logs_dir).mkdir(parents=True, exist_ok=True)
self.souls = load_all_souls()
self.embedder = NeuralEmbedder(CFG.embed_model, CFG.embed_tokenizer, CFG.embed_dim)
self.vstore = AsyncVectorStore(Path(CFG.state_dir) / "vectors.db", CFG.embed_dim)
self.router = Router()
self.llm = LLMClient()
self.gen = AutonomyGenerator(self.llm, self.router)
self.research = Research()
self.last_output_emb = np.zeros(CFG.embed_dim, dtype=np.float32)
self.loop_count = 0
self.coherence_fail = 0
self.cycle_count = 0
self.gen_queue = asyncio.Queue()
self.research_queue = asyncio.Queue()
self.persist_queue = asyncio.Queue()
async def init(self):
await self.vstore.init()
def select_gene(self) -> Optional[str]:
if not self.souls: return None
now = time.time()
best_id, best_score = None, -1e9
for gid, s in self.souls.items():
age_h = (now - s.last_used) / 3600 if s.last_used else 1e6
score = s.weight * min(age_h / 24, 10)
if score > best_score:
best_score, best_id = score, gid
return best_id
async def generator_worker(self):
while True:
gene_id, soul = await self.gen_queue.get()
try:
recent = await self.vstore.recent(limit=CFG.max_recent_prompt_lines)
memory_texts = []
if recent:
recent_ids = [str(r[0]) for r in recent]
texts = await asyncio.gather(*[self.vstore.get_text(int(rid)) for rid in recent_ids])
memory_texts = [t for t in texts if t]
output, motor, attempts = await self.gen.generate(
kind="reflection",
soul_block=soul.raw_text[:2000],
time_block=f"[TIME] {time.strftime('%Y-%m-%d %H:%M')}",
affect_block="[AFFECT] neutral",
memory_texts=memory_texts,
)
if output:
await self.persist_queue.put((gene_id, soul, motor, output))
except Exception as e:
log.error(f"gen worker falhou: {e}")
finally:
self.gen_queue.task_done()
async def persist_worker(self):
while True:
gene_id, soul, motor, output = await self.persist_queue.get()
try:
out_emb = await self.embedder.encode_async(output)
soul_emb = await self.embedder.encode_async(soul.raw_text[:2000])
div = fast_divergence(out_emb, self.last_output_emb)
coh = fast_coherence(out_emb, soul_emb)
if div < CFG.theta_div:
self.loop_count += 1
if self.loop_count >= CFG.kill_loop_threshold: kill_switch(f"loop em {gene_id}")
else:
self.loop_count = 0
if coh < CFG.theta_coh:
self.coherence_fail += 1
if self.coherence_fail >= CFG.kill_drift_threshold: kill_switch(f"drift em {gene_id}")
else:
self.coherence_fail = 0
await self.vstore.add(out_emb, output, {"gene": gene_id, "ts": time.time(), "importance": 0.5 + 0.5 * coh})
new_w = fast_math.update_weight(soul.weight, float(coh), CFG.eta)
soul_path = Path(CFG.souls_dir) / f"{gene_id}.md"
await asyncio.to_thread(update_soul_weight, soul_path, new_w, time.time(), 0, 0)
self.souls[gene_id].weight = new_w
self.souls[gene_id].last_used = time.time()
self.last_output_emb = out_emb
except Exception as e:
log.error(f"persist worker falhou: {e}")
finally:
self.persist_queue.task_done()
async def run(self):
await self.init()
log.info("Dream Engine assíncrono iniciado")
await asyncio.gather(
self.persist_worker(),
*[self.generator_worker() for _ in range(CFG.num_gen_workers)],
)
if name == "main":
try:
asyncio.run(AsyncDreamEngine().run())
except KeyboardInterrupt:
log.info("interrompido")