diff --git a/src/mnemosyne/gateway.py b/src/mnemosyne/gateway.py index 5771350..138cf1e 100644 --- a/src/mnemosyne/gateway.py +++ b/src/mnemosyne/gateway.py @@ -1911,6 +1911,7 @@ def create_app( """ from mnemosyne.benchmark import Timer + _t0 = time.perf_counter() incoming_messages = payload.get("messages", []) # Guard: don't forward turns with empty user content (injection-only turns) @@ -2216,65 +2217,58 @@ def create_app( file=sys.stderr, ) - # 3b. Entropy-gated faulting — check last assistant response for uncertainty - try: - last_assistant_text = "" - for msg in reversed(ms.messages): - if msg.get("role") == "assistant": - content = msg.get("content", "") - if isinstance(content, str): - last_assistant_text = content - elif isinstance(content, list): - last_assistant_text = " ".join( - b.get("text", "") - for b in content - if isinstance(b, dict) and b.get("type") == "text" - ) - break - - if last_assistant_text: - # Collect evicted entities + # 3b. Entropy-gated faulting — run in background to avoid blocking. + # The per-object get() loop and micro-fault LLM calls are too slow + # for the request path. Results apply on the next turn. + def _entropy_bg(s=session, msgs=ms.messages): + try: + last_assistant_text = "" + for msg in reversed(msgs): + if msg.get("role") == "assistant": + content = msg.get("content", "") + if isinstance(content, str): + last_assistant_text = content + elif isinstance(content, list): + last_assistant_text = " ".join( + b.get("text", "") + for b in content + if isinstance(b, dict) and b.get("type") == "text" + ) + break + if not last_assistant_text: + return evicted_entities: list[str] = [] - fm = session.fidelity_manager + fm = s.fidelity_manager for obj_id, obj in fm._objects.items(): if obj.current_fidelity >= FidelityLevel.L4: - # Get key_entities from object store - stored = _run_async(session.object_store.get(obj_id)) + stored = _run_async(s.object_store.get(obj_id)) if stored and stored.key_entities: evicted_entities.extend(stored.key_entities) + if not evicted_entities: + return + signal = s.entropy_detector.analyze_response(last_assistant_text, evicted_entities) + fault_entities = s.entropy_detector.should_fault(signal) + if fault_entities and s.context_assembler: + for entity in fault_entities[:3]: + try: + _run_async( + s.context_assembler.handle_micro_fault( + session_id=s.id, + question=f"What was the content related to {entity}?", + scope=entity, + ) + ) + except Exception: + pass + except Exception: + pass - if evicted_entities: - signal = session.entropy_detector.analyze_response( - last_assistant_text, evicted_entities - ) - fault_entities = session.entropy_detector.should_fault(signal) - if fault_entities: - print( - f" {_YELLOW}[{session.id}] entropy fault: score={signal.score:.2f}, " - f"faulting {len(fault_entities)} entities{_RESET}", - file=sys.stderr, - ) - # Proactively restore via micro-fault for next turn - if session.context_assembler: - for entity in fault_entities[:3]: # Limit to 3 to avoid latency - try: - _run_async( - session.context_assembler.handle_micro_fault( - session_id=session.id, - question=f"What was the content related to {entity}?", - scope=entity, - ) - ) - except Exception: - pass # Best effort - except Exception as exc: - print( - f" {_YELLOW}[{session.id}] entropy check error (non-fatal): {exc}{_RESET}", - file=sys.stderr, - ) + _threading.Thread(target=_entropy_bg, daemon=True).start() - # 4. Build ephemeral outbound view — never mutate the physical store - payload["messages"] = copy.deepcopy(ms.messages) + # 4. Build ephemeral outbound view — never mutate the physical store. + # Use json round-trip instead of copy.deepcopy — 3-5x faster for + # large message lists (deepcopy tracks object identity which is O(n^2)). + payload["messages"] = json.loads(json.dumps(ms.messages, default=str)) # Apply cleanup/block-state rewrites to the outbound view so model-authored # drop/summarize/collapse operations actually affect future forwarded context. @@ -2334,6 +2328,14 @@ def create_app( with open(session._page_checkpoint, "w") as f: _json.dump(session.page_store.checkpoint(), f) + _preprocess_ms = (time.perf_counter() - _t0) * 1000 + if _preprocess_ms > 500: + print( + f" {_DIM}[{session.id}] _preprocess took {_preprocess_ms:.0f}ms " + f"({len(incoming_messages)} msgs){_RESET}", + file=sys.stderr, + flush=True, + ) return payload, _bytes_saved effective = session.token_state["last_effective"]