aae8d0a379
In-progress work committed to unblock the config import/export plan: - prompts: FILE_REFERENCES section for workspace-relative file citation - backends: resolve quoted virtual absolute paths onto the sandbox workspace - middleware: read_file_images middleware; message_budget extensions - image_gen/model_registry: image model 'enabled' flag refactor - memory/launch, gateway/background_runs, tools/image follow-ons - scripts: dev_backend.sh, release.sh - tests for the above
466 lines
16 KiB
Python
466 lines
16 KiB
Python
"""EvoMemory LangGraph launch adapter."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import logging
|
|
import uuid
|
|
from collections.abc import Callable
|
|
from pathlib import Path
|
|
from typing import cast
|
|
|
|
from ..config import MemoryControls, get_effective_config
|
|
from ..gateway.background_runs import (
|
|
BackgroundRun,
|
|
BackgroundRunHooks,
|
|
BackgroundRunPayload,
|
|
BackgroundRunRequest,
|
|
alaunch_background_run,
|
|
launch_background_run,
|
|
)
|
|
from ..langgraph_dev.sdk import messages_input
|
|
from ..model_registry.schemas import ModelRef, ReasoningEffort
|
|
from .observations import build_observation_linker_index_context
|
|
from .scheduler import ObservationLinkerContext
|
|
from .source_context import MemorySourceContext, _trajectory_for_prompt
|
|
from .types import MemorySourceType
|
|
from .worker_activity import (
|
|
MemoryOutputDelta,
|
|
MemoryOutputSnapshot,
|
|
ObservationRelationSnapshot,
|
|
forget_memory_worker,
|
|
forget_observation_linker,
|
|
mark_memory_worker_finished,
|
|
mark_memory_worker_started,
|
|
mark_observation_linker_finished,
|
|
mark_observation_linker_started,
|
|
snapshot_memory_outputs,
|
|
snapshot_observation_relations,
|
|
)
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
SUBAGENT_MEMORY_WORKER_GRAPH_ID = "evomemory-subagent-worker"
|
|
TURN_MEMORY_WORKER_GRAPH_ID = "evomemory-turn-worker"
|
|
OBSERVATION_LINKER_GRAPH_ID = "evomemory-observation-linker"
|
|
|
|
MemoryWorkerFinishedHook = Callable[[BackgroundRun, MemoryOutputDelta | None], None]
|
|
MemoryWorkerAbortedHook = Callable[[BackgroundRun, MemoryOutputDelta | None], None]
|
|
|
|
|
|
def _observation_linking_enabled() -> bool:
|
|
return MemoryControls.from_config(get_effective_config()).observations_enabled
|
|
|
|
|
|
def _memory_worker_graph_id(source_type: MemorySourceType) -> str:
|
|
match source_type:
|
|
case MemorySourceType.TURN:
|
|
return TURN_MEMORY_WORKER_GRAPH_ID
|
|
case MemorySourceType.SUBAGENT:
|
|
return SUBAGENT_MEMORY_WORKER_GRAPH_ID
|
|
case _:
|
|
raise ValueError(f"Unsupported memory source type: {source_type!r}")
|
|
|
|
|
|
def _memory_worker_user_prompt(context: MemorySourceContext) -> str:
|
|
match context.source_type:
|
|
case MemorySourceType.TURN:
|
|
return (
|
|
"Review this completed orchestrator turn.\n\n"
|
|
f"Source agent: {context.source_agent}\n"
|
|
f"Source session: {context.session_id}\n\n"
|
|
f"Turn trajectory:\n{_trajectory_for_prompt(context.trajectory)}"
|
|
)
|
|
case MemorySourceType.SUBAGENT:
|
|
return (
|
|
"Review this completed subagent run.\n\n"
|
|
f"Source agent: {context.source_agent}\n"
|
|
f"Source session: {context.session_id}\n\n"
|
|
f"Trajectory:\n{_trajectory_for_prompt(context.trajectory)}"
|
|
)
|
|
case _:
|
|
raise ValueError(f"Unsupported memory source type: {context.source_type!r}")
|
|
|
|
|
|
def _runs_create_kwargs(payload: BackgroundRunPayload) -> BackgroundRunPayload:
|
|
try:
|
|
from EvoScientist.llm.patches import _merge_runs_config_kwargs
|
|
except Exception:
|
|
return payload
|
|
return cast("BackgroundRunPayload", _merge_runs_config_kwargs(dict(payload)))
|
|
|
|
|
|
def _worker_workspace_dir(workspace_dir: str | Path) -> str:
|
|
return str(Path(workspace_dir).expanduser().resolve())
|
|
|
|
|
|
def _memory_worker_metadata(context: MemorySourceContext) -> dict[str, str]:
|
|
metadata = {
|
|
"run_kind": f"evomemory_{context.source_type.value}_worker",
|
|
"source_session_id": context.session_id,
|
|
"source_agent": context.source_agent,
|
|
"project_id": context.project_id,
|
|
"trajectory_digest": context.trajectory_digest,
|
|
"workspace_dir": _worker_workspace_dir(context.workspace_dir),
|
|
}
|
|
from ..usage.callback import usage_tracking_requested
|
|
|
|
if usage_tracking_requested():
|
|
metadata["usage_scope"] = "memory"
|
|
if usage_tracking_requested() and context.turn_id:
|
|
metadata["turn_id"] = context.turn_id
|
|
return metadata
|
|
|
|
|
|
def _source_thread_model_selection(
|
|
session_id: str,
|
|
) -> tuple[ModelRef, ReasoningEffort | None, int] | None:
|
|
"""Read the source conversation thread's explicit model selection.
|
|
|
|
Returns ``(primary, reasoning_effort, revision)`` for an explicit
|
|
ThreadModelSelection, or ``None`` when the thread is unreadable, has no
|
|
selection, or is set to ``inherit`` — in those cases the worker falls
|
|
back to the section 8.1 lazy local snapshot (registry default). A legacy
|
|
``auxiliary`` key in stored metadata is tolerated and dropped, matching
|
|
the BFF validator.
|
|
"""
|
|
from langgraph_sdk import get_sync_client
|
|
|
|
from ..langgraph_dev.sdk import (
|
|
configured_langgraph_dev_url,
|
|
langgraph_dev_headers,
|
|
)
|
|
|
|
client = get_sync_client(
|
|
url=configured_langgraph_dev_url(),
|
|
headers=langgraph_dev_headers(None),
|
|
)
|
|
thread = client.threads.get(session_id)
|
|
metadata = (
|
|
thread.get("metadata") if isinstance(thread, dict) else getattr(thread, "metadata", None)
|
|
)
|
|
if not isinstance(metadata, dict):
|
|
return None
|
|
raw = metadata.get("model_selection")
|
|
if not isinstance(raw, dict):
|
|
return None
|
|
primary_raw = raw.get("primary")
|
|
if not isinstance(primary_raw, dict):
|
|
return None
|
|
provider_id = primary_raw.get("provider_id")
|
|
model_key = primary_raw.get("model_key")
|
|
if (
|
|
not isinstance(provider_id, str)
|
|
or not provider_id
|
|
or not isinstance(model_key, str)
|
|
or not model_key
|
|
):
|
|
return None
|
|
effort_raw = raw.get("reasoning_effort")
|
|
effort: ReasoningEffort | None = (
|
|
effort_raw if effort_raw in ("low", "medium", "high") else None
|
|
)
|
|
revision_raw = metadata.get("model_selection_revision")
|
|
revision = revision_raw if isinstance(revision_raw, int) and revision_raw >= 0 else 0
|
|
return (ModelRef(provider_id=provider_id, model_key=model_key), effort, revision)
|
|
|
|
|
|
def _worker_snapshot_id(context: MemorySourceContext, worker_thread_id: str) -> str | None:
|
|
"""Freeze the source conversation's model selection for the worker run.
|
|
|
|
The worker runs on its own thread, so the conversation's snapshot cannot
|
|
be reused (bindings are per-thread); a fresh snapshot with the same
|
|
selection is created and bound to the worker thread instead. Returns
|
|
``None`` to keep the section 8.1 lazy-default behavior when the source
|
|
thread has no explicit selection or any step fails — memory work must
|
|
never fail just because selection inheritance did.
|
|
"""
|
|
try:
|
|
selection = _source_thread_model_selection(context.session_id)
|
|
except Exception:
|
|
logger.warning(
|
|
"Memory worker: could not read source thread %s model selection; "
|
|
"falling back to the registry default",
|
|
context.session_id,
|
|
exc_info=True,
|
|
)
|
|
return None
|
|
if selection is None:
|
|
return None
|
|
primary, reasoning_effort, revision = selection
|
|
|
|
from ..model_registry.runtime import get_snapshot_runtime
|
|
from ..model_registry.snapshots import SnapshotCreateRequest
|
|
|
|
try:
|
|
runtime = get_snapshot_runtime()
|
|
creation = runtime.snapshots.create( SnapshotCreateRequest(
|
|
run_request_id=uuid.uuid4().hex,
|
|
thread_id=worker_thread_id,
|
|
deployment_id=runtime.local_deployment_id,
|
|
model_selection_revision=revision,
|
|
primary=primary,
|
|
reasoning_effort=reasoning_effort,
|
|
)
|
|
)
|
|
except Exception:
|
|
logger.warning(
|
|
"Memory worker: could not freeze the source selection into a run "
|
|
"snapshot; falling back to the registry default",
|
|
exc_info=True,
|
|
)
|
|
return None
|
|
return creation.snapshot.snapshot_id
|
|
|
|
|
|
def _memory_worker_run_payload(
|
|
*,
|
|
context: MemorySourceContext,
|
|
thread_id: str,
|
|
) -> BackgroundRunPayload:
|
|
"""Build the LangGraph SDK run payload for a memory worker."""
|
|
metadata = _memory_worker_metadata(context)
|
|
configurable = {
|
|
"thread_id": thread_id,
|
|
"evomemory_source_session_id": context.session_id,
|
|
"evomemory_source_agent": context.source_agent,
|
|
"evomemory_project_id": context.project_id,
|
|
"evomemory_trajectory_digest": context.trajectory_digest,
|
|
}
|
|
snapshot_id = _worker_snapshot_id(context, thread_id)
|
|
if snapshot_id is not None:
|
|
configurable["runtime_snapshot_id"] = snapshot_id
|
|
from ..usage.callback import usage_tracking_requested
|
|
|
|
if usage_tracking_requested() and context.turn_id:
|
|
configurable["evomemory_source_turn_id"] = context.turn_id
|
|
payload: BackgroundRunPayload = {
|
|
"assistant_id": _memory_worker_graph_id(context.source_type),
|
|
"input": messages_input(_memory_worker_user_prompt(context)),
|
|
"metadata": metadata,
|
|
"config": {"configurable": configurable},
|
|
}
|
|
return _runs_create_kwargs(payload)
|
|
|
|
|
|
def memory_worker_launch_request(
|
|
context: MemorySourceContext,
|
|
) -> BackgroundRunRequest:
|
|
"""Build the background run request for a memory worker."""
|
|
metadata = _memory_worker_metadata(context)
|
|
|
|
def run_payload(thread_id: str) -> BackgroundRunPayload:
|
|
return _memory_worker_run_payload(context=context, thread_id=thread_id)
|
|
|
|
return BackgroundRunRequest(
|
|
graph_id=_memory_worker_graph_id(context.source_type),
|
|
run_payload=run_payload,
|
|
thread_metadata=metadata,
|
|
name="EvoMemory worker",
|
|
)
|
|
|
|
|
|
def _observation_linker_user_prompt(context: ObservationLinkerContext) -> str:
|
|
payload = {
|
|
"project_id": context.project_id,
|
|
"new_observation_ids": sorted(context.observation_ids),
|
|
}
|
|
prompt = (
|
|
"Link newly recorded observations when there is a strong reusable "
|
|
"relationship.\n\n"
|
|
f"{json.dumps(payload, ensure_ascii=False, indent=2, sort_keys=True)}"
|
|
)
|
|
observation_index = build_observation_linker_index_context(
|
|
memory_dir=context.memory_dir,
|
|
project_id=context.project_id,
|
|
exclude_ids=context.observation_ids,
|
|
)
|
|
if observation_index:
|
|
prompt += f"\n\n{observation_index}"
|
|
return prompt
|
|
|
|
|
|
def _observation_linker_metadata(
|
|
context: ObservationLinkerContext,
|
|
) -> dict[str, str]:
|
|
return {
|
|
"run_kind": "evomemory_observation_linker",
|
|
"project_id": context.project_id,
|
|
"observation_count": str(len(context.observation_ids)),
|
|
"workspace_dir": str(context.workspace_dir.expanduser().resolve()),
|
|
}
|
|
|
|
|
|
def _observation_linker_run_payload(
|
|
*,
|
|
context: ObservationLinkerContext,
|
|
thread_id: str,
|
|
) -> BackgroundRunPayload:
|
|
payload: BackgroundRunPayload = {
|
|
"assistant_id": OBSERVATION_LINKER_GRAPH_ID,
|
|
"input": messages_input(_observation_linker_user_prompt(context)),
|
|
"metadata": _observation_linker_metadata(context),
|
|
"config": {
|
|
"configurable": {
|
|
"thread_id": thread_id,
|
|
"evomemory_project_id": context.project_id,
|
|
"evomemory_observation_ids": json.dumps(
|
|
list(context.observation_ids),
|
|
ensure_ascii=False,
|
|
),
|
|
}
|
|
},
|
|
}
|
|
return _runs_create_kwargs(payload)
|
|
|
|
|
|
def observation_linker_launch_request(
|
|
context: ObservationLinkerContext,
|
|
) -> BackgroundRunRequest:
|
|
"""Build the background run request for the observation linker."""
|
|
|
|
def run_payload(thread_id: str) -> BackgroundRunPayload:
|
|
return _observation_linker_run_payload(
|
|
context=context,
|
|
thread_id=thread_id,
|
|
)
|
|
|
|
return BackgroundRunRequest(
|
|
graph_id=OBSERVATION_LINKER_GRAPH_ID,
|
|
run_payload=run_payload,
|
|
thread_metadata=_observation_linker_metadata(context),
|
|
name="EvoMemory observation linker",
|
|
)
|
|
|
|
|
|
def _observation_linker_launch_hooks(memory_dir: str | Path) -> BackgroundRunHooks:
|
|
before_relations: dict[str, ObservationRelationSnapshot] = {}
|
|
|
|
def on_before_run(_thread_id: str) -> None:
|
|
before_relations["value"] = snapshot_observation_relations(memory_dir)
|
|
|
|
def on_started(run: BackgroundRun) -> None:
|
|
mark_observation_linker_started(
|
|
thread_id=run.thread_id,
|
|
run_id=run.run_id,
|
|
before_relations=before_relations.get("value"),
|
|
)
|
|
|
|
def on_finished(run: BackgroundRun) -> None:
|
|
mark_observation_linker_finished(
|
|
run.thread_id,
|
|
run.run_id,
|
|
memory_dir=memory_dir,
|
|
)
|
|
|
|
def on_aborted(run: BackgroundRun) -> None:
|
|
forget_observation_linker(run.thread_id, run.run_id)
|
|
|
|
return BackgroundRunHooks(
|
|
on_before_run=on_before_run,
|
|
on_started=on_started,
|
|
on_finished=on_finished,
|
|
on_aborted=on_aborted,
|
|
on_watcher_start_failed=on_aborted,
|
|
)
|
|
|
|
|
|
def _memory_worker_launch_hooks(
|
|
memory_dir: str | Path,
|
|
*,
|
|
on_worker_finished: MemoryWorkerFinishedHook | None = None,
|
|
on_worker_aborted: MemoryWorkerAbortedHook | None = None,
|
|
) -> BackgroundRunHooks:
|
|
before_outputs: dict[str, MemoryOutputSnapshot] = {}
|
|
|
|
def on_before_run(_thread_id: str) -> None:
|
|
before_outputs["value"] = snapshot_memory_outputs(memory_dir)
|
|
|
|
def on_started(run: BackgroundRun) -> None:
|
|
mark_memory_worker_started(
|
|
thread_id=run.thread_id,
|
|
run_id=run.run_id,
|
|
memory_dir=memory_dir,
|
|
before_outputs=before_outputs.get("value"),
|
|
)
|
|
|
|
def on_finished(run: BackgroundRun) -> None:
|
|
delta = mark_memory_worker_finished(run.thread_id, run.run_id)
|
|
if on_worker_finished is not None:
|
|
on_worker_finished(run, delta)
|
|
|
|
def on_aborted(run: BackgroundRun) -> None:
|
|
delta = mark_memory_worker_finished(run.thread_id, run.run_id)
|
|
if on_worker_aborted is not None:
|
|
on_worker_aborted(run, delta)
|
|
|
|
def on_status_unknown(run: BackgroundRun) -> None:
|
|
forget_memory_worker(run.thread_id, run.run_id)
|
|
|
|
return BackgroundRunHooks(
|
|
on_before_run=on_before_run,
|
|
on_started=on_started,
|
|
on_finished=on_finished,
|
|
on_aborted=on_aborted,
|
|
on_status_unknown=on_status_unknown,
|
|
on_watcher_start_failed=on_aborted,
|
|
)
|
|
|
|
|
|
def launch_memory_worker(
|
|
context: MemorySourceContext,
|
|
*,
|
|
on_worker_finished: MemoryWorkerFinishedHook | None = None,
|
|
on_worker_aborted: MemoryWorkerAbortedHook | None = None,
|
|
) -> BackgroundRun | None:
|
|
"""Launch one synchronous EvoMemory worker for a source context."""
|
|
return launch_background_run(
|
|
memory_worker_launch_request(context),
|
|
hooks=_memory_worker_launch_hooks(
|
|
context.memory_dir,
|
|
on_worker_finished=on_worker_finished,
|
|
on_worker_aborted=on_worker_aborted,
|
|
),
|
|
)
|
|
|
|
|
|
async def alaunch_memory_worker(
|
|
context: MemorySourceContext,
|
|
*,
|
|
on_worker_finished: MemoryWorkerFinishedHook | None = None,
|
|
on_worker_aborted: MemoryWorkerAbortedHook | None = None,
|
|
) -> BackgroundRun | None:
|
|
"""Launch one asynchronous EvoMemory worker for a source context."""
|
|
return await alaunch_background_run(
|
|
memory_worker_launch_request(context),
|
|
hooks=_memory_worker_launch_hooks(
|
|
context.memory_dir,
|
|
on_worker_finished=on_worker_finished,
|
|
on_worker_aborted=on_worker_aborted,
|
|
),
|
|
)
|
|
|
|
|
|
def launch_observation_linker(
|
|
context: ObservationLinkerContext,
|
|
) -> BackgroundRun | None:
|
|
"""Launch one synchronous observation-linking pass."""
|
|
if not _observation_linking_enabled():
|
|
return None
|
|
return launch_background_run(
|
|
observation_linker_launch_request(context),
|
|
hooks=_observation_linker_launch_hooks(context.memory_dir),
|
|
)
|
|
|
|
|
|
async def alaunch_observation_linker(
|
|
context: ObservationLinkerContext,
|
|
) -> BackgroundRun | None:
|
|
"""Launch one asynchronous observation-linking pass."""
|
|
if not _observation_linking_enabled():
|
|
return None
|
|
return await alaunch_background_run(
|
|
observation_linker_launch_request(context),
|
|
hooks=_observation_linker_launch_hooks(context.memory_dir),
|
|
)
|