Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 34 additions & 5 deletions agentscope-core/src/main/java/io/agentscope/core/ReActAgent.java
Original file line number Diff line number Diff line change
Expand Up @@ -1630,13 +1630,15 @@ private List<ConfirmResult> extractConfirmResults(List<Msg> msgs) {
* modified) one from the result, set state to {@link ToolCallState#ALLOWED}, and
* register any attached {@link PermissionRule}s with the engine.</li>
* <li>{@code confirmed == false}: write a DENIED {@link ToolResultBlock} to context so
* the tool will no longer be pending on resume.</li>
* the tool will no longer be pending on resume, and emit the same
* {@link ToolResultStartEvent} / {@link ToolResultTextDeltaEvent} /
* {@link ToolResultEndEvent} sequence used for auto-denied tools.</li>
* </ul>
*/
private void applyConfirmResults(List<ConfirmResult> results) {
// Replace ASKING ToolUseBlocks with possibly-modified ones from the user, and
// promote them to ALLOWED. Collect denied ones for separate handling.
List<ToolUseBlock> deniedToolCalls = new ArrayList<>();
List<Map.Entry<ToolUseBlock, String>> deniedEntries = new ArrayList<>();
Map<String, ToolUseBlock> replacements = new HashMap<>();
Map<String, ToolCallState> stateUpdates = new HashMap<>();
for (ConfirmResult r : results) {
Expand All @@ -1655,20 +1657,47 @@ private void applyConfirmResults(List<ConfirmResult> results) {
}
}
} else {
deniedToolCalls.add(target);
deniedEntries.add(Map.entry(target, resolveUserDenyMessage(r)));
}
}
applyToolUseBlockReplacements(replacements);
for (ToolUseBlock denied : deniedToolCalls) {
if (deniedEntries.isEmpty()) {
return;
}
// Correlate all deny tool-result events for this resume with one reply id, matching
// the auto-deny path in runToolBatch.
String replyId = UUID.randomUUID().toString().replace("-", "");
for (Map.Entry<ToolUseBlock, String> entry : deniedEntries) {
ToolUseBlock denied = entry.getKey();
String denyMessage = entry.getValue();
ToolResultBlock deniedResult =
ToolResultBlock.text("Permission denied by user")
ToolResultBlock.text(denyMessage)
.withIdAndName(denied.getId(), denied.getName())
.withState(ToolResultState.DENIED);
Msg deniedMsg =
ToolResultMessageBuilder.buildToolResultMsg(
deniedResult, denied, getName());
state.contextMutable().add(deniedMsg);
publishEvent(new ToolResultStartEvent(replyId, denied.getId(), denied.getName()));
publishEvent(
new ToolResultTextDeltaEvent(
replyId, denied.getId(), denied.getName(), denyMessage));
publishEvent(
new ToolResultEndEvent(
replyId, denied.getId(), denied.getName(), ToolResultState.DENIED));
}
}

/**
* Resolve the tool-result text for a user deny. Prefer {@link ConfirmResult#getMessage()}
* when present; otherwise keep the historical default.
*/
private static String resolveUserDenyMessage(ConfirmResult result) {
String message = result.getMessage();
if (message != null && !message.isBlank()) {
return message;
}
return "Permission denied by user";
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,26 +27,55 @@
* <p>When confirmed, the caller may supply a modified {@link #toolCall} (allowing the user to
* tweak input) and/or new {@link #rules} that the {@code PermissionEngine} should remember for
* future calls — e.g. "always allow this command going forward".
*
* <p>When denied ({@code confirmed == false}), an optional {@link #message} can override the
* default tool-result text ({@code "Permission denied by user"}). Prefer {@link #withMessage} or
* the 4-arg constructor — there is no 3-arg {@code (confirmed, toolCall, message)} overload so
* {@code new ConfirmResult(c, t, null)} stays unambiguous against the rules overload.
*/
public class ConfirmResult {

private final boolean confirmed;
private final ToolUseBlock toolCall;
private final List<PermissionRule> rules;
private final String message;

@JsonCreator
public ConfirmResult(
@JsonProperty("confirmed") boolean confirmed,
@JsonProperty("toolCall") ToolUseBlock toolCall,
@JsonProperty("rules") List<PermissionRule> rules) {
@JsonProperty("rules") List<PermissionRule> rules,
@JsonProperty("message") String message) {
this.confirmed = confirmed;
this.toolCall = toolCall;
this.rules = rules;
this.message = message;
}

/** Convenience constructor without a custom deny message. */
public ConfirmResult(boolean confirmed, ToolUseBlock toolCall, List<PermissionRule> rules) {
this(confirmed, toolCall, rules, null);
}

/** Convenience constructor without rules. */
/** Convenience constructor without rules or a custom deny message. */
public ConfirmResult(boolean confirmed, ToolUseBlock toolCall) {
this(confirmed, toolCall, null);
this(confirmed, toolCall, null, null);
}

/**
* Factory for a confirm result with a custom deny message and no rules.
*
* <p>Use this (or the 4-arg constructor) instead of a 3-arg message overload so that {@code
* new ConfirmResult(confirmed, toolCall, null)} remains unambiguously the rules constructor.
*
* @param confirmed whether the user approved the tool call
* @param toolCall the (possibly modified) tool call being decided
* @param message custom tool-result text when denying; blank/null falls back to the default
* @return a new {@link ConfirmResult}
*/
public static ConfirmResult withMessage(
boolean confirmed, ToolUseBlock toolCall, String message) {
return new ConfirmResult(confirmed, toolCall, null, message);
}

public boolean isConfirmed() {
Expand All @@ -66,4 +95,15 @@ public ToolUseBlock getToolCall() {
public List<PermissionRule> getRules() {
return rules;
}

/**
* Optional text used as the denied tool-result content when {@link #confirmed} is false.
*
* <p>When null or blank, the agent falls back to {@code "Permission denied by user"}.
*
* @return custom deny message, or null
*/
public String getMessage() {
return message;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,8 @@
import io.agentscope.core.event.RequestStopEvent;
import io.agentscope.core.event.RequireUserConfirmEvent;
import io.agentscope.core.event.ToolResultEndEvent;
import io.agentscope.core.event.ToolResultStartEvent;
import io.agentscope.core.event.ToolResultTextDeltaEvent;
import io.agentscope.core.message.ContentBlock;
import io.agentscope.core.message.GenerateReason;
import io.agentscope.core.message.Msg;
Expand Down Expand Up @@ -200,10 +202,12 @@ private static int countOf(List<AgentEvent> events, Class<?> type) {
}

private static Msg confirmMsg(boolean confirmed, ToolUseBlock toolCall) {
return confirmMsg(new ConfirmResult(confirmed, toolCall));
}

private static Msg confirmMsg(ConfirmResult result) {
Map<String, Object> meta = new HashMap<>();
meta.put(
Msg.METADATA_CONFIRM_RESULTS,
List.of(new ConfirmResult(confirmed, toolCall, null)));
meta.put(Msg.METADATA_CONFIRM_RESULTS, List.of(result));
return Msg.builder()
.name("user")
.role(MsgRole.USER)
Expand All @@ -212,6 +216,23 @@ private static Msg confirmMsg(boolean confirmed, ToolUseBlock toolCall) {
.build();
}

private static ToolUseBlock pendingAskingTool(ReActAgent agent) {
for (int i = agent.getAgentState().getContext().size() - 1; i >= 0; i--) {
Msg m = agent.getAgentState().getContext().get(i);
if (m.getRole() == MsgRole.ASSISTANT) {
return m.getContentBlocks(ToolUseBlock.class).get(0);
}
}
throw new AssertionError("expected an assistant message with a pending tool call");
}

private static String toolResultText(ToolResultBlock block) {
return block.getOutput().stream()
.filter(b -> b instanceof TextBlock)
.map(b -> ((TextBlock) b).getText())
.reduce("", String::concat);
}

@Test
void askingToolPausesFirstCallAndExecutesOnConfirmedSecondCall() {
ChatModelBase model =
Expand Down Expand Up @@ -295,29 +316,130 @@ void askingToolResumeWithDeniedConfirmResultProducesDeniedToolResult() {
assertNotNull(first);
assertEquals(GenerateReason.PERMISSION_ASKING, first.getGenerateReason());

Msg lastAssistant = null;
for (int i = agent.getAgentState().getContext().size() - 1; i >= 0; i--) {
Msg m = agent.getAgentState().getContext().get(i);
if (m.getRole() == MsgRole.ASSISTANT) {
lastAssistant = m;
break;
}
}
ToolUseBlock pending = lastAssistant.getContentBlocks(ToolUseBlock.class).get(0);
ToolUseBlock pending = pendingAskingTool(agent);

// Second call → deny
Msg second = agent.call(List.of(confirmMsg(false, pending))).block();
assertNotNull(second);

// Context should contain a DENIED ToolResultBlock for tc1
boolean foundDenied =
// Context should contain a DENIED ToolResultBlock for tc1 with the default message
ToolResultBlock denied =
agent.getAgentState().getContext().stream()
.flatMap(m -> m.getContentBlocks(ToolResultBlock.class).stream())
.anyMatch(
tr ->
"tc1".equals(tr.getId())
&& tr.getState() == ToolResultState.DENIED);
assertTrue(foundDenied, "expected a DENIED ToolResultBlock for the rejected tool");
.filter(tr -> "tc1".equals(tr.getId()))
.findFirst()
.orElse(null);
assertNotNull(denied, "expected a DENIED ToolResultBlock for the rejected tool");
assertEquals(ToolResultState.DENIED, denied.getState());
assertEquals("Permission denied by user", toolResultText(denied));
}

@Test
void askingToolResumeWithCustomDenyMessageUsesConfirmResultMessage() {
ChatModelBase model =
new ScriptedModel(
List.of(
() -> Flux.just(toolUseResponse("tc1", "ask", "x")),
() -> Flux.just(textResponse("done"))));
ReActAgent agent = buildAgent(model, toolkitWith(new AskingTool("ask")));

Msg first = agent.call(List.of()).block();
assertNotNull(first);
assertEquals(GenerateReason.PERMISSION_ASKING, first.getGenerateReason());

ToolUseBlock pending = pendingAskingTool(agent);
String reason = "User rejected: do not delete production data";
Msg second =
agent.call(List.of(confirmMsg(ConfirmResult.withMessage(false, pending, reason))))
.block();
assertNotNull(second);

ToolResultBlock denied =
agent.getAgentState().getContext().stream()
.flatMap(m -> m.getContentBlocks(ToolResultBlock.class).stream())
.filter(tr -> "tc1".equals(tr.getId()))
.findFirst()
.orElse(null);
assertNotNull(denied);
assertEquals(ToolResultState.DENIED, denied.getState());
assertEquals(reason, toolResultText(denied));
}

@Test
void askingToolResumeWithBlankDenyMessageFallsBackToDefault() {
ChatModelBase model =
new ScriptedModel(
List.of(
() -> Flux.just(toolUseResponse("tc1", "ask", "x")),
() -> Flux.just(textResponse("done"))));
ReActAgent agent = buildAgent(model, toolkitWith(new AskingTool("ask")));

Msg first = agent.call(List.of()).block();
assertNotNull(first);
assertEquals(GenerateReason.PERMISSION_ASKING, first.getGenerateReason());

ToolUseBlock pending = pendingAskingTool(agent);
Msg second =
agent.call(List.of(confirmMsg(ConfirmResult.withMessage(false, pending, " "))))
.block();
assertNotNull(second);

ToolResultBlock denied =
agent.getAgentState().getContext().stream()
.flatMap(m -> m.getContentBlocks(ToolResultBlock.class).stream())
.filter(tr -> "tc1".equals(tr.getId()))
.findFirst()
.orElse(null);
assertNotNull(denied);
assertEquals(ToolResultState.DENIED, denied.getState());
assertEquals("Permission denied by user", toolResultText(denied));
}

@Test
void askingToolResumeWithDenyEmitsToolResultEvents() {
ChatModelBase model =
new ScriptedModel(
List.of(
() -> Flux.just(toolUseResponse("tc1", "ask", "x")),
() -> Flux.just(textResponse("done"))));
ReActAgent agent = buildAgent(model, toolkitWith(new AskingTool("ask")));

// First stream: pause on ASK
List<AgentEvent> firstEvents = agent.streamEvents(List.of()).collectList().block();
assertNotNull(firstEvents);
assertTrue(indexOf(firstEvents, RequireUserConfirmEvent.class) >= 0);

ToolUseBlock pending = pendingAskingTool(agent);
String reason = "blocked by policy";

// Second stream: deny — should surface the same Start/Delta/End sequence as auto-deny
List<AgentEvent> resumeEvents =
agent.streamEvents(
List.of(
confirmMsg(
ConfirmResult.withMessage(false, pending, reason))))
.collectList()
.block();
assertNotNull(resumeEvents);

int iStart = indexOf(resumeEvents, ToolResultStartEvent.class);
int iDelta = indexOf(resumeEvents, ToolResultTextDeltaEvent.class);
int iEnd = indexOf(resumeEvents, ToolResultEndEvent.class);
assertTrue(iStart >= 0, "ToolResultStartEvent must be emitted on manual deny");
assertTrue(iDelta > iStart, "ToolResultTextDeltaEvent must follow Start");
assertTrue(iEnd > iDelta, "ToolResultEndEvent must follow TextDelta");

ToolResultStartEvent start = (ToolResultStartEvent) resumeEvents.get(iStart);
assertEquals("tc1", start.getToolCallId());
assertEquals("ask", start.getToolCallName());

ToolResultTextDeltaEvent delta = (ToolResultTextDeltaEvent) resumeEvents.get(iDelta);
assertEquals("tc1", delta.getToolCallId());
assertEquals(reason, delta.getDelta());

ToolResultEndEvent end = (ToolResultEndEvent) resumeEvents.get(iEnd);
assertEquals("tc1", end.getToolCallId());
assertEquals(ToolResultState.DENIED, end.getState());
}

@Test
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import reactor.core.Disposable;
Expand Down Expand Up @@ -673,20 +674,40 @@ private Mono<Msg> execLocalSync(
new AgentStartEvent(spawned.sessionId(), null, spawned.agentId())
.withSource(sourcePath));

// Emit AGENT_END at most once. Prefer the success path (flatMap) so the
// bookend is delivered *before* the Mono value reaches
// execWithTimeoutPromotion's CompletableFuture bridge — otherwise
// whenComplete can advance the parent and complete the streamEvents sink
// before doFinally runs, dropping the child AGENT_END under load (CI).
AtomicBoolean endEmitted = new AtomicBoolean(false);
Runnable emitChildEnd =
() -> {
if (endEmitted.compareAndSet(false, true)) {
parentEmitter.emit(
new AgentEndEvent(null).withSource(sourcePath));
}
};

return manager.invokeAgent(agent, sessionId, userId, prompt, parentCtx)
.contextWrite(
c ->
c.put(
AgentEventEmitter.FORWARDING_CONTEXT_KEY,
taggedEmitter))
// doFinally, not doOnTerminate: the latter skips cancel, so a
// parent cancel would leave the AgentStartEvent above unmatched
// and consumers would render this subagent as running forever.
.flatMap(
msg -> {
emitChildEnd.run();
return Mono.just(msg);
})
.doOnError(err -> emitChildEnd.run())
// Cancel: flatMap/doOnError do not run; still need the bookend so
// consumers do not render the subagent as running forever.
.doFinally(
signal ->
parentEmitter.emit(
new AgentEndEvent(null)
.withSource(sourcePath)));
signal -> {
if (signal == SignalType.CANCEL) {
emitChildEnd.run();
}
});
}

// ── Path 2: stream() (deprecated) — SubagentEventBus forwarding ──
Expand Down
Loading
Loading