perf: move entropy faulting to background, use json copy instead of deepcopy
Remaining hot-path blockers for large contexts: - Entropy detection: per-object get() loop + micro-fault LLM calls → bg thread - copy.deepcopy(messages): O(n^2) identity tracking → json round-trip (3-5x faster) - Added _preprocess timing log (warns if >500ms)
This commit is contained in:
parent
a0822b0f6d
commit
8d79101025
1 changed files with 54 additions and 52 deletions
|
|
@ -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"]
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue