Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
a352a40
chore(hook): 初始化 hook 能力分支
lyfmt Jun 21, 2026
7718b79
feat(hook): 添加工具 hook 契约
lyfmt Jun 21, 2026
097bb88
fix(hook): 统一 hook 上下文输入快照
lyfmt Jun 21, 2026
91d33e0
test(hook): 放宽输入快照测试语义
lyfmt Jun 21, 2026
5a8170e
fix(hook): 固化嵌套输入快照
lyfmt Jun 21, 2026
e7c03a0
fix(hook): 支持 null 输入快照
lyfmt Jun 21, 2026
45a292b
test(hook): 补充嵌套 null 和 set 快照覆盖
lyfmt Jun 21, 2026
1ff273f
feat(hook): 接入工具 hook 拦截器
lyfmt Jun 21, 2026
7520ab6
fix(hook): 兼容空工具输入
lyfmt Jun 21, 2026
9eb6763
feat(hook): 在启动装配中接入工具 hook
lyfmt Jun 21, 2026
ca9169b
test(hook): 强化 boot hook 接线覆盖
lyfmt Jun 21, 2026
8d8c17c
feat(hook): 添加 hook 生命周期事件契约
lyfmt Jun 22, 2026
fdeebb9
feat(hook): 发布工具 hook 生命周期事件
lyfmt Jun 22, 2026
90825cf
feat(hook): 接入 hook 生命周期事件总线
lyfmt Jun 22, 2026
2739e6f
fix(hook): 容错缺失审计元数据
lyfmt Jun 22, 2026
e47d3e2
test(hook): 校验 hook 审计字段一致性
lyfmt Jun 22, 2026
a6e9595
feat(hook): 添加 turn hook 契约
lyfmt Jun 22, 2026
0f721a7
feat(hook): 接入 turn hook 运行时
lyfmt Jun 22, 2026
84309a5
feat(hook): 自动接线 turn hook
lyfmt Jun 22, 2026
38a7734
test(hook): 校验 after turn hook 替换状态
lyfmt Jun 22, 2026
82ce207
fix(hook): 收窄 after turn hook 语义
lyfmt Jun 22, 2026
ac658d3
docs(hook): 修正 after turn hook 注释
lyfmt Jun 22, 2026
0535797
fix(hook): 保持 after turn 失败 fork point
lyfmt Jun 22, 2026
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
104 changes: 79 additions & 25 deletions lypi-agent-core/src/main/java/cn/lypi/agent/DefaultTurnExecutor.java
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,10 @@
import cn.lypi.contracts.context.ToolCallContentBlock;
import cn.lypi.contracts.event.ErrorEvent;
import cn.lypi.contracts.event.TurnStartEvent;
import cn.lypi.contracts.hook.AfterTurnHookContext;
import cn.lypi.contracts.hook.BeforeTurnHookContext;
import cn.lypi.contracts.hook.BeforeTurnHookResult;
import cn.lypi.contracts.hook.TurnHookRuntime;
import cn.lypi.contracts.model.AssistantEventStream;
import cn.lypi.contracts.model.AssistantError;
import cn.lypi.contracts.model.AssistantStart;
Expand Down Expand Up @@ -44,8 +48,13 @@ public final class DefaultTurnExecutor implements TurnExecutor {
private final ContextBudgetEstimator budgetEstimator;
private final TurnContinuationGuard continuationGuard;
private final TurnEventPublisher eventPublisher;
private final TurnHookRuntime turnHooks;

public DefaultTurnExecutor(AgentCoreRuntimePorts ports, TurnIds ids, Clock clock) {
this(ports, ids, clock, TurnHookRuntime.noop());
}

public DefaultTurnExecutor(AgentCoreRuntimePorts ports, TurnIds ids, Clock clock, TurnHookRuntime turnHooks) {
this.ports = ports;
this.ids = ids;
this.clock = clock;
Expand All @@ -55,6 +64,7 @@ public DefaultTurnExecutor(AgentCoreRuntimePorts ports, TurnIds ids, Clock clock
this.budgetEstimator = new ContextBudgetEstimator();
this.continuationGuard = new TurnContinuationGuard(ports.sessionManager());
this.eventPublisher = new TurnEventPublisher(ports.eventBus(), clock);
this.turnHooks = turnHooks == null ? TurnHookRuntime.noop() : turnHooks;
}

@Override
Expand All @@ -73,30 +83,43 @@ private TurnState executeWithTurnId(TurnRequest request, String turnId) {
request.parentEntryId().ifPresent(parentEntryId -> ports.sessionManager().switchLeaf(parentEntryId));
Instant startedAt = clock.instant();
ports.eventBus().publish(new TurnStartEvent(request.sessionId(), turnId, startedAt, startedAt));
Optional<String> unsafeReason = continuationGuard.unsafeContinuationReason(currentLeafId());
if (unsafeReason.isPresent()) {
ports.eventBus().publish(new ErrorEvent(
request.sessionId(),
unsafeReason.orElseThrow(),
"当前分支停在 assistant 工具调用消息上,不能直接追加用户消息。请选择上一条用户消息或工具结果之后继续。",
clock.instant()
));
return failedState(turnId, request.sessionId(), null, List.of(), 0, startedAt, currentLeafId());
}

AgentMessage user = messageFactory.userMessage(ids.newMessageId(), request.userInput());
String contextLeafId = appendNewMessage(request.sessionId(), user);
newMessages.add(user);

ContextSnapshot context = null;
int toolRound = 0;
String contextLeafId = currentLeafId();
try {
BeforeTurnHookResult beforeHook = turnHooks.beforeTurn(new BeforeTurnHookContext(request, turnId, ports.cwd()));
if (beforeHook != null && beforeHook.blocked()) {
String message = beforeHook.message() == null ? "turn hook blocked" : beforeHook.message();
ports.eventBus().publish(new ErrorEvent(
request.sessionId(),
"turn-hook-blocked",
message,
clock.instant()
));
return failedState(request, turnId, null, List.of(), 0, startedAt, contextLeafId);
}
Optional<String> unsafeReason = continuationGuard.unsafeContinuationReason(currentLeafId());
if (unsafeReason.isPresent()) {
ports.eventBus().publish(new ErrorEvent(
request.sessionId(),
unsafeReason.orElseThrow(),
"当前分支停在 assistant 工具调用消息上,不能直接追加用户消息。请选择上一条用户消息或工具结果之后继续。",
clock.instant()
));
return failedState(request, turnId, null, List.of(), 0, startedAt, currentLeafId());
}

AgentMessage user = messageFactory.userMessage(ids.newMessageId(), request.userInput());
contextLeafId = appendNewMessage(request.sessionId(), user);
newMessages.add(user);

context = buildContext(request, Optional.of(contextLeafId));
AgentMessage assistant = runModel(request, context);
contextLeafId = appendStartedMessage(request.sessionId(), assistant);
newMessages.add(assistant);
if (isAssistantError(assistant, request)) {
return failedState(turnId, request.sessionId(), context, newMessages, toolRound, startedAt, contextLeafId);
return failedState(request, turnId, context, newMessages, toolRound, startedAt, contextLeafId);
}

while (!request.abortSignal().aborted()) {
Expand All @@ -106,9 +129,9 @@ private TurnState executeWithTurnId(TurnRequest request, String turnId) {
"incomplete-tool-call",
"模型返回的工具调用参数未完成,已终止本轮执行。"
);
appendNewMessage(request.sessionId(), error);
contextLeafId = appendNewMessage(request.sessionId(), error);
newMessages.add(error);
return failedState(turnId, request.sessionId(), context, newMessages, toolRound, startedAt, contextLeafId);
return failedState(request, turnId, context, newMessages, toolRound, startedAt, contextLeafId);
}
List<ToolUseRequest> toolRequests = toolCallMapper.requestsFrom(assistant);
if (toolRequests.isEmpty()) {
Expand All @@ -134,7 +157,7 @@ private TurnState executeWithTurnId(TurnRequest request, String turnId) {
contextLeafId = appendStartedMessage(request.sessionId(), assistant);
newMessages.add(assistant);
if (isAssistantError(assistant, request)) {
return failedState(turnId, request.sessionId(), context, newMessages, toolRound, startedAt, contextLeafId);
return failedState(request, turnId, context, newMessages, toolRound, startedAt, contextLeafId);
}
}
} catch (RuntimeException failure) {
Expand All @@ -145,13 +168,12 @@ private TurnState executeWithTurnId(TurnRequest request, String turnId) {
);
appendNewMessage(request.sessionId(), handled.message());
newMessages.add(handled.message());
return failedState(turnId, request.sessionId(), context, newMessages, toolRound, startedAt, currentLeafId());
return failedState(request, turnId, context, newMessages, toolRound, startedAt, currentLeafId());
}

TurnStatus status = request.abortSignal().aborted() ? TurnStatus.ABORTED : TurnStatus.COMPLETED;
TurnState state = new TurnState(turnId, request.sessionId(), context, List.copyOf(newMessages), toolRound, status);
eventPublisher.publishTurnEnd(request.sessionId(), turnId, status, startedAt, toolRound, contextLeafId);
return state;
return finishTurn(request, state, startedAt, contextLeafId);
}

private boolean isAssistantError(AgentMessage assistant, TurnRequest request) {
Expand All @@ -171,8 +193,8 @@ private String currentLeafId() {
}

private TurnState failedState(
TurnRequest request,
String turnId,
String sessionId,
ContextSnapshot context,
List<AgentMessage> newMessages,
int toolRound,
Expand All @@ -181,14 +203,46 @@ private TurnState failedState(
) {
TurnState state = new TurnState(
turnId,
sessionId,
request.sessionId(),
context,
List.copyOf(newMessages),
toolRound,
TurnStatus.FAILED
);
eventPublisher.publishTurnEnd(sessionId, turnId, TurnStatus.FAILED, startedAt, toolRound, leafEntryId);
return state;
return finishTurn(request, state, startedAt, leafEntryId);
}

private TurnState finishTurn(TurnRequest request, TurnState state, Instant startedAt, String leafEntryId) {
TurnState finalState = state;
try {
turnHooks.afterTurn(new AfterTurnHookContext(request, state, ports.cwd()));
} catch (RuntimeException failure) {
AgentCoreExceptionHandler.Failure handled = exceptionHandler.handle(
request.sessionId(),
ids.newMessageId(),
failure
);
appendNewMessage(request.sessionId(), handled.message());
List<AgentMessage> messages = new ArrayList<>(state.newMessages());
messages.add(handled.message());
finalState = new TurnState(
state.turnId(),
state.sessionId(),
state.context(),
List.copyOf(messages),
state.currentToolRound(),
TurnStatus.FAILED
);
}
eventPublisher.publishTurnEnd(
request.sessionId(),
finalState.turnId(),
finalState.status(),
startedAt,
finalState.currentToolRound(),
leafEntryId
);
return finalState;
}

private ContextSnapshot buildContext(TurnRequest request, Optional<String> leafEntryId) {
Expand Down
Loading
Loading