diff --git a/solstone/convey/chat.py b/solstone/convey/chat.py index 032d82590..1b0653524 100644 --- a/solstone/convey/chat.py +++ b/solstone/convey/chat.py @@ -29,6 +29,7 @@ from solstone.apps.chat.copy import ( CHAT_CLOSER_SUPPORT_SEND_FAILED, CHAT_CLOSER_TALENT_ERRORED_FORMAT, CHAT_CLOSER_TALENT_ERRORED_GENERIC, + CHAT_LIVENESS_TASK_FORMAT, CHAT_OFFER_SUPPORT_DECLINE, CHAT_OFFER_SUPPORT_PROMPT, CHAT_SUPPORT_ATTACH_FILED_FORMAT, @@ -73,7 +74,6 @@ MAX_ACTIVE_TALENTS = 2 _WATCHDOG_TIMEOUTS = {"chat": 30, "talent": 180} _DEFAULT_WATCHDOG_SECONDS = 180 _RESERVED_USE_ID_CAP = 256 -MAX_ACTIVE_REASON = "max active — waiting for one to finish" _state_lock = threading.Lock() _runtime_lock = threading.Lock() @@ -171,21 +171,10 @@ def post_chat() -> Any: 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 = _enqueue_trigger_locked( - trigger, - location, - ) - queued = True + response_use_id, queued, start_info = _activate_or_enqueue_trigger_locked( + trigger, + location, + ) if start_info is not None: spawn_result = _spawn_chat_generate(start_info) @@ -557,7 +546,7 @@ def _on_cortex_finish(message: dict[str, Any]) -> None: if not use_id: return - next_info: dict[str, Any] | None = None + next_actions: list[dict[str, Any] | None] = [] finish_payload: dict[str, Any] | None = None error_payload: dict[str, Any] | None = None @@ -579,7 +568,7 @@ def _on_cortex_finish(message: dict[str, Any]) -> None: _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) + next_actions.append(_build_spawn_info_locked(logical_use_id)) else: _evict_thinking_locked(use_id) append_chat_event( @@ -593,7 +582,7 @@ def _on_cortex_finish(message: dict[str, Any]) -> None: "use_id": logical_use_id, "reason": "provider_response_invalid", } - next_info = _clear_current_locked() + next_actions.append(_clear_current_locked()) else: message_text = parsed["message"] or "" requested_target = ( @@ -669,6 +658,12 @@ def _on_cortex_finish(message: dict[str, Any]) -> None: requested_target = None requested_task = None # consent in {"pending", "confirmed"}: allow the spawn, no offer. + if requested_target: + message_text = _dispatch_ack_text( + requested_target, + requested_task, + message_text, + ) thinking = _drain_thinking_locked(use_id, message) sol_message_fields: dict[str, Any] = { "use_id": logical_use_id, @@ -685,6 +680,9 @@ def _on_cortex_finish(message: dict[str, Any]) -> None: sol_message_fields["offer"] = offer if draft is not None: sol_message_fields["draft"] = draft + origin = trigger.get("origin") + if origin is not None and not requested_target: + sol_message_fields["origin"] = origin append_chat_event( "sol_message", **sol_message_fields, @@ -692,39 +690,14 @@ def _on_cortex_finish(message: dict[str, Any]) -> None: _current_chat_state["retry_count"] = 0 _set_current_raw_use_locked(logical_use_id, None) if requested_target: - active_talent_count = _active_talent_count_for_today_locked() - if active_talent_count >= MAX_ACTIVE_TALENTS: - _current_chat_state["trigger"] = { - "type": "synthetic-max-active", - "reason": MAX_ACTIVE_REASON, - } - synthetic_use_id = _reserve_use_id_locked() - _set_current_raw_use_locked(logical_use_id, synthetic_use_id) - next_info = _build_spawn_info_locked(logical_use_id) - else: - talent_use_id = _reserve_use_id_locked() - _active_talents[talent_use_id] = { - "chat_use_id": logical_use_id, - "target": requested_target, - "task": requested_task, - "location": dict(_current_chat_state["location"]), - } - append_chat_event( - "talent_spawned", - use_id=talent_use_id, - name=requested_target, - task=requested_task, - started_at=int(talent_use_id), - ) - next_info = { - "kind": "talent", - "logical_use_id": logical_use_id, - "target": requested_target, - "use_id": talent_use_id, - "task": requested_task, - "context": parsed["talent_request"].get("context") or {}, - "location": dict(_current_chat_state["location"]), - } + dispatch_job = _build_dispatch_job_locked( + logical_use_id, + requested_target, + requested_task, + parsed["talent_request"].get("context") or {}, + ) + next_actions.append(_clear_current_locked()) + next_actions.append(_spawn_or_defer_dispatch_locked(dispatch_job)) else: if not message_text: provider = str(message.get("provider") or "") @@ -746,16 +719,18 @@ def _on_cortex_finish(message: dict[str, Any]) -> None: } if not message_text: _evict_thinking_locked(use_id) - next_info = _clear_current_locked() + next_actions.append(_clear_current_locked()) elif use_id in _active_talents: summary = str(message.get("result") or "").strip() - next_info = _handle_talent_terminal_locked( - use_id, - "talent_finished", - "summary", - summary, - terminal_message=message, + next_actions.extend( + _handle_talent_terminal_locked( + use_id, + "talent_finished", + "summary", + summary, + terminal_message=message, + ) ) elif _is_superseded_raw_use_id_locked(use_id): logger.debug( @@ -773,7 +748,7 @@ def _on_cortex_finish(message: dict[str, Any]) -> None: "no matching active chat-generate or talent", ) - _run_next_action(next_info) + _run_next_actions(next_actions) if finish_payload is not None: _emit_finish(finish_payload["use_id"], finish_payload["message"]) if error_payload is not None: @@ -785,7 +760,7 @@ def _on_cortex_error(message: dict[str, Any]) -> None: if not use_id: return - next_info: dict[str, Any] | None = None + next_actions: list[dict[str, Any] | None] = [] error_payload: dict[str, Any] | None = None with _state_lock: @@ -811,17 +786,19 @@ def _on_cortex_error(message: dict[str, Any]) -> None: "provider": provider, "detail": detail, } - next_info = _clear_current_locked() + next_actions.append(_clear_current_locked()) elif use_id in _active_talents: reason = str(message.get("error") or "unknown") reason_code = message.get("reason_code") or None _evict_thinking_locked(use_id) - next_info = _handle_talent_terminal_locked( - use_id, - "talent_errored", - "reason", - reason, - reason_code=reason_code, + next_actions.extend( + _handle_talent_terminal_locked( + use_id, + "talent_errored", + "reason", + reason, + reason_code=reason_code, + ) ) elif _is_superseded_raw_use_id_locked(use_id): logger.debug( @@ -839,7 +816,7 @@ def _on_cortex_error(message: dict[str, Any]) -> None: "no matching active chat-generate or talent", ) - _run_next_action(next_info) + _run_next_actions(next_actions) if error_payload is not None: _emit_error( error_payload["use_id"], @@ -857,11 +834,15 @@ def _handle_talent_terminal_locked( *, reason_code: str | None = None, terminal_message: dict[str, Any] | None = None, -) -> dict[str, Any] | None: +) -> list[dict[str, Any] | None]: _cancel_watchdog_locked(use_id) talent_state = _active_talents.pop(use_id) logical_use_id = str(talent_state["chat_use_id"]) talent_name = str(talent_state["target"]) + origin = { + "logical_use_id": logical_use_id, + "ask": str(talent_state.get("ask") or ""), + } trigger = _talent_terminal_trigger( kind, use_id, @@ -869,6 +850,7 @@ def _handle_talent_terminal_locked( result_field_name, result_value, reason_code=reason_code, + origin=origin, ) event_fields: dict[str, Any] = { "use_id": use_id, @@ -882,16 +864,11 @@ def _handle_talent_terminal_locked( if thinking is not None: event_fields["thinking"] = thinking append_chat_event(kind, **event_fields) - if _current_chat_use_id != logical_use_id or _current_chat_state is None: - return None - - _current_chat_state["trigger"] = trigger - _set_current_raw_use_locked( - logical_use_id, - _reserve_use_id_locked(), + _, _, synth_action = _activate_or_enqueue_trigger_locked( + trigger, + dict(talent_state["location"]), ) - _current_chat_state["retry_count"] = 0 - return _build_spawn_info_locked(logical_use_id) + return [synth_action, _promote_deferred_spawn_locked(_today_day())] def _run_next_action(action: dict[str, Any] | None) -> None: @@ -918,6 +895,11 @@ def _run_next_action(action: dict[str, Any] | None) -> None: ) +def _run_next_actions(actions: list[dict[str, Any] | None]) -> None: + for action in actions: + _run_next_action(action) + + def _spawn_chat_generate(action: dict[str, Any]) -> ChatSpawnResult: logger.info( "starting chat generate logical=%s raw=%s trigger=%s", @@ -1001,30 +983,36 @@ def _spawn_talent(action: dict[str, Any]) -> bool: def _handle_talent_spawn_failure(action: dict[str, Any]) -> None: - next_info: dict[str, Any] | None = None + next_actions: list[dict[str, Any] | None] = [] with _state_lock: - _cancel_watchdog_locked(str(action["use_id"])) - _active_talents.pop(str(action["use_id"]), None) + use_id = str(action["use_id"]) + _cancel_watchdog_locked(use_id) + talent_state = _active_talents.pop(use_id, None) + logical_use_id = str( + (talent_state or {}).get("chat_use_id") or action["logical_use_id"] + ) + talent_name = str((talent_state or {}).get("target") or action["target"]) + ask = str((talent_state or {}).get("ask") or "") + location = dict((talent_state or {}).get("location") or action["location"]) append_chat_event( "talent_errored", - use_id=action["use_id"], - name=action["target"], + use_id=use_id, + name=talent_name, reason="unknown", ) - 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": action["target"], - "reason": "unknown", - } - _set_current_raw_use_locked( - str(action["logical_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) + trigger = _talent_terminal_trigger( + "talent_errored", + use_id, + talent_name, + "reason", + "unknown", + origin={"logical_use_id": logical_use_id, "ask": ask}, + ) + _, _, synth_action = _activate_or_enqueue_trigger_locked(trigger, location) + next_actions.extend( + [synth_action, _promote_deferred_spawn_locked(_today_day())] + ) + _run_next_actions(next_actions) def _handle_chat_failure( @@ -1058,6 +1046,7 @@ def _recover_active_talents_locked(day: str) -> None: events = read_chat_events(day) latest_owner_message: dict[str, Any] | None = None latest_sol_message: dict[str, Any] | None = None + queued_events: dict[str, dict[str, Any]] = {} spawned: dict[str, dict[str, Any]] = {} latest_parent_kind: str | None = None @@ -1071,37 +1060,69 @@ def _recover_active_talents_locked(day: str) -> None: latest_sol_message = event latest_parent_kind = "sol_message" continue + if kind == "talent_queued": + use_id = str(event.get("use_id") or "") + if use_id: + queued_events[use_id] = event + continue if kind == "talent_spawned": use_id = str(event.get("use_id") or "") if not use_id: continue - if latest_sol_message is None or latest_owner_message is None: + queued_event = queued_events.get(use_id) + if queued_event is None and ( + latest_sol_message is None or latest_owner_message is None + ): logger.warning( "skipping active-talent recovery for %s: no parent chat turn", use_id, ) continue - chat_use_id = str(latest_sol_message.get("use_id") or "") + chat_use_id = str( + (queued_event or {}).get("chat_use_id") + or (latest_sol_message or {}).get("use_id") + or "" + ) if not chat_use_id: logger.warning( "skipping active-talent recovery for %s: sol_message missing use_id", use_id, ) continue + location_source = ( + queued_event.get("location") + if queued_event is not None + and isinstance(queued_event.get("location"), dict) + else None + ) spawned[use_id] = { "chat_use_id": chat_use_id, "target": str(event.get("name") or ""), "task": str(event.get("task") or ""), "trigger": latest_parent_kind or "sol_message", - "location": _normalize_location( - latest_owner_message.get("app"), - latest_owner_message.get("path"), - latest_owner_message.get("facet"), + "location": ( + _normalize_location( + location_source.get("app"), + location_source.get("path"), + location_source.get("facet"), + ) + if location_source is not None + else _normalize_location( + (latest_owner_message or {}).get("app"), + (latest_owner_message or {}).get("path"), + (latest_owner_message or {}).get("facet"), + ) + ), + "ask": str( + (queued_event or {}).get("ask") + or (latest_owner_message or {}).get("text") + or "" ), } continue if kind in {"talent_finished", "talent_errored"}: spawned.pop(str(event.get("use_id") or ""), None) + queued_events.pop(str(event.get("use_id") or ""), None) for use_id, state in spawned.items(): # recovery blind spot: pre-crash reservations are not seen here @@ -1120,28 +1141,42 @@ def _recover_active_talents_locked(day: str) -> None: def _recover_chat_if_needed() -> None: day = _today_day() - start_info: dict[str, Any] | None = None + start_actions: list[dict[str, Any]] = [] with _state_lock: _recover_active_talents_locked(day) + while _active_talent_count_for_today_locked() < MAX_ACTIVE_TALENTS: + promotion = _promote_deferred_spawn_locked(day) + if promotion is None: + break + start_actions.append(promotion) if _current_chat_use_id is not None: - return - unresolved = find_unresponded_trigger(day) - if unresolved is None: - 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) + unresolved = None + else: + unresolved = find_unresponded_trigger(day) + if unresolved is not None: + location = _location_for_trigger(day, unresolved) + trigger = _trigger_from_stream_event(day, unresolved) + _, _, start_info = _activate_or_enqueue_trigger_locked(trigger, location) + if start_info is not None: + start_actions.append(start_info) - if start_info is not None: - spawn_result = _spawn_chat_generate(start_info) - if not spawn_result.ok: - _handle_chat_failure( - start_info["logical_use_id"], - spawn_result.reason, - detail=spawn_result.detail, - ) + _run_next_actions(start_actions) + + +def _activate_or_enqueue_trigger_locked( + trigger: dict[str, Any], + location: dict[str, str], +) -> tuple[str, bool, dict[str, Any] | None]: + if _current_chat_use_id is None: + logical_use_id = _reserve_use_id_locked() + return ( + logical_use_id, + False, + _activate_current_locked(logical_use_id, trigger, location), + ) + queued_use_id = _enqueue_trigger_locked(trigger, location) + return queued_use_id, True, None def _activate_current_locked( @@ -1175,6 +1210,149 @@ def _build_spawn_info_locked(logical_use_id: str) -> dict[str, Any]: } +def _dispatch_ack_text(target: str, task: str | None, message_text: str) -> str: + text = message_text.strip() + if text: + return text + return CHAT_LIVENESS_TASK_FORMAT.format( + label=chat_stream._talent_label(target, "running"), + task=task or "", + ).strip() + + +def _current_trigger_ask_locked() -> str: + if _current_chat_state is None: + return "" + trigger = _current_chat_state.get("trigger") or {} + message = trigger.get("message") + if message: + return str(message) + origin = trigger.get("origin") + if isinstance(origin, dict) and origin.get("ask"): + return str(origin["ask"]) + return "" + + +def _build_dispatch_job_locked( + logical_use_id: str, + target: str, + task: str | None, + context: dict[str, Any], +) -> dict[str, Any]: + assert _current_chat_state is not None + return { + "use_id": _reserve_use_id_locked(), + "chat_use_id": logical_use_id, + "target": target, + "task": task, + "context": dict(context), + "location": dict(_current_chat_state["location"]), + "ask": _current_trigger_ask_locked(), + } + + +def _spawn_or_defer_dispatch_locked(job: dict[str, Any]) -> dict[str, Any] | None: + if _active_talent_count_for_today_locked() >= MAX_ACTIVE_TALENTS: + append_chat_event( + "talent_queued", + use_id=job["use_id"], + name=job["target"], + task=job["task"], + queued_at=now_ms(), + chat_use_id=job["chat_use_id"], + ask=job["ask"], + context=dict(job["context"]), + location=dict(job["location"]), + ) + return None + return _register_talent_spawn_locked(job, started_at=int(str(job["use_id"]))) + + +def _register_talent_spawn_locked( + job: dict[str, Any], + *, + started_at: int, +) -> dict[str, Any]: + use_id = str(job["use_id"]) + chat_use_id = str(job["chat_use_id"]) + _reserved_use_ids[use_id] = None + _reserved_use_ids[chat_use_id] = None + _active_talents[use_id] = { + "chat_use_id": chat_use_id, + "target": str(job["target"]), + "task": job["task"], + "location": dict(job["location"]), + "ask": str(job.get("ask") or ""), + } + append_chat_event( + "talent_spawned", + use_id=use_id, + name=str(job["target"]), + task=job["task"], + started_at=started_at, + ) + return { + "kind": "talent", + "logical_use_id": chat_use_id, + "target": str(job["target"]), + "use_id": use_id, + "task": job["task"], + "context": dict(job["context"]), + "location": dict(job["location"]), + } + + +def _promote_deferred_spawn_locked(day: str) -> dict[str, Any] | None: + if _active_talent_count_for_today_locked() >= MAX_ACTIVE_TALENTS: + return None + event = _oldest_unpromoted_queued_talent(day) + if event is None: + return None + job = { + "use_id": str(event["use_id"]), + "chat_use_id": str(event["chat_use_id"]), + "target": str(event["name"]), + "task": event.get("task"), + "context": dict(event.get("context") or {}), + "location": _normalize_location( + (event.get("location") or {}).get("app") + if isinstance(event.get("location"), dict) + else "", + (event.get("location") or {}).get("path") + if isinstance(event.get("location"), dict) + else "", + (event.get("location") or {}).get("facet") + if isinstance(event.get("location"), dict) + else "", + ), + "ask": str(event.get("ask") or ""), + } + return _register_talent_spawn_locked(job, started_at=now_ms()) + + +def _oldest_unpromoted_queued_talent(day: str) -> dict[str, Any] | None: + queued: dict[str, dict[str, Any]] = {} + for event in read_chat_events(day): + use_id = str(event.get("use_id") or "") + if not use_id: + continue + kind = event.get("kind") + if kind == "talent_queued": + queued[use_id] = event + continue + if kind in {"talent_spawned", "talent_finished", "talent_errored"}: + queued.pop(use_id, None) + if not queued: + return None + return sorted( + queued.values(), + key=lambda event: ( + int(event.get("queued_at", 0) or 0), + str(event.get("use_id") or ""), + ), + )[0] + + def _enqueue_trigger_locked( trigger: dict[str, Any], location: dict[str, str], @@ -1326,7 +1504,7 @@ def _is_routeable_cortex_use_id_locked(use_id: str) -> bool: def _on_watchdog_timeout(use_id: str, kind: str, logical_use_id: str) -> None: - next_info: dict[str, Any] | None = None + next_actions: list[dict[str, Any] | None] = [] should_emit = False with _state_lock: @@ -1351,7 +1529,7 @@ def _on_watchdog_timeout(use_id: str, kind: str, logical_use_id: str) -> None: provider="", detail="", ) - next_info = _clear_current_locked() + next_actions.append(_clear_current_locked()) should_emit = True elif kind == "talent": talent_state = _active_talents.get(use_id) @@ -1367,33 +1545,21 @@ def _on_watchdog_timeout(use_id: str, kind: str, logical_use_id: str) -> None: logical_use_id, ) _evict_thinking_locked(use_id) - append_chat_event( - "talent_errored", - use_id=use_id, - name=str(talent_state["target"]), - reason="talent took too long", - ) - _active_talents.pop(use_id, None) - append_chat_event( - "chat_error", - reason="chat_timeout", - use_id=logical_use_id, - provider="", - detail="", + next_actions.extend( + _handle_talent_terminal_locked( + use_id, + "talent_errored", + "reason", + "talent took too long", + ) ) - if ( - _current_chat_use_id == logical_use_id - and _current_chat_state is not None - and not _current_chat_state.get("raw_use_id") - ): - next_info = _clear_current_locked() - should_emit = True + should_emit = False else: return if should_emit: _emit_error(logical_use_id, "chat_timeout") - _run_next_action(next_info) + _run_next_actions(next_actions) def _active_talent_count_for_today_locked() -> int: @@ -1762,7 +1928,7 @@ def _location_for_trigger(day: str, trigger: dict[str, Any]) -> dict[str, str]: return _normalize_location("", "", "") -def _trigger_from_stream_event(event: dict[str, Any]) -> dict[str, Any]: +def _trigger_from_stream_event(day: str, event: dict[str, Any]) -> dict[str, Any]: kind = event.get("kind") if kind == "owner_message": return {"type": "owner_message", "message": event.get("text", "")} @@ -1783,6 +1949,7 @@ def _trigger_from_stream_event(event: dict[str, Any]) -> dict[str, Any]: event.get("name", "exec"), "summary", event.get("summary", ""), + origin=_reconstruct_origin_for_terminal(day, event), ) if kind == "talent_errored": return _talent_terminal_trigger( @@ -1792,10 +1959,59 @@ def _trigger_from_stream_event(event: dict[str, Any]) -> dict[str, Any]: "reason", event.get("reason", ""), reason_code=event.get("reason_code"), + origin=_reconstruct_origin_for_terminal(day, event), ) raise ValueError(f"unsupported trigger event: {kind}") +def _reconstruct_origin_for_terminal( + day: str, + terminal_event: dict[str, Any], +) -> dict[str, str] | None: + terminal_use_id = str(terminal_event.get("use_id") or "") + if not terminal_use_id: + return None + + latest_owner_message = "" + latest_dispatch_origin: dict[str, str] | None = None + origins_by_talent_use_id: dict[str, dict[str, str]] = {} + + for event in read_chat_events(day): + kind = event.get("kind") + use_id = str(event.get("use_id") or "") + + if event is terminal_event or ( + kind == terminal_event.get("kind") + and use_id == terminal_use_id + and event.get("ts") == terminal_event.get("ts") + ): + return origins_by_talent_use_id.get(terminal_use_id) + + if kind == "owner_message": + latest_owner_message = str(event.get("text") or "") + continue + + if kind == "sol_message": + if event.get("requested_target") is not None: + latest_dispatch_origin = { + "logical_use_id": str(event.get("use_id") or ""), + "ask": latest_owner_message, + } + continue + + if kind == "talent_queued" and use_id: + origins_by_talent_use_id[use_id] = { + "logical_use_id": str(event.get("chat_use_id") or ""), + "ask": str(event.get("ask") or ""), + } + continue + + if kind == "talent_spawned" and use_id and latest_dispatch_origin is not None: + origins_by_talent_use_id.setdefault(use_id, latest_dispatch_origin) + + return origins_by_talent_use_id.get(terminal_use_id) + + def _talent_terminal_trigger( kind: str, use_id: Any, @@ -1804,6 +2020,7 @@ def _talent_terminal_trigger( result_value: Any, *, reason_code: str | None = None, + origin: dict[str, str] | None = None, ) -> dict[str, Any]: trigger = { "type": kind, @@ -1813,6 +2030,8 @@ def _talent_terminal_trigger( } if reason_code: trigger["reason_code"] = reason_code + if origin is not None: + trigger["origin"] = dict(origin) return trigger diff --git a/solstone/convey/chat_stream.py b/solstone/convey/chat_stream.py index 674a0b272..73afc07dc 100644 --- a/solstone/convey/chat_stream.py +++ b/solstone/convey/chat_stream.py @@ -28,9 +28,10 @@ _CHAT_STREAM = "chat" _SEGMENT_WINDOW_MS = 300_000 _APPENDED_CHAT_PATHS: dict[int, Path] = {} # owner_message may carry optional `source`; extras flow through unchanged. -# sol_message may carry optional `thinking`, `offer`, `draft`, `sources`, and -# `answer_state`; result may carry optional `ticket_id`, `error`, `ambiguous`, and -# `cancelled`; talent_finished may carry optional `thinking`. Extras flow through +# sol_message may carry optional `thinking`, `offer`, `draft`, `sources`, +# `answer_state`, and folded-turn `origin`; result may carry optional `ticket_id`, +# `error`, `ambiguous`, and `cancelled`; talent_finished may carry optional +# `thinking`. talent_queued records a deferred spawn intent. Extras flow through # unchanged and are not part of the required-field tuples below. _VALID_KINDS = { "owner_message": ("text", "app", "path", "facet"), @@ -42,6 +43,16 @@ _VALID_KINDS = { "requested_task", ), "talent_spawned": ("use_id", "name", "task", "started_at"), + "talent_queued": ( + "use_id", + "name", + "task", + "queued_at", + "chat_use_id", + "ask", + "context", + "location", + ), "talent_finished": ("use_id", "name", "summary"), "talent_errored": ("use_id", "name", "reason"), "reflection_ready": ("day", "url"), @@ -233,6 +244,7 @@ def reduce_chat_state(day: str) -> dict[str, Any]: active_talents: dict[str, dict[str, Any]] = {} completed_talents: list[dict[str, Any]] = [] errored_talents: list[dict[str, Any]] = [] + queued_talents: dict[str, dict[str, Any]] = {} chat_error: dict[str, Any] | None = None queue_depth = 0 @@ -252,13 +264,24 @@ def reduce_chat_state(day: str) -> dict[str, Any]: "requested_task": event["requested_task"], "offer": event.get("offer"), "draft": event.get("draft"), + "origin": event.get("origin"), "sources": event.get("sources", []), "answer_state": event.get("answer_state", "answered"), } chat_error = None continue + if kind == "talent_queued": + queued_talents[str(event["use_id"])] = { + "use_id": event["use_id"], + "name": event["name"], + "task": event["task"], + "queued_at": event["queued_at"], + } + continue + if kind == "talent_spawned": + queued_talents.pop(str(event["use_id"]), None) active_talents[str(event["use_id"])] = { "use_id": event["use_id"], "name": event["name"], @@ -269,6 +292,7 @@ def reduce_chat_state(day: str) -> dict[str, Any]: continue if kind == "talent_finished": + queued_talents.pop(str(event["use_id"]), None) started = active_talents.pop(str(event["use_id"]), None) completed_talents.append( { @@ -283,6 +307,7 @@ def reduce_chat_state(day: str) -> dict[str, Any]: continue if kind == "talent_errored": + queued_talents.pop(str(event["use_id"]), None) active_talents.pop(str(event["use_id"]), None) errored_talents.append( { @@ -314,6 +339,13 @@ def reduce_chat_state(day: str) -> dict[str, Any]: str(talent["use_id"]), ), ), + "queued_talents": sorted( + queued_talents.values(), + key=lambda talent: ( + int(talent.get("queued_at", 0) or 0), + str(talent["use_id"]), + ), + ), "completed_talents": completed_talents, "errored_talents": errored_talents, "chat_error": chat_error, diff --git a/solstone/convey/sol_initiated/copy.py b/solstone/convey/sol_initiated/copy.py index 55be28ad6..89c59edf1 100644 --- a/solstone/convey/sol_initiated/copy.py +++ b/solstone/convey/sol_initiated/copy.py @@ -12,7 +12,6 @@ SURFACE_CONVEY = "convey" SOL_PINGED_OFFLINE_TOOLTIP = "sol-pinged but offline — refresh" TRIGGER_LABEL_SOL_INITIATED = "sol_initiated" -SYNTHETIC_TRIGGER_LABEL = "synthetic" THROTTLE_MUTE_WINDOW = "mute-window" THROTTLE_RATE_FLOOR = "rate-floor" diff --git a/solstone/talent/chat_context.py b/solstone/talent/chat_context.py index 9ca4e066c..f6cf1b6e4 100644 --- a/solstone/talent/chat_context.py +++ b/solstone/talent/chat_context.py @@ -12,7 +12,6 @@ from typing import Any from solstone.convey.chat_stream import read_chat_tail, reduce_chat_state from solstone.convey.sol_initiated.copy import ( KIND_SOL_CHAT_REQUEST, - SYNTHETIC_TRIGGER_LABEL, TRIGGER_LABEL_SOL_INITIATED, ) @@ -254,9 +253,6 @@ def _render_trigger_context( _append_terminal_trigger_context(lines, trigger_kind, payload) elif trigger_kind == "talent_errored": _append_terminal_trigger_context(lines, trigger_kind, payload) - elif trigger_kind == "synthetic-max-active": - if payload.get("reason"): - lines.append(f"- Reason: {payload['reason']}") else: if payload: for key, value in payload.items(): @@ -268,8 +264,6 @@ def _render_trigger_context( def _prompt_trigger_kind(trigger_kind: str | None) -> str: if trigger_kind == KIND_SOL_CHAT_REQUEST: return TRIGGER_LABEL_SOL_INITIATED - if trigger_kind == "synthetic-max-active": - return SYNTHETIC_TRIGGER_LABEL return str(trigger_kind or "") diff --git a/solstone/think/chat_cli.py b/solstone/think/chat_cli.py index d48979055..f41e0523c 100644 --- a/solstone/think/chat_cli.py +++ b/solstone/think/chat_cli.py @@ -138,13 +138,30 @@ def _render_post_error(exc: ConveyClientError) -> str: return "\n".join(lines) +def _origin_logical_use_id(message: dict) -> str: + origin = message.get("origin") + if not isinstance(origin, dict): + return "" + return str(origin.get("logical_use_id") or "") + + +def _is_fold_terminal_message(message: dict, use_id: str) -> bool: + return ( + message.get("requested_target") is None + and _origin_logical_use_id(message) == use_id + ) + + def _session_terminal(client: ConveyClient, use_id: str) -> dict | None: try: data = client.request("GET", "/api/chat/session") except ConveyClientError: return None latest = (data or {}).get("latest_sol_message") or {} - if str(latest.get("use_id") or "") != use_id: + if str(latest.get("use_id") or "") != use_id and not _is_fold_terminal_message( + latest, + use_id, + ): return None if latest.get("requested_target") is not None: return None @@ -322,6 +339,21 @@ def main() -> None: render_event_progress(msg) return + if ( + logical_use_id is not None + and event_name == "sol_message" + and _is_fold_terminal_message(msg, logical_use_id) + ): + with lock: + state["last_event_at"] = time.monotonic() + set_terminal( + { + "kind": "finish", + "result": str(msg.get("text") or ""), + } + ) + return + if logical_use_id is not None and event_name == "talent_finished": with lock: state["last_event_at"] = time.monotonic() diff --git a/tests/baselines/api/chat/session.json b/tests/baselines/api/chat/session.json index 8bb80ed29..5265b85b7 100644 --- a/tests/baselines/api/chat/session.json +++ b/tests/baselines/api/chat/session.json @@ -4,5 +4,6 @@ "completed_talents": [], "errored_talents": [], "latest_sol_message": null, - "queue_depth": 0 + "queue_depth": 0, + "queued_talents": [] } diff --git a/tests/test_chat_cli.py b/tests/test_chat_cli.py index 0fed9d1fe..8300cb38a 100644 --- a/tests/test_chat_cli.py +++ b/tests/test_chat_cli.py @@ -103,12 +103,14 @@ def _session_finish( use_id: str = USE_ID, text: str = "Recovered answer", requested_target: str | None = None, + origin: dict[str, Any] | None = None, ) -> dict[str, Any]: return { "latest_sol_message": { "use_id": use_id, "text": text, "requested_target": requested_target, + "origin": origin, } } @@ -326,6 +328,18 @@ def test_session_terminal_finish_only_current_use_id() -> None: ) is None ) + assert chat_cli._session_terminal( + FakeClient( + [ + _session_finish( + use_id=FOREIGN_USE_ID, + text="Folded answer", + origin={"logical_use_id": USE_ID, "ask": "hello"}, + ) + ] + ), + USE_ID, + ) == {"kind": "finish", "result": "Folded answer"} assert ( chat_cli._session_terminal( FakeClient([_session_finish(requested_target="exec")]), @@ -371,6 +385,34 @@ def test_main_live_sse_finish_prints_answer_without_fallback_warning( assert client.calls == [("POST", "/api/chat", {"message": "hello"})] +def test_main_live_sse_origin_fold_prints_answer(monkeypatch, capsys) -> None: + client = FakeClient([_post_success()]) + _install_main_fakes( + monkeypatch, + argv=["hello"], + client=client, + sse_chunks=[ + _frame( + { + "tract": "chat", + "event": "sol_message", + "use_id": FOREIGN_USE_ID, + "text": "Folded answer", + "requested_target": None, + "origin": {"logical_use_id": USE_ID, "ask": "hello"}, + } + ) + ], + ) + + chat_cli.main() + + captured = capsys.readouterr() + assert captured.out == "Folded answer\n" + assert captured.err == "" + assert client.calls == [("POST", "/api/chat", {"message": "hello"})] + + def test_main_live_sse_error_prints_terminal_error_without_session_poll( monkeypatch, capsys ) -> None: diff --git a/tests/test_chat_runtime.py b/tests/test_chat_runtime.py index ed40fb61d..8ac72f6ec 100644 --- a/tests/test_chat_runtime.py +++ b/tests/test_chat_runtime.py @@ -11,7 +11,11 @@ import pytest from flask import Flask from solstone.apps.chat.copy import CHAT_CLOSER_SUPPORT_SEND_FAILED -from solstone.convey.chat_stream import append_chat_event, read_chat_events +from solstone.convey.chat_stream import ( + append_chat_event, + read_chat_events, + reduce_chat_state, +) def _reset_chat_state(chat_module) -> None: @@ -104,7 +108,7 @@ def _append_recoverable_talent_events( ) -def test_chat_result_with_two_active_talents_retriggers_with_max_active_reason( +def test_chat_result_with_two_active_talents_queues_deferred_spawn( tmp_path, monkeypatch ): import solstone.convey.chat as chat @@ -129,7 +133,8 @@ def test_chat_result_with_two_active_talents_retriggers_with_max_active_reason( actions: list[dict] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -165,18 +170,29 @@ def test_chat_result_with_two_active_talents_retriggers_with_max_active_reason( } ) - assert actions - assert actions[-1]["kind"] == "chat" - assert actions[-1]["trigger"] == { - "type": "synthetic-max-active", - "reason": "max active — waiting for one to finish", - } + assert actions == [] sol_messages = [ e for e in read_chat_events(chat._today_day()) if e["kind"] == "sol_message" ] assert sol_messages[-1]["requested_target"] == "exec" assert sol_messages[-1]["requested_task"] == "research it" + queued = [ + e for e in read_chat_events(chat._today_day()) if e["kind"] == "talent_queued" + ] + assert queued[-1]["name"] == "exec" + assert queued[-1]["task"] == "research it" + assert queued[-1]["chat_use_id"] == "1713620000100" + assert queued[-1]["ask"] == "help" + assert queued[-1]["context"] == {"k": "v"} + assert reduce_chat_state(chat._today_day())["queued_talents"] == [ + { + "use_id": queued[-1]["use_id"], + "name": "exec", + "task": "research it", + "queued_at": queued[-1]["queued_at"], + } + ] def test_post_talent_finished_request_is_forced_terminal(tmp_path, monkeypatch): @@ -188,7 +204,8 @@ def test_post_talent_finished_request_is_forced_terminal(tmp_path, monkeypatch): actions: list[dict | None] = [] finishes: list[tuple[str, str]] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", @@ -226,7 +243,7 @@ def test_post_talent_finished_request_is_forced_terminal(tmp_path, monkeypatch): } ) - assert actions == [None] + assert actions == [] expected_text = ( "Here's what I have so far: Found three relevant notes. " "Want me to try a different angle?" @@ -253,7 +270,8 @@ def test_post_talent_errored_request_is_forced_terminal(tmp_path, monkeypatch): actions: list[dict | None] = [] finishes: list[tuple[str, str]] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", @@ -298,7 +316,7 @@ def test_post_talent_errored_request_is_forced_terminal(tmp_path, monkeypatch): "I couldn't finish that lookup — talent timed out waiting for provider " "response. Want to try a different angle, or rephrase the question?" ) - assert actions == [None] + assert actions == [] events = read_chat_events(chat._today_day()) sol_messages = [event for event in events if event["kind"] == "sol_message"] assert len(sol_messages) == 1 @@ -323,7 +341,8 @@ def test_post_support_talent_errored_request_uses_send_failed_closer( actions: list[dict | None] = [] finishes: list[tuple[str, str]] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", @@ -366,7 +385,7 @@ def test_post_support_talent_errored_request_uses_send_failed_closer( } ) - assert actions == [None] + assert actions == [] events = read_chat_events(chat._today_day()) sol_messages = [event for event in events if event["kind"] == "sol_message"] assert len(sol_messages) == 1 @@ -437,7 +456,8 @@ def test_owner_message_request_still_spawns_talent(tmp_path, monkeypatch): actions: list[dict] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -486,10 +506,6 @@ def test_owner_message_request_still_spawns_talent(tmp_path, monkeypatch): [ None, {}, - { - "type": "synthetic-max-active", - "reason": "max active — waiting for one to finish", - }, ], ) def test_non_post_talent_triggers_allow_dispatch(tmp_path, monkeypatch, trigger): @@ -500,7 +516,8 @@ def test_non_post_talent_triggers_allow_dispatch(tmp_path, monkeypatch, trigger) actions: list[dict] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -554,7 +571,8 @@ def test_cortex_finish_and_error_append_exec_terminal_events_by_use_id( actions: list[dict] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -577,6 +595,7 @@ def test_cortex_finish_and_error_append_exec_terminal_events_by_use_id( "target": "exec", "task": "summarize", "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "ask": "help", } chat._on_cortex_finish({"use_id": "1713623000001", "result": "done"}) @@ -584,7 +603,14 @@ def test_cortex_finish_and_error_append_exec_terminal_events_by_use_id( 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" + assert actions == [] + with chat._state_lock: + queued = chat._queued_triggers[-1] + assert queued["trigger"]["type"] == "talent_finished" + assert queued["trigger"]["origin"] == { + "logical_use_id": "1713623000000", + "ask": "help", + } _reset_chat_state(chat) actions.clear() @@ -602,6 +628,7 @@ def test_cortex_finish_and_error_append_exec_terminal_events_by_use_id( "target": "exec", "task": "summarize", "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "ask": "help", } chat._on_cortex_error( @@ -615,22 +642,30 @@ def test_cortex_finish_and_error_append_exec_terminal_events_by_use_id( 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" - assert actions[-1]["trigger"]["reason_code"] == "wall_clock_exceeded" + assert actions == [] + with chat._state_lock: + queued = chat._queued_triggers[-1] + assert queued["trigger"]["type"] == "talent_errored" + assert queued["trigger"]["reason"] == "boom" + assert queued["trigger"]["reason_code"] == "wall_clock_exceeded" + assert queued["trigger"]["origin"] == { + "logical_use_id": "1713624000000", + "ask": "help", + } def test_talent_errored_trigger_recovers_reason_code(): import solstone.convey.chat as chat trigger = chat._trigger_from_stream_event( + "20260420", { "kind": "talent_errored", "use_id": "1713624500001", "name": "support", "reason": "Traceback (most recent call last)", "reason_code": "wall_clock_exceeded", - } + }, ) assert trigger["type"] == "talent_errored" @@ -868,6 +903,7 @@ def test_recover_active_talents_repopulates_from_chat_stream(tmp_path, monkeypat "task": "research it", "trigger": "sol_message", "location": {"app": "home", "path": "/app/home", "facet": "work"}, + "ask": "Help me with this", } assert talent_use_id in chat._watchdog_timers assert len(timers) == 1 @@ -891,7 +927,8 @@ def test_late_talent_finish_after_recovery_routes_to_chat_continuation( actions: list[dict | None] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -918,8 +955,17 @@ def test_late_talent_finish_after_recovery_routes_to_chat_continuation( e for e in read_chat_events(chat._today_day()) if e["kind"] == "talent_finished" ] assert finished_events[-1]["use_id"] == talent_use_id - assert actions[-1]["logical_use_id"] == chat_use_id - assert actions[-1]["trigger"]["type"] == "talent_finished" + assert actions == [] + with chat._state_lock: + queued = chat._queued_triggers[-1] + assert queued["trigger"]["type"] == "talent_finished" + assert queued["trigger"]["origin"] == { + "logical_use_id": chat_use_id, + "ask": "Help me with this", + } + next_action = chat._clear_current_locked() + assert next_action["logical_use_id"] == queued["use_id"] + assert next_action["trigger"]["type"] == "talent_finished" def test_recovery_is_idempotent_for_active_talents(tmp_path, monkeypatch): @@ -999,7 +1045,8 @@ def test_chat_generate_schema_violation_retries_once_then_chat_errors( actions: list[dict | None] = [] emitted_errors: list[tuple[str, str]] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -1048,7 +1095,8 @@ def test_chat_generate_absorbs_prose_context_and_dispatches_talent( actions: list[dict | None] = [] emitted_errors: list[tuple[str, str]] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -1101,7 +1149,8 @@ def test_superseded_raw_finish_after_retry_is_dropped_without_warning( actions: list[dict | None] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -1162,7 +1211,8 @@ def test_superseded_raw_error_after_followup_rotation_is_dropped_without_warning actions: list[dict | None] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -1191,7 +1241,9 @@ def test_superseded_raw_error_after_followup_rotation_is_dropped_without_warning chat._on_cortex_finish({"use_id": "1713625200002", "result": "summary"}) with chat._state_lock: - followup_use_id = str(chat._current_chat_state["raw_use_id"]) + assert chat._current_chat_state["raw_use_id"] is None + queued = chat._queued_triggers[-1] + assert queued["trigger"]["type"] == "talent_finished" events_before = list(read_chat_events(chat._today_day())) with caplog.at_level("DEBUG"): @@ -1203,24 +1255,7 @@ def test_superseded_raw_error_after_followup_rotation_is_dropped_without_warning ) assert "unrouteable cortex event" not in caplog.text assert read_chat_events(chat._today_day()) == events_before - assert actions[0]["trigger"]["type"] == "talent_finished" - - chat._on_cortex_finish( - { - "use_id": followup_use_id, - "result": '{"message":"wrapped up","notes":"ok","talent_request":null}', - } - ) - - sol_messages = [ - event - for event in read_chat_events(chat._today_day()) - if event["kind"] == "sol_message" - ] - assert ( - sol_messages[-1]["text"] - == "Here's what I have so far: wrapped up Want me to try a different angle?" - ) + assert actions == [] def test_reserved_unknown_raw_use_id_still_warns(tmp_path, monkeypatch, caplog): @@ -1725,7 +1760,7 @@ def test_chat_watchdog_times_out_current_chat_generate(tmp_path, monkeypatch): assert raw_use_id not in chat._watchdog_timers -def test_chat_watchdog_times_out_active_talent_and_clears_blocked_chat( +def test_chat_watchdog_times_out_active_talent_and_queues_fold_when_chat_busy( tmp_path, monkeypatch ): import solstone.convey.chat as chat @@ -1734,14 +1769,10 @@ def test_chat_watchdog_times_out_active_talent_and_clears_blocked_chat( _reset_chat_state(chat) timers = _install_fake_timers(monkeypatch) - emitted_errors: list[tuple[str, str]] = [] monkeypatch.setattr( "solstone.convey.chat._emit_cortex_event", lambda *args, **kwargs: None ) - monkeypatch.setattr( - "solstone.convey.chat._emit_error", - lambda use_id, reason: emitted_errors.append((use_id, reason)), - ) + monkeypatch.setattr("solstone.convey.chat._emit_error", lambda *args: None) monkeypatch.setattr( "solstone.convey.utils.spawn_agent", lambda *args, **kwargs: kwargs["use_id"] ) @@ -1760,6 +1791,7 @@ def test_chat_watchdog_times_out_active_talent_and_clears_blocked_chat( "target": "exec", "task": "summarize", "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "ask": "help", } chat._run_next_action( @@ -1777,28 +1809,38 @@ def test_chat_watchdog_times_out_active_talent_and_clears_blocked_chat( assert "1713629000001" in chat._watchdog_timers timers[-1].fire() - errors = [ + chat_errors = [ event for event in read_chat_events(chat._today_day()) if event["kind"] == "chat_error" ] - assert emitted_errors == [("1713629000000", "chat_timeout")] - assert errors[-1]["use_id"] == "1713629000000" - assert errors[-1]["reason"] == "chat_timeout" + talent_errors = [ + event + for event in read_chat_events(chat._today_day()) + if event["kind"] == "talent_errored" + ] + assert chat_errors == [] + assert talent_errors[-1]["use_id"] == "1713629000001" + assert talent_errors[-1]["reason"] == "talent took too long" with chat._state_lock: assert "1713629000001" not in chat._active_talents - assert chat._current_chat_use_id is None - assert chat._current_chat_state is None + assert chat._current_chat_use_id == "1713629000000" + assert chat._queued_triggers[-1]["trigger"]["type"] == "talent_errored" + assert chat._queued_triggers[-1]["trigger"]["origin"] == { + "logical_use_id": "1713629000000", + "ask": "help", + } assert "1713629000001" not in chat._watchdog_timers -def test_chat_watchdog_marks_timed_out_talent_result_as_errored(tmp_path, monkeypatch): +def test_chat_watchdog_timeout_folds_talent_error_with_origin(tmp_path, monkeypatch): import solstone.convey.chat as chat _setup_journal(tmp_path, monkeypatch) _reset_chat_state(chat) timers = _install_fake_timers(monkeypatch) + actions: list[dict] = [] monkeypatch.setattr( "solstone.convey.chat._emit_cortex_event", lambda *args, **kwargs: None ) @@ -1806,7 +1848,8 @@ def test_chat_watchdog_marks_timed_out_talent_result_as_errored(tmp_path, monkey "solstone.convey.chat._emit_error", lambda *args, **kwargs: None ) monkeypatch.setattr( - "solstone.convey.utils.spawn_agent", lambda *args, **kwargs: kwargs["use_id"] + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) with chat._state_lock: @@ -1822,36 +1865,18 @@ def test_chat_watchdog_marks_timed_out_talent_result_as_errored(tmp_path, monkey ) with chat._state_lock: - chat._current_chat_use_id = logical_use_id - chat._current_chat_state = { - "raw_use_id": None, - "raw_use_ids_seen": set(), - "trigger": {"type": "owner_message", "message": "help"}, - "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, - "retry_count": 0, - } chat._active_talents[talent_use_id] = { "chat_use_id": logical_use_id, "target": "exec", "task": "summarize", "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "ask": "why did it hang?", } - - chat._run_next_action( - { - "kind": "talent", - "logical_use_id": logical_use_id, - "target": "exec", - "use_id": talent_use_id, - "task": "summarize", - "context": {}, - "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, - } - ) + chat._arm_watchdog_locked(talent_use_id, "talent", logical_use_id) timers[-1].fire() - parent_errors = [ + chat_errors = [ event for event in read_chat_events(chat._today_day()) if event["kind"] == "chat_error" @@ -1861,10 +1886,34 @@ def test_chat_watchdog_marks_timed_out_talent_result_as_errored(tmp_path, monkey for event in read_chat_events(chat._today_day()) if event["kind"] == "talent_errored" ] - assert parent_errors[-1]["use_id"] == logical_use_id - assert parent_errors[-1]["reason"] == "chat_timeout" + assert chat_errors == [] assert talent_errors[-1]["use_id"] == talent_use_id assert talent_errors[-1]["reason"] == "talent took too long" + assert actions and actions[-1]["kind"] == "chat" + assert actions[-1]["logical_use_id"] != logical_use_id + assert actions[-1]["trigger"]["origin"] == { + "logical_use_id": logical_use_id, + "ask": "why did it hang?", + } + + chat._on_cortex_finish( + { + "use_id": actions[-1]["raw_use_id"], + "result": {"message": "ignored", "notes": "ok", "talent_request": None}, + } + ) + + folded = [ + event + for event in read_chat_events(chat._today_day()) + if event["kind"] == "sol_message" + ][-1] + assert folded["requested_target"] is None + assert folded["origin"] == { + "logical_use_id": logical_use_id, + "ask": "why did it hang?", + } + assert folded["text"].startswith("I couldn't finish that lookup") def test_cortex_finish_logs_warning_for_unrouteable_use_id( @@ -2354,7 +2403,8 @@ def test_read_dispatch_spawns_read_talent(tmp_path, monkeypatch): actions: list[dict | None] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -2408,7 +2458,8 @@ def test_read_finish_retriggers_chat_like_exec(tmp_path, monkeypatch): actions: list[dict | None] = [] monkeypatch.setattr( - "solstone.convey.chat._run_next_action", lambda action: actions.append(action) + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, ) monkeypatch.setattr( "solstone.convey.chat._emit_finish", lambda *args, **kwargs: None @@ -2431,6 +2482,7 @@ def test_read_finish_retriggers_chat_like_exec(tmp_path, monkeypatch): "target": "read", "task": "Reflect on the week", "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "ask": "help", } chat._on_cortex_finish({"use_id": "1713627000001", "result": "A reflective note"}) @@ -2439,5 +2491,12 @@ def test_read_finish_retriggers_chat_like_exec(tmp_path, monkeypatch): e for e in read_chat_events(chat._today_day()) if e["kind"] == "talent_finished" ] assert finished_events[-1]["name"] == "read" - assert actions[-1]["trigger"]["type"] == "talent_finished" - assert actions[-1]["trigger"]["name"] == "read" + assert actions == [] + with chat._state_lock: + queued = chat._queued_triggers[-1] + assert queued["trigger"]["type"] == "talent_finished" + assert queued["trigger"]["name"] == "read" + assert queued["trigger"]["origin"] == { + "logical_use_id": "1713627000000", + "ask": "help", + } diff --git a/tests/test_chat_stream.py b/tests/test_chat_stream.py index b59370133..508d3d965 100644 --- a/tests/test_chat_stream.py +++ b/tests/test_chat_stream.py @@ -530,6 +530,7 @@ def test_reduce_chat_state_extracts_latest_sol_and_active_talents( "requested_task": "compare drafts", "offer": None, "draft": None, + "origin": None, "sources": [], "answer_state": "answered", } @@ -566,6 +567,7 @@ def test_reduce_chat_state_extracts_latest_sol_and_active_talents( "detail": "", } assert reduced["queue_depth"] == 0 + assert reduced["queued_talents"] == [] def test_reduce_chat_state_enriches_talent_labels(tmp_path, monkeypatch): @@ -901,6 +903,66 @@ def test_reduce_chat_state_returns_last_queue_depth(tmp_path, monkeypatch): assert reduce_chat_state("20260420")["queue_depth"] == 1 +def test_reduce_chat_state_surfaces_sol_message_origin(tmp_path, monkeypatch): + _setup_journal(tmp_path, monkeypatch) + start = _ms(2026, 4, 20, 12, 0, 0) + origin = {"logical_use_id": "chat-dispatch", "ask": "look this up"} + + append_chat_event( + "sol_message", + ts=start, + use_id="chat-fold", + text="folded answer", + notes="", + requested_target=None, + requested_task=None, + origin=origin, + ) + + reduced = reduce_chat_state("20260420") + assert reduced["latest_sol_message"]["origin"] == origin + + +def test_reduce_chat_state_tracks_queued_talents_until_spawn(tmp_path, monkeypatch): + _setup_journal(tmp_path, monkeypatch) + start = _ms(2026, 4, 20, 12, 0, 0) + + append_chat_event( + "talent_queued", + ts=start, + use_id="talent-queued", + name="exec", + task="research", + queued_at=start, + chat_use_id="chat-dispatch", + ask="research this", + context={"scope": "today"}, + location={"app": "sol", "path": "/app/sol", "facet": "work"}, + ) + + assert reduce_chat_state("20260420")["queued_talents"] == [ + { + "use_id": "talent-queued", + "name": "exec", + "task": "research", + "queued_at": start, + } + ] + + append_chat_event( + "talent_spawned", + ts=start + 1_000, + use_id="talent-queued", + name="exec", + task="research", + started_at=start + 1_000, + ) + + reduced = reduce_chat_state("20260420") + assert reduced["queued_talents"] == [] + assert reduced["active_talents"][0]["use_id"] == "talent-queued" + + def test_append_reflection_ready_event(tmp_path, monkeypatch): _setup_journal(tmp_path, monkeypatch) ts = _ms(2026, 4, 20, 12, 0, 0) @@ -971,6 +1033,86 @@ def test_find_unresponded_trigger_talent_finished(tmp_path, monkeypatch): assert trigger["summary"] == "done" +def test_find_unresponded_trigger_after_dispatch_ack_and_spawn(tmp_path, monkeypatch): + _setup_journal(tmp_path, monkeypatch) + start = _ms(2026, 4, 20, 12, 0, 0) + + append_chat_event( + "owner_message", + ts=start, + text="look this up", + app="sol", + path="/chat", + facet="work", + ) + append_chat_event( + "sol_message", + ts=start + 1_000, + use_id="chat-dispatch", + text="working", + notes="", + requested_target="exec", + requested_task="research", + ) + append_chat_event( + "talent_spawned", + ts=start + 2_000, + use_id="talent-1", + name="exec", + task="research", + started_at=start + 2_000, + ) + append_chat_event( + "talent_finished", + ts=start + 3_000, + use_id="talent-1", + name="exec", + summary="done", + ) + + trigger = find_unresponded_trigger("20260420") + assert trigger is not None + assert trigger["kind"] == "talent_finished" + assert trigger["use_id"] == "talent-1" + + +def test_talent_queued_is_not_an_unresponded_trigger(tmp_path, monkeypatch): + _setup_journal(tmp_path, monkeypatch) + start = _ms(2026, 4, 20, 12, 0, 0) + + append_chat_event( + "owner_message", + ts=start, + text="look this up", + app="sol", + path="/chat", + facet="work", + ) + append_chat_event( + "sol_message", + ts=start + 1_000, + use_id="chat-dispatch", + text="working", + notes="", + requested_target="exec", + requested_task="research", + ) + append_chat_event( + "talent_queued", + ts=start + 2_000, + use_id="talent-queued", + name="exec", + task="research", + queued_at=start + 2_000, + chat_use_id="chat-dispatch", + ask="look this up", + context={}, + location={"app": "sol", "path": "/chat", "facet": "work"}, + ) + + assert find_unresponded_trigger("20260420") is None + + def test_find_unresponded_trigger_resolved(tmp_path, monkeypatch): _setup_journal(tmp_path, monkeypatch) start = _ms(2026, 4, 20, 12, 0, 0) diff --git a/tests/test_convey_chat.py b/tests/test_convey_chat.py index c674e39a9..55cf6b185 100644 --- a/tests/test_convey_chat.py +++ b/tests/test_convey_chat.py @@ -253,6 +253,7 @@ def test_cortex_thinking_reaches_talent_finished(chat_client, monkeypatch): "target": "exec", "task": "research", "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "ask": "help", } chat._handle_callosum_message( @@ -428,6 +429,275 @@ def test_non_outbound_routes_are_not_gated(chat_client, monkeypatch, target): assert "offer" not in sol_message +def test_dispatch_clears_turn_and_records_origin_ask(chat_client, monkeypatch): + import solstone.convey.chat as chat + + actions: list[dict] = [] + monkeypatch.setattr( + "solstone.convey.chat._emit_cortex_event", lambda *_args, **_kwargs: None + ) + monkeypatch.setattr( + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, + ) + _set_current_chat(chat, "logical-chat", "raw-chat") + + chat._on_cortex_finish( + { + "use_id": "raw-chat", + "result": _talent_route_result("exec", "do the work"), + } + ) + + with chat._state_lock: + assert chat._current_chat_state is None + assert chat._current_chat_use_id is None + talent_state = next(iter(chat._active_talents.values())) + assert talent_state["ask"] == "help" + assert actions[-1]["kind"] == "talent" + events = read_chat_events(date.today().strftime("%Y%m%d")) + sol_message = next(event for event in events if event["kind"] == "sol_message") + assert sol_message["requested_target"] == "exec" + assert sol_message["requested_task"] == "do the work" + assert "origin" not in sol_message + + +def test_second_owner_message_starts_after_dispatch_not_queued( + chat_client, monkeypatch +): + import solstone.convey.chat as chat + + starts: list[dict] = [] + monkeypatch.setattr( + "solstone.think.identity.ensure_identity_directory", lambda: None + ) + monkeypatch.setattr( + "solstone.convey.chat._spawn_chat_generate", + lambda action: starts.append(action) or ChatSpawnResult(ok=True), + ) + monkeypatch.setattr("solstone.convey.chat._spawn_talent", lambda _action: True) + monkeypatch.setattr( + "solstone.convey.chat._emit_cortex_event", lambda *_args, **_kwargs: None + ) + monkeypatch.setattr( + "solstone.convey.chat._arm_watchdog_locked", lambda *_args, **_kwargs: None + ) + + first = _post_chat_message(chat_client, "first") + assert first.status_code == 200 + assert first.get_json()["queued"] is False + chat._on_cortex_finish( + { + "use_id": starts[0]["raw_use_id"], + "result": _talent_route_result("exec", "do the work"), + } + ) + + second = _post_chat_message(chat_client, "second") + + assert second.status_code == 200 + assert second.get_json()["queued"] is False + assert [start["trigger"]["message"] for start in starts] == ["first", "second"] + + +def test_empty_dispatch_ack_uses_nonraising_liveness_backstop(chat_client, monkeypatch): + import solstone.convey.chat as chat + + monkeypatch.setattr( + "solstone.convey.chat._emit_cortex_event", lambda *_args, **_kwargs: None + ) + monkeypatch.setattr("solstone.convey.chat._run_next_action", lambda _action: None) + _set_current_chat(chat, "logical-chat", "raw-chat") + + chat._on_cortex_finish( + { + "use_id": "raw-chat", + "result": { + "message": "", + "notes": "ok", + "talent_request": {"target": "exec", "task": "research"}, + }, + } + ) + + sol_message = next( + event + for event in read_chat_events(date.today().strftime("%Y%m%d")) + if event["kind"] == "sol_message" + ) + assert sol_message["text"] == "Making that change… research" + assert chat._dispatch_ack_text("unknown", None, "") == "unknown" + + +def test_talent_finish_folds_fresh_turn_with_origin(chat_client, monkeypatch): + import solstone.convey.chat as chat + + actions: list[dict] = [] + monkeypatch.setattr( + "solstone.convey.chat._emit_cortex_event", lambda *_args, **_kwargs: None + ) + monkeypatch.setattr( + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, + ) + with chat._state_lock: + chat._active_talents["talent-raw"] = { + "chat_use_id": "dispatch-chat", + "target": "exec", + "task": "research", + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "ask": "what happened?", + } + + chat._on_cortex_finish({"use_id": "talent-raw", "result": "summary"}) + + assert actions and actions[-1]["kind"] == "chat" + assert actions[-1]["logical_use_id"] != "dispatch-chat" + assert actions[-1]["trigger"]["origin"] == { + "logical_use_id": "dispatch-chat", + "ask": "what happened?", + } + chat._on_cortex_finish( + { + "use_id": actions[-1]["raw_use_id"], + "result": { + "message": "Here is the answer with enough detail to stand alone.", + "notes": "ok", + "talent_request": None, + }, + } + ) + + sol_messages = [ + event + for event in read_chat_events(date.today().strftime("%Y%m%d")) + if event["kind"] == "sol_message" + ] + assert sol_messages[-1]["origin"] == { + "logical_use_id": "dispatch-chat", + "ask": "what happened?", + } + assert sol_messages[-1]["requested_target"] is None + + +@pytest.mark.parametrize("spawn_failure", (False, True)) +def test_talent_failure_paths_fold_with_origin( + chat_client, + monkeypatch, + spawn_failure, +): + import solstone.convey.chat as chat + + actions: list[dict] = [] + monkeypatch.setattr( + "solstone.convey.chat._emit_cortex_event", lambda *_args, **_kwargs: None + ) + monkeypatch.setattr( + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, + ) + with chat._state_lock: + chat._active_talents["talent-raw"] = { + "chat_use_id": "dispatch-chat", + "target": "exec", + "task": "research", + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "ask": "what failed?", + } + + if spawn_failure: + chat._handle_talent_spawn_failure( + { + "kind": "talent", + "logical_use_id": "dispatch-chat", + "target": "exec", + "use_id": "talent-raw", + "task": "research", + "context": {}, + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + } + ) + else: + chat._on_cortex_error({"use_id": "talent-raw", "error": "boom"}) + + assert actions and actions[-1]["kind"] == "chat" + assert actions[-1]["trigger"]["origin"] == { + "logical_use_id": "dispatch-chat", + "ask": "what failed?", + } + assert actions[-1]["trigger"]["type"] == "talent_errored" + chat._on_cortex_finish( + { + "use_id": actions[-1]["raw_use_id"], + "result": {"message": "ignored", "notes": "ok", "talent_request": None}, + } + ) + + folded = [ + event + for event in read_chat_events(date.today().strftime("%Y%m%d")) + if event["kind"] == "sol_message" + ][-1] + assert folded["origin"] == { + "logical_use_id": "dispatch-chat", + "ask": "what failed?", + } + assert folded["answer_state"] == "failed" + assert folded["text"].startswith("I couldn't finish that lookup") + + +def test_deferred_spawn_promotes_when_talent_slot_frees(chat_client, monkeypatch): + import solstone.convey.chat as chat + + actions: list[dict] = [] + monkeypatch.setattr( + "solstone.convey.chat._emit_cortex_event", lambda *_args, **_kwargs: None + ) + monkeypatch.setattr( + "solstone.convey.chat._run_next_action", + lambda action: actions.append(action) if action is not None else None, + ) + day = date.today().strftime("%Y%m%d") + for use_id in ("active-1", "active-2"): + append_chat_event( + "talent_spawned", + use_id=use_id, + name="exec", + task="existing", + started_at=int(time.time() * 1000), + ) + with chat._state_lock: + for use_id in ("active-1", "active-2"): + chat._active_talents[use_id] = { + "chat_use_id": f"chat-{use_id}", + "target": "exec", + "task": "existing", + "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "ask": "old ask", + } + _set_current_chat(chat, "logical-chat", "raw-chat") + + chat._on_cortex_finish( + { + "use_id": "raw-chat", + "result": _talent_route_result("exec", "deferred"), + } + ) + + queued = reduce_chat_state(day)["queued_talents"] + assert len(queued) == 1 + queued_use_id = queued[0]["use_id"] + chat._on_cortex_finish({"use_id": "active-1", "result": "done"}) + + reduced = reduce_chat_state(day) + assert reduced["queued_talents"] == [] + assert any( + talent["use_id"] == queued_use_id for talent in reduced["active_talents"] + ) + assert actions[-1]["kind"] == "talent" + assert actions[-1]["use_id"] == queued_use_id + + def test_clean_support_finish_with_pending_draft_emits_marker(chat_client, monkeypatch): import solstone.convey.chat as chat @@ -1742,9 +2012,7 @@ def test_thinking_buffer_evicted_on_retry_rotation(chat_client, monkeypatch): assert sol_message["thinking"]["content"] == "new thought" -def test_thinking_buffer_clear_across_synthetic_max_active_rotation( - chat_client, monkeypatch -): +def test_thinking_buffer_clear_across_at_cap_defer(chat_client, monkeypatch): import solstone.convey.chat as chat monkeypatch.setattr( @@ -1787,17 +2055,29 @@ def test_thinking_buffer_clear_across_synthetic_max_active_rotation( ) with chat._state_lock: - raw_new = str(chat._current_chat_state["raw_use_id"]) + assert chat._current_chat_state is None + assert chat._current_chat_use_id is None assert "raw-old" not in chat._thinking_buffers - chat._handle_callosum_message( + events = read_chat_events(date.today().strftime("%Y%m%d")) + sol_message = next(event for event in events if event["kind"] == "sol_message") + assert sol_message["text"] == "checking" + assert sol_message["requested_target"] == "exec" + assert sol_message["requested_task"] == "research" + assert sol_message["thinking"]["content"] == "old thought" + queued = next(event for event in events if event["kind"] == "talent_queued") + assert queued["name"] == "exec" + assert queued["task"] == "research" + assert queued["chat_use_id"] == "logical-chat" + assert queued["ask"] == "help" + assert queued["context"] == {} + assert reduce_chat_state(date.today().strftime("%Y%m%d"))["queued_talents"] == [ { - "tract": "cortex", - "event": "thinking", - "use_id": raw_new, - "summary": "synthetic thought", + "use_id": queued["use_id"], + "name": "exec", + "task": "research", + "queued_at": queued["queued_at"], } - ) - assert chat._thinking_buffers[raw_new] == ["synthetic thought"] + ] def test_late_thinking_arrival_drops_without_mutating_events( @@ -1839,6 +2119,7 @@ def test_talent_error_evicts_but_does_not_attach_thinking(chat_client, monkeypat "target": "exec", "task": "research", "location": {"app": "sol", "path": "/app/sol", "facet": "work"}, + "ask": "help", } chat._handle_callosum_message( { @@ -2271,6 +2552,64 @@ def test_chat_session_retries_unresolved_trigger_when_idle(chat_client, monkeypa assert approvals == [None] +def test_chat_session_reconstructs_origin_for_unresponded_terminal( + chat_client, + monkeypatch, +): + + day = "20260420" + monkeypatch.setattr("solstone.convey.chat._today_day", lambda: day) + start = _ms(2026, 4, 20, 12, 0, 0) + append_chat_event( + "owner_message", + ts=start, + text="look this up", + app="sol", + path="/app/sol", + facet="work", + ) + append_chat_event( + "sol_message", + ts=start + 1_000, + use_id="dispatch-chat", + text="working", + notes="", + requested_target="exec", + requested_task="research", + ) + append_chat_event( + "talent_spawned", + ts=start + 2_000, + use_id="talent-raw", + name="exec", + task="research", + started_at=start + 2_000, + ) + append_chat_event( + "talent_finished", + ts=start + 3_000, + use_id="talent-raw", + name="exec", + summary="done", + ) + + starts: list[dict] = [] + monkeypatch.setattr( + "solstone.convey.chat._spawn_chat_generate", + lambda action: starts.append(action) or ChatSpawnResult(ok=True), + ) + + response = chat_client.get("/api/chat/session") + + assert response.status_code == 200 + assert len(starts) == 1 + assert starts[0]["trigger"]["type"] == "talent_finished" + assert starts[0]["trigger"]["origin"] == { + "logical_use_id": "dispatch-chat", + "ask": "look this up", + } + + def test_chat_session_retries_again_when_spawn_fails_and_trigger_remains_unresolved( chat_client, monkeypatch ):