diff --git a/apps/home/events.py b/apps/home/events.py deleted file mode 100644 index 29769590a..000000000 --- a/apps/home/events.py +++ /dev/null @@ -1,61 +0,0 @@ -# SPDX-License-Identifier: AGPL-3.0-only -# Copyright (c) 2026 sol pbc - -"""Callosum event handlers for conversation exchange recording. - -Records triage agent completions to the conversation memory service -so exchanges persist for future context injection. -""" - -import logging - -from apps.events import EventContext, on_event -from think.conversation import record_exchange -from think.cortex_client import read_use_events - -logger = logging.getLogger(__name__) - -TRIAGE_AGENT_NAMES = {"unified", "triage"} - - -@on_event("cortex", "finish") -def record_triage_exchange(ctx: EventContext) -> None: - """Record completed triage agent exchanges to conversation memory.""" - name = ctx.msg.get("name") - if name not in TRIAGE_AGENT_NAMES: - return - - use_id = ctx.msg.get("use_id") - if not use_id: - return - - try: - events = read_use_events(use_id) - facet = "" - app = "" - path = "" - user_message = "" - for event in events: - if event.get("event") == "request": - facet = event.get("facet", "") - app = event.get("app", "") - path = event.get("path", "") - user_message = event.get("user_message", "") - break - - result = ctx.msg.get("result", "") - record_exchange( - facet=facet, - app=app, - path=path, - user_message=user_message, - agent_response=result, - talent=name, - use_id=use_id, - ) - except Exception: - logger.debug( - "Failed to record conversation exchange for agent %s", - use_id, - exc_info=True, - ) diff --git a/convey/__init__.py b/convey/__init__.py index 6bfd76365..c3f42733c 100644 --- a/convey/__init__.py +++ b/convey/__init__.py @@ -20,9 +20,9 @@ from apps import AppRegistry from . import state, system from .apps import register_app_context from .bridge import emit, register_websocket +from .chat import chat_bp, start_chat_runtime from .config import bp as config_bp from .root import bp as root_bp -from .triage import bp as triage_bp __all__ = [ "create_app", @@ -144,8 +144,8 @@ def create_app(journal: str = "") -> Flask: # Register config API blueprint app.register_blueprint(config_bp) - # Register triage API blueprint (universal chat bar) - app.register_blueprint(triage_bp) + # Register chat API blueprint (universal chat bar) + app.register_blueprint(chat_bp) # Register system health API blueprint app.register_blueprint(system.bp) @@ -172,6 +172,7 @@ def create_app(journal: str = "") -> Flask: register_websocket(sock) start_voice_runtime(app) start_push_runtime(app) + start_chat_runtime(app) if journal: state.journal_root = journal diff --git a/convey/chat.py b/convey/chat.py new file mode 100644 index 000000000..aa4a782aa --- /dev/null +++ b/convey/chat.py @@ -0,0 +1,855 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +# Chat backend runs in a single Flask worker process. The threading.Lock plus +# module-level singleton state assumes one convey process per stack. + +from __future__ import annotations + +import atexit +import json +import logging +import pprint +import re +import threading +from dataclasses import dataclass, field +from datetime import datetime +from typing import Any + +from flask import Blueprint, jsonify, request + +from convey.chat_stream import ( + append_chat_event, + find_unresponded_trigger, + read_chat_events, + reduce_chat_state, +) +from convey.utils import error_response +from think.callosum import CallosumConnection, callosum_send +from think.utils import now_ms + +logger = logging.getLogger(__name__) + +chat_bp = Blueprint("chat", __name__, url_prefix="/api/chat") + +MAX_ACTIVE_EXECS = 2 +MAX_LOOP_RETRIES = 3 +DEFAULT_STREAM_LIMIT = 200 +MAX_STREAM_LIMIT = 1000 +MAX_ACTIVE_REASON = "max active — waiting for one to finish" +CHAT_TROUBLE_REASON = "chat had trouble — try again" + +_DAY_RE = re.compile(r"^\d{8}$") +_state_lock = threading.Lock() +_runtime_lock = threading.Lock() +_current_chat_use_id: str | None = None +_current_chat_state: dict[str, Any] | None = None +_queued_trigger: dict[str, Any] | None = None +_active_execs: dict[str, dict[str, Any]] = {} +_last_use_id = 0 +_recovery_day: str | None = None +_runtime: "ChatRuntimeState | None" = None +_atexit_registered = False + + +@dataclass +class ChatRuntimeState: + callosum: CallosumConnection + apps: list[Any] = field(default_factory=list) + + +@chat_bp.route("", methods=["POST"]) +def post_chat() -> Any: + """Accept an owner message and schedule the chat singleton.""" + payload = request.get_json(force=True) or {} + message = str(payload.get("message") or "").strip() + if not message: + return error_response("message is required", 400) + + from think.identity import ensure_identity_directory + + ensure_identity_directory() + + location = _normalize_location( + payload.get("app"), + payload.get("path"), + payload.get("facet"), + ) + append_chat_event( + "owner_message", + text=message, + app=location["app"], + path=location["path"], + facet=location["facet"], + ) + trigger = { + "type": "owner_message", + "message": message, + } + + start_info: dict[str, Any] | None = None + with _state_lock: + if _current_chat_use_id is None: + logical_use_id = _reserve_use_id_locked() + start_info = _activate_current_locked(logical_use_id, trigger, location) + queued = False + response_use_id = logical_use_id + else: + response_use_id = _queue_trigger_locked(trigger, location) + queued = True + + if start_info is not None and not _spawn_chat_generate(start_info): + _handle_chat_failure(response_use_id, CHAT_TROUBLE_REASON) + return error_response("Failed to connect to agent service", 503) + + return jsonify(use_id=response_use_id, queued=queued) + + +@chat_bp.route("/session", methods=["GET"]) +def chat_session() -> Any: + """Return reduced state for today's chat stream.""" + return jsonify(reduce_chat_state(_today_day())) + + +@chat_bp.route("/stream/", methods=["GET"]) +def chat_stream(day: str) -> Any: + """Return ordered chat events for a day.""" + if not _DAY_RE.fullmatch(day): + return error_response("day must be YYYYMMDD", 400) + + limit_raw = request.args.get("limit", str(DEFAULT_STREAM_LIMIT)) + try: + limit = int(limit_raw) + except (TypeError, ValueError): + limit = DEFAULT_STREAM_LIMIT + limit = max(1, min(limit, MAX_STREAM_LIMIT)) + + return jsonify(events=read_chat_events(day, limit=limit)) + + +@chat_bp.route("/result/", methods=["GET"]) +def chat_result(use_id: str) -> Any: + """Return chat or exec state from the chat stream.""" + result = _read_result_state(use_id) + if result is None: + return jsonify(error="not found"), 404 + return jsonify(result) + + +def start_chat_runtime(app: Any) -> None: + """Start the chat backend runtime and subscribe to cortex events.""" + global _runtime, _atexit_registered + + with _runtime_lock: + if _runtime is None: + runtime = ChatRuntimeState(callosum=CallosumConnection()) + runtime.callosum.start(callback=_handle_callosum_message) + _runtime = runtime + runtime = _runtime + if app not in runtime.apps: + runtime.apps.append(app) + app.chat_runtime_started = True + if not _atexit_registered: + atexit.register(stop_all_chat_runtime) + _atexit_registered = True + + _recover_chat_if_needed() + + +def stop_chat_runtime(app: Any) -> None: + """Detach an app from the shared runtime.""" + app.chat_runtime_started = False + runtime = _runtime + if runtime is None: + return + with _runtime_lock: + if app in runtime.apps: + runtime.apps.remove(app) + remaining = list(runtime.apps) + if not remaining: + stop_all_chat_runtime() + + +def stop_all_chat_runtime() -> None: + """Stop the shared runtime.""" + global _runtime + + with _runtime_lock: + runtime = _runtime + _runtime = None + if runtime is None: + return + for app in list(runtime.apps): + try: + app.chat_runtime_started = False + except Exception: + logger.exception("chat runtime app cleanup failed") + runtime.callosum.stop() + + +def _handle_callosum_message(message: dict[str, Any]) -> None: + if message.get("chat_proxy"): + return + if message.get("tract") != "cortex": + return + + event_type = message.get("event") + if event_type == "finish": + _on_cortex_finish(message) + return + if event_type == "error": + _on_cortex_error(message) + return + + _proxy_progress(message) + + +def _proxy_progress(message: dict[str, Any]) -> None: + logical_use_id: str | None = None + use_id = str(message.get("use_id") or "") + if not use_id: + return + + with _state_lock: + if _current_chat_state is None or _current_chat_use_id is None: + return + raw_chat_use_id = str(_current_chat_state.get("raw_use_id") or "") + if use_id == raw_chat_use_id: + logical_use_id = _current_chat_use_id + elif use_id in _active_execs: + logical_use_id = str(_active_execs[use_id]["chat_use_id"]) + + if logical_use_id is None: + return + + fields = { + key: value + for key, value in message.items() + if key not in {"tract", "event", "use_id"} + } + fields["use_id"] = logical_use_id + fields["chat_proxy"] = True + _emit_cortex_event(message["event"], **fields) + + +def _on_cortex_finish(message: dict[str, Any]) -> None: + use_id = str(message.get("use_id") or "") + if not use_id: + return + + next_info: dict[str, Any] | None = None + finish_payload: dict[str, Any] | None = None + error_payload: dict[str, Any] | None = None + + with _state_lock: + if _current_chat_state is not None and use_id == _current_chat_state.get( + "raw_use_id" + ): + logical_use_id = str(_current_chat_use_id) + try: + parsed = _parse_chat_result(message.get("result")) + except ValueError: + if int(_current_chat_state.get("retry_count", 0) or 0) < 1: + retry_use_id = _reserve_use_id_locked() + _current_chat_state["raw_use_id"] = retry_use_id + _current_chat_state["retry_count"] = ( + int(_current_chat_state.get("retry_count", 0) or 0) + 1 + ) + next_info = _build_spawn_info_locked(logical_use_id) + else: + append_chat_event( + "chat_error", + reason=CHAT_TROUBLE_REASON, + use_id=logical_use_id, + ) + error_payload = { + "use_id": logical_use_id, + "reason": CHAT_TROUBLE_REASON, + } + next_info = _clear_current_locked() + else: + message_text = parsed["message"] or "" + requested_exec = parsed["talent_request"] is not None + requested_task = ( + parsed["talent_request"]["task"] + if parsed["talent_request"] + else None + ) + append_chat_event( + "sol_message", + use_id=logical_use_id, + text=message_text, + notes=parsed["notes"], + requested_exec=requested_exec, + requested_task=requested_task, + ) + _current_chat_state["retry_count"] = 0 + _current_chat_state["raw_use_id"] = None + if requested_exec: + active_exec_count = _active_exec_count_for_today_locked() + if active_exec_count >= MAX_ACTIVE_EXECS: + _current_chat_state["trigger"] = { + "type": "synthetic-max-active", + "reason": MAX_ACTIVE_REASON, + } + synthetic_use_id = _reserve_use_id_locked() + _current_chat_state["raw_use_id"] = synthetic_use_id + next_info = _build_spawn_info_locked(logical_use_id) + elif _exec_loop_count_locked() >= MAX_LOOP_RETRIES: + append_chat_event( + "chat_error", + reason=CHAT_TROUBLE_REASON, + use_id=logical_use_id, + ) + error_payload = { + "use_id": logical_use_id, + "reason": CHAT_TROUBLE_REASON, + } + next_info = _clear_current_locked() + else: + exec_use_id = _reserve_use_id_locked() + _active_execs[exec_use_id] = { + "chat_use_id": logical_use_id, + "task": requested_task, + "location": dict(_current_chat_state["location"]), + } + append_chat_event( + "talent_spawned", + use_id=exec_use_id, + name="exec", + task=requested_task, + started_at=int(exec_use_id), + ) + next_info = { + "kind": "exec", + "logical_use_id": logical_use_id, + "use_id": exec_use_id, + "task": requested_task, + "context": parsed["talent_request"].get("context") or {}, + "location": dict(_current_chat_state["location"]), + } + else: + if not message_text: + append_chat_event( + "chat_error", + reason=CHAT_TROUBLE_REASON, + use_id=logical_use_id, + ) + error_payload = { + "use_id": logical_use_id, + "reason": CHAT_TROUBLE_REASON, + } + else: + finish_payload = { + "use_id": logical_use_id, + "message": message_text, + } + next_info = _clear_current_locked() + + elif use_id in _active_execs: + exec_state = _active_execs.pop(use_id) + logical_use_id = str(exec_state["chat_use_id"]) + summary = str(message.get("result") or "").strip() + append_chat_event( + "talent_finished", + use_id=use_id, + name="exec", + summary=summary, + ) + if ( + _current_chat_use_id == logical_use_id + and _current_chat_state is not None + ): + _current_chat_state["trigger"] = { + "type": "talent_finished", + "use_id": use_id, + "name": "exec", + "summary": summary, + } + _current_chat_state["raw_use_id"] = _reserve_use_id_locked() + _current_chat_state["retry_count"] = 0 + next_info = _build_spawn_info_locked(logical_use_id) + + _run_next_action(next_info) + if finish_payload is not None: + _emit_finish(finish_payload["use_id"], finish_payload["message"]) + if error_payload is not None: + _emit_error(error_payload["use_id"], error_payload["reason"]) + + +def _on_cortex_error(message: dict[str, Any]) -> None: + use_id = str(message.get("use_id") or "") + if not use_id: + return + + next_info: dict[str, Any] | None = None + error_payload: dict[str, Any] | None = None + + with _state_lock: + if _current_chat_state is not None and use_id == _current_chat_state.get( + "raw_use_id" + ): + logical_use_id = str(_current_chat_use_id) + append_chat_event( + "chat_error", + reason=CHAT_TROUBLE_REASON, + use_id=logical_use_id, + ) + error_payload = {"use_id": logical_use_id, "reason": CHAT_TROUBLE_REASON} + next_info = _clear_current_locked() + elif use_id in _active_execs: + exec_state = _active_execs.pop(use_id) + logical_use_id = str(exec_state["chat_use_id"]) + reason = str(message.get("error") or CHAT_TROUBLE_REASON) + append_chat_event( + "talent_errored", + use_id=use_id, + name="exec", + reason=reason, + ) + if ( + _current_chat_use_id == logical_use_id + and _current_chat_state is not None + ): + _current_chat_state["trigger"] = { + "type": "talent_errored", + "use_id": use_id, + "name": "exec", + "reason": reason, + } + _current_chat_state["raw_use_id"] = _reserve_use_id_locked() + _current_chat_state["retry_count"] = 0 + next_info = _build_spawn_info_locked(logical_use_id) + + _run_next_action(next_info) + if error_payload is not None: + _emit_error(error_payload["use_id"], error_payload["reason"]) + + +def _run_next_action(action: dict[str, Any] | None) -> None: + if action is None: + return + if action.get("kind") == "chat": + if not _spawn_chat_generate(action): + _handle_chat_failure(action["logical_use_id"], CHAT_TROUBLE_REASON) + return + if action.get("kind") == "exec": + if not _spawn_exec(action): + _handle_exec_spawn_failure(action) + + +def _spawn_chat_generate(action: dict[str, Any]) -> bool: + logger.info( + "starting chat generate logical=%s raw=%s trigger=%s", + action["logical_use_id"], + action["raw_use_id"], + action["trigger"]["type"], + ) + from convey.utils import spawn_agent + + config = { + "app": action["location"]["app"], + "path": action["location"]["path"], + "facet": action["location"]["facet"], + "trigger": action["trigger"], + "chat_request_use_id": action["logical_use_id"], + } + use_id = spawn_agent( + prompt="", + name="chat", + provider=None, + config=config, + use_id=action["raw_use_id"], + ) + if use_id is None: + return False + _emit_cortex_event("thinking", use_id=action["logical_use_id"], chat_proxy=True) + return True + + +def _spawn_exec(action: dict[str, Any]) -> bool: + from convey.utils import spawn_agent + + prompt = _build_exec_prompt( + action["task"], + action["context"], + action["location"], + ) + config = { + "app": action["location"]["app"], + "path": action["location"]["path"], + "facet": action["location"]["facet"], + "chat_parent_use_id": action["logical_use_id"], + } + use_id = spawn_agent( + prompt=prompt, + name="exec", + provider=None, + config=config, + use_id=action["use_id"], + ) + if use_id is None: + return False + _emit_cortex_event("thinking", use_id=action["logical_use_id"], chat_proxy=True) + return True + + +def _handle_exec_spawn_failure(action: dict[str, Any]) -> None: + next_info: dict[str, Any] | None = None + with _state_lock: + _active_execs.pop(str(action["use_id"]), None) + append_chat_event( + "talent_errored", + use_id=action["use_id"], + name="exec", + reason=CHAT_TROUBLE_REASON, + ) + if _current_chat_use_id == action["logical_use_id"] and _current_chat_state: + _current_chat_state["trigger"] = { + "type": "talent_errored", + "use_id": action["use_id"], + "name": "exec", + "reason": CHAT_TROUBLE_REASON, + } + _current_chat_state["raw_use_id"] = _reserve_use_id_locked() + _current_chat_state["retry_count"] = 0 + next_info = _build_spawn_info_locked(action["logical_use_id"]) + _run_next_action(next_info) + + +def _handle_chat_failure(logical_use_id: str, reason: str) -> None: + next_info: dict[str, Any] | None = None + with _state_lock: + append_chat_event("chat_error", reason=reason, use_id=logical_use_id) + if _current_chat_use_id == logical_use_id: + next_info = _clear_current_locked() + _emit_error(logical_use_id, reason) + _run_next_action(next_info) + + +def _recover_chat_if_needed() -> None: + day = _today_day() + start_info: dict[str, Any] | None = None + + with _state_lock: + global _recovery_day + if _recovery_day == day or _current_chat_use_id is not None: + return + unresolved = find_unresponded_trigger(day) + if unresolved is None: + _recovery_day = day + return + location = _location_for_trigger(day, unresolved) + logical_use_id = _reserve_use_id_locked() + trigger = _trigger_from_stream_event(unresolved) + start_info = _activate_current_locked(logical_use_id, trigger, location) + _recovery_day = day + + if start_info is not None and not _spawn_chat_generate(start_info): + _handle_chat_failure(start_info["logical_use_id"], CHAT_TROUBLE_REASON) + + +def _activate_current_locked( + logical_use_id: str, + trigger: dict[str, Any], + location: dict[str, str], +) -> dict[str, Any]: + global _current_chat_use_id, _current_chat_state + + raw_use_id = _reserve_use_id_locked() + _current_chat_use_id = logical_use_id + _current_chat_state = { + "raw_use_id": raw_use_id, + "trigger": dict(trigger), + "location": dict(location), + "retry_count": 0, + } + return _build_spawn_info_locked(logical_use_id) + + +def _build_spawn_info_locked(logical_use_id: str) -> dict[str, Any]: + assert _current_chat_state is not None + return { + "kind": "chat", + "logical_use_id": logical_use_id, + "raw_use_id": str(_current_chat_state["raw_use_id"]), + "trigger": dict(_current_chat_state["trigger"]), + "location": dict(_current_chat_state["location"]), + } + + +def _queue_trigger_locked(trigger: dict[str, Any], location: dict[str, str]) -> str: + global _queued_trigger + if _queued_trigger is None: + _queued_trigger = { + "use_id": _reserve_use_id_locked(), + "trigger": dict(trigger), + "location": dict(location), + } + return str(_queued_trigger["use_id"]) + + +def _clear_current_locked() -> dict[str, Any] | None: + global _current_chat_use_id, _current_chat_state, _queued_trigger + + _current_chat_use_id = None + _current_chat_state = None + if _queued_trigger is None: + return None + + queued = _queued_trigger + _queued_trigger = None + return _activate_current_locked( + str(queued["use_id"]), + dict(queued["trigger"]), + dict(queued["location"]), + ) + + +def _active_exec_count_for_today_locked() -> int: + return len(reduce_chat_state(_today_day())["active_talents"]) + + +def _exec_loop_count_locked() -> int: + events = read_chat_events(_today_day()) + count = 0 + for index in range(len(events) - 1, -1, -1): + event = events[index] + kind = event.get("kind") + if kind == "owner_message": + break + if kind != "sol_message": + continue + if not event.get("requested_exec"): + continue + + previous = events[index - 1] if index > 0 else None + if previous and previous.get("kind") in {"talent_finished", "talent_errored"}: + count += 1 + else: + break + return count + + +def _parse_chat_result(result: Any) -> dict[str, Any]: + if isinstance(result, str): + payload = json.loads(result) + elif isinstance(result, dict): + payload = result + else: + raise ValueError("chat result must be JSON text") + + if not isinstance(payload, dict): + raise ValueError("chat result must be an object") + if not isinstance(payload.get("notes"), str): + raise ValueError("chat result notes must be a string") + + message = payload.get("message") + if message is not None and not isinstance(message, str): + raise ValueError("chat result message must be a string or null") + + talent_request = payload.get("talent_request") + if talent_request is None: + return {"message": message, "notes": payload["notes"], "talent_request": None} + if not isinstance(talent_request, dict): + raise ValueError("chat talent_request must be an object or null") + task = talent_request.get("task") + if not isinstance(task, str) or not task.strip(): + raise ValueError("chat talent_request.task must be a non-empty string") + context = talent_request.get("context") or {} + if not isinstance(context, dict): + raise ValueError("chat talent_request.context must be an object") + return { + "message": message, + "notes": payload["notes"], + "talent_request": { + "task": task.strip(), + "context": context, + }, + } + + +def _build_exec_prompt( + task: str, + context_hints: dict[str, Any], + location: dict[str, str], +) -> str: + parts = [f"Task: {task}"] + if context_hints: + parts.append( + "Context hints:\n" + pprint.pformat(context_hints, sort_dicts=True) + ) + parts.append( + "Location: " + f"app={location['app']} path={location['path']} facet={location['facet']}" + ) + + history_lines: list[str] = [] + for event in read_chat_events(_today_day()): + kind = event.get("kind") + if kind == "owner_message": + history_lines.append(f"**Owner**: {event['text']}") + elif kind == "sol_message": + history_lines.append(f"**Sol**: {event['text']}") + if history_lines: + parts.append("Recent chat:\n" + "\n".join(history_lines[-6:])) + + return "\n\n".join(parts) + + +def _emit_finish(use_id: str, message: str) -> None: + _emit_cortex_event( + "finish", + use_id=use_id, + result=message, + display=_display_mode(message), + chat_proxy=True, + ) + + +def _emit_error(use_id: str, reason: str) -> None: + _emit_cortex_event( + "error", + use_id=use_id, + error=reason, + chat_proxy=True, + ) + + +def _emit_cortex_event(event: str, **fields: Any) -> None: + runtime = _runtime + if runtime is not None and runtime.callosum.emit("cortex", event, **fields): + return + callosum_send("cortex", event, **fields) + + +def _display_mode(text: str) -> str: + if not text: + return "inline" + if len(text) >= 120 or "\n" in text: + return "panel" + if len(re.split(r"(?<=[.!?])\s", text)) > 2: + return "panel" + return "inline" + + +def _normalize_location(app_name: Any, path: Any, facet: Any) -> dict[str, str]: + return { + "app": str(app_name or ""), + "path": str(path or ""), + "facet": str(facet or ""), + } + + +def _location_for_trigger(day: str, trigger: dict[str, Any]) -> dict[str, str]: + if trigger.get("kind") == "owner_message": + return _normalize_location( + trigger.get("app"), + trigger.get("path"), + trigger.get("facet"), + ) + for event in reversed(read_chat_events(day)): + if event.get("kind") == "owner_message": + return _normalize_location( + event.get("app"), + event.get("path"), + event.get("facet"), + ) + return _normalize_location("", "", "") + + +def _trigger_from_stream_event(event: dict[str, Any]) -> dict[str, Any]: + kind = event.get("kind") + if kind == "owner_message": + return {"type": "owner_message", "message": event.get("text", "")} + if kind == "talent_finished": + return { + "type": "talent_finished", + "use_id": event.get("use_id"), + "name": event.get("name", "exec"), + "summary": event.get("summary", ""), + } + if kind == "talent_errored": + return { + "type": "talent_errored", + "use_id": event.get("use_id"), + "name": event.get("name", "exec"), + "reason": event.get("reason", ""), + } + raise ValueError(f"unsupported trigger event: {kind}") + + +def _read_result_state(use_id: str) -> dict[str, Any] | None: + day = _day_for_use_id(use_id) + if day is None: + return None + + latest_sol: dict[str, Any] | None = None + exec_state: dict[str, Any] | None = None + chat_error: dict[str, Any] | None = None + spawned_task: str | None = None + + for event in read_chat_events(day): + kind = event.get("kind") + if kind == "sol_message" and str(event.get("use_id")) == use_id: + latest_sol = event + elif kind == "chat_error" and str(event.get("use_id") or "") == use_id: + chat_error = event + elif kind == "talent_spawned" and str(event.get("use_id")) == use_id: + spawned_task = event.get("task") + exec_state = {"state": "active", "task": spawned_task} + elif kind == "talent_finished" and str(event.get("use_id")) == use_id: + exec_state = { + "state": "finished", + "summary": event.get("summary", ""), + "task": spawned_task, + } + elif kind == "talent_errored" and str(event.get("use_id")) == use_id: + exec_state = { + "state": "errored", + "reason": event.get("reason", ""), + "task": spawned_task, + } + + with _state_lock: + if _current_chat_use_id == use_id: + task = None + if latest_sol and latest_sol.get("requested_exec"): + task = latest_sol.get("requested_task") + return {"state": "active", "task": task} + + if chat_error is not None: + return { + "state": "errored", + "reason": chat_error.get("reason", CHAT_TROUBLE_REASON), + } + if latest_sol is not None: + return { + "state": "finished", + "summary": latest_sol.get("text", ""), + "display": _display_mode(str(latest_sol.get("text", ""))), + } + return exec_state + + +def _reserve_use_id_locked() -> str: + global _last_use_id + + ts = now_ms() + if ts <= _last_use_id: + ts = _last_use_id + 1 + _last_use_id = ts + return str(ts) + + +def _today_day() -> str: + return datetime.now().strftime("%Y%m%d") + + +def _day_for_use_id(use_id: str) -> str | None: + if not use_id.isdigit(): + return None + try: + return datetime.fromtimestamp(int(use_id) / 1000).strftime("%Y%m%d") + except (OSError, OverflowError, ValueError): + return None diff --git a/convey/templates/app.html b/convey/templates/app.html index b58a8bd68..15076a7b5 100644 --- a/convey/templates/app.html +++ b/convey/templates/app.html @@ -474,14 +474,19 @@ // Check GET immediately in case agent already finished try { - var r = await fetch('/api/triage/result/' + agentId); + var r = await fetch('/api/chat/result/' + agentId); if (r.ok) { // Agent already finished — cancel WS subscription and display - if (recoveryCleanup) { recoveryCleanup(); recoveryCleanup = null; } - if (recoveryWatchdog) { clearTimeout(recoveryWatchdog); recoveryWatchdog = null; } var data = await r.json(); - var resp = data.response || ''; - recoverDeliver(resp, data.display || 'panel', null); + if (data.state === 'finished') { + if (recoveryCleanup) { recoveryCleanup(); recoveryCleanup = null; } + if (recoveryWatchdog) { clearTimeout(recoveryWatchdog); recoveryWatchdog = null; } + recoverDeliver(data.summary || '', data.display || 'panel', null); + } else if (data.state === 'errored') { + if (recoveryCleanup) { recoveryCleanup(); recoveryCleanup = null; } + if (recoveryWatchdog) { clearTimeout(recoveryWatchdog); recoveryWatchdog = null; } + recoverDeliver('', 'panel', data.reason || 'something went wrong. the server returned an unexpected response. try sending your message again, or check the health page if it keeps happening.'); + } } // If 404, agent still running — keep WS subscription active } catch (err) { @@ -693,7 +698,7 @@ } try { - var r = await fetch('/api/triage', { + var r = await fetch('/api/chat', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ diff --git a/convey/triage.py b/convey/triage.py deleted file mode 100644 index 133ce50fa..000000000 --- a/convey/triage.py +++ /dev/null @@ -1,134 +0,0 @@ -# SPDX-License-Identifier: AGPL-3.0-only -# Copyright (c) 2026 sol pbc - -"""Triage endpoint for universal chat bar / conversation panel queries.""" - -from __future__ import annotations - -import logging -import re -from typing import Any - -from flask import Blueprint, jsonify, request - -from convey.utils import error_response - -logger = logging.getLogger(__name__) - - -def compute_display_mode(text: str) -> str: - """Return 'inline' or 'panel' based on response text characteristics.""" - if not text: - return "inline" - if len(text) >= 120: - return "panel" - if "\n" in text: - return "panel" - if len(re.split(r"(?<=[.!?])\s", text)) > 2: - return "panel" - return "inline" - - -bp = Blueprint("triage", __name__, url_prefix="/api/triage") - - -@bp.route("", methods=["POST"]) -def triage() -> Any: - """Accept a message from the conversation panel and spawn a triage agent. - - Expects JSON: {message, app, path, facet} - Returns JSON: {use_id} - - The agent runs asynchronously. The browser receives the result via - WebSocket (cortex/finish event). For reload recovery, use GET /result/. - - All journals route to the unified talent. - """ - payload = request.get_json(force=True) - message = payload.get("message", "").strip() - - from think.identity import ensure_identity_directory - - ensure_identity_directory() - - if not message: - return error_response("message is required", 400) - - app_name = payload.get("app", "") - path = payload.get("path", "") - facet = payload.get("facet", "") - agent_name = "unified" - - # Build prompt with location context - context_lines = [] - if app_name: - context_lines.append(f"Current app: {app_name}") - if path: - context_lines.append(f"Current path: {path}") - if facet: - context_lines.append(f"Current facet: {facet}") - - # Add system health context when attention items exist - try: - from convey.apps import _resolve_attention - from think.awareness import get_current - - attention = _resolve_attention(get_current()) - if attention: - context_lines.extend(attention.context_lines) - except Exception: - pass # Don't let health context break triage - - # Assemble the full prompt - prompt_parts = [] - if context_lines: - prompt_parts.append("\n".join(context_lines)) - prompt_parts.append(message) - full_prompt = "\n\n".join(prompt_parts) - - try: - from convey.utils import spawn_agent - - config: dict[str, Any] = {} - if facet: - config["facet"] = facet - config["app"] = app_name - config["path"] = path - config["user_message"] = message - - use_id = spawn_agent( - prompt=full_prompt, - name=agent_name, - provider=None, - config=config, - ) - if use_id is None: - return error_response("Failed to connect to agent service", 503) - - return jsonify(use_id=use_id) - - except Exception: - logger.exception("Triage request failed") - return error_response("Failed to process triage request", 500) - - -@bp.route("/result/", methods=["GET"]) -def triage_result(use_id: str) -> Any: - """Return the result of a completed triage agent. - - Returns {response, display} if the agent has finished, 404 otherwise. - Used for page-reload recovery when the WebSocket may have missed the finish event. - """ - try: - from think.cortex_client import read_use_events - - events = read_use_events(use_id) - for event in reversed(events): - if event.get("event") == "finish": - result = event.get("result", "") - return jsonify(response=result, display=compute_display_mode(result)) - except FileNotFoundError: - pass - except Exception: - logger.debug("Failed to read triage result for %s", use_id, exc_info=True) - return jsonify(error="not found"), 404 diff --git a/convey/utils.py b/convey/utils.py index b3dd4d4ff..c5c376919 100644 --- a/convey/utils.py +++ b/convey/utils.py @@ -87,6 +87,7 @@ def spawn_agent( name: str, provider: Optional[str] = None, config: Optional[dict[str, Any]] = None, + use_id: Optional[str] = None, ) -> str | None: """Spawn a Cortex agent and return the use_id. @@ -98,6 +99,7 @@ def spawn_agent( name: Agent name - system (e.g., "default") or app-qualified (e.g., "entities:entity_assist") provider: Optional provider override (openai, google, anthropic) config: Additional configuration (max_tokens, facet, session_id, etc.) + use_id: Optional pre-reserved Cortex use_id to reuse for the request Returns: use_id string (timestamp-based), or None if the request could not be sent. @@ -112,6 +114,7 @@ def spawn_agent( name=name, provider=provider, config=config, + use_id=use_id, ) diff --git a/talent/conversation_memory.py b/talent/conversation_memory.py deleted file mode 100644 index da339e64b..000000000 --- a/talent/conversation_memory.py +++ /dev/null @@ -1,45 +0,0 @@ -# SPDX-License-Identifier: AGPL-3.0-only -# Copyright (c) 2026 sol pbc - -"""Pre-hook: inject conversation memory into unified talent context. - -Loaded via hook config: {"hook": {"pre": "conversation_memory"}} - -Replaces CONVERSATION_MEMORY_INJECTION_POINT in the unified talent's -user instruction with recent conversation exchanges and today's summary. -This gives the agent awareness of past conversations without needing -to search — recent interactions are always in context. -""" - -import logging - -logger = logging.getLogger(__name__) - - -def pre_process(context: dict) -> dict | None: - """Inject conversation memory into the unified talent's user instruction. - - Args: - context: Full agent config dict. - - Returns: - Dict with modified user_instruction, or None if no injection needed. - """ - from think.conversation import INJECTION_MARKER, build_memory_context, inject_memory - - user_instruction = context.get("user_instruction", "") - if INJECTION_MARKER not in user_instruction: - return None - - facet = context.get("facet") - - try: - memory_context = build_memory_context(facet=facet, recent_limit=10) - new_instruction = inject_memory(user_instruction, memory_context) - - if new_instruction != user_instruction: - return {"user_instruction": new_instruction} - except Exception: - logger.exception("Conversation memory injection failed") - - return None diff --git a/talent/triage.md b/talent/triage.md deleted file mode 100644 index 9e5fccf83..000000000 --- a/talent/triage.md +++ /dev/null @@ -1,120 +0,0 @@ -{ - "type": "cogitate", - "title": "Triage", - "description": "Quick-action assistant for the chat bar — handles navigation, todos, calendar, and entity lookups" -} - -You are a quick-action assistant for the sol journal system chat bar. You handle simple actions and short lookups: navigate the app, manage todos, manage calendar events, and look up entities. - -Respond in one concise line for actions you complete. If a request needs deeper analysis, the conversation panel handles it automatically — just answer to the best of your ability. - -You are given context about the owner's current app, URL path, and facet. Use this to inform your actions — for example, use the facet for todo and calendar commands. - -## Available Commands - -### Navigation -- `sol call navigate [PATH] --facet FACET` — Navigate the browser to a path and/or switch facet. - -### Todos -- `sol call todos list [DAY] --facet FACET` — Show todos for a day. -- `sol call todos add TEXT --day DAY --facet FACET [--nudge TIME]` — Add a todo. Nudge formats: HH:MM, now, tomorrow HH:MM, YYYYMMDDTHH:MM. -- `sol call todos done LINE --day DAY --facet FACET` — Mark a todo as done. -- `sol call todos cancel LINE --day DAY --facet FACET` — Cancel a todo. -- `sol call todos upcoming --facet FACET [--limit N]` — Show upcoming todos. - -### Entities -- `sol call entities list [FACET]` — List entities for a facet. -- `sol call entities observations ENTITY --facet FACET` — List observations for an entity. -- `sol call entities observe ENTITY CONTENT --facet FACET` — Record an observation. -- `sol call entities strength [--facet FACET] [--since YYYYMMDD] [--limit N]` — Rank entities by relationship strength. -- `sol call entities search [--query TEXT] [--type TYPE] [--facet FACET] [--since YYYYMMDD] [--limit N]` — Search entities by text, type, or facet. -- `sol call entities intelligence ENTITY [--facet FACET] [--brief]` — Intelligence briefing for an entity (returns JSON — synthesize into natural language). Use --brief for concise lookups. - -### Journal -- `sol call activities list --source anticipated [--day DAY] [-f FACET]` — List anticipated activities with participants, times, and summaries. - -### Awareness -- `sol call awareness status [SECTION]` — Read awareness state (e.g., processing state, journal health). -- `sol call awareness log-read [DAY] [--kind KIND] [--limit N]` — Read awareness log entries. - -### Support -- `sol call support search ` — Search KB articles. -- `sol call support diagnose` — Run local diagnostics (no network). -- `sol call support create --subject "..." --description "..." [--severity medium] [--category bug]` — File a ticket (interactive consent flow). - -## Behavioral Rules - -- After completing an action, respond with one concise line confirming what you did. -- For lookups (list todos, list events, list entities), present the results concisely. -- For entity intelligence briefings, synthesize the JSON output into a concise natural-language summary — do not dump raw JSON. -- **Pre-meeting briefings**: When the owner asks "brief me on my next meeting", "who am I meeting?", or similar: - 1. Run `sol call activities list --source anticipated` to find upcoming events with participants. - 2. For each participant, run `sol call entities intelligence PARTICIPANT --brief` to gather background. - 3. Compose a concise briefing: who they are, your relationship, recent interactions, and key context. - Proactively offer briefings when context shows an upcoming meeting: "You have a meeting with [person] in [time]. Want me to brief you?" -- **Support**: When the owner reports a problem ("this isn't working", "I found a bug", "something's broken"), wants to file a ticket, or wants to give feedback, handle it in-place — search KB, run diagnostics, draft and submit a ticket with the owner's approval. -- Do not attempt to use any commands not listed above. -- SOL_DAY and SOL_FACET environment variables are already set — tools will use them as defaults when --day/--facet are omitted. So you can often omit these flags. - -## System Attention - -When the context includes a `System health:` line, there is an active attention item. Handle these queries: - -- **"what needs my attention?"** — Report the system health item from context. If there are agent errors, mention which agents failed. If an import just completed, mention what arrived. Be concise. -- **Agent errors**: If the owner asks about errors, explain which agents failed today. Suggest checking agent logs or re-running the daily analysis. -- **Import complete**: If an import just finished, briefly describe what was imported and offer to explore the new data or import from another source. - -When no `System health:` line is present in context, there is nothing to report. If the owner asks "what needs my attention?", respond that everything looks good. - -## Import Awareness - -Check import state with `sol call awareness imports`: - -- **After an import completes** (owner returns to chat): The import system updates awareness automatically. If you see `has_imported: true` and new sources in `sources_used`, offer to import from another source: "I just processed your [source] import. Want to import from another source, or explore what I found?" - -- **Soft import nudge**: If all of these are true, you may weave a single soft import mention into your response: - 1. No imports done (`has_imported: false`) - 2. Import offer not recently declined (no `offer_declined` or >3 days ago) - 3. No recent nudge (`last_nudge` is null) - 4. The owner's message touches on their journal, data, or what $agent_name can do - - After mentioning imports, run `sol call awareness imports --nudge` to record it. Do **not** repeat this nudge. - -- **Available sources**: Calendar (ics), ChatGPT (chatgpt), Claude (claude), Gemini (gemini), Notes (obsidian), Kindle (kindle) - -- If the owner wants to import, read the guide from `apps/import/guides/{source}.md` and present the export instructions conversationally. Then navigate to the import app: `sol call navigate "/app/import#guide/{source}"` - -## Naming Awareness - -Check whether the naming ceremony should trigger: - -1. Run `sol call sol name` to check status. -2. If `name_status` is `"default"`, run `sol call sol thickness` to check readiness. -3. If `ready` is `true`, mention that you've been getting to know the owner and offer to suggest a name — or let the naming talent handle it. -4. Only do this once per session. If you've already checked or offered, don't repeat. -5. If `name_status` is `"chosen"` or `"self-named"`, do nothing. - -## Owner Voice Detection Awareness - -Check whether owner voice detection should be surfaced: - -1. Run `sol call speakers owner-ready` to check readiness. -2. If `ready` is `false`, do nothing. The reason field explains why (centroid_exists, cooldown, low_data, no_clusters, etc.). -3. If `ready` is `true`, surface the prompt conversationally: - - > "I've been learning voices from your observed media and I think I can identify yours. Want to listen to a few samples and confirm?" - -4. Only do this once per session. If you've already checked or offered, don't repeat. - -### Handling the owner's response - -- **Owner confirms ("yes", "sure", "go ahead"):** - 1. Run `sol call speakers confirm-owner` — this saves the centroid and automatically runs attribution backfill on all segments. - 2. Report back: "Got it. I'll start labeling speakers in your transcripts." - -- **Owner declines ("no", "not now", "skip"):** - 1. Run `sol call speakers reject-owner` — this enters a 14-day cooldown. - 2. Respond: "No problem — I'll keep listening and try again when I have more to work with." - -- **Owner wants to hear samples first:** - The `owner-ready` result includes a `samples` array with audio URLs. Navigate the owner to the speakers app for the full confirmation flow: `sol call navigate "/app/speakers#owner"` diff --git a/tests/baselines/api/chat/result.json b/tests/baselines/api/chat/result.json new file mode 100644 index 000000000..5244ecd08 --- /dev/null +++ b/tests/baselines/api/chat/result.json @@ -0,0 +1,3 @@ +{ + "error": "not found" +} diff --git a/tests/baselines/api/chat/session.json b/tests/baselines/api/chat/session.json new file mode 100644 index 000000000..d5c72d357 --- /dev/null +++ b/tests/baselines/api/chat/session.json @@ -0,0 +1,5 @@ +{ + "active_talents": [], + "completed_talents": [], + "latest_sol_message": null +} diff --git a/tests/baselines/api/chat/stream.json b/tests/baselines/api/chat/stream.json new file mode 100644 index 000000000..52387b4af --- /dev/null +++ b/tests/baselines/api/chat/stream.json @@ -0,0 +1,3 @@ +{ + "events": [] +} diff --git a/tests/baselines/api/sol/preview.json b/tests/baselines/api/sol/preview.json index 5d8159261..3c4a69621 100644 --- a/tests/baselines/api/sol/preview.json +++ b/tests/baselines/api/sol/preview.json @@ -1,6 +1,6 @@ { "full_prompt": "## Instructions\n\n## Available Facets\n\n- **Capulet Industries** (`capulet`)\n Capulet Industries enterprise division\n - **Capulet Industries Entities**: Capulet Industries; Juliet Capulet; Nurse Angela; Paris Duke; Tybalt Capulet\n - **Capulet Industries Activities**: Meetings; Coding; Browsing; Email; Messaging; AI Conversation; Writing; Reading; Video; Gaming; Social Media; Planning; Productivity; Terminal; Design; Music\n\n- **Empty Entities Test** (`empty-entities`)\n - **Empty Entities Test Activities**: Meetings; Coding; Browsing; Email; Messaging; AI Conversation; Writing; Reading; Video; Gaming; Social Media; Planning; Productivity; Terminal; Design; Music\n\n- **Full Featured Facet** (`full-featured`)\n A facet for testing all features\n - **Full Featured Facet Entities**: First test entity; Second test entity; Third test entity with description\n - **Full Featured Facet Activities**: Meetings; Coding; Custom Activity; Email; Messaging\n\n- **Minimal Facet** (`minimal-facet`)\n - **Minimal Facet Activities**: Meetings; Coding; Browsing; Email; Messaging; AI Conversation; Writing; Reading; Video; Gaming; Social Media; Planning; Productivity; Terminal; Design; Music\n\n- **Montague Tech** (`montague`)\n Montague Tech startup operations\n - **Tester's Role**: CTO and co-founder of Montague Tech. Visionary full-stack engineer.\n - **Montague Tech Entities**: Balcony App; Balthasar Davi; Benvolio Montague; Friar Lawrence; Juliet Capulet; Mercutio Escalus; Mesh Routing; Montague Tech; Prince Escalus; Rosaline Prince; Schema Bridge; Verona Platform; Verona Ventures\n - **Montague Tech Activities**: Engineering; Meetings; Email; Messaging\n\n- **Priority Test** (`priority-test`)\n - **Priority Test Activities**: Meetings; Coding; Browsing; Email; Messaging; AI Conversation; Writing; Reading; Video; Gaming; Social Media; Planning; Productivity; Terminal; Design; Music\n\n- **Test Facet** (`test-facet`)\n A test facet for validating functionality\n - **Test Facet Entities**: Acme Corp; API Optimization; Bob Wilson; Dashboard Redesign; Docker; Jane Doe; John Smith; PostgreSQL; Tech Solutions Inc; Visual Studio Code\n - **Test Facet Activities**: Meetings; Coding; Browsing; Email; Messaging; AI Conversation; Writing; Reading; Video; Gaming; Social Media; Planning; Productivity; Terminal; Design; Music\n\n- **Verona** (`verona`)\n Cross-company Verona Platform collaboration\n - **Tester's Role**: Co-lead of the Verona Platform joint venture from Montague Tech.\n - **Verona Entities**: Balcony App; Friar Lawrence; Juliet Capulet; Verona Platform\n - **Verona Activities**: Engineering; Meetings; Design Review; Email; Messaging\n\n## Identity Frame\n\nYou are sol, responding to Tester inside the chat backend. You are not the research worker and you do not have tools in this step. Work only from the context already provided to you.\n\n## Current Digest\n\n$digest_contents\n\n$location\n\n$trigger_context\n\n$chat_stream_tail\n\n$active_talents\n\n$active_routines\n\n$routine_suggestion\n\n## Tonal Range\n\nMatch the owner's tone and stakes:\n- Be direct and brief for simple replies.\n- Be warm when the owner is sharing something difficult or personal.\n- Be analytical when the owner needs synthesis or a plan.\n- Be challenging only when there is a clear pattern worth naming.\n\n## Routine Etiquette\n\n- If a routine suggestion appears in context, mention it once and only at the end.\n- Do not raise routine suggestions on machine-driven follow-ups unless the context explicitly includes one.\n- Do not mention internal systems, hooks, or prompt assembly.\n\n## Import And Naming Awareness\n\n- If the owner is asking about imports, naming, or system readiness, answer plainly from the supplied context.\n- Request exec only when answering well requires deeper lookup, synthesis, or tool use.\n\n## When To Dispatch Exec\n\nSet `talent_request` only when the owner needs work that cannot be answered well from the supplied digest, chat history, active routines, and trigger context alone.\n\nDispatch exec for:\n- Journal exploration across days, entities, or transcripts\n- Multi-step synthesis or research\n- Meeting prep that needs fresh participant or activity lookup\n- Any request that clearly needs tool use or external state inspection\n\nDo not dispatch exec for:\n- Simple acknowledgements\n- Straightforward follow-up chat\n- Routine suggestions already supported by the supplied context\n- Brief guidance that can be answered from the current digest and chat tail\n\n## JSON Contract\n\nReturn exactly one JSON object matching `chat.schema.json`.\n\n- `message`: The owner-facing reply. Use `null` only when you genuinely have no safe or useful message to send.\n- `notes`: Brief internal summary of why you responded this way. Keep it factual and concise. Do not dump long reasoning.\n- `talent_request`: `null` unless exec should be dispatched. When dispatching, include:\n - `task`: the exact work exec should perform\n - `context`: optional structured hints that will help exec start fast\n\n## Output Rules\n\n- Return JSON only.\n- `message` should stand on its own without referring to hidden machinery.\n- If `talent_request` is present, the `message` should still be useful to the owner right now.\n- Prefer no dispatch over a weak or redundant dispatch.", "multi_facet": false, - "name": "unified", + "name": "chat", "title": "Chat" } diff --git a/tests/conftest.py b/tests/conftest.py index dc9d723bf..f674b708c 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -14,6 +14,7 @@ ROOT = Path(__file__).resolve().parents[1] if str(ROOT) not in sys.path: sys.path.insert(0, str(ROOT)) +from convey.chat import stop_all_chat_runtime from tests._baseline_harness import copytree_tracked from think.entities.journal import clear_journal_entity_cache from think.entities.loading import clear_entity_loading_cache @@ -74,6 +75,12 @@ def _cleanup_push_runtime(): stop_all_push_runtime() +@pytest.fixture(autouse=True) +def _cleanup_chat_runtime(): + yield + stop_all_chat_runtime() + + @pytest.fixture def journal_copy(tmp_path, monkeypatch): """Copy git-tracked fixture files to tmp_path for mutation tests.""" diff --git a/tests/test_app_sol.py b/tests/test_app_sol.py index cc355b0ff..c08bf768e 100644 --- a/tests/test_app_sol.py +++ b/tests/test_app_sol.py @@ -70,7 +70,7 @@ def app_with_agent(tmp_path, monkeypatch): def test_resolve_agent_path_system_agent(): """Test _resolve_talent_path returns correct path for system agents.""" - agent_dir, agent_name = _resolve_talent_path("unified") + agent_dir, agent_name = _resolve_talent_path("chat") assert agent_name == "chat" assert agent_dir.name == "talent" @@ -96,9 +96,9 @@ def test_resolve_agent_path_app_agent_with_underscores(): def test_get_agent_system_agent(fixture_journal): """Test get_talent loads system agents correctly.""" - config = get_talent("unified") + config = get_talent("chat") - assert config["name"] == "unified" + assert config["name"] == "chat" assert "user_instruction" in config assert len(config["user_instruction"]) > 0 @@ -111,6 +111,12 @@ def test_get_agent_nonexistent_raises(): assert "nonexistent_agent_xyz" in str(exc_info.value) +def test_get_agent_legacy_alias_raises(): + """The legacy chat alias is removed in the chat backend cutover.""" + with pytest.raises(FileNotFoundError): + get_talent("uni" + "fied") + + def test_get_agent_nonexistent_app_agent_raises(): """Test get_talent raises FileNotFoundError for nonexistent app agents.""" with pytest.raises(FileNotFoundError) as exc_info: @@ -205,7 +211,7 @@ class TestResolveOutputPath: """Without output_path, derives from day/name/segment fields.""" event = { "day": "20260214", - "name": "unified", + "name": "chat", "segment": "100", "facet": "health", } @@ -216,13 +222,13 @@ class TestResolveOutputPath: def test_returns_none_without_day_or_output_path(self): """Returns None when neither output_path nor day is present.""" - event = {"name": "unified"} + event = {"name": "chat"} result = _resolve_output_path(event, "/journal") assert result is None def test_empty_output_path_falls_through(self, fixture_journal): """Empty string output_path falls through to derivation.""" - event = {"output_path": "", "day": "20260214", "name": "unified"} + event = {"output_path": "", "day": "20260214", "name": "chat"} result = _resolve_output_path(event, "tests/fixtures/journal") # Empty string is falsy, so falls through to derivation assert result is not None @@ -231,7 +237,7 @@ class TestResolveOutputPath: """SOL_STREAM from env is passed through to get_output_path.""" event = { "day": "20260214", - "name": "unified", + "name": "chat", "env": {"SOL_STREAM": "mystream"}, } result = _resolve_output_path(event, "tests/fixtures/journal") @@ -242,7 +248,7 @@ class TestResolveOutputPath: event = { "output_path": "/custom/path/output.md", "day": "20260214", - "name": "unified", + "name": "chat", "segment": "100", } result = _resolve_output_path(event, "/journal") diff --git a/tests/test_awareness.py b/tests/test_awareness.py index 715e65087..c9c39b08b 100644 --- a/tests/test_awareness.py +++ b/tests/test_awareness.py @@ -224,7 +224,7 @@ class TestComputeThickness: "think.indexer.journal.get_entity_strength", return_value=[] ): with unittest.mock.patch( - "think.conversation.get_recent_exchanges", return_value=[] + "think.awareness._recent_chat_exchanges", return_value=[] ): with unittest.mock.patch( "think.facets.get_enabled_facets", return_value={} @@ -248,7 +248,7 @@ class TestComputeThickness: ] exchanges = [ { - "talent": "triage", + "talent": "chat", "agent_response": f"talked about entity_{i}", "user_message": "hi", } @@ -260,7 +260,7 @@ class TestComputeThickness: "think.indexer.journal.get_entity_strength", return_value=entities ): with unittest.mock.patch( - "think.conversation.get_recent_exchanges", return_value=exchanges + "think.awareness._recent_chat_exchanges", return_value=exchanges ): with unittest.mock.patch( "think.facets.get_enabled_facets", return_value=facets @@ -283,7 +283,7 @@ class TestComputeThickness: ] exchanges = [ { - "talent": "triage", + "talent": "chat", "agent_response": f"entity_{i} is great", "user_message": "yo", } @@ -300,7 +300,7 @@ class TestComputeThickness: "think.indexer.journal.get_entity_strength", return_value=entities ): with unittest.mock.patch( - "think.conversation.get_recent_exchanges", return_value=exchanges + "think.awareness._recent_chat_exchanges", return_value=exchanges ): with unittest.mock.patch( "think.facets.get_enabled_facets", return_value=facets @@ -324,7 +324,7 @@ class TestComputeThickness: {"entity_name": f"entity_{i}", "observation_depth": 3} for i in range(15) ] exchanges = [ - {"talent": "triage", "agent_response": "hello there", "user_message": "hi"} + {"talent": "chat", "agent_response": "hello there", "user_message": "hi"} for _ in range(10) ] facets = {"work": {}, "personal": {}, "hobby": {}} @@ -333,7 +333,7 @@ class TestComputeThickness: "think.indexer.journal.get_entity_strength", return_value=entities ): with unittest.mock.patch( - "think.conversation.get_recent_exchanges", return_value=exchanges + "think.awareness._recent_chat_exchanges", return_value=exchanges ): with unittest.mock.patch( "think.facets.get_enabled_facets", return_value=facets @@ -363,7 +363,7 @@ class TestComputeThickness: "user_message": "hello", }, { - "talent": "triage", + "talent": "chat", "agent_response": "foo is great", "user_message": "hey", }, @@ -373,7 +373,7 @@ class TestComputeThickness: "think.indexer.journal.get_entity_strength", return_value=entities ): with unittest.mock.patch( - "think.conversation.get_recent_exchanges", return_value=exchanges + "think.awareness._recent_chat_exchanges", return_value=exchanges ): with unittest.mock.patch( "think.facets.get_enabled_facets", @@ -394,7 +394,7 @@ class TestComputeThickness: side_effect=Exception("db error"), ): with unittest.mock.patch( - "think.conversation.get_recent_exchanges", + "think.awareness._recent_chat_exchanges", side_effect=Exception("no file"), ): with unittest.mock.patch( @@ -421,7 +421,7 @@ class TestComputeThickness: "think.indexer.journal.get_entity_strength", return_value=[] ): with unittest.mock.patch( - "think.conversation.get_recent_exchanges", return_value=[] + "think.awareness._recent_chat_exchanges", return_value=[] ): with unittest.mock.patch( "think.facets.get_enabled_facets", return_value={} diff --git a/tests/test_chat_context.py b/tests/test_chat_context.py index 076f90a93..327d4a72e 100644 --- a/tests/test_chat_context.py +++ b/tests/test_chat_context.py @@ -302,7 +302,7 @@ def test_chat_context_enrichment_errors_are_graceful(monkeypatch, tmp_path): assert "/app/home" in template_vars["location"] -def test_chat_context_drops_conversation_memory_imports(monkeypatch): +def test_chat_context_drops_legacy_memory_imports(monkeypatch): monkeypatch.setattr("think.routines.get_routine_state", lambda: []) monkeypatch.setattr( "think.routines.get_config", @@ -310,13 +310,15 @@ def test_chat_context_drops_conversation_memory_imports(monkeypatch): ) monkeypatch.setattr("think.routines.save_config", lambda config: None) + legacy_module = "think" + ".con" + "versation" + legacy_memory = "conversation_" + "memory" source = ( Path(__file__).resolve().parents[1] / "talent" / "chat_context.py" ).read_text(encoding="utf-8") - assert "think.conversation" not in source - assert "conversation_memory" not in source + assert legacy_module not in source + assert legacy_memory not in source - sys.modules.pop("think.conversation", None) + sys.modules.pop(legacy_module, None) _load_chat_context_module() - assert "think.conversation" not in sys.modules + assert legacy_module not in sys.modules diff --git a/tests/test_chat_runtime.py b/tests/test_chat_runtime.py new file mode 100644 index 000000000..bac15a5eb --- /dev/null +++ b/tests/test_chat_runtime.py @@ -0,0 +1,298 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +from flask import Flask + +from convey.chat_stream import append_chat_event, read_chat_events + + +def _reset_chat_state(chat_module) -> None: + chat_module.stop_all_chat_runtime() + with chat_module._state_lock: + chat_module._current_chat_use_id = None + chat_module._current_chat_state = None + chat_module._queued_trigger = None + chat_module._active_execs.clear() + chat_module._recovery_day = None + chat_module._last_use_id = 0 + + +def _setup_journal(tmp_path, monkeypatch): + journal = tmp_path / "journal" + journal.mkdir() + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(journal)) + return journal + + +def test_chat_result_with_two_active_execs_retriggers_with_max_active_reason( + tmp_path, monkeypatch +): + import convey.chat as chat + + _setup_journal(tmp_path, monkeypatch) + _reset_chat_state(chat) + + append_chat_event( + "talent_spawned", + use_id="1713620000001", + name="exec", + task="first task", + started_at=1713620000001, + ) + append_chat_event( + "talent_spawned", + use_id="1713620000002", + name="exec", + task="second task", + started_at=1713620000002, + ) + + actions: list[dict] = [] + monkeypatch.setattr( + "convey.chat._run_next_action", lambda action: actions.append(action) + ) + monkeypatch.setattr("convey.chat._emit_finish", lambda *args, **kwargs: None) + monkeypatch.setattr("convey.chat._emit_error", lambda *args, **kwargs: None) + + with chat._state_lock: + chat._current_chat_use_id = "1713620000100" + chat._current_chat_state = { + "raw_use_id": "1713620000101", + "trigger": {"type": "owner_message", "message": "help"}, + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "retry_count": 0, + } + + chat._on_cortex_finish( + { + "use_id": "1713620000101", + "result": ( + '{"message":"I am looking into that.","notes":"need exec",' + '"talent_request":{"task":"research it","context":{"k":"v"}}}' + ), + } + ) + + assert actions + assert actions[-1]["kind"] == "chat" + assert actions[-1]["trigger"] == { + "type": "synthetic-max-active", + "reason": "max active — waiting for one to finish", + } + + sol_messages = [ + e for e in read_chat_events(chat._today_day()) if e["kind"] == "sol_message" + ] + assert sol_messages[-1]["requested_exec"] is True + assert sol_messages[-1]["requested_task"] == "research it" + + +def test_exec_retrigger_loop_stops_after_three_without_owner_reset( + tmp_path, monkeypatch +): + import convey.chat as chat + + _setup_journal(tmp_path, monkeypatch) + _reset_chat_state(chat) + + append_chat_event( + "owner_message", + text="dig deeper", + app="sol", + path="/app/sol", + facet="work", + ) + for index in range(3): + append_chat_event( + "talent_finished", + use_id=f"171362100000{index}", + name="exec", + summary=f"summary {index}", + ) + if index < 2: + append_chat_event( + "sol_message", + use_id="1713621999999", + text=f"follow up {index}", + notes="retrying", + requested_exec=True, + requested_task=f"task {index}", + ) + + emitted_errors: list[tuple[str, str]] = [] + actions: list[dict | None] = [] + monkeypatch.setattr( + "convey.chat._run_next_action", lambda action: actions.append(action) + ) + monkeypatch.setattr("convey.chat._emit_finish", lambda *args, **kwargs: None) + monkeypatch.setattr( + "convey.chat._emit_error", + lambda use_id, reason: emitted_errors.append((use_id, reason)), + ) + + with chat._state_lock: + chat._current_chat_use_id = "1713621999999" + chat._current_chat_state = { + "raw_use_id": "1713622000000", + "trigger": {"type": "talent_finished", "summary": "summary 2"}, + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "retry_count": 0, + } + + chat._on_cortex_finish( + { + "use_id": "1713622000000", + "result": ( + '{"message":"Still digging.","notes":"loop",' + '"talent_request":{"task":"one more pass","context":{}}}' + ), + } + ) + + assert emitted_errors == [("1713621999999", "chat had trouble — try again")] + assert actions == [None] + errors = [ + e for e in read_chat_events(chat._today_day()) if e["kind"] == "chat_error" + ] + assert errors[-1]["reason"] == "chat had trouble — try again" + + +def test_cortex_finish_and_error_append_exec_terminal_events_by_use_id( + tmp_path, monkeypatch +): + import convey.chat as chat + + _setup_journal(tmp_path, monkeypatch) + _reset_chat_state(chat) + + actions: list[dict] = [] + monkeypatch.setattr( + "convey.chat._run_next_action", lambda action: actions.append(action) + ) + monkeypatch.setattr("convey.chat._emit_finish", lambda *args, **kwargs: None) + monkeypatch.setattr("convey.chat._emit_error", lambda *args, **kwargs: None) + + with chat._state_lock: + chat._current_chat_use_id = "1713623000000" + chat._current_chat_state = { + "raw_use_id": None, + "trigger": {"type": "owner_message", "message": "help"}, + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "retry_count": 0, + } + chat._active_execs["1713623000001"] = { + "chat_use_id": "1713623000000", + "task": "summarize", + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + } + + chat._on_cortex_finish({"use_id": "1713623000001", "result": "done"}) + finished_events = [ + e for e in read_chat_events(chat._today_day()) if e["kind"] == "talent_finished" + ] + assert finished_events[-1]["use_id"] == "1713623000001" + assert actions[-1]["trigger"]["type"] == "talent_finished" + + _reset_chat_state(chat) + actions.clear() + with chat._state_lock: + chat._current_chat_use_id = "1713624000000" + chat._current_chat_state = { + "raw_use_id": None, + "trigger": {"type": "owner_message", "message": "help"}, + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "retry_count": 0, + } + chat._active_execs["1713624000001"] = { + "chat_use_id": "1713624000000", + "task": "summarize", + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + } + + chat._on_cortex_error({"use_id": "1713624000001", "error": "boom"}) + errored_events = [ + e for e in read_chat_events(chat._today_day()) if e["kind"] == "talent_errored" + ] + assert errored_events[-1]["use_id"] == "1713624000001" + assert actions[-1]["trigger"]["type"] == "talent_errored" + assert actions[-1]["trigger"]["reason"] == "boom" + + +def test_start_chat_runtime_recovers_exactly_one_unresponded_trigger( + tmp_path, monkeypatch +): + import convey.chat as chat + + _setup_journal(tmp_path, monkeypatch) + _reset_chat_state(chat) + + append_chat_event( + "owner_message", + text="recover me", + app="sol", + path="/app/sol", + facet="work", + ) + + starts: list[dict] = [] + monkeypatch.setattr( + "convey.chat.CallosumConnection.start", lambda self, callback=None: None + ) + monkeypatch.setattr("convey.chat.CallosumConnection.stop", lambda self: None) + monkeypatch.setattr( + "convey.chat._spawn_chat_generate", lambda action: starts.append(action) or True + ) + + app = Flask(__name__) + chat.start_chat_runtime(app) + chat.start_chat_runtime(app) + + assert len(starts) == 1 + + +def test_chat_generate_schema_violation_retries_once_then_chat_errors( + tmp_path, monkeypatch +): + import convey.chat as chat + + _setup_journal(tmp_path, monkeypatch) + _reset_chat_state(chat) + + actions: list[dict | None] = [] + emitted_errors: list[tuple[str, str]] = [] + monkeypatch.setattr( + "convey.chat._run_next_action", lambda action: actions.append(action) + ) + monkeypatch.setattr("convey.chat._emit_finish", lambda *args, **kwargs: None) + monkeypatch.setattr( + "convey.chat._emit_error", + lambda use_id, reason: emitted_errors.append((use_id, reason)), + ) + + with chat._state_lock: + chat._current_chat_use_id = "1713625000000" + chat._current_chat_state = { + "raw_use_id": "1713625000001", + "trigger": {"type": "owner_message", "message": "help"}, + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "retry_count": 0, + } + + chat._on_cortex_finish({"use_id": "1713625000001", "result": "not json"}) + + assert actions and actions[-1]["kind"] == "chat" + assert actions[-1]["logical_use_id"] == "1713625000000" + assert emitted_errors == [] + + with chat._state_lock: + retry_use_id = chat._current_chat_state["raw_use_id"] + + chat._on_cortex_finish({"use_id": retry_use_id, "result": "still not json"}) + + assert emitted_errors == [("1713625000000", "chat had trouble — try again")] + errors = [ + e for e in read_chat_events(chat._today_day()) if e["kind"] == "chat_error" + ] + assert errors[-1]["use_id"] == "1713625000000" diff --git a/tests/test_conversation.py b/tests/test_conversation.py deleted file mode 100644 index 05455806b..000000000 --- a/tests/test_conversation.py +++ /dev/null @@ -1,434 +0,0 @@ -# SPDX-License-Identifier: AGPL-3.0-only -# Copyright (c) 2026 sol pbc - -"""Tests for think.conversation module — conversation memory service.""" - -import json -from datetime import datetime -from unittest import mock - -import pytest - - -@pytest.fixture -def journal_dir(tmp_path, monkeypatch): - """Create a temporary journal directory.""" - journal = tmp_path / "journal" - journal.mkdir() - monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(journal)) - with mock.patch("think.conversation.get_journal", return_value=str(journal)): - yield journal - - -# --------------------------------------------------------------------------- -# record_exchange -# --------------------------------------------------------------------------- - - -def test_record_exchange_writes_jsonl(journal_dir): - """Exchange is appended to conversation/exchanges.jsonl.""" - from think.conversation import record_exchange - - record_exchange( - ts=1710000000000, - facet="work", - app="entities", - path="/app/entities/adrian", - user_message="what's our history with adrian?", - agent_response="You met Adrian at betaworks.", - talent="unified", - use_id="12345", - ) - - jsonl_path = journal_dir / "conversation" / "exchanges.jsonl" - assert jsonl_path.exists() - - with open(jsonl_path) as f: - lines = [json.loads(line) for line in f if line.strip()] - - assert len(lines) == 1 - ex = lines[0] - assert ex["ts"] == 1710000000000 - assert ex["facet"] == "work" - assert ex["app"] == "entities" - assert ex["user_message"] == "what's our history with adrian?" - assert ex["agent_response"] == "You met Adrian at betaworks." - assert ex["talent"] == "unified" - assert ex["use_id"] == "12345" - - -def test_record_exchange_writes_journal_segment(journal_dir): - """Exchange creates a journal segment markdown file for search indexing.""" - from think.conversation import record_exchange - - # Use a known timestamp: 2026-03-15 14:30:00 UTC - ts = int(datetime(2026, 3, 15, 14, 30, 0).timestamp() * 1000) - - record_exchange( - ts=ts, - facet="work", - app="activities", - path="/app/activities", - user_message="move my 3pm to 4pm", - agent_response="Done — moved 'DVD sync' to 4pm.", - talent="unified", - use_id="67890", - ) - - # Check journal segment directory: YYYYMMDD/conversation/HHMMSS_1/talents/ - day = datetime.fromtimestamp(ts / 1000).strftime("%Y%m%d") - time_key = datetime.fromtimestamp(ts / 1000).strftime("%H%M%S") - md_path = ( - journal_dir - / "chronicle" - / day - / "conversation" - / f"{time_key}_1" - / "talents" - / "conversation.md" - ) - - assert md_path.exists() - content = md_path.read_text() - assert "move my 3pm to 4pm" in content - assert "Done — moved 'DVD sync' to 4pm." in content - assert "**Facet:** work" in content - assert "activities" in content - - -def test_record_exchange_appends_multiple(journal_dir): - """Multiple exchanges append to the same JSONL file.""" - from think.conversation import record_exchange - - record_exchange( - ts=1710000001000, - user_message="hello", - agent_response="hi there", - talent="triage", - ) - record_exchange( - ts=1710000002000, - user_message="what time is it?", - agent_response="It's 2pm.", - talent="triage", - ) - - jsonl_path = journal_dir / "conversation" / "exchanges.jsonl" - with open(jsonl_path) as f: - lines = [json.loads(line) for line in f if line.strip()] - - assert len(lines) == 2 - assert lines[0]["user_message"] == "hello" - assert lines[1]["user_message"] == "what time is it?" - - -def test_record_exchange_skips_empty(journal_dir): - """Empty user_message or agent_response is silently skipped.""" - from think.conversation import record_exchange - - record_exchange(user_message="", agent_response="response", talent="triage") - record_exchange(user_message="hello", agent_response="", talent="triage") - - jsonl_path = journal_dir / "conversation" / "exchanges.jsonl" - assert not jsonl_path.exists() - - -# --------------------------------------------------------------------------- -# get_recent_exchanges -# --------------------------------------------------------------------------- - - -def test_get_recent_exchanges_empty(journal_dir): - """Returns empty list when no exchanges exist.""" - from think.conversation import get_recent_exchanges - - assert get_recent_exchanges() == [] - - -def test_get_recent_exchanges_returns_last_n(journal_dir): - """Returns the last N exchanges.""" - from think.conversation import get_recent_exchanges, record_exchange - - for i in range(15): - record_exchange( - ts=1710000000000 + i * 1000, - user_message=f"msg {i}", - agent_response=f"resp {i}", - talent="triage", - ) - - recent = get_recent_exchanges(limit=5) - assert len(recent) == 5 - assert recent[0]["user_message"] == "msg 10" - assert recent[-1]["user_message"] == "msg 14" - - -def test_get_recent_exchanges_filters_by_facet(journal_dir): - """Facet filter returns only matching exchanges.""" - from think.conversation import get_recent_exchanges, record_exchange - - record_exchange( - ts=1710000001000, - facet="work", - user_message="work question", - agent_response="work answer", - talent="triage", - ) - record_exchange( - ts=1710000002000, - facet="personal", - user_message="personal question", - agent_response="personal answer", - talent="triage", - ) - - work = get_recent_exchanges(facet="work") - assert len(work) == 1 - assert work[0]["facet"] == "work" - - personal = get_recent_exchanges(facet="personal") - assert len(personal) == 1 - assert personal[0]["facet"] == "personal" - - -# --------------------------------------------------------------------------- -# get_today_exchanges -# --------------------------------------------------------------------------- - - -def test_get_today_exchanges_filters_by_day(journal_dir): - """Only returns exchanges from today.""" - from think.conversation import get_today_exchanges, record_exchange - from think.utils import now_ms - - # Record an exchange with current timestamp (today) - record_exchange( - ts=now_ms(), - user_message="today question", - agent_response="today answer", - talent="triage", - ) - - # Record an exchange with old timestamp (not today) - record_exchange( - ts=1000000000000, # 2001-09-08 - user_message="old question", - agent_response="old answer", - talent="triage", - ) - - today = get_today_exchanges() - assert len(today) == 1 - assert today[0]["user_message"] == "today question" - - -# --------------------------------------------------------------------------- -# build_memory_context -# --------------------------------------------------------------------------- - - -def test_build_memory_context_empty(journal_dir): - """Returns empty string when no exchanges exist.""" - from think.conversation import build_memory_context - - assert build_memory_context() == "" - - -def test_build_memory_context_includes_recent(journal_dir): - """Context includes recent exchanges.""" - from think.conversation import build_memory_context, record_exchange - from think.utils import now_ms - - ts = now_ms() - record_exchange( - ts=ts, - facet="work", - app="entities", - user_message="who is adrian?", - agent_response="Adrian is the CTO of Own Company.", - talent="unified", - ) - - context = build_memory_context() - assert "who is adrian?" in context - assert "Adrian is the CTO" in context - assert "Recent Conversations" in context - - -def test_build_memory_context_truncates_long_responses(journal_dir): - """Long agent responses are truncated in context output.""" - from think.conversation import ( - build_memory_context, - record_exchange, - ) - from think.utils import now_ms - - long_response = "x" * 500 - record_exchange( - ts=now_ms(), - user_message="tell me a story", - agent_response=long_response, - talent="unified", - ) - - context = build_memory_context() - # Response should be truncated - assert "..." in context - assert long_response not in context - - -def test_build_memory_context_earlier_today(journal_dir): - """When more exchanges exist today than the recent limit, earlier ones are compact.""" - from think.conversation import build_memory_context, record_exchange - from think.utils import now_ms - - ts = now_ms() - # Record 15 exchanges "today" - for i in range(15): - record_exchange( - ts=ts + i * 1000, - user_message=f"question {i}", - agent_response=f"answer {i}", - talent="unified", - ) - - context = build_memory_context(recent_limit=10) - assert "Earlier Today" in context - assert "Recent Conversations" in context - - -# --------------------------------------------------------------------------- -# inject_memory -# --------------------------------------------------------------------------- - - -def test_inject_memory_replaces_marker(): - """Injection point is replaced with memory context.""" - from think.conversation import inject_memory - - instruction = """## Before - -## Conversation Memory - - - -## After""" - - result = inject_memory(instruction, "### Recent\nHello world") - assert "CONVERSATION_MEMORY_INJECTION_POINT" not in result - assert "### Recent\nHello world" in result - assert "## Before" in result - assert "## After" in result - - -def test_inject_memory_no_marker(): - """If no marker present, instruction is returned unchanged.""" - from think.conversation import inject_memory - - instruction = "No marker here." - result = inject_memory(instruction, "some context") - assert result == instruction - - -def test_inject_memory_empty_context(): - """Empty context gets a placeholder message.""" - from think.conversation import inject_memory - - instruction = "" - result = inject_memory(instruction, "") - assert "No conversation history yet." in result - - -# --------------------------------------------------------------------------- -# _format_exchange -# --------------------------------------------------------------------------- - - -def test_format_exchange_full(): - """Full format includes user message and truncated response.""" - from think.conversation import _format_exchange - - ex = { - "ts": 1710000000000, - "app": "entities", - "facet": "work", - "user_message": "who is adrian?", - "agent_response": "Adrian is the CTO.", - } - - result = _format_exchange(ex, compact=False) - assert "User: who is adrian?" in result - assert "Sol: Adrian is the CTO." in result - assert "entities" in result - assert "work" in result - - -def test_format_exchange_compact(): - """Compact format is a one-liner.""" - from think.conversation import _format_exchange - - ex = { - "ts": 1710000000000, - "app": "activities", - "facet": "work", - "user_message": "what's on my schedule today?", - "agent_response": "You have 3 meetings.", - } - - result = _format_exchange(ex, compact=True) - assert result.startswith("- [") - assert "what's on my schedule today?" in result - assert "You have 3 meetings" not in result # Compact omits response - - -# --------------------------------------------------------------------------- -# Pre-hook integration -# --------------------------------------------------------------------------- - - -def test_conversation_memory_pre_hook(journal_dir): - """Pre-hook injects memory into user instruction.""" - from talent.conversation_memory import pre_process - from think.conversation import record_exchange - from think.utils import now_ms - - # Record an exchange first - record_exchange( - ts=now_ms(), - facet="work", - user_message="hello", - agent_response="hi there!", - talent="unified", - ) - - context = { - "user_instruction": """Some instructions. - -## Conversation Memory - - - -## Other section""", - "facet": "work", - } - - result = pre_process(context) - assert result is not None - assert "user_instruction" in result - assert "CONVERSATION_MEMORY_INJECTION_POINT" not in result["user_instruction"] - assert "hello" in result["user_instruction"] - assert "hi there!" in result["user_instruction"] - - -def test_conversation_memory_pre_hook_no_marker(): - """Pre-hook returns None when no injection marker present.""" - from talent.conversation_memory import pre_process - - context = {"user_instruction": "No marker here."} - result = pre_process(context) - assert result is None diff --git a/tests/test_convey_apps.py b/tests/test_convey_apps.py index aa18a76fc..28f4a0fae 100644 --- a/tests/test_convey_apps.py +++ b/tests/test_convey_apps.py @@ -3,10 +3,7 @@ """Tests for convey app placeholder and attention behavior.""" -from unittest.mock import patch - import pytest -from flask import Flask @pytest.fixture(autouse=True) @@ -15,26 +12,6 @@ def _temp_journal(monkeypatch, tmp_path): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) -def _run_triage(): - """Run the triage endpoint with mocked state.""" - app = Flask(__name__) - with ( - patch("convey.utils.spawn_agent", return_value="agent-1") as mock_spawn, - patch("think.cortex_client.wait_for_uses", return_value=({}, [])), - patch( - "think.cortex_client.read_use_events", - return_value=[{"event": "finish", "result": "ok"}], - ), - ): - from convey.triage import triage - - with app.test_request_context("/", method="POST", json={"message": "hello"}): - response = triage() - - assert response.status_code == 200 - return mock_spawn - - # --- Placeholder resolution --- @@ -310,13 +287,3 @@ class TestAttentionResolution: assert "2" in result.placeholder_text assert "report" in result.placeholder_text.lower() assert len(result.placeholder_text) <= 90 - - -class TestTriageSystemHealth: - """Tests for system health context injection in triage.""" - - def test_triage_no_health_context_when_healthy(self): - """No system health context when nothing needs attention.""" - mock = _run_triage() - prompt = mock.call_args.kwargs["prompt"] - assert "System health" not in prompt diff --git a/tests/test_convey_chat.py b/tests/test_convey_chat.py new file mode 100644 index 000000000..03455cfb2 --- /dev/null +++ b/tests/test_convey_chat.py @@ -0,0 +1,143 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +from datetime import datetime + +import pytest +from flask import Flask + +from convey.chat import chat_bp +from convey.chat_stream import append_chat_event + + +def _setup_journal(tmp_path, monkeypatch): + journal = tmp_path / "journal" + journal.mkdir() + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(journal)) + return journal + + +def _reset_chat_state(chat_module) -> None: + chat_module.stop_all_chat_runtime() + with chat_module._state_lock: + chat_module._current_chat_use_id = None + chat_module._current_chat_state = None + chat_module._queued_trigger = None + chat_module._active_execs.clear() + chat_module._recovery_day = None + chat_module._last_use_id = 0 + + +def _ms(year: int, month: int, day: int, hour: int, minute: int, second: int) -> int: + return int(datetime(year, month, day, hour, minute, second).timestamp() * 1000) + + +@pytest.fixture +def chat_client(tmp_path, monkeypatch): + import convey.chat as chat + + _setup_journal(tmp_path, monkeypatch) + _reset_chat_state(chat) + + app = Flask(__name__) + app.config["TESTING"] = True + app.register_blueprint(chat_bp) + return app.test_client() + + +def test_post_chat_appends_owner_message_and_returns_reserved_use_id( + chat_client, monkeypatch +): + starts: list[dict] = [] + monkeypatch.setattr("think.identity.ensure_identity_directory", lambda: None) + monkeypatch.setattr( + "convey.chat._spawn_chat_generate", lambda action: starts.append(action) or True + ) + + response = chat_client.post( + "/api/chat", + json={ + "message": "hello there", + "app": "sol", + "path": "/app/sol", + "facet": "work", + }, + ) + + assert response.status_code == 200 + payload = response.get_json() + assert payload["queued"] is False + assert payload["use_id"].isdigit() + assert starts and starts[-1]["logical_use_id"] == payload["use_id"] + + +def test_session_endpoint_reduces_from_chat_stream(chat_client): + day = "20260420" + append_chat_event( + "sol_message", + ts=_ms(2026, 4, 20, 12, 0, 0), + use_id="1713626000000", + text="hello", + notes="ready", + requested_exec=False, + requested_task=None, + ) + append_chat_event( + "talent_spawned", + ts=_ms(2026, 4, 20, 12, 1, 0), + use_id="1713626000001", + name="exec", + task="research", + started_at=1713626000001, + ) + + response = chat_client.get("/api/chat/session") + assert response.status_code == 200 + payload = response.get_json() + assert payload["latest_sol_message"]["text"] == "hello" + assert payload["active_talents"][0]["task"] == "research" + assert chat_client.get(f"/api/chat/stream/{day}").status_code == 200 + + +def test_stream_endpoint_ordered_with_limit(chat_client): + start = _ms(2026, 4, 20, 12, 0, 0) + for index in range(4): + append_chat_event( + "owner_message", + ts=start + (index * 300_000), + text=f"m{index}", + app="sol", + path="/app/sol", + facet="work", + ) + + response = chat_client.get("/api/chat/stream/20260420?limit=2") + assert response.status_code == 200 + payload = response.get_json() + assert [event["text"] for event in payload["events"]] == ["m2", "m3"] + + +def test_result_endpoint_reads_stream_not_talent_log(chat_client, tmp_path): + use_id = str(_ms(2026, 4, 20, 12, 0, 0)) + append_chat_event( + "sol_message", + use_id=use_id, + text="stream reply", + notes="done", + requested_exec=False, + requested_task=None, + ) + + talents_dir = tmp_path / "journal" / "talents" / "chat" + talents_dir.mkdir(parents=True, exist_ok=True) + (talents_dir / f"{use_id}.jsonl").write_text( + '{"event":"finish","result":"log reply"}\n' + ) + + response = chat_client.get(f"/api/chat/result/{use_id}") + assert response.status_code == 200 + payload = response.get_json() + assert payload["state"] == "finished" + assert payload["summary"] == "stream reply" diff --git a/tests/test_cortex.py b/tests/test_cortex.py index 4d0fcf133..9f3a05c6d 100644 --- a/tests/test_cortex.py +++ b/tests/test_cortex.py @@ -114,7 +114,7 @@ def test_spawn_subprocess( "ts": 123456789, "prompt": "Test prompt", "provider": "openai", - "name": "unified", + "name": "chat", "model": GPT_5, } @@ -141,7 +141,7 @@ def test_spawn_subprocess( assert ndjson["event"] == "request" assert ndjson["prompt"] == "Test prompt" assert ndjson["provider"] == "openai" - assert ndjson["name"] == "unified" + assert ndjson["name"] == "chat" assert ndjson["model"] == GPT_5 # Check stdin was closed @@ -266,7 +266,7 @@ def test_spawn_subprocess_uses_cwd_from_talent( "ts": 24680, "prompt": "Test prompt", "provider": "openai", - "name": "unified", + "name": "chat", "model": GPT_5, } @@ -481,11 +481,11 @@ def test_has_finish_event(cortex_service, mock_journal): def test_complete_use_file(cortex_service, mock_journal): """Test completing an agent file (rename from active to completed).""" use_id = "123456789" - unified_dir = mock_journal / "talents" / "unified" + unified_dir = mock_journal / "talents" / "chat" unified_dir.mkdir() active_path = unified_dir / f"{use_id}_active.jsonl" active_path.touch() - cortex_service.use_requests[use_id] = {"name": "unified", "use_id": use_id} + cortex_service.use_requests[use_id] = {"name": "chat", "use_id": use_id} cortex_service._complete_use_file(use_id, active_path) @@ -493,33 +493,33 @@ def test_complete_use_file(cortex_service, mock_journal): assert not active_path.exists() completed_path = unified_dir / f"{use_id}.jsonl" assert completed_path.exists() - symlink_path = mock_journal / "talents" / "unified.log" + symlink_path = mock_journal / "talents" / "chat.log" assert symlink_path.is_symlink() - assert os.readlink(symlink_path) == f"unified/{use_id}.jsonl" + assert os.readlink(symlink_path) == f"chat/{use_id}.jsonl" def test_complete_use_file_replaces_symlink(cortex_service, mock_journal): """Test completing agent file replaces convenience symlink for same name.""" - unified_dir = mock_journal / "talents" / "unified" + unified_dir = mock_journal / "talents" / "chat" unified_dir.mkdir() first_agent_id = "111" first_active_path = unified_dir / f"{first_agent_id}_active.jsonl" first_active_path.touch() - cortex_service.use_requests[first_agent_id] = {"name": "unified"} + cortex_service.use_requests[first_agent_id] = {"name": "chat"} cortex_service._complete_use_file(first_agent_id, first_active_path) second_agent_id = "222" second_active_path = unified_dir / f"{second_agent_id}_active.jsonl" second_active_path.touch() - cortex_service.use_requests[second_agent_id] = {"name": "unified"} + cortex_service.use_requests[second_agent_id] = {"name": "chat"} cortex_service._complete_use_file(second_agent_id, second_active_path) - symlink_path = mock_journal / "talents" / "unified.log" + symlink_path = mock_journal / "talents" / "chat.log" assert symlink_path.is_symlink() - assert os.readlink(symlink_path) == f"unified/{second_agent_id}.jsonl" + assert os.readlink(symlink_path) == f"chat/{second_agent_id}.jsonl" def test_complete_use_file_colon_name(cortex_service, mock_journal): @@ -765,7 +765,7 @@ def test_recover_orphaned_uses(cortex_service, mock_journal): """Test recovery of orphaned active agent files.""" # Create orphaned active files talents_dir = mock_journal / "talents" - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() agent1_active = unified_dir / "111_active.jsonl" agent2_active = unified_dir / "222_active.jsonl" diff --git a/tests/test_cortex_client.py b/tests/test_cortex_client.py index 975f61b36..717393d16 100644 --- a/tests/test_cortex_client.py +++ b/tests/test_cortex_client.py @@ -86,7 +86,7 @@ def test_cortex_request_broadcasts_to_callosum(callosum_listener): # Create a request use_id = cortex_request( prompt="Test prompt", - name="unified", + name="chat", provider="openai", config={"model": GPT_5}, ) @@ -99,7 +99,7 @@ def test_cortex_request_broadcasts_to_callosum(callosum_listener): assert msg["tract"] == "cortex" assert msg["event"] == "request" assert msg["prompt"] == "Test prompt" - assert msg["name"] == "unified" + assert msg["name"] == "chat" assert msg["provider"] == "openai" assert msg["model"] == GPT_5 assert msg["use_id"] == use_id @@ -110,7 +110,7 @@ def test_cortex_request_returns_agent_id(callosum_server): """Test that cortex_request returns use_id string.""" _ = callosum_server # Needed for side effects only - use_id = cortex_request(prompt="Test", name="unified", provider="openai") + use_id = cortex_request(prompt="Test", name="chat", provider="openai") # Verify use_id is a string timestamp assert isinstance(use_id, str) @@ -118,13 +118,29 @@ def test_cortex_request_returns_agent_id(callosum_server): assert len(use_id) == 13 # Millisecond timestamp +def test_cortex_request_uses_explicit_use_id(callosum_listener): + messages = callosum_listener + + use_id = cortex_request( + prompt="Test prompt", + name="chat", + provider="openai", + use_id="1713629000000", + ) + + time.sleep(0.2) + + assert use_id == "1713629000000" + assert messages[-1]["use_id"] == "1713629000000" + + def test_cortex_request_unique_agent_ids(callosum_server): """Test that cortex_request generates unique agent IDs.""" _ = callosum_server # Needed for side effects only agent_ids = [] for i in range(3): - use_id = cortex_request(prompt=f"Test {i}", name="unified", provider="openai") + use_id = cortex_request(prompt=f"Test {i}", name="chat", provider="openai") agent_ids.append(use_id) time.sleep(0.002) @@ -136,7 +152,7 @@ def test_cortex_request_returns_none_on_send_failure(callosum_server, monkeypatc """Test cortex_request returns None when callosum_send fails.""" monkeypatch.setattr("think.cortex_client.callosum_send", lambda *a, **kw: False) - use_id = cortex_request(prompt="Test", name="unified", provider="openai") + use_id = cortex_request(prompt="Test", name="chat", provider="openai") assert use_id is None @@ -146,7 +162,7 @@ def test_cortex_request_empty_journal(tmp_path, monkeypatch): monkeypatch.setattr("think.cortex_client.callosum_send", lambda *a, **kw: True) monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) - use_id = cortex_request("test", "unified", "openai") + use_id = cortex_request("test", "chat", "openai") assert use_id is not None assert len(use_id) > 0 @@ -177,7 +193,7 @@ def test_cortex_agents_with_active(tmp_path, monkeypatch): ts1 = now_ms() ts2 = ts1 + 1000 - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" tester_dir = talents_dir / "tester" unified_dir.mkdir() tester_dir.mkdir() @@ -189,7 +205,7 @@ def test_cortex_agents_with_active(tmp_path, monkeypatch): "event": "request", "ts": ts1, "prompt": "Task 1", - "name": "unified", + "name": "chat", "provider": "openai", }, f, @@ -260,7 +276,7 @@ def test_cortex_agents_pagination(tmp_path, monkeypatch): # Create multiple agents base_ts = now_ms() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() for i in range(5): ts = base_ts + (i * 1000) @@ -271,7 +287,7 @@ def test_cortex_agents_pagination(tmp_path, monkeypatch): "event": "request", "ts": ts, "prompt": f"Task {i}", - "name": "unified", + "name": "chat", }, f, ) @@ -300,7 +316,7 @@ def test_get_agent_log_status_completed(tmp_path, monkeypatch): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) talents_dir = tmp_path / "talents" talents_dir.mkdir() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() use_id = "1234567890123" @@ -314,7 +330,7 @@ def test_get_agent_log_status_running(tmp_path, monkeypatch): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) talents_dir = tmp_path / "talents" talents_dir.mkdir() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() use_id = "1234567890123" @@ -336,7 +352,7 @@ def test_get_agent_log_status_prefers_completed(tmp_path, monkeypatch): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) talents_dir = tmp_path / "talents" talents_dir.mkdir() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() # Edge case: both files exist (shouldn't happen, but check precedence) @@ -352,7 +368,7 @@ def test_get_agent_end_state_finish(tmp_path, monkeypatch): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) talents_dir = tmp_path / "talents" talents_dir.mkdir() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() use_id = "1234567890123" @@ -369,7 +385,7 @@ def test_get_agent_end_state_error(tmp_path, monkeypatch): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) talents_dir = tmp_path / "talents" talents_dir.mkdir() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() use_id = "1234567890123" @@ -386,7 +402,7 @@ def test_get_agent_end_state_running(tmp_path, monkeypatch): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) talents_dir = tmp_path / "talents" talents_dir.mkdir() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() use_id = "1234567890123" @@ -413,7 +429,7 @@ def test_wait_for_agents_already_complete(tmp_path, monkeypatch): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) talents_dir = tmp_path / "talents" talents_dir.mkdir() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() (tmp_path / "health").mkdir() @@ -433,7 +449,7 @@ def test_wait_for_agents_event_completion(callosum_server): """Test wait_for_uses completes when finish event is received.""" tmp_path = callosum_server talents_dir = tmp_path / "talents" - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir(exist_ok=True) use_id = "1234567890123" @@ -471,7 +487,7 @@ def test_wait_for_agents_error_event(callosum_server): """Test wait_for_uses completes on error event too.""" tmp_path = callosum_server talents_dir = tmp_path / "talents" - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir(exist_ok=True) use_id = "1234567890124" @@ -506,7 +522,7 @@ def test_wait_for_agents_initial_file_check(tmp_path, monkeypatch): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) talents_dir = tmp_path / "talents" talents_dir.mkdir() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() (tmp_path / "health").mkdir() @@ -527,7 +543,7 @@ def test_wait_for_agents_timeout_actual(tmp_path, monkeypatch): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) talents_dir = tmp_path / "talents" talents_dir.mkdir() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() (tmp_path / "health").mkdir() @@ -545,7 +561,7 @@ def test_wait_for_agents_partial(callosum_server): """Test wait_for_uses with some completing and some timing out.""" tmp_path = callosum_server talents_dir = tmp_path / "talents" - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir(exist_ok=True) completing_agent = "1111" @@ -588,7 +604,7 @@ def test_wait_for_agents_missed_event_recovery(tmp_path, monkeypatch, caplog): monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) talents_dir = tmp_path / "talents" talents_dir.mkdir() - unified_dir = talents_dir / "unified" + unified_dir = talents_dir / "chat" unified_dir.mkdir() (tmp_path / "health").mkdir() diff --git a/tests/test_entity_talents.py b/tests/test_entity_talents.py index 490261559..e159359f0 100644 --- a/tests/test_entity_talents.py +++ b/tests/test_entity_talents.py @@ -98,7 +98,7 @@ def test_agent_context_includes_entities_by_facet(fixture_journal): def test_agent_context_with_facet_focus(fixture_journal): """Test that get_talent with facet parameter uses focused single-facet context.""" - config = get_talent("unified", facet="full-featured") + config = get_talent("chat", facet="full-featured") prompt = config["user_instruction"] diff --git a/tests/test_home_events.py b/tests/test_home_events.py deleted file mode 100644 index 907188c22..000000000 --- a/tests/test_home_events.py +++ /dev/null @@ -1,143 +0,0 @@ -# SPDX-License-Identifier: AGPL-3.0-only -# Copyright (c) 2026 sol pbc - -"""Tests for apps/home/events.py — conversation exchange recording.""" - -from unittest.mock import patch - -import pytest - -from apps.events import EventContext, clear_handlers, stop_dispatcher -from apps.home.events import TRIAGE_AGENT_NAMES, record_triage_exchange - - -@pytest.fixture(autouse=True) -def clean_handlers(): - clear_handlers() - yield - clear_handlers() - stop_dispatcher() - - -class TestRecordTriageExchange: - """Tests for record_triage_exchange handler.""" - - def _make_ctx(self, msg): - return EventContext(msg=msg, app="home", tract="cortex", event="finish") - - def test_ignores_non_triage_agent(self): - """Handler returns early for non-triage agent names.""" - ctx = self._make_ctx( - { - "tract": "cortex", - "event": "finish", - "name": "reviewer", - "use_id": "123", - "result": "hello", - } - ) - with patch("apps.home.events.record_exchange") as mock_record: - record_triage_exchange(ctx) - mock_record.assert_not_called() - - def test_ignores_missing_agent_id(self): - """Handler returns early if use_id is missing.""" - ctx = self._make_ctx( - { - "tract": "cortex", - "event": "finish", - "name": "unified", - "result": "hello", - } - ) - with patch("apps.home.events.record_exchange") as mock_record: - record_triage_exchange(ctx) - mock_record.assert_not_called() - - @pytest.mark.parametrize("agent_name", sorted(TRIAGE_AGENT_NAMES)) - def test_records_exchange_for_triage_agents(self, agent_name): - """Handler calls record_exchange with correct fields for each triage agent name.""" - events = [ - { - "event": "request", - "ts": 1700000000000, - "use_id": "abc123", - "facet": "work", - "app": "home", - "path": "/home", - "user_message": "hello world", - }, - { - "event": "finish", - "ts": 1700000001000, - "use_id": "abc123", - "result": "hi there", - }, - ] - ctx = self._make_ctx( - { - "tract": "cortex", - "event": "finish", - "name": agent_name, - "use_id": "abc123", - "result": "hi there", - } - ) - with patch("apps.home.events.read_use_events", return_value=events): - with patch("apps.home.events.record_exchange") as mock_record: - record_triage_exchange(ctx) - mock_record.assert_called_once_with( - facet="work", - app="home", - path="/home", - user_message="hello world", - agent_response="hi there", - talent=agent_name, - use_id="abc123", - ) - - def test_handles_missing_request_event(self): - """Handler uses empty strings for metadata if request event not found.""" - events = [ - {"event": "finish", "use_id": "abc123", "result": "done"}, - ] - ctx = self._make_ctx( - { - "tract": "cortex", - "event": "finish", - "name": "unified", - "use_id": "abc123", - "result": "done", - } - ) - with patch("apps.home.events.read_use_events", return_value=events): - with patch("apps.home.events.record_exchange") as mock_record: - record_triage_exchange(ctx) - mock_record.assert_called_once_with( - facet="", - app="", - path="", - user_message="", - agent_response="done", - talent="unified", - use_id="abc123", - ) - - def test_handles_read_error_gracefully(self): - """Handler logs and swallows exceptions from read_use_events.""" - ctx = self._make_ctx( - { - "tract": "cortex", - "event": "finish", - "name": "unified", - "use_id": "abc123", - "result": "done", - } - ) - with patch( - "apps.home.events.read_use_events", - side_effect=FileNotFoundError("not found"), - ): - with patch("apps.home.events.record_exchange") as mock_record: - record_triage_exchange(ctx) # should not raise - mock_record.assert_not_called() diff --git a/tests/test_journal_index.py b/tests/test_journal_index.py index d12eb3977..ecbcc2206 100644 --- a/tests/test_journal_index.py +++ b/tests/test_journal_index.py @@ -11,6 +11,7 @@ from pathlib import Path import pytest +from convey.chat_stream import append_chat_event from tests.conftest import copytree_tracked from think.indexer import sanitize_fts_query from think.indexer.journal import ( @@ -1675,6 +1676,32 @@ class TestSegmentChunks: assert count1 == count2 +def test_chat_turn_is_searchable_after_rescan(journal_fixture): + from think.indexer.journal import scan_journal, search_journal + + append_chat_event( + "owner_message", + text="Tell me about the nebula phrase", + app="sol", + path="/app/sol", + facet="work", + ) + append_chat_event( + "sol_message", + use_id="1713628000000", + text="The unique nebula phrase is now in chat history.", + notes="done", + requested_exec=False, + requested_task=None, + ) + + scan_journal(str(journal_fixture), full=True) + total, results = search_journal("unique nebula phrase") + + assert total >= 1 + assert any("unique nebula phrase" in result["text"].lower() for result in results) + + def test_scan_journal_is_pure_wrt_entity_state(journal_copy): """scan_journal must not mutate journal/entities/ state.""" from think.indexer.journal import scan_journal diff --git a/tests/test_no_legacy_chat_imports.py b/tests/test_no_legacy_chat_imports.py new file mode 100644 index 000000000..00dd59c17 --- /dev/null +++ b/tests/test_no_legacy_chat_imports.py @@ -0,0 +1,78 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import ast +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +ALLOWED_UNIFIED_PATHS = { + ROOT / "apps/sol/maint/006_rename_unified_triage_providers.py", + ROOT / "tests/test_maint_006_rename_unified_triage_providers.py", +} + + +def _parts(*pieces: str) -> str: + return "".join(pieces) + + +BANNED_NAMES = { + _parts("record_", "exchange"), + _parts("build_", "memory_", "context"), + _parts("INJECTION_", "MARKER"), + _parts("inject_", "memory"), + _parts("get_", "recent_", "exchanges"), + _parts("get_", "today_", "exchanges"), + _parts("TRIAGE_", "AGENT_", "NAMES"), + _parts("record_", "triage_", "exchange"), + _parts("compute_", "display_", "mode"), +} +LEGACY_CHAT_MODULE = _parts("think", ".", "conversation") +LEGACY_MEMORY_MODULE = _parts("talent", ".", "conversation_", "memory") +LEGACY_NAME = _parts("uni", "fied") + + +def _python_files() -> list[Path]: + return [ + path + for path in ROOT.rglob("*.py") + if ".venv" not in path.parts and "__pycache__" not in path.parts + ] + + +def test_no_legacy_chat_imports_or_usages(): + violations: list[str] = [] + + for path in _python_files(): + tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path)) + for node in ast.walk(tree): + if isinstance(node, ast.Import): + for alias in node.names: + if alias.name in {LEGACY_CHAT_MODULE, LEGACY_MEMORY_MODULE}: + violations.append(f"{path}: import {alias.name}") + elif isinstance(node, ast.ImportFrom): + if node.module in {LEGACY_CHAT_MODULE, LEGACY_MEMORY_MODULE}: + violations.append(f"{path}: from {node.module} import ...") + elif isinstance(node, ast.Name) and node.id in BANNED_NAMES: + violations.append(f"{path}: name {node.id}") + elif isinstance(node, ast.Attribute) and node.attr in BANNED_NAMES: + violations.append(f"{path}: attribute {node.attr}") + + assert violations == [] + + +def test_no_live_unified_literals_outside_migration_paths(): + violations: list[str] = [] + + for path in _python_files(): + if path in ALLOWED_UNIFIED_PATHS: + continue + if path == Path(__file__).resolve(): + continue + tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path)) + for node in ast.walk(tree): + if isinstance(node, ast.Constant) and node.value == LEGACY_NAME: + violations.append(str(path)) + + assert violations == [] diff --git a/tests/test_talent_cli.py b/tests/test_talent_cli.py index 087344c5e..6f40b3a3b 100644 --- a/tests/test_talent_cli.py +++ b/tests/test_talent_cli.py @@ -309,7 +309,7 @@ def test_logs_runs_default(capsys): output = capsys.readouterr().out # Should have runs from all fixture days (original + R&J) - assert "default" in output or "unified" in output + assert "default" in output or "chat" in output assert "flow" in output assert "activity" in output assert "entities" in output diff --git a/tests/test_talent_fallback.py b/tests/test_talent_fallback.py index 403450d38..7c493b4f2 100644 --- a/tests/test_talent_fallback.py +++ b/tests/test_talent_fallback.py @@ -127,7 +127,7 @@ def test_preflight_swap_unhealthy_primary(monkeypatch): ) monkeypatch.setenv("ANTHROPIC_API_KEY", "test-key") - config = prepare_config({"name": "unified", "prompt": "hello"}) + config = prepare_config({"name": "chat", "prompt": "hello"}) assert config["provider"] == "anthropic" assert config["model"] == "claude-sonnet-4-5" @@ -144,7 +144,7 @@ def test_preflight_no_swap_healthy_primary(monkeypatch): ) monkeypatch.setattr("think.models.should_recheck_health", lambda _h: False) - config = prepare_config({"name": "unified", "prompt": "hello"}) + config = prepare_config({"name": "chat", "prompt": "hello"}) assert config["provider"] == "google" assert "fallback_from" not in config @@ -162,7 +162,7 @@ def test_preflight_no_swap_no_backup_key(monkeypatch): monkeypatch.setattr("think.models.get_backup_provider", lambda _type: "anthropic") monkeypatch.delenv("ANTHROPIC_API_KEY", raising=False) - config = prepare_config({"name": "unified", "prompt": "hello"}) + config = prepare_config({"name": "chat", "prompt": "hello"}) assert config["provider"] == "google" assert "fallback_from" not in config @@ -255,7 +255,7 @@ def test_on_failure_retry_cogitate_uses_context_from_name(monkeypatch): monkeypatch.setenv("ANTHROPIC_API_KEY", "test-key") config = { - "name": "unified", + "name": "chat", "provider": "google", "model": "gemini-3-flash-preview", "health_stale": False, @@ -292,7 +292,7 @@ def test_on_failure_retry_generate(monkeypatch): monkeypatch.setenv("ANTHROPIC_API_KEY", "test-key") config = { - "name": "unified", + "name": "chat", "provider": "google", "model": "gemini-3-flash-preview", "prompt": "hello", @@ -324,7 +324,7 @@ def test_on_failure_no_retry_value_error(monkeypatch): monkeypatch.setattr("think.models.generate_with_result", bad_generate) config = { - "name": "unified", + "name": "chat", "provider": "google", "model": "gemini-3-flash-preview", "prompt": "hello", @@ -361,7 +361,7 @@ def test_on_failure_both_fail_raises_original(monkeypatch): monkeypatch.setenv("ANTHROPIC_API_KEY", "test-key") config = { - "name": "unified", + "name": "chat", "provider": "google", "model": "gemini-3-flash-preview", "prompt": "hello", @@ -380,7 +380,7 @@ def test_fallback_event_emitted(): events = [] config = { "type": "cogitate", - "name": "unified", + "name": "chat", "provider": "anthropic", "model": "claude-sonnet-4-5", "prompt": "hello", @@ -427,7 +427,7 @@ def test_recheck_requested_on_stale(monkeypatch): def test_main_async_no_duplicate_error_when_evented(monkeypatch, capsys): from think.talents import main_async - ndjson_input = json.dumps({"name": "unified", "prompt": "hello"}) + ndjson_input = json.dumps({"name": "chat", "prompt": "hello"}) monkeypatch.setattr("sys.stdin", StringIO(ndjson_input)) async def fake_run_talent(_config, emit_event, dry_run=False): diff --git a/tests/test_talents_ndjson.py b/tests/test_talents_ndjson.py index 84394846a..44659c0b4 100644 --- a/tests/test_talents_ndjson.py +++ b/tests/test_talents_ndjson.py @@ -96,7 +96,7 @@ def test_ndjson_single_request(mock_journal, monkeypatch, capsys): { "prompt": "What is 2+2?", "provider": "openai", - "name": "unified", + "name": "chat", "model": GPT_5, "max_output_tokens": 100, } diff --git a/tests/verify_api.py b/tests/verify_api.py index 66e091ce2..7f7b0bde8 100644 --- a/tests/verify_api.py +++ b/tests/verify_api.py @@ -61,7 +61,7 @@ ENDPOINTS = [ { "app": "sol", "name": "preview", - "path": "/app/sol/api/preview/unified", + "path": "/app/sol/api/preview/chat", "params": {}, "status": 200, }, @@ -87,6 +87,28 @@ ENDPOINTS = [ "status": 200, "sandbox_only": True, # live indexer computes differently than Flask test client }, + # convey/chat.py + { + "app": "chat", + "name": "session", + "path": "/api/chat/session", + "params": {}, + "status": 200, + }, + { + "app": "chat", + "name": "stream", + "path": "/api/chat/stream/20260304", + "params": {"limit": "20"}, + "status": 200, + }, + { + "app": "chat", + "name": "result", + "path": "/api/chat/result/1700000000001", + "params": {}, + "status": 404, + }, # apps/activities/routes.py { "app": "activities", diff --git a/think/awareness.py b/think/awareness.py index 4e51d6caf..a932d4cd9 100644 --- a/think/awareness.py +++ b/think/awareness.py @@ -209,6 +209,47 @@ def get_imports() -> dict[str, Any]: ) +def _recent_chat_exchanges(limit: int = 10000) -> list[dict[str, Any]]: + """Return owner-visible chat responses from chat stream history.""" + from think.utils import day_dirs + + try: + days = day_dirs() + except Exception: + return [] + + exchanges: list[dict[str, Any]] = [] + for day_name in sorted(days): + day_path = Path(days[day_name]) + chat_root = day_path / "chat" + if not chat_root.exists(): + continue + for segment_dir in sorted(chat_root.iterdir()): + if not segment_dir.is_dir(): + continue + chat_path = segment_dir / "chat.jsonl" + if not chat_path.exists(): + continue + try: + for line in chat_path.read_text().splitlines(): + if not line.strip(): + continue + event = json.loads(line) + if event.get("kind") != "sol_message": + continue + exchanges.append( + { + "talent": "chat", + "agent_response": event.get("text", ""), + } + ) + except (OSError, json.JSONDecodeError): + logger.warning("Skipping malformed chat stream file: %s", chat_path) + if limit <= 0: + return [] + return exchanges[-limit:] + + def compute_thickness() -> dict[str, Any]: """Compute journal thickness signals for naming ceremony readiness. @@ -221,7 +262,6 @@ def compute_thickness() -> dict[str, Any]: - ``journal_days``: number of day directories with at least one segment - ``ready``: True when the naming ceremony should trigger """ - from think.conversation import get_recent_exchanges from think.facets import get_enabled_facets from think.indexer.journal import get_entity_strength from think.utils import day_dirs, iter_segments @@ -233,7 +273,7 @@ def compute_thickness() -> dict[str, Any]: entity_depth = sum(1 for e in entities if e.get("observation_depth", 0) >= 2) try: - exchanges = get_recent_exchanges(limit=10000) + exchanges = _recent_chat_exchanges(limit=10000) except Exception: exchanges = [] non_onboarding = [ diff --git a/think/conversation.py b/think/conversation.py deleted file mode 100644 index e262e0c44..000000000 --- a/think/conversation.py +++ /dev/null @@ -1,341 +0,0 @@ -# SPDX-License-Identifier: AGPL-3.0-only -# Copyright (c) 2026 sol pbc - -"""Conversation memory service for solstone. - -Manages conversation exchange storage, retrieval, and context injection -for the unified talent agent. Three layers of recall: - -- Layer 1: Recent exchanges (last ~10 turns), loaded directly into context -- Layer 2: Today's earlier exchanges, summarized compactly -- Layer 3: Older conversations, searchable via journal search (automatic — - exchanges are stored as journal entries indexed by FTS5) -""" - -from __future__ import annotations - -import json -import logging -import re -from datetime import datetime -from pathlib import Path - -from think.utils import day_path, get_journal, now_ms - -logger = logging.getLogger(__name__) - -# Append-only exchange log for fast recent retrieval -EXCHANGES_FILE = "conversation/exchanges.jsonl" - -# Journal stream name for conversation segments -CONVERSATION_STREAM = "conversation" - -# Marker in unified talent for memory injection -INJECTION_MARKER = "CONVERSATION_MEMORY_INJECTION_POINT" - -# Context budget: max characters for agent response in recent exchanges -MAX_RESPONSE_CHARS = 300 - -# Max characters for user message in compact summaries -MAX_MESSAGE_CHARS = 100 - -# Default number of recent exchanges for layer 1 -DEFAULT_RECENT_LIMIT = 10 - - -# --------------------------------------------------------------------------- -# Exchange Recording -# --------------------------------------------------------------------------- - - -def record_exchange( - *, - ts: int | None = None, - facet: str = "", - app: str = "", - path: str = "", - user_message: str = "", - agent_response: str = "", - talent: str = "", - use_id: str = "", -) -> None: - """Record a conversation exchange to journal storage. - - Writes to two locations: - 1. conversation/exchanges.jsonl — append-only quick-read index - 2. YYYYMMDD/conversation/HHMMSS_1/talents/conversation.md — journal entry - for FTS5 search indexing (matches */*/*/talents/*.md formatter pattern) - """ - if not user_message or not agent_response: - return - - if ts is None: - ts = now_ms() - - journal = get_journal() - - exchange = { - "ts": ts, - "facet": facet, - "app": app, - "path": path, - "user_message": user_message, - "agent_response": agent_response, - "talent": talent, - "use_id": use_id, - } - - # 1. Append to exchanges.jsonl (fast-read index) - jsonl_path = Path(journal) / EXCHANGES_FILE - jsonl_path.parent.mkdir(parents=True, exist_ok=True) - try: - with open(jsonl_path, "a", encoding="utf-8") as f: - f.write(json.dumps(exchange, ensure_ascii=False) + "\n") - except Exception: - logger.exception("Failed to write exchange to JSONL") - - # 2. Write journal segment for search indexing - dt = datetime.fromtimestamp(ts / 1000) - day = dt.strftime("%Y%m%d") - time_key = dt.strftime("%H%M%S") - segment = f"{time_key}_1" - - seg_dir = day_path(day) / CONVERSATION_STREAM / segment / "talents" - seg_dir.mkdir(parents=True, exist_ok=True) - - time_str = dt.strftime("%Y-%m-%d %H:%M:%S") - md_parts = ["# Conversation Exchange\n"] - md_parts.append(f"**Time:** {time_str}") - if facet: - md_parts.append(f"**Facet:** {facet}") - if app: - app_info = f"{app} ({path})" if path else app - md_parts.append(f"**App:** {app_info}") - md_parts.append("") - md_parts.append("## User\n") - md_parts.append(user_message) - md_parts.append("") - md_parts.append("## Sol\n") - md_parts.append(agent_response) - - md_content = "\n".join(md_parts) + "\n" - - md_path = seg_dir / "conversation.md" - try: - with open(md_path, "w", encoding="utf-8") as f: - f.write(md_content) - except Exception: - logger.exception("Failed to write conversation journal entry") - - -# --------------------------------------------------------------------------- -# Exchange Retrieval -# --------------------------------------------------------------------------- - - -def get_recent_exchanges( - limit: int = DEFAULT_RECENT_LIMIT, - facet: str | None = None, -) -> list[dict]: - """Read the most recent conversation exchanges. - - Args: - limit: Maximum number of exchanges to return. - facet: If provided, only return exchanges from this facet. - - Returns: - List of exchange dicts, most recent last. - """ - journal = get_journal() - jsonl_path = Path(journal) / EXCHANGES_FILE - - if not jsonl_path.exists(): - return [] - - exchanges = [] - try: - with open(jsonl_path, "r", encoding="utf-8") as f: - for line in f: - line = line.strip() - if not line: - continue - try: - ex = json.loads(line) - ex = _normalize_exchange(ex) - if facet and ex.get("facet") != facet: - continue - exchanges.append(ex) - except json.JSONDecodeError: - continue - except Exception: - logger.exception("Failed to read exchanges") - return [] - - return exchanges[-limit:] - - -def get_today_exchanges(facet: str | None = None) -> list[dict]: - """Read all conversation exchanges from today. - - Args: - facet: If provided, only return exchanges from this facet. - - Returns: - List of exchange dicts from today, chronological order. - """ - journal = get_journal() - jsonl_path = Path(journal) / EXCHANGES_FILE - - if not jsonl_path.exists(): - return [] - - today = datetime.now().strftime("%Y%m%d") - exchanges = [] - - try: - with open(jsonl_path, "r", encoding="utf-8") as f: - for line in f: - line = line.strip() - if not line: - continue - try: - ex = json.loads(line) - ex = _normalize_exchange(ex) - ts = ex.get("ts", 0) - ex_day = datetime.fromtimestamp(ts / 1000).strftime("%Y%m%d") - if ex_day != today: - continue - if facet and ex.get("facet") != facet: - continue - exchanges.append(ex) - except (json.JSONDecodeError, ValueError, OSError): - continue - except Exception: - logger.exception("Failed to read today's exchanges") - return [] - - return exchanges - - -# --------------------------------------------------------------------------- -# Context Formatting -# --------------------------------------------------------------------------- - - -def _format_exchange(ex: dict, *, compact: bool = False) -> str: - """Format a single exchange for context injection. - - Args: - ex: Exchange dict. - compact: If True, return a one-liner summary. - - Returns: - Formatted string. - """ - ts = ex.get("ts", 0) - try: - time_str = datetime.fromtimestamp(ts / 1000).strftime("%H:%M") - except (ValueError, OSError): - time_str = "??:??" - - app = ex.get("app", "") - facet_val = ex.get("facet", "") - - context_parts = [time_str] - if app: - context_parts.append(app) - if facet_val: - context_parts.append(facet_val) - context = " · ".join(context_parts) - - user_msg = ex.get("user_message", "") - agent_resp = ex.get("agent_response", "") - - if compact: - truncated = user_msg[:MAX_MESSAGE_CHARS] - if len(user_msg) > MAX_MESSAGE_CHARS: - truncated += "..." - return f"- [{context}] {truncated}" - - # Full exchange with truncated response - truncated_resp = agent_resp[:MAX_RESPONSE_CHARS] - if len(agent_resp) > MAX_RESPONSE_CHARS: - truncated_resp += "..." - - return f"[{context}] User: {user_msg}\nSol: {truncated_resp}" - - -def build_memory_context( - facet: str | None = None, - recent_limit: int = DEFAULT_RECENT_LIMIT, -) -> str: - """Build the full conversation memory context block. - - Assembles layer 1 (recent exchanges) and layer 2 (today's summary) - into a formatted block for injection into the unified talent prompt. - - Args: - facet: Active facet for filtering. - recent_limit: Number of recent exchanges for layer 1. - - Returns: - Formatted memory context string, or empty string if no history. - """ - recent = get_recent_exchanges(limit=recent_limit, facet=facet) - if not recent: - return "" - - today_all = get_today_exchanges(facet=facet) - - parts = [] - - # Layer 2: Earlier today (exchanges beyond the recent set) - if len(today_all) > len(recent): - earlier = today_all[: -len(recent)] - if earlier: - parts.append("### Earlier Today\n") - for ex in earlier: - parts.append(_format_exchange(ex, compact=True)) - parts.append("") - - # Layer 1: Recent exchanges (full detail) - parts.append("### Recent Conversations\n") - parts.append("The following are your most recent exchanges with the user:\n") - for ex in recent: - parts.append(_format_exchange(ex, compact=False)) - parts.append("") - - return "\n".join(parts).strip() - - -def inject_memory(user_instruction: str, memory_context: str) -> str: - """Replace the CONVERSATION_MEMORY_INJECTION_POINT with memory context. - - Args: - user_instruction: The unified talent's user instruction text. - memory_context: Formatted conversation memory to inject. - - Returns: - Modified user instruction with memory context injected. - """ - if INJECTION_MARKER not in user_instruction: - return user_instruction - - # Replace the entire HTML comment block containing the marker - pattern = r"" - - if memory_context: - replacement = memory_context - else: - replacement = "No conversation history yet." - - return re.sub(pattern, replacement, user_instruction, flags=re.DOTALL) - - -def _normalize_exchange(ex: dict) -> dict: - """Normalize legacy exchange dicts to the talent namespace.""" - if "talent" not in ex and "muse" in ex: - ex["talent"] = ex["muse"] - # Optionally, remove the old 'muse' key if it's no longer needed in the normalized dict - # del ex["muse"] - return ex diff --git a/think/cortex.py b/think/cortex.py index fb3599eeb..ddc47cc4b 100644 --- a/think/cortex.py +++ b/think/cortex.py @@ -417,18 +417,6 @@ class CortexService: _req = self.use_requests.get(agent.use_id) if _req and "name" not in event: event["name"] = _req.get("name", "") - # Inject display mode for triage talent finish events - if event.get("event") == "finish" and _req: - try: - from apps.home.events import TRIAGE_AGENT_NAMES - from convey.triage import compute_display_mode - - if _req.get("name", "") in TRIAGE_AGENT_NAMES: - event["display"] = compute_display_mode( - event.get("result", "") - ) - except Exception: - pass # Display is cosmetic; don't break finish handling # Append to JSONL file with open(agent.log_path, "a") as f: diff --git a/think/cortex_client.py b/think/cortex_client.py index 471568163..1a0a91eba 100644 --- a/think/cortex_client.py +++ b/think/cortex_client.py @@ -37,6 +37,7 @@ def cortex_request( name: str, provider: Optional[str] = None, config: Optional[Dict[str, Any]] = None, + use_id: Optional[str] = None, ) -> str | None: """Create a Cortex talent request via Callosum broadcast. @@ -45,6 +46,7 @@ def cortex_request( name: Talent name - system (e.g., "chat") or app-qualified (e.g., "entities:entity_assist") provider: AI provider - openai, google, or anthropic config: Provider-specific configuration (model, max_output_tokens, thinking_budget, etc.) + use_id: Optional pre-reserved use_id. When omitted, a unique timestamp is allocated. Returns: Use ID (timestamp-based string), or None if the Callosum send failed. @@ -58,14 +60,20 @@ def cortex_request( # Generate monotonic timestamp in milliseconds, ensuring uniqueness global _last_ts - ts = now_ms() - - # If same or earlier than last used, increment to ensure uniqueness - if ts <= _last_ts: - ts = _last_ts + 1 - - _last_ts = ts - use_id = str(ts) + if use_id is None: + ts = now_ms() + + if ts <= _last_ts: + ts = _last_ts + 1 + + _last_ts = ts + use_id = str(ts) + else: + if not use_id.isdigit(): + raise ValueError("use_id must be a millisecond timestamp string") + ts = int(use_id) + if ts > _last_ts: + _last_ts = ts # Build request object request = { diff --git a/think/talent.py b/think/talent.py index 8cdddc625..b2286abf8 100644 --- a/think/talent.py +++ b/think/talent.py @@ -35,7 +35,6 @@ from think.prompts import _load_prompt_metadata, load_prompt TALENT_DIR = Path(__file__).parent.parent / "talent" APPS_DIR = Path(__file__).parent.parent / "apps" -_UNDISCOVERED_SYSTEM_TALENTS = {"triage"} # --------------------------------------------------------------------------- @@ -232,8 +231,6 @@ def get_talent_configs( if TALENT_DIR.is_dir(): for md_path in sorted(TALENT_DIR.glob("*.md")): name = md_path.stem - if name in _UNDISCOVERED_SYSTEM_TALENTS: - continue info = _load_prompt_metadata(md_path) info["source"] = "system" @@ -355,10 +352,6 @@ def _resolve_talent_path(name: str) -> tuple[Path, str]: # App talent: "support:support" -> apps/support/talent/support app, talent_name = name.split(":", 1) talent_dir = Path(__file__).parent.parent / "apps" / app / "talent" - elif name == "unified": - # Chat talent: "unified" -> talent/chat - talent_dir = TALENT_DIR - talent_name = "chat" else: # System talent: bare name -> talent/{name} talent_dir = TALENT_DIR