feat(flows): builder clarifies ambiguous inputs instead of guessing (orchestrator posture) (#4630)

This commit is contained in:
Cyrus Gray
2026-07-07 19:40:43 +05:30
committed by GitHub
parent 2ce5e6027f
commit b2f6206290
5 changed files with 352 additions and 24 deletions
@@ -1,4 +1,4 @@
import { fireEvent, render, screen } from '@testing-library/react';
import { act, fireEvent, render, screen } from '@testing-library/react';
import { beforeEach, describe, expect, it, vi } from 'vitest';
import type { WorkflowGraph, WorkflowNode } from '../../lib/flows/types';
@@ -45,7 +45,7 @@ describe('WorkflowCopilotPanel', () => {
hookState.toolTimeline = [];
hookState.liveResponse = '';
hookState.error = null;
hookState.send = vi.fn().mockResolvedValue(undefined);
hookState.send = vi.fn().mockResolvedValue({ proposed: false });
hookState.clearProposal = vi.fn();
});
@@ -76,6 +76,65 @@ describe('WorkflowCopilotPanel', () => {
expect(arg.request.graph).toEqual(baseGraph);
});
it('carries the original ask forward across a clarifying-question turn, then drops it once a proposal lands', async () => {
hookState.send = vi
.fn()
// Turn 1: the agent asks a clarifying question instead of proposing.
.mockResolvedValueOnce({ proposed: false })
// Turn 2: the user's answer resolves it and a proposal lands.
.mockResolvedValueOnce({ proposed: true })
// Turn 3 (and any further calls): a normal revise turn, already resolved.
.mockResolvedValue({ proposed: true });
render(
<WorkflowCopilotPanel
graph={baseGraph}
onProposal={vi.fn()}
onAccept={vi.fn()}
onReject={vi.fn()}
onClose={vi.fn()}
/>
);
fireEvent.change(screen.getByPlaceholderText('flows.copilot.placeholder'), {
target: { value: 'post a daily summary to slack' },
});
await act(async () => {
fireEvent.click(screen.getByTestId('send-message-button'));
// Flush the microtasks `submit` awaits before it records `pendingAskRef`.
await Promise.resolve();
await Promise.resolve();
});
expect(hookState.send).toHaveBeenCalledTimes(1);
fireEvent.change(screen.getByPlaceholderText('flows.copilot.placeholder'), {
target: { value: '#eng' },
});
await act(async () => {
fireEvent.click(screen.getByTestId('send-message-button'));
await Promise.resolve();
await Promise.resolve();
});
expect(hookState.send).toHaveBeenCalledTimes(2);
const secondArg = hookState.send.mock.calls[1][0];
// The follow-up must carry the ORIGINAL ask forward — a bare "#eng" alone
// would strand the agent with no idea what it was asked to build (the
// current graph is still blank/unchanged since no proposal has landed).
expect(secondArg.request.mode).toBe('revise');
expect(secondArg.request.instruction).toContain('post a daily summary to slack');
expect(secondArg.request.instruction).toContain('#eng');
// Turn 3, after a proposal has landed: the graph itself now carries the
// state, so the original ask must NOT be repeated.
fireEvent.change(screen.getByPlaceholderText('flows.copilot.placeholder'), {
target: { value: 'also add a filter step' },
});
fireEvent.click(screen.getByTestId('send-message-button'));
expect(hookState.send).toHaveBeenCalledTimes(3);
const thirdArg = hookState.send.mock.calls[2][0];
expect(thirdArg.request.instruction).toBe('also add a filter step');
});
it('renders the conversation transcript (user + agent turns)', () => {
hookState.messages = [
{ id: 'm1', content: 'add a Slack step', sender: 'user' },
@@ -253,4 +312,44 @@ describe('WorkflowCopilotPanel', () => {
);
expect(hookState.send).toHaveBeenCalledTimes(1);
});
it('carries the build seed description forward when the auto-sent build turn asks a clarifying question instead of proposing', async () => {
hookState.send = vi
.fn()
// The auto-sent build turn asks a question rather than proposing.
.mockResolvedValueOnce({ proposed: false })
// The user's free-text answer then resolves it.
.mockResolvedValueOnce({ proposed: true });
render(
<WorkflowCopilotPanel
graph={baseGraph}
flowId="flow-1"
onProposal={vi.fn()}
onAccept={vi.fn()}
onReject={vi.fn()}
onClose={vi.fn()}
buildSeed={{ description: 'post a daily summary to slack' }}
/>
);
// Flush the microtasks the seed effect awaits before recording
// `pendingAskRef` from the resolved `{ proposed: false }`.
await act(async () => {
await Promise.resolve();
await Promise.resolve();
});
expect(hookState.send).toHaveBeenCalledTimes(1);
fireEvent.change(screen.getByPlaceholderText('flows.copilot.placeholder'), {
target: { value: '#eng' },
});
fireEvent.click(screen.getByTestId('send-message-button'));
expect(hookState.send).toHaveBeenCalledTimes(2);
const secondArg = hookState.send.mock.calls[1][0];
// The follow-up must carry the build seed's original description forward,
// not just the bare "#eng" answer.
expect(secondArg.request.instruction).toContain('post a daily summary to slack');
expect(secondArg.request.instruction).toContain('#eng');
});
});
@@ -19,6 +19,7 @@
* Invariant: the copilot only PROPOSES. Accept applies to the UNSAVED local
* draft (no `flows_update`); persistence stays behind the canvas's own Save.
*/
import createDebug from 'debug';
import { useCallback, useEffect, useRef, useState } from 'react';
import { BubbleMarkdown } from '../../features/conversations/components/AgentMessageBubble';
@@ -31,6 +32,8 @@ import type { WorkflowProposal } from '../../store/chatRuntimeSlice';
import ChatComposer from '../chat/ChatComposer';
import Button from '../ui/Button';
const log = createDebug('app:flows:copilot-panel');
/**
* Context for a repair turn opened from a failed run's inspector ("Fix with
* agent"). Maps directly onto a `repair`-mode builder request.
@@ -133,13 +136,42 @@ export default function WorkflowCopilotPanel({
}
}, [proposal, onProposal]);
// Holds the ORIGINAL ask when a turn ends without a proposal — i.e. the
// agent asked a genuinely-ambiguous clarifying question (the prompt's
// "bucket 3" branch) and stopped rather than revising. `submit` always
// sends `mode: 'revise'` with the CURRENT graph, but while a question is
// still open that graph hasn't changed yet, so a bare follow-up answer
// ("#eng") would be the agent's ENTIRE context for the next turn — the
// original request ("post a daily summary to Slack") would be lost and the
// turn renders as "Revise it as follows: #eng" against a stale/blank draft.
// Prepending the unresolved ask keeps that context alive across the Q&A
// round-trip; it's cleared once a turn actually proposes (the graph itself
// then carries the state, so later revises don't need it).
const pendingAskRef = useRef<string | null>(null);
// Sets/clears `pendingAskRef` after a turn settles, logging the decision
// (stable prefix + thread correlation, never the raw ask/answer text — that
// may carry user-authored content).
const updatePendingAsk = useCallback(
(proposed: boolean, ask: string) => {
log(
'pendingAsk: %s thread=%s',
proposed ? 'cleared (proposal landed)' : 'set (still open)',
threadId
);
pendingAskRef.current = proposed ? null : ask;
},
[threadId]
);
// Auto-send the repair turn once when opened from a failed run.
const repairSentRef = useRef(false);
useEffect(() => {
if (!repairSeed || repairSentRef.current) return;
repairSentRef.current = true;
void send({
displayText: t('flows.copilot.repairDisplay'),
const instruction = t('flows.copilot.repairDisplay');
send({
displayText: instruction,
request: {
mode: 'repair',
instruction: '',
@@ -148,8 +180,10 @@ export default function WorkflowCopilotPanel({
error: repairSeed.error ?? null,
failingNodeIds: repairSeed.failingNodeIds ?? [],
},
}).then(({ proposed }) => {
updatePendingAsk(proposed, instruction);
});
}, [repairSeed, send, t]);
}, [repairSeed, send, t, updatePendingAsk]);
// Auto-send the build turn once when opened from the prompt bar's
// instant-create path: the user's description becomes the first user turn on
@@ -162,16 +196,22 @@ export default function WorkflowCopilotPanel({
useEffect(() => {
if (!buildSeed || buildSentRef.current) return;
buildSentRef.current = true;
void send({
send({
displayText: buildSeed.description,
request: flowId
? { mode: 'build', instruction: buildSeed.description, graph, flowId }
: { mode: 'revise', instruction: buildSeed.description, graph, flowId },
}).then(({ proposed }) => {
// Not proposed => the seed turn asked a clarifying question instead of
// building. Carry the original description forward so the user's
// free-text answer (via `submit` below) doesn't strand the agent with
// no idea what it was asked to build.
updatePendingAsk(proposed, buildSeed.description);
});
// `graph`/`flowId` are read once for the seed turn — later edits must not
// re-fire it (guarded by the ref regardless).
// eslint-disable-next-line react-hooks/exhaustive-deps
}, [buildSeed, send]);
}, [buildSeed, send, updatePendingAsk]);
// Keep the transcript pinned to the newest message / streamed activity.
// `scrollTo` is optional-chained: jsdom (tests) doesn't implement it.
@@ -184,12 +224,17 @@ export default function WorkflowCopilotPanel({
const trimmed = (raw ?? text).trim();
if (!trimmed || sending) return;
setText('');
await send({
const priorAsk = pendingAskRef.current;
const instruction = priorAsk
? `${priorAsk}\n\n(This is my answer to your question above: ${trimmed})`
: trimmed;
const { proposed } = await send({
displayText: trimmed,
request: { mode: 'revise', instruction: trimmed, graph, flowId },
request: { mode: 'revise', instruction, graph, flowId },
});
updatePendingAsk(proposed, instruction);
},
[text, sending, send, graph, flowId]
[text, sending, send, graph, flowId, updatePendingAsk]
);
const handleInputKeyDown = useCallback(
+84 -3
View File
@@ -3,6 +3,7 @@ import { beforeEach, describe, expect, it, vi } from 'vitest';
import type { BuilderTurnResult } from '../services/api/flowsApi';
import type { WorkflowProposal } from '../store/chatRuntimeSlice';
import type { ThreadMessage } from '../types/thread';
import { useWorkflowBuilderChat } from './useWorkflowBuilderChat';
// The hook now runs the builder server-side via `openhuman.flows_build`.
@@ -41,6 +42,15 @@ vi.mock('../store/chatRuntimeSlice', () => ({
setWorkflowProposalForThread: (p: unknown) => ({ type: 'setProposal', p }),
}));
// The hook reads the live store directly (not the stale closed-over selector
// value) to dedup against a message the streamed `chat_done` path may have
// already appended for this exact turn — see the `assistantText` fallback
// branch. Controllable per test via `rawStoreState.thread.messagesByThreadId`.
const rawStoreState = vi.hoisted(() => ({
thread: { messagesByThreadId: {} as Record<string, { sender: string; content: string }[]> },
}));
vi.mock('../store', () => ({ store: { getState: () => rawStoreState } }));
const okResult = (over: Partial<BuilderTurnResult> = {}): BuilderTurnResult => ({
proposal: null,
assistantText: 'done',
@@ -55,6 +65,7 @@ describe('useWorkflowBuilderChat', () => {
selectorState.messagesByThreadId = {};
selectorState.toolTimelineByThread = {};
selectorState.streamingAssistantByThread = {};
rawStoreState.thread.messagesByThreadId = {};
dispatch.mockReset().mockImplementation((action: { type: string }) => {
if (action.type === 'createNewThread') {
return { unwrap: () => Promise.resolve({ id: 'builder-1' }) };
@@ -113,7 +124,13 @@ describe('useWorkflowBuilderChat', () => {
);
});
it('appends only the user turn locally — the runtime owns the agent reply', async () => {
it('appends the user turn locally — the runtime normally owns the agent reply', async () => {
// Simulate the streamed path already having delivered this exact text via
// `chat_done` (the normal case when streaming is wired) so the fallback
// branch below can prove it does NOT double the bubble.
rawStoreState.thread.messagesByThreadId = {
'builder-1': [{ sender: 'agent', content: 'Here is your workflow.' }],
};
buildWorkflow.mockResolvedValue(okResult({ assistantText: 'Here is your workflow.' }));
const { result } = renderHook(() => useWorkflowBuilderChat());
await act(async () => {
@@ -128,11 +145,75 @@ describe('useWorkflowBuilderChat', () => {
// The web channel never persists user messages, so the hook appends the
// user turn itself...
expect(appended.some(a => a.p?.message?.sender === 'user')).toBe(true);
// ...but NOT the agent reply — `ChatRuntimeProvider` appends that on the
// streamed `chat_done`, so appending here too would double it.
// ...but NOT the agent reply when it was already streamed — appending
// here too would double it.
expect(appended.some(a => a.p?.message?.sender === 'agent')).toBe(false);
});
it('surfaces a clarifying question as an assistant message when the builder returns plain text with no proposal (fallback)', async () => {
buildWorkflow.mockResolvedValue(
okResult({
proposal: null,
error: null,
assistantText: 'Which Slack channel — #eng or #sales?',
})
);
const { result } = renderHook(() => useWorkflowBuilderChat());
await act(async () => {
await result.current.send({
displayText: 'post a daily summary to slack',
request: { mode: 'create', instruction: 'post a daily summary to slack' },
});
});
const appendedAgentMessages = dispatch.mock.calls
.map(([a]) => a as { type: string; p?: { threadId?: string; message?: ThreadMessage } })
.filter(a => a.type === 'addMessageLocal' && a.p?.message?.sender === 'agent');
expect(appendedAgentMessages).toHaveLength(1);
expect(appendedAgentMessages[0]?.p?.message?.content).toBe(
'Which Slack channel — #eng or #sales?'
);
expect(appendedAgentMessages[0]?.p?.threadId).toBe('builder-1');
// No proposal was surfaced for this turn.
expect(dispatch.mock.calls.some(([a]) => (a as { type: string }).type === 'setProposal')).toBe(
false
);
});
it('does not double-append when a proposal is returned alongside assistant text', async () => {
const proposal: WorkflowProposal = {
name: 'Digest',
graph: { nodes: [], edges: [] },
requireApproval: true,
summary: { trigger: 'schedule', steps: [] },
};
buildWorkflow.mockResolvedValue(
okResult({ proposal, assistantText: "I've built this — review below." })
);
const { result } = renderHook(() => useWorkflowBuilderChat());
await act(async () => {
await result.current.send({
displayText: 'hi',
request: { mode: 'create', instruction: 'x' },
});
});
// A proposal result still sets the proposal, unchanged...
expect(dispatch).toHaveBeenCalledWith(
expect.objectContaining({ type: 'setProposal', p: { threadId: 'builder-1', proposal } })
);
// ...and does NOT also append an agent chat message (the proposal branch
// is exclusive of the assistant-text fallback branch).
expect(
dispatch.mock.calls.some(
([a]) =>
(a as { type: string; p?: { message?: { sender?: string } } }).type ===
'addMessageLocal' &&
(a as { p?: { message?: { sender?: string } } }).p?.message?.sender === 'agent'
)
).toBe(false);
});
it('reuses the same dedicated thread across sends (creates it once)', async () => {
const { result } = renderHook(() => useWorkflowBuilderChat());
await act(async () => {
+52 -7
View File
@@ -26,6 +26,7 @@ import createDebug from 'debug';
import { useCallback, useMemo, useState } from 'react';
import { type BuilderTurnRequest, buildWorkflow } from '../services/api/flowsApi';
import { store } from '../store';
import {
clearWorkflowProposalForThread,
setWorkflowProposalForThread,
@@ -79,8 +80,16 @@ export interface UseWorkflowBuilderChat {
liveResponse: string;
/** Last send error (thread create / RPC failure), or `null`. */
error: string | null;
/** Send a builder turn, creating the dedicated thread on first use. */
send: (params: WorkflowBuilderSendParams) => Promise<void>;
/**
* Send a builder turn, creating the dedicated thread on first use. Resolves
* with `proposed: true` iff this turn's `flows_build` call returned a
* proposal — `false` for a clarifying question, an error, or a call that
* never ran (already sending / offline). Callers that loop a conversation
* (the copilot's free-text follow-ups) use this to know whether the turn's
* instruction is still "unresolved" and must be carried into the next turn
* — see `WorkflowCopilotPanel`'s `pendingAskRef`.
*/
send: (params: WorkflowBuilderSendParams) => Promise<{ proposed: boolean }>;
/** Clear the current proposal (e.g. after Accept/Reject) without persisting. */
clearProposal: () => void;
}
@@ -141,16 +150,17 @@ export function useWorkflowBuilderChat(seedThreadId?: string | null): UseWorkflo
async ({ displayText, request }: WorkflowBuilderSendParams) => {
if (localSending) {
log('send: ignored — a turn is already dispatching');
return;
return { proposed: false };
}
if (socketStatus !== 'connected') {
log('send: blocked — socket not connected (%s)', socketStatus);
setError('offline');
return;
return { proposed: false };
}
setLocalSending(true);
setError(null);
let targetThreadId = threadId;
let proposed = false;
try {
if (!targetThreadId) {
log('send: creating dedicated builder thread');
@@ -178,9 +188,10 @@ export function useWorkflowBuilderChat(seedThreadId?: string | null): UseWorkflo
// text/thinking/tool events + a terminal `chat_done` keyed by it. The
// GLOBAL `ChatRuntimeProvider` owns that transcript — it appends the
// final assistant message on `chat_done` and fills the streaming/tool
// slices as the turn runs so this hook must NOT also append the agent
// reply (doing so would double it). We still await the blocking result
// for its `proposal`/`error` fallback.
// slices as the turn runs, so in the normal (streaming-wired) case this
// hook must NOT also append the agent reply (doing so would double
// it) — see the dedup check below. We still await the blocking result
// for its `proposal`/`error`/`assistantText` fallback.
log('send: running flows_build thread=%s mode=%s', targetThreadId, request.mode);
const result = await buildWorkflow(request, targetThreadId);
@@ -190,11 +201,44 @@ export function useWorkflowBuilderChat(seedThreadId?: string | null): UseWorkflo
// `pendingWorkflowProposalsByThread` from the tool result; re-writing the
// same value here is idempotent and covers a missed socket event / CLI.
if (result.proposal) {
proposed = true;
dispatch(
setWorkflowProposalForThread({ threadId: targetThreadId, proposal: result.proposal })
);
} else if (result.error) {
setError(result.error);
} else if (result.assistantText?.trim()) {
// Neither a proposal nor an error: the agent replied with plain
// text instead of proposing this turn — most commonly a clarifying
// question (the "ask" branch of the clarify/verify posture). When
// streaming is wired (the normal case) `ChatRuntimeProvider` already
// appended this exact text on the turn's `chat_done` — the Rust
// side (`finalize_flow_stream`) delivers it unconditionally,
// independent of whether a proposal was made — so re-appending here
// would double the bubble. Read the live store (not the stale
// closed-over `messages`) to check whether that already landed;
// only append when it hasn't, which is the actual fallback case
// (streaming not wired: CLI / tests / a missed socket event).
const latest = store.getState().thread.messagesByThreadId[targetThreadId] ?? [];
const lastMessage = latest[latest.length - 1];
const alreadyStreamed =
lastMessage?.sender === 'agent' && lastMessage.content === result.assistantText;
log(
'send: assistantText fallback thread=%s alreadyStreamed=%s',
targetThreadId,
alreadyStreamed
);
if (!alreadyStreamed) {
const assistantMessage: ThreadMessage = {
id: `msg_${globalThis.crypto.randomUUID()}`,
content: result.assistantText,
type: 'text',
extraMetadata: {},
sender: 'agent',
createdAt: new Date().toISOString(),
};
dispatch(addMessageLocal({ threadId: targetThreadId, message: assistantMessage }));
}
}
} catch (err) {
const msg = err instanceof Error ? err.message : String(err);
@@ -203,6 +247,7 @@ export function useWorkflowBuilderChat(seedThreadId?: string | null): UseWorkflo
} finally {
setLocalSending(false);
}
return { proposed };
},
[dispatch, localSending, socketStatus, threadId]
);
@@ -361,7 +361,65 @@ Prefer `retry` + `on_error: "route"` for flaky network/tool steps, and
## Style
Be concise. Ask a clarifying question only when the trigger or a critical step is
genuinely ambiguous — otherwise make a sensible proposal and let the user refine
it. Always end by proposing (or revising) the workflow; describe what it does in
one or two plain sentences alongside the proposal.
Be concise. Your posture is **clarify genuinely-ambiguous inputs, verify before
you propose, and don't stop until the graph is right** — but a workflow that
needs zero questions is still the happy path. Don't let "ask when truly
unsure" turn into "ask about everything": most requests carry enough signal
to build immediately.
### The ask-vs-just-build rule
Once `get_tool_contract` hands you a node's `required_args`, sort each one
into exactly one bucket before you write the node:
1. **WIRED** — an upstream node's output already produces the value. Bind it
(`=nodes.<id>.item.json.<field>`, per "the envelope" above) and move on —
no question, nothing to state.
2. **INFERABLE** — the request implies the value even though nothing
upstream produces it:
- "to me" / "message me" / "DM me" → the user's OWN Slack/Discord/etc. DM
target, never a public channel. **Never default a personal request to
`#general`** — that's a different destination than the user asked for,
not a safe guess.
- Exactly one connected account for the toolkit the step needs → that
account (`list_flow_connections` / `composio_list_connections` tell
you this; don't ask "which Gmail?" when there's only one).
- An unambiguous, low-stakes default implied by the ask ("daily" → a
sensible `schedule` hour if none was named).
Fill these in yourself, then **name the choice in your final summary**
(below) so the user can correct it in one message if you guessed wrong.
3. **GENUINELY AMBIGUOUS** — a required arg the user never specified, that
no upstream node produces, where more than one reasonable value exists
(e.g. "post to Slack" with several channels connected and no hint which).
**Ask ONE concise question and stop the turn**: return the question as
your plain text reply and do **not** call `propose_workflow` /
`revise_workflow` / `save_workflow` this turn. Wait for the user's answer
on the next turn before building further.
Ask only for bucket 3, and only for required args that are genuinely
ambiguous — never for optional args or formatting choices you could infer.
Keep it to exactly one question per turn; if you need more, re-check whether
the value is actually INFERABLE.
### The verify loop — don't stop at "it compiles"
`dry_run_workflow` isn't a formality you run once. Treat a flagged result
(`"ok": false`, a `null_resolutions` entry, an `agent_prompt_nulls` entry, or
a rejected contract) as unfinished work: fix the binding/schema/slug it
names, `dry_run_workflow` again, and repeat until it comes back clean. Only
then call `propose_workflow` / `save_workflow`. Don't hand back a proposal
you haven't verified just because the turn has run long — the user would
rather wait one more tool call than review a graph that silently does
nothing.
### Say what you inferred
In the proposal's summary (or your closing reply if you asked a question
instead), name every INFERABLE choice in half a sentence — "sending as a DM
to you", "using your only connected Gmail account", "running every morning
at 8am since none was specified". This is what makes bucket 2 safe to skip
asking about: the guess stays visible and one message away from being
corrected, never silently locked in.
Always end a building turn with either a proposal (or revision), or — only
for bucket 3 — a single clarifying question. Never both, never neither.