イベントは、エージェントの作業中に起きたことを通知します。アイテムは、後から取得できる保存済みのメッセージやツール呼び出しです。イベントを使ってアプリケーションをリアルタイムに更新し、アイテムを使って保存された履歴を表示します。
アプリケーションは入力イベントを送信して、メッセージの送信、ターンのキャンセル、ツールの結果の返却を行います。エージェントは、出力やセッションの変更を通知するイベントを送信します。入力の送信については、セッションの実行と継続を参照してください。
ターンの初期イベントをアプリケーションで受信できるように、作業を送信する前に購読を開始します。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 を使って、テキストの更新を同じコンテンツパートに対応付けます。たとえば、以下の省略したイベント例は、1 つのパートを更新します。
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 で絞り込みます。各サブエージェントには、独自のアイテム履歴とターンごとのアイテム取得用エンドポイントがあります。
一覧取得エンドポイントは、1 回に 1 ページを返します。さらに結果を取得するには、SDK のページネーションヘルパーまたは after カーソルを使います。1 ページにターンのすべてのアイテムが含まれるとは限りません。アイテムを古い順に読み取るには、order: "asc" を使います。
ストリームは、受信できなかったイベントを再送しません。アプリケーションの表示を復元するには、次の手順を実行します。
- 新しいストリームを開き、受信するイベントをバッファーに保存します。
- ストリームの接続を維持したまま、セッションと保存済みのアイテムを取得します。
- アイテム ID をキーとして、それらのアイテムからローカルの状態を復元します。
item_id を使って、バッファーに保存したアイテムの更新を適用します。取得した履歴ですでに最終状態に達しているアイテムについては、更新を破棄します。
- リアルタイムのイベント処理を再開します。
output_text.done イベントを使うと、一時的なテキストバッファーを完全なテキストで置き換えられます。保存済みのアイテムから完了した作業内容を復元できますが、受信できなかった途中のイベントをすべて復元できるわけではありません。