事件报告智能体工作过程中发生的情况。条目是已保存的消息和工具调用,您可以在之后获取这些内容。使用事件实时更新您的应用,使用条目显示已保存的历史记录。
您的应用通过发送输入事件来提交消息、取消轮次或返回工具结果。智能体发送的事件则报告输出和会话的变化。有关发送输入的说明,请参阅运行和继续会话。
请在发送任务前订阅事件流,以便您的应用接收到该轮次的早期事件。传入您的 API 客户端、对话的会话 ID 和事件处理函数:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38// Pass your saved session ID to this helper.
async function streamSession(client, sessionId, handleEvent) {
const events = await client.beta.agents.sessions.events.stream(sessionId);
try {
for await (const event of events) {
await handleEvent(event);
switch (event.type) {
case "agent.session.idle":
continue;
case "error":
throw new Error(event.error.message);
case "agent.session.failed":
case "agent.session.environment.failed":
throw new Error(`Agent lifecycle failure: ${event.type}`);
case "agent.session.turn.failed":
if (event.turn.subagent_id === null) {
throw new Error(
`${event.type}: ${event.turn.error?.message ?? ""}`
);
}
break;
case "agent.session.turn.cancelled":
if (event.turn.subagent_id === null) {
throw new Error("The agent turn was cancelled");
}
break;
case "agent.session.turn.completed":
if (event.turn.subagent_id === null) return;
break;
}
}
throw new Error(
"Stream closed before a turn ended. Retrieve the saved state."
);
} finally {
events.controller.abort();
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23# Pass your saved session ID to this helper.
def stream_session(client: OpenAI, session_id: str, handle_event):
with client.beta.agents.sessions.events.stream(session_id) as events:
for event in events:
handle_event(event)
match event.type:
case "agent.session.idle":
continue
case "error":
raise RuntimeError(event.error.message)
case "agent.session.failed" | "agent.session.environment.failed":
raise RuntimeError(f"Agent lifecycle failure: {event.type}")
case "agent.session.turn.failed":
if event.turn.subagent_id is None:
detail = event.turn.error.message if event.turn.error else ""
raise RuntimeError(f"{event.type}: {detail}")
case "agent.session.turn.cancelled":
if event.turn.subagent_id is None:
raise RuntimeError("The agent turn was cancelled")
case "agent.session.turn.completed":
if event.turn.subagent_id is None:
return
raise RuntimeError("Stream closed before a turn ended. Retrieve the saved state.")
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29// Pass your saved session ID to this helper.
func streamSession(ctx context.Context, client *openai.Client, sessionID string, handleEvent func(openai.AgentSessionEventUnion)) error {
events := client.Beta.Agents.Sessions.Events.StreamStreaming(ctx, sessionID)
defer events.Close()
for events.Next() {
event := events.Current()
handleEvent(event)
switch event.Type {
case "agent.session.idle":
continue
case "error":
return fmt.Errorf("agent error: %s", event.RawJSON())
case "agent.session.failed", "agent.session.environment.failed":
return fmt.Errorf("agent lifecycle failure: %s", event.RawJSON())
case "agent.session.turn.failed", "agent.session.turn.cancelled":
if event.Turn.SubagentID == "" {
return fmt.Errorf("agent turn did not complete: %s", event.RawJSON())
}
case "agent.session.turn.completed":
if event.Turn.SubagentID == "" {
return nil
}
}
}
if err := events.Err(); err != nil {
return err
}
return fmt.Errorf("stream closed before a turn ended; retrieve the saved state")
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30// Pass your saved session ID to this helper.
public static void streamSession(
OpenAIClient client, String sessionId, Consumer<AgentSessionEvent> handleEvent) {
try (StreamResponse<AgentSessionEvent> events =
client.beta().agents().sessions().events().streamStreaming(sessionId)) {
var iterator = events.stream().iterator();
while (iterator.hasNext()) {
var event = iterator.next();
handleEvent.accept(event);
if (event.idle().isPresent()) {
continue;
}
if (event.error().isPresent()) {
throw new IllegalStateException("Agent error: " + event);
}
if (event.failed().isPresent() || event.environmentFailed().isPresent()) {
throw new IllegalStateException("Agent lifecycle failure: " + event);
}
if (event.turnFailed().filter(e -> e.turn().subagentId().isEmpty()).isPresent()
|| event.turnCancelled().filter(e -> e.turn().subagentId().isEmpty()).isPresent()) {
throw new IllegalStateException("Agent turn did not complete: " + event);
}
if (event.turnCompleted().filter(e -> e.turn().subagentId().isEmpty()).isPresent()) {
return;
}
}
throw new IllegalStateException(
"Stream closed before a turn ended. Retrieve the saved state.");
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26# Pass your saved session ID to this helper.
def stream_session(client, session_id, &handle_event)
events = client.beta.agents.sessions.events.stream_streaming(session_id)
begin
events.each do |event|
handle_event.call(event)
case event.type.to_s
when "agent.session.idle"
next
when "error"
raise event.error.message
when "agent.session.failed", "agent.session.environment.failed"
raise "Agent lifecycle failure: #{event.type}"
when "agent.session.turn.failed"
raise "#{event.type}: #{event.turn.error&.message}" if event.turn.subagent_id.nil?
when "agent.session.turn.cancelled"
raise "The agent turn was cancelled" if event.turn.subagent_id.nil?
when "agent.session.turn.completed"
return nil if event.turn.subagent_id.nil?
end
end
raise "Stream closed before a turn ended. Retrieve the saved state."
ensure
events.close
end
end
1
2
3
4
5curl -N \
"https://api.openai.com/v1/agents/sessions/$session_id/events?stream=true" \
-H "OpenAI-Beta: agents=v1" \
-H "Authorization: Bearer $OPENAI_API_KEY" \
-H "Accept: text/event-stream"
辅助函数会将每个事件传递给您的处理程序,然后检查常见的事件类型。收到 agent.session.idle 时,它会继续处理,并在根轮次完成时返回。如果根轮次失败或被取消、会话或环境发生故障,或收到 error 事件,它会抛出错误。子智能体的轮次事件不会结束流。您的处理程序决定如何显示输出;调用方负责处理辅助函数抛出的错误。如果流在轮次结束前关闭,辅助函数会抛出错误。请参阅恢复断开的流。
订阅后发送消息
此版本接收一条消息,并在打开流后提交该消息:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46// Pass your saved session ID and message to this helper.
async function sendAndStream(client, sessionId, text, handleEvent) {
const events = await client.beta.agents.sessions.events.stream(sessionId);
try {
await client.beta.agents.sessions.events.create(sessionId, {
events: [
{
type: "agent.session.input.message",
input: [{ role: "user", content: [{ type: "input_text", text }] }],
},
],
});
for await (const event of events) {
await handleEvent(event);
switch (event.type) {
case "agent.session.idle":
continue;
case "error":
throw new Error(event.error.message);
case "agent.session.failed":
case "agent.session.environment.failed":
throw new Error(`Agent lifecycle failure: ${event.type}`);
case "agent.session.turn.failed":
if (event.turn.subagent_id === null) {
throw new Error(
`${event.type}: ${event.turn.error?.message ?? ""}`
);
}
break;
case "agent.session.turn.cancelled":
if (event.turn.subagent_id === null) {
throw new Error("The agent turn was cancelled");
}
break;
case "agent.session.turn.completed":
if (event.turn.subagent_id === null) return;
break;
}
}
throw new Error(
"Stream closed before a turn ended. Retrieve the saved state."
);
} finally {
events.controller.abort();
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37# Pass your saved session ID and message to this helper.
def send_and_stream(client: OpenAI, session_id: str, text, handle_event):
with client.beta.agents.sessions.events.stream(session_id) as events:
client.beta.agents.sessions.events.create(
session_id,
events=[
{
"type": "agent.session.input.message",
"input": [
{
"role": "user",
"content": [{"type": "input_text", "text": text}],
}
],
}
],
)
for event in events:
handle_event(event)
match event.type:
case "agent.session.idle":
continue
case "error":
raise RuntimeError(event.error.message)
case "agent.session.failed" | "agent.session.environment.failed":
raise RuntimeError(f"Agent lifecycle failure: {event.type}")
case "agent.session.turn.failed":
if event.turn.subagent_id is None:
detail = event.turn.error.message if event.turn.error else ""
raise RuntimeError(f"{event.type}: {detail}")
case "agent.session.turn.cancelled":
if event.turn.subagent_id is None:
raise RuntimeError("The agent turn was cancelled")
case "agent.session.turn.completed":
if event.turn.subagent_id is None:
return
raise RuntimeError("Stream closed before a turn ended. Retrieve the saved state.")
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56// Pass your saved session ID and message to this helper.
func sendAndStream(ctx context.Context, client *openai.Client, sessionID string, text string, handleEvent func(openai.AgentSessionEventUnion)) error {
events := client.Beta.Agents.Sessions.Events.StreamStreaming(ctx, sessionID)
defer events.Close()
if err := events.Err(); err != nil {
return err
}
err := client.Beta.Agents.Sessions.Events.New(ctx,
sessionID,
openai.BetaAgentSessionEventNewParams{
Events: []openai.AgentSessionInputParamUnion{
{
OfParamAgentSessionInputMessage: &openai.AgentSessionInputParamAgentSessionInputMessage{
Input: []openai.AgentSessionInputMessageParam{
{
Content: []openai.InputContentParamUnion{
{
OfParamInputText: &openai.InputContentParamInputText{
Text: text,
},
},
},
},
},
},
},
},
})
if err != nil {
return err
}
for events.Next() {
event := events.Current()
handleEvent(event)
switch event.Type {
case "agent.session.idle":
continue
case "error":
return fmt.Errorf("agent error: %s", event.RawJSON())
case "agent.session.failed", "agent.session.environment.failed":
return fmt.Errorf("agent lifecycle failure: %s", event.RawJSON())
case "agent.session.turn.failed", "agent.session.turn.cancelled":
if event.Turn.SubagentID == "" {
return fmt.Errorf("agent turn did not complete: %s", event.RawJSON())
}
case "agent.session.turn.completed":
if event.Turn.SubagentID == "" {
return nil
}
}
}
if err := events.Err(); err != nil {
return err
}
return fmt.Errorf("stream closed before a turn ended; retrieve the saved state")
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46// Pass your saved session ID and message to this helper.
public static void sendAndStream(
OpenAIClient client, String sessionId, String text, Consumer<AgentSessionEvent> handleEvent) {
try (StreamResponse<AgentSessionEvent> events =
client.beta().agents().sessions().events().streamStreaming(sessionId)) {
client
.beta()
.agents()
.sessions()
.events()
.create(
EventCreateParams.builder()
.sessionId(sessionId)
.addEvent(
AgentSessionInputParam.AgentSessionInputMessage.builder()
.addInput(
AgentSessionInputMessageParam.builder()
.addInputTextContent(text)
.build())
.build())
.build());
var iterator = events.stream().iterator();
while (iterator.hasNext()) {
var event = iterator.next();
handleEvent.accept(event);
if (event.idle().isPresent()) {
continue;
}
if (event.error().isPresent()) {
throw new IllegalStateException("Agent error: " + event);
}
if (event.failed().isPresent() || event.environmentFailed().isPresent()) {
throw new IllegalStateException("Agent lifecycle failure: " + event);
}
if (event.turnFailed().filter(e -> e.turn().subagentId().isEmpty()).isPresent()
|| event.turnCancelled().filter(e -> e.turn().subagentId().isEmpty()).isPresent()) {
throw new IllegalStateException("Agent turn did not complete: " + event);
}
if (event.turnCompleted().filter(e -> e.turn().subagentId().isEmpty()).isPresent()) {
return;
}
}
throw new IllegalStateException(
"Stream closed before a turn ended. Retrieve the saved state.");
}
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45# Pass your saved session ID and message to this helper.
def send_and_stream(client, session_id, text, &handle_event)
events = client.beta.agents.sessions.events.stream_streaming(session_id)
begin
client.beta.agents.sessions.events.create(
session_id,
events: [
{
type: "agent.session.input.message",
input: [
{
role: "user",
content: [
{
type: "input_text",
text: text
}
]
}
]
}
]
)
events.each do |event|
handle_event.call(event)
case event.type.to_s
when "agent.session.idle"
next
when "error"
raise event.error.message
when "agent.session.failed", "agent.session.environment.failed"
raise "Agent lifecycle failure: #{event.type}"
when "agent.session.turn.failed"
raise "#{event.type}: #{event.turn.error&.message}" if event.turn.subagent_id.nil?
when "agent.session.turn.cancelled"
raise "The agent turn was cancelled" if event.turn.subagent_id.nil?
when "agent.session.turn.completed"
return nil if event.turn.subagent_id.nil?
end
end
raise "Stream closed before a turn ended. Retrieve the saved state."
ensure
events.close
end
end
根据事件的 type 决定您的应用应执行的操作:
- 显示文本: 将
agent.session.turn.output_text.delta 追加到相应的内容部分。收到 agent.session.turn.output_text.done 时,用该部分的完整文本替换它。增量事件可能不会出现。
- 跟踪工作进度: 会话、轮次和条目事件会报告进度。检查是否收到
agent.session.turn.completed、agent.session.turn.failed 或 agent.session.turn.cancelled,以确定轮次的结果。
- 提供所需输入: 收到
agent.session.requires_action 时,获取会话并检查 required_actions。您的代码可能需要返回函数结果或连接环境。
仅凭会话空闲或流已关闭,并不能确定任务成功。轮次完成也不保证每个工具都执行成功。请检查智能体的输出。
使用 item_id、output_index 和 content_index 将文本更新关联到同一个内容部分。例如,以下简化后的事件会更新同一个部分:
1234567{
"type": "agent.session.turn.output_text.delta",
"item_id": "msg_789",
"output_index": 0,
"content_index": 0,
"delta": "Acme competes"
}
1234567{
"type": "agent.session.turn.output_text.done",
"item_id": "msg_789",
"output_index": 0,
"content_index": 0,
"text": "Acme competes on price and distribution."
}
每个事件都有自己的 event_id。相同的 item_id 标识同一个已保存条目,其中包含消息的内容、状态和阶段。请参阅获取已保存的工作记录。
有关所有事件类型和字段,请参阅流式事件参考资料。这些流式事件与 Webhook 不同。有关子智能体活动和命令归属的信息,请参阅观察委派过程。
使用您的应用对话状态中保存的会话 ID,获取已保存的工作记录:
- 会话条目: 通过列出条目,获取根智能体在各个轮次中的消息和工具调用。
- 轮次: 通过列出轮次浏览会话中的工作。按 ID 获取轮次,以检查其状态、时间戳、用量和错误。
- 单个轮次中的条目: 对于根智能体的轮次,按
turn_id 筛选会话条目。每个子智能体都有自己的条目历史记录和按轮次获取条目的端点。
列表端点每次返回一页结果。使用 SDK 分页辅助函数或 after 游标获取更多结果。单页结果可能不包含某个轮次的所有条目。使用 order: "asc" 按从旧到新的顺序读取条目。
流不会重放错过的事件。要恢复您的应用视图,请执行以下操作:
- 打开新的流,并缓冲传入的事件。
- 保持流连接,同时获取会话及其已保存的条目。
- 以条目 ID 为键,根据这些条目恢复本地状态。
- 使用
item_id 应用已缓冲的条目更新。如果获取的历史记录显示某个条目已达到最终状态,则丢弃该条目的更新。
- 恢复处理实时事件。
output_text.done 事件可用完整文本替换临时文本缓冲区中的内容。已保存的条目可帮助您恢复已完成的工作,但无法恢复您错过的每一个中间事件。