diff --git a/backend/app/gateway/AGENTS.md b/backend/app/gateway/AGENTS.md index 09fe21bf6..186cff6e1 100644 --- a/backend/app/gateway/AGENTS.md +++ b/backend/app/gateway/AGENTS.md @@ -159,7 +159,7 @@ startup gate rejects process-local memory and JSONL event stores when - Redis `StreamBridge` keys use a rolling retained-buffer TTL (`stream_bridge.stream_ttl_seconds`, refreshed on `publish()` / `publish_end()`) as a leak safety net, not as a run timeout. Startup and lease-driven periodic orphan recovery share one Gateway stream-terminalization path: after `RunManager` durably marks a run `error` with `stop_reason=orphan_recovered`, Gateway publishes `END_SENTINEL` and schedules stream cleanup. The periodic store scan, per-row status writes, and Gateway callback run as one supervised single-flight task, so a slow pass is skipped at the next interval instead of piling up or pausing the sole lease-renewal loop. Store retries have bounded attempts/backoff; an individual operation still relies on the database driver/pool timeout. `RunManager.shutdown()` gives active user runs priority within its shared deadline, then drains or cancels orphan recovery. Gateway tracks delayed recovered-stream cleanups and converts unfinished delays to immediate deletes before closing the bridge; the Redis TTL remains the outage safety net. Only startup recovery, before the runtime yields to requests, projects the latest affected thread to `error`; periodic recovery deliberately avoids that non-atomic projection because `ThreadMetaStore` has no `latest_run_id` conditional-update contract. Store-only SSE and `/wait` consumers wait for the bridge's real END marker after an ordinary durable terminal status, because status persistence can precede tail events. The explicit `orphan_recovered` signal is the only heartbeat fallback: its publisher is known to be gone, so it supplies the liveness boundary if END publication fails or the retained key expires. Malformed `Last-Event-ID` reconnect values live-tail new Redis events rather than replaying the retained buffer. Keep cross-component recovery orchestration in Gateway through the generic `RunManager.on_orphans_recovered` callback; do not introduce a harness-to-app dependency. Callback failure warnings include every recovered `run_id` so operators can identify rows whose Gateway-side terminalization needs inspection. - Thread-scoped run creation accepts an optional `Idempotency-Key` header on create, stream, and wait. Gateway hashes the caller key with the authenticated owner and `thread_id` before passing it to `RunManager`, whose persistence index is process-wide; never pass an unscoped external key to that index. The same scoped key is shared across `/runs`, `/runs/stream`, and `/runs/wait`; a reused admission whose stored `input` or `assistant_id` differs from the retry returns 409. `/wait` must not treat `task is None` as completion: `store_only` records without a cross-process bridge return durable `status`/`error` instead of serializing the current checkpoint; otherwise wait on the bridge. An idempotent reuse must not serialize the latest thread checkpoint as this run's result — a later run on the same thread may have advanced the head — so reused `/wait` returns durable `status`/`error`. Capture that reuse decision before awaiting completion; `idempotency_reused` is sticky on the shared cached record and an overlapping retry must not suppress the original creating request's checkpoint. After observing completion, refresh store-backed `status`/`error` before returning them — a hydrated peer record still holds admission-time fields. A creating-endpoint retry of a terminal record whose stream is gone emits SSE `gap`/`stream_replay_gap` with `recovery: reload_durable_state` rather than a bare `end`; observer joins of that same record still emit `end`. That gap is opt-in via `sse_consumer(..., emit_gap_on_missing_stream=True)` from thread-scoped `/runs/stream` on this request's reuse — do not key it off `apply_on_disconnect` or the sticky `idempotency_reused` flag. Default `sse_consumer` callers, including stateless `/api/runs/stream`, still emit `end`. A reused still-running `store_only` record on a process-local bridge returns 409 from `/stream` with no `Retry-After`, matching `join`. Missing headers preserve ordinary non-idempotent admission. Stateless `/api/runs/*` stays outside this contract because a request without an explicit thread creates a fresh temporary thread before admission. - Thread-scoped run creation accepts `checkpoint` / `checkpoint_id`; Gateway validates the checkpoint belongs to the request thread before writing `checkpoint_id` / `checkpoint_ns` into `config.configurable` for LangGraph branching. In `delta` checkpoint mode the worker rewrites that fork into a linear head write before the graph starts (see "A delta-mode run cannot fork" under Checkpoint Channel Modes), because delta state for a fork replays the abandoned sibling's writes. -- Thread-scoped Gateway runs evaluate an active `ThreadState.goal` after the visible turn completes. `runtime/goal.py` asks a non-thinking evaluator model to judge only visible conversation evidence and return a typed blocker; the evaluator model is created once per run and reused across hidden continuation checks. The evaluator runs after the graph root's tracing scope has already closed, so `create_goal_evaluator_model`/`evaluate_goal_completion` attach their own model-level tracing callbacks (`attach_tracing=True`) and inject Langfuse trace metadata (`thread_id`/`user_id`/`deerflow_trace_id`) directly onto the `ainvoke` call — the same standalone-caller pattern as `oneshot_llm.run_oneshot_llm` and `MemoryUpdater` (see Tracing System below). Satisfied goals are cleared; every non-satisfied evaluation — continuable or stand-down — is persisted with `last_evaluation` (the blocker, reason, and evidence summary; outcomes that stop the loop additionally record a `stand_down_reason` for observability), but only `goal_not_met_yet` evaluations are streamed as hidden `HumanMessage` continuations, and only when a durable assistant end-of-turn checkpoint exists, the run has not been aborted, the run has not hit its token-budget hard stop (`stop_reason == "token_capped"`; the goal stands down with `stand_down_reason: "token_capped"`), the thread did not change during evaluation, and the no-progress breaker has not fired. The continuation cap is 8 — a hard maximum in the `0`–`8` range; callers requesting more are clamped (`set_goal`/TUI) or rejected with 422 (`PUT /goal`). The no-progress breaker keys on the latest visible assistant evidence (not the evaluator's free-text reason, which an LLM rewords every turn), so two consecutive continuations that add no new visible assistant output stop the loop after 2 attempts. Model-response cleanup helpers such as think-block stripping and code-fence stripping live in `deerflow.utils.llm_text` so `runtime/goal.py` and Gateway suggestion parsing share the same JSON-prep behavior. +- Thread-scoped Gateway runs evaluate an active `ThreadState.goal` after the visible turn completes. `runtime/goal.py` asks a non-thinking evaluator model to judge only visible conversation evidence and return a typed blocker; the evaluator model is created once per run and reused across hidden continuation checks. The evaluator runs after the graph root's tracing scope has already closed, so `create_goal_evaluator_model`/`evaluate_goal_completion` attach their own model-level tracing callbacks (`attach_tracing=True`) and inject Langfuse trace metadata (`thread_id`/`user_id`/`deerflow_trace_id`) directly onto the `ainvoke` call — the same standalone-caller pattern as `oneshot_llm.run_oneshot_llm` and `MemoryUpdater` (see Tracing System below). Satisfied goals are cleared; every non-satisfied evaluation — continuable or stand-down — is persisted with `last_evaluation` (the blocker, reason, and evidence summary; outcomes that stop the loop additionally record a `stand_down_reason` for observability), but only `goal_not_met_yet` evaluations are streamed as hidden `HumanMessage` continuations, and only when the turn did not end on an unanswered Human Input Card (`ask_clarification` or the sandbox network prompt; the evaluator never sees those ToolMessages, so the goal stands down with `blocker: "needs_user_input"` without calling it), a durable assistant end-of-turn checkpoint exists, the run has not been aborted, the run has not hit its token-budget hard stop (`stop_reason == "token_capped"`; the goal stands down with `stand_down_reason: "token_capped"`), the thread did not change during evaluation, and the no-progress breaker has not fired. The continuation cap is 8 — a hard maximum in the `0`–`8` range; callers requesting more are clamped (`set_goal`/TUI) or rejected with 422 (`PUT /goal`). The no-progress breaker keys on the latest visible assistant evidence (not the evaluator's free-text reason, which an LLM rewords every turn), so two consecutive continuations that add no new visible assistant output stop the loop after 2 attempts. Model-response cleanup helpers such as think-block stripping and code-fence stripping live in `deerflow.utils.llm_text` so `runtime/goal.py` and Gateway suggestion parsing share the same JSON-prep behavior. - Run event stream changes must keep producer code, `deerflow/constants.py`, `runtime/events/catalog.py`, `contracts/run_event_stream_contract.json`, `backend/docs/RUN_EVENT_STREAM.md`, and `tests/test_run_event_stream_contract.py` in sync. The dependency-free constants module owns the persisted envelope limits (`event_type` 32 characters, `category` 16) and cross-layer workspace event identity; the catalog owns validated runtime definitions and categories. Dynamic middleware tags are limited to 21 characters after the `middleware:` prefix. The JSON contract owns payload schemas, backend-specific storage semantics, legacy aliases, and compatibility rules; conformance tests require both views and all producer groups to agree. `run.end.content` remains opaque and may retain nested Python values in memory while JSONL/database stores stringify non-JSON nested values, so consumers must not assume backend-identical nested output representations. Proxied through nginx: `/api/langgraph/*` → Gateway LangGraph-compatible runtime, all other `/api/*` → Gateway REST APIs. diff --git a/backend/packages/harness/deerflow/runtime/runs/worker.py b/backend/packages/harness/deerflow/runtime/runs/worker.py index aaead6fb1..89dd4da3f 100644 --- a/backend/packages/harness/deerflow/runtime/runs/worker.py +++ b/backend/packages/harness/deerflow/runtime/runs/worker.py @@ -1820,6 +1820,23 @@ def _has_durable_goal_turn_receipt(checkpoint_tuple: Any, messages: list[Any]) - return _message_type(visible_messages[-1]) == "ai" +def _ends_on_human_input_request(messages: list[Any]) -> bool: + """Return true when the turn ended on a Human Input Card the user has not answered. + + ``ask_clarification`` and the sandbox network prompt put the request in a + ToolMessage artifact and end the graph there, so it sits in the trailing run of + tool results. The goal evaluator only reads human and AI text and never sees it. + """ + for message in reversed(messages): + if _message_type(message) != "tool": + return False + artifact = message.get("artifact") if isinstance(message, dict) else getattr(message, "artifact", None) + human_input = artifact.get("human_input") if isinstance(artifact, Mapping) else None + if isinstance(human_input, Mapping) and human_input.get("kind") == "human_input_request": + return True + return False + + def _stand_down_reason(goal: GoalState, evaluation: GoalEvaluation, no_progress_count: int) -> str | None: if evaluation["satisfied"]: return None @@ -1977,6 +1994,19 @@ async def _prepare_goal_continuation_input( conversation_signature_before = visible_conversation_signature(messages) evidence_signature = latest_visible_assistant_signature(messages) + if _ends_on_human_input_request(messages): + # The agent asked the user something. Continuing would tell it to keep + # going while the question is still open on screen. + evaluation = GoalEvaluation( + satisfied=False, + blocker="needs_user_input", + reason="The turn ended on a question to the user that has not been answered.", + evidence_summary="", + ) + no_progress_count = compute_no_progress_count(goal, evaluation, evidence_signature=evidence_signature) + await _persist(goal, evaluation, no_progress_count, stand_down_reason=_stand_down_reason(goal, evaluation, no_progress_count)) + return None + if not _has_durable_goal_turn_receipt(checkpoint_tuple, messages): evaluation = GoalEvaluation( satisfied=False, diff --git a/backend/tests/test_goal_worker.py b/backend/tests/test_goal_worker.py index 6bae94f53..803c2aa31 100644 --- a/backend/tests/test_goal_worker.py +++ b/backend/tests/test_goal_worker.py @@ -3,7 +3,7 @@ import copy import pytest from deerflow_extension_api import ExtensionData -from langchain_core.messages import AIMessage, HumanMessage +from langchain_core.messages import AIMessage, HumanMessage, ToolMessage from langgraph.checkpoint.base import empty_checkpoint, uuid6 from langgraph.checkpoint.memory import InMemorySaver @@ -345,6 +345,117 @@ async def test_goal_worker_stands_down_after_the_run_hit_its_token_budget(monkey assert latest_goal["last_evaluation"]["stand_down_reason"] == "token_capped" +@pytest.mark.asyncio +@pytest.mark.parametrize("assistant_text", ["Let me check one thing first.", ""]) +@pytest.mark.parametrize("with_parallel_tool_result", [False, True]) +async def test_goal_worker_waits_for_an_unanswered_clarification(monkeypatch, assistant_text, with_parallel_tool_result): + """The clarification question lives in a ToolMessage the evaluator never reads.""" + checkpointer = InMemorySaver() + thread_id = "clarification-goal-thread" + ask = {"name": "ask_clarification", "args": {"question": "Drop the legacy table?"}, "id": "call-ask"} + look = {"name": "bash", "args": {"command": "ls"}, "id": "call-ls"} + messages = [ + HumanMessage(content="Migrate the orders database."), + AIMessage(content=assistant_text, tool_calls=[look, ask] if with_parallel_tool_result else [ask]), + ToolMessage( + content="Drop the legacy table?", + tool_call_id="call-ask", + name="ask_clarification", + artifact={"human_input": {"version": 1, "kind": "human_input_request", "source": "ask_clarification", "request_id": "req-1", "question": "Drop the legacy table?"}}, + ), + ] + if with_parallel_tool_result: + messages.append(ToolMessage(content="orders_v1 orders_v2", tool_call_id="call-ls", name="bash")) + await _seed_goal_thread(checkpointer, thread_id=thread_id, goal_text="Finish the migration", messages=messages) + evaluator_calls = [] + + async def fake_evaluate_goal_completion(_goal, _messages, **_kwargs): + evaluator_calls.append(_messages) + return GoalEvaluation(satisfied=False, blocker="goal_not_met_yet", reason="Not migrated yet.", evidence_summary="") + + monkeypatch.setattr(worker, "evaluate_goal_completion", fake_evaluate_goal_completion) + + continuation = await worker._prepare_goal_continuation_input( + accessor=_full_accessor(checkpointer), + bridge=_CollectingBridge(), + checkpointer=checkpointer, + thread_id=thread_id, + run_id="run-clarification", + model_name="test-model", + app_config=None, + ) + + assert continuation is None + assert evaluator_calls == [] + latest_goal = await read_thread_goal(checkpointer, thread_id) + assert latest_goal is not None + assert latest_goal["continuation_count"] == 0 + assert latest_goal["last_evaluation"]["blocker"] == "needs_user_input" + assert latest_goal["last_evaluation"]["stand_down_reason"] == "blocked:needs_user_input" + + +@pytest.mark.asyncio +@pytest.mark.parametrize("with_later_tool_result", [False, True]) +async def test_goal_worker_resumes_after_the_user_answers_clarification(monkeypatch, with_later_tool_result): + """An old card remains in history, but cannot block the answered turn.""" + checkpointer = InMemorySaver() + thread_id = "answered-clarification-goal-thread" + messages = [ + HumanMessage(content="Migrate the orders database."), + AIMessage(content="", tool_calls=[{"name": "ask_clarification", "args": {"question": "Keep the legacy table?"}, "id": "call-ask"}]), + ToolMessage( + content="Keep the legacy table?", + tool_call_id="call-ask", + name="ask_clarification", + artifact={"human_input": {"kind": "human_input_request"}}, + ), + ] + await _seed_goal_thread(checkpointer, thread_id=thread_id, goal_text="Finish the migration", messages=messages) + evaluator_calls = [] + + async def fake_evaluate_goal_completion(_goal, evaluated_messages, **_kwargs): + evaluator_calls.append(evaluated_messages) + return GoalEvaluation(satisfied=False, blocker="goal_not_met_yet", reason="Validation remains.", evidence_summary="Migrated while retaining the legacy table.") + + monkeypatch.setattr(worker, "evaluate_goal_completion", fake_evaluate_goal_completion) + kwargs = { + "accessor": _full_accessor(checkpointer), + "bridge": _CollectingBridge(), + "checkpointer": checkpointer, + "thread_id": thread_id, + "model_name": "test-model", + "app_config": None, + } + assert await worker._prepare_goal_continuation_input(**kwargs, run_id="run-question") is None + assert evaluator_calls == [] + waiting_goal = await read_thread_goal(checkpointer, thread_id) + assert waiting_goal["last_evaluation"]["stand_down_reason"] == "blocked:needs_user_input" + + messages.extend( + [ + HumanMessage(content="Yes, keep the legacy table."), + AIMessage( + content="Migrated the schema and kept the legacy table. Validation remains.", + tool_calls=[{"name": "bash", "args": {"command": "ls"}, "id": "call-ls"}] if with_later_tool_result else [], + ), + ] + ) + if with_later_tool_result: + messages.append(ToolMessage(content="orders_v1 orders_v2", tool_call_id="call-ls", name="bash")) + await _write_messages(checkpointer, thread_id=thread_id, messages=messages) + + continuation = await worker._prepare_goal_continuation_input(**kwargs, run_id="run-answer") + + assert len(evaluator_calls) == 1 + assert any(isinstance(message, HumanMessage) and message.content == "Yes, keep the legacy table." for message in evaluator_calls[0]) + assert continuation is not None + assert continuation["messages"][0].additional_kwargs["hide_from_ui"] is True + latest_goal = await read_thread_goal(checkpointer, thread_id) + assert latest_goal["continuation_count"] == 1 + assert latest_goal["last_evaluation"]["blocker"] == "goal_not_met_yet" + assert "stand_down_reason" not in latest_goal["last_evaluation"] + + @pytest.mark.asyncio async def test_goal_worker_clears_a_satisfied_goal_even_after_the_run_hit_its_token_budget(monkeypatch): checkpointer = InMemorySaver()