WebSocket Mode

WebSocket Mode

Responses API는 오래 실행되고 도구 호출이 많은 워크플로를 위한 WebSocket 모드를 지원해요. 지연을 줄이는 것 외에도 stream_id를 사용하면 WebSocket 멀티플렉싱이 가능해져요. /v1/responses에 대한 하나의 영구 연결로 병렬 대화를 실행하고 기존 대화를 새 스트림으로 포크(fork)할 수 있어요. 각 턴은 새 입력 항목과 previous_response_id만 보내면 계속할 수 있어요.

출처: 문서

본문

WebSocket 모드는 Zero Data Retention (ZDR)과 store=false 모두와 호환돼요.

WebSocket 모드를 쓰는 이유

WebSocket 모드는 워크플로가 많은 모델-도구 왕복을 포함할 때(예: 반복된 도구 호출이 있는 에이전틱 코딩이나 오케스트레이션 루프) 가장 유용해요.

연결이 계속 열려 있고 각 턴이 증분 입력만 보내기 때문에, WebSocket 모드는 턴마다의 계속(continuation) 오버헤드를 줄이고 긴 체인에서 엔드투엔드 지연을 개선해요. 도구 호출이 20개 이상인 롤아웃에서는 엔드투엔드 실행이 대략 40%까지 빨라지는 것을 확인했어요.

연결하고 응답 만들기

WebSocket 의존성을 Python은 pip install "openai[realtime]>=3.8.0", JavaScript는 npm install openai@^7.10.0 ws, Ruby는 gem install openai async-websocket으로 설치하세요.

WebSocket 모드에서는 클라이언트가 response.create 이벤트를 보내 각 턴을 시작해요. 페이로드는 일반 Responses create 본문과 동일하지만, stream과 background 같은 전송 특정 필드는 사용하지 않아요.

import OpenAI from "openai";
import { ResponsesWS } from "openai/resources/responses/ws";

const client = new OpenAI();

const ws = new ResponsesWS(client);
try {
  ws.send({
    type: "response.create",
    stream_id: "main",
    model: "gpt-6-astra",
    store: false,
    input: [
      {
        type: "message",
        role: "user",
        content: [{ type: "input_text", text: "Find fizz_buzz()" }],
      },
    ],
    tools: [],
  });
  let completed = false;
  for await (const event of ws) {
    if (event.type === "error") throw event.error;
    if (event.type !== "message") continue;
    const message = event.message;
    if (message.type === "response.output_text.delta") {
      process.stdout.write(message.delta);
    } else if (message.type === "response.completed") {
      completed = true;
      break;
    } else if (
      message.type === "response.failed" ||
      message.type === "response.incomplete"
    ) {
      throw new Error(JSON.stringify(message));
    }
  }
  if (!completed)
    throw new Error("Connection closed before the response finished.");
} finally {
  ws.close();
}
from openai import OpenAI

client = OpenAI()

with client.responses.connect() as connection:
    connection.response.create(
        stream_id="main",
        model="gpt-6-astra",
        store=False,
        input=[
            {
                "type": "message",
                "role": "user",
                "content": [{"type": "input_text", "text": "Find fizz_buzz()"}],
            }
        ],
        tools=[],
    )
    for event in connection:
        if event.type == "response.completed":
            print(event.response.output_text)
            break
        if event.type in {"response.failed", "response.incomplete", "error"}:
            raise RuntimeError(event.to_json())
require "async"
require "openai"

def wait_for_response(connection)
  while (event = connection.receive)
    case event.type.to_s
    when "response.completed" then return event.response
    when "response.failed", "response.incomplete", "error"
      raise "Response failed: #{event.to_json}"
    end
  end
  raise "Connection closed before the response finished"
end

client = OpenAI::Client.new
Sync do |task|
  task.with_timeout(120) do
    client.responses.connect(request_options: { timeout: 10 }) do |connection|
      connection.response.create(
        stream_id: "main", model: "gpt-6-astra", store: false,
        input: [
          {
            role: "user",
            content: "Find fizz_buzz()"
          }
        ], tools: []
      )
      puts(wait_for_response(connection).output_text)
    end
  end
end

클라이언트는 선택적으로 generate: false로 response.create를 보내 요청 상태를 예열할 수 있어요. 이는 다음 턴에 보낼 도구, 지시, 그리고/또는 커스텀 메시지를 이미 알고 있을 때 유용해요. generate: false는 모델 출력을 반환하지 않지만 요청 상태를 준비해서 다음 생성 턴이 더 빨리 시작될 수 있게 해요. 예열 요청은 응답 ID를 반환하며, 여기서 previous_response_id로 이어갈 수 있어요(응답 체인의 이후 턴에도). 다음 섹션에서는 previous_response_id와 증분 입력을 사용해 세션을 계속하는 방법을 설명해요.

증분 입력으로 계속

응답이 아직 실행 중일 때 사용자 지시를 추가하려면 Mid-turn steering을 사용하세요. Steering은 완료된 작업을 보존하고 새 지시를 연속(continuation)에 포함해요. 일반적인 턴 사이 연속과 도구 결과에는 다음 response.create 패턴을 사용하세요.

실행을 계속하려면 다음으로 response.create를 보내세요.

  • previous_response_id를 이전 응답 ID로 설정.
  • input에 새 항목만 포함 (예: 도구 출력과 다음 사용자 메시지).
import OpenAI from "openai";
import { ResponsesWS } from "openai/resources/responses/ws";

const client = new OpenAI();
const model = "gpt-6-astra";

const tools = [
  {
    type: "function",
    name: "get_test_results",
    description: "Return a local demo test result.",
    parameters: { type: "object", properties: {}, additionalProperties: false },
    strict: true,
  },
];

async function waitForResponse(ws) {
  for await (const event of ws) {
    if (event.type === "error") throw event.error;
    if (event.type !== "message") continue;
    const message = event.message;
    if (message.type === "response.output_text.delta") {
      process.stdout.write(message.delta);
    } else if (message.type === "response.completed") {
      return message.response;
    } else if (
      message.type === "response.failed" ||
      message.type === "response.incomplete"
    ) {
      throw new Error(JSON.stringify(message));
    }
  }
  throw new Error("Connection closed before the response finished.");
}

const ws = new ResponsesWS(client);
try {
  ws.send({
    type: "response.create",
    stream_id: "main",
    model,
    store: false,
    input: "Find the failing test and suggest a fix.",
    tools,
    tool_choice: { type: "function", name: "get_test_results" },
    parallel_tool_calls: false,
  });
  const first = await waitForResponse(ws);
  const call = first.output.find((item) => item.type === "function_call");
  if (!call || call.name !== "get_test_results") {
    throw new Error("Expected a get_test_results function call.");
  }
  const result = {
    test: "test_fizz_buzz",
    failure: 'Expected "FizzBuzz" for 15, got "Fizz".',
  };

  // Continue on the same socket with the actual response and tool-call IDs.
  ws.send({
    type: "response.create",
    stream_id: "main",
    model,
    store: false,
    previous_response_id: first.id,
    input: [
      {
        type: "function_call_output",
        call_id: call.call_id,
        output: JSON.stringify(result),
      },
      { role: "user", content: "Now optimize it." },
    ],
    tools,
    tool_choice: "none",
  });
  await waitForResponse(ws);
} finally {
  ws.close();
}
import json

from openai import OpenAI
from openai.resources.responses.responses import ResponsesConnection
from openai.types.responses import FunctionToolParam, Response

client = OpenAI()
model = "gpt-6-astra"
tools: list[FunctionToolParam] = [
    {
        "type": "function",
        "name": "get_test_results",
        "description": "Read the demo test results.",
        "parameters": {
            "type": "object",
            "properties": {},
            "required": [],
            "additionalProperties": False,
        },
        "strict": True,
    }
]


def get_test_results():
    # Demo data. Replace this function with your test runner.
    return {
        "test": "test_fizz_buzz",
        "failure": 'Expected "FizzBuzz" for 15, got "Fizz".',
    }


def wait_for_response(connection: ResponsesConnection) -> Response:
    for event in connection:
        if event.type == "response.completed":
            return event.response
        if event.type in {"response.failed", "response.incomplete", "error"}:
            raise RuntimeError(event.to_json())
    raise RuntimeError("Connection closed before the response finished.")


with client.responses.connect() as connection:
    connection.response.create(
        stream_id="main",
        model=model,
        store=False,
        input="Find the failing test and suggest a fix.",
        tools=tools,
        tool_choice={"type": "function", "name": "get_test_results"},
        parallel_tool_calls=False,
    )
    response = wait_for_response(connection)
    call = next(item for item in response.output if item.type == "function_call")
    if call.name != "get_test_results" or json.loads(call.arguments) != {}:
        raise ValueError("Expected a get_test_results call with no arguments")

    # Continue on the same connection using the actual response and tool-call IDs.
    connection.response.create(
        stream_id="main",
        model=model,
        store=False,
        previous_response_id=response.id,
        input=[
            {
                "type": "function_call_output",
                "call_id": call.call_id,
                "output": json.dumps(get_test_results()),
            },
            {"role": "user", "content": "Now optimize it."},
        ],
        tools=tools,
        tool_choice="none",
    )
    print(wait_for_response(connection).output_text)
require "async"
require "openai"
require "json"

def wait_for_response(connection)
  while (event = connection.receive)
    case event.type.to_s
    when "response.completed" then return event.response
    when "response.failed", "response.incomplete", "error"
      raise "Response failed: #{event.to_json}"
    end
  end
  raise "Connection closed before the response finished"
end

tools = [
  {
    type: "function",
    name: "get_test_results",
    description: "Read the demo test results.",
    parameters: {
      type: "object",
      properties: {},
      required: [],
      additionalProperties: false
    },
    strict: true
  }
]

client = OpenAI::Client.new
Sync do |task|
  task.with_timeout(120) do
    client.responses.connect(request_options: { timeout: 10 }) do |connection|
      connection.response.create(
        stream_id: "main", model: "gpt-6-astra", store: false,
        input: "Find the failing test and suggest a fix.", tools: tools,
        tool_choice: {
          type: "function",
          name: "get_test_results"
        }, parallel_tool_calls: false
      )
      response = wait_for_response(connection)
      call = response.output.grep(OpenAI::Responses::ResponseFunctionToolCall).first
      unless call && call.name == "get_test_results" && JSON.parse(call.arguments) == {}
        raise "Expected a get_test_results call with no arguments"
      end

      # Demo data. Replace this with your test runner.
      result = {
        test: "test_fizz_buzz",
        failure: 'Expected "FizzBuzz" for 15, got "Fizz".'
      }
      connection.response.create(
        stream_id: "main", model: "gpt-6-astra", store: false,
        previous_response_id: response.id,
        input: [
          {
            type: "function_call_output",
            call_id: call.call_id,
            output: JSON.generate(result)
          },
          {
            role: "user",
            content: "Now optimize it."
          }
        ],
        tools: tools, tool_choice: "none"
      )
      puts(wait_for_response(connection).output_text)
    end
  end
end

연속(continuation)이 작동하는 방식

WebSocket 모드는 HTTP 모드와 같은 previous_response_id 체이닝 의미론을 사용하지만, 활성 소켓에 더 낮은 지연의 연속 경로를 추가해요.

활성 WebSocket 연결에서 서비스는 최근 previous-response 상태를 연결 로컬 인메모리 캐시에 보관해요. stream_id를 사용하면 각 레인(lane)이 가장 최근 캐시 응답을 유지하므로, 서비스가 연결 로컬 상태를 재사용할 수 있어 그 레인의 최신 응답에서 계속하는 것이 빠르다는 뜻이에요. 서비스는 previous-response 상태를 메모리에만 보관하고 디스크에 기록하지 않기 때문에, store=false와 Zero Data Retention (ZDR)과 호환되는 방식으로 WebSocket 모드를 사용할 수 있어요.

previous_response_id가 인메모리 캐시에 없으면 동작은 응답을 저장하는지 여부에 따라 달라져요.

  • store=true면 서비스는 가능할 때 영구 상태에서 더 오래된 응답 ID를 복원할 수 있어요. 연속은 여전히 동작하지만 인메모리 지연 이점은 사라져요.
  • store=false(ZDR 포함)면 영구 대체가 없어요. ID가 캐시되지 않으면 요청은 previous_response_not_found를 반환해요.

같은 레인 연속이 4xx나 5xx를 반환하면 서비스는 참조된 previous_response_id를 연결 로컬 캐시에서 퇴출해요. 오류를 반환하는 크로스-레인 포크는 공유 부모를 보존하므로 원본 레인이 계속될 수 있어요.

컴팩션과 새 응답 만들기

컴팩션(compaction)을 사용한다면 두 가지 다른 연속 패턴이 있어요.

서버 측 컴팩션 (context_management)

서버 측 컴팩션(compact_threshold가 있는 context_management)을 활성화하면 컴팩션은 일반 /responses 생성 중에 일어나요. WebSocket 모드에서는 평소처럼 계속하면 돼요. 최신 previous_response_id와 새 입력 항목만 보내 다음 response.create를 보내세요.

독립형 /responses/compact

독립형 /responses/compact 엔드포인트는 응답 ID가 아니라 새 컴팩트 입력 창을 반환해요. 컴팩션 후에는 WebSocket 연결에서 컴팩트 창을 input으로 사용해(다음 사용자/도구 항목과 함께) 새 응답을 만드세요.

previous_response_id를 생략하거나 null로 설정해 새 체인을 시작하세요. 컴팩트 출력을 그대로 전달하세요. 반환된 창을 정리(prune)하지 마세요.

import { toResponseInputItems } from "openai/lib/responses/ResponseInputItems";

// Compact your current window with an HTTP request.
const compacted = await client.responses.compact({
  model: "gpt-6-astra",
  input: longInputItems,
});
const nextInput = toResponseInputItems(compacted.output);
nextInput.push({
  type: "message",
  role: "user",
  content: [{ type: "input_text", text: "Continue from here." }],
});

// Start a new response on the WebSocket using the compacted window.
const ws = new ResponsesWS(client);
try {
  ws.send({
    type: "response.create",
    stream_id: "main",
    model: "gpt-6-astra",
    store: false,
    input: nextInput,
    tools: [],
  });
  let completed = false;
  for await (const event of ws) {
    if (event.type === "error") throw event.error;
    if (event.type !== "message") continue;
    const message = event.message;
    if (message.type === "response.output_text.delta") {
      process.stdout.write(message.delta);
    } else if (message.type === "response.completed") {
      completed = true;
      break;
    } else if (
      message.type === "response.failed" ||
      message.type === "response.incomplete"
    ) {
      throw new Error(JSON.stringify(message));
    }
  }
  if (!completed)
    throw new Error("Connection closed before the response finished.");
} finally {
  ws.close();
}
from typing import cast

from openai import OpenAI
from openai.types.responses import ResponseInputParam

# Compact your current window (HTTP call).
compacted = client.responses.compact(
    model="gpt-6-astra",
    input=long_input_items_array,
)
next_input = cast(
    ResponseInputParam,
    [item.to_dict() for item in compacted.output],
)
next_input.append(
    {
        "type": "message",
        "role": "user",
        "content": [{"type": "input_text", "text": "Continue from here."}],
    }
)

# Start a new response on the WebSocket using the compacted window.
with client.responses.connect() as connection:
    connection.response.create(
        stream_id="main",
        model="gpt-6-astra",
        store=False,
        input=next_input,
        tools=[],
    )
    for event in connection:
        if event.type == "response.completed":
            print(event.response.output_text)
            break
        if event.type in {"response.failed", "response.incomplete", "error"}:
            raise RuntimeError(event.to_json())
require "async"
require "openai"

def wait_for_response(connection)
  while (event = connection.receive)
    case event.type.to_s
    when "response.completed" then return event.response
    when "response.failed", "response.incomplete", "error"
      raise "Response failed: #{event.to_json}"
    end
  end
  raise "Connection closed before the response finished"
end

client = OpenAI::Client.new
compacted = client.responses.compact(
  model: "gpt-6-astra",
  input: [
    {
      role: :user,
      content: "Find the failing test."
    }
  ]
)
next_input = compacted.output.map(&:to_h)
next_input << {
  role: :user,
  content: "Continue from here."
}

Sync do |task|
  task.with_timeout(120) do
    client.responses.connect(request_options: { timeout: 10 }) do |connection|
      connection.response.create(
        stream_id: "main", model: "gpt-6-astra", store: false,
        input: next_input, tools: []
      )
      puts(wait_for_response(connection).output_text)
    end
  end
end

대화를 병렬로 실행

stream_id 파라미터를 사용해 같은 연결에서 병렬 대화를 유지할 수 있어요. 서로 다른 stream_id 값으로 독립적인 response.create 이벤트를 연달아 보내세요. 서버가 하나의 연결에서 그들을 동시에 실행할 수 있어요. 그들의 이벤트는 서로 섞일 수 있으므로 리더 루프 하나를 유지하고 각 이벤트를 stream_id로 라우팅하세요.

stream_id는 하나의 WebSocket 연결에서 정렬된 레인을 이름 붙입니다. stream_id와 previous_response_id를 구분하세요.

  • stream_id는 이벤트가 어디로 가는지, 어떤 요청이 선입선출(first-in, first-out) 순서로 실행되는지를 통제해요.
  • previous_response_id는 대화의 계보(lineage)를 통제해요.

그 구분은 두 가지 유용한 패턴을 가능하게 해요.

one WebSocket connection
├─ stream_id="planner"   draft a deployment plan
└─ stream_id="research"  list deployment risks

같은 stream_id를 가진 요청은 선입선출로 유지되고 겹치지 않아요. 서로 다른 stream_id를 가진 요청은 동시에 실행될 수 있어요.

연결당 한계

  • 연결은 이름 붙은 레인과 기본 레인을 합쳐 최대 16개의 활성 진행 중 응답을 가질 수 있어요. 연결은 더 많은 response.create 이벤트를 받아들이고 활성 응답이 끝날 때까지 대기시켜요.
  • 연결은 최대 32개의 서로 다른 이름 붙은 stream_id 값을 받아들여요. 암시적 기본 레인은 이 이름 붙은 스트림 한계에 계산되지 않아요. 한계에 도달하면 기존 stream_id를 재사용하거나 새 연결을 여세요.

대화를 새 스트림으로 포크

완료된 응답에서 분기하려면 새 stream_id와 함께 그 ID를 previous_response_id로 보내세요. 그 응답이 사용 가능한 동안 새 스트림은 그 컨텍스트를 상속하고 원본 스트림은 계속 갈 수 있어요. 포크가 시작된 후에는 다른 스트림 ID를 사용하므로 두 분기가 동시에 실행될 수 있어요.

store=false(ZDR 포함)에서는 크로스-레인 포크가 부모가 연결 로컬 캐시에 남아 있는 것에 의존해요. 포크가 큐에 있는 동안 소스 레인이 진행하거나 실패하면 부모가 포크가 시작되기 전에 퇴출될 수 있고, 포크는 previous_response_not_found를 반환해요. 소스 레인을 진행하기 전에 포크 레인이 response.in_progress를 방출할 때까지 기다리거나, previous_response_id를 null로 설정하고 전체 입력 컨텍스트를 재생하며 재시도하세요.

main:   resp_1 ──▶ resp_2 ──▶ resp_3
                       ╲
critic:                 resp_4 ──▶ resp_5

previous_response_id 없이 stream_id를 재사용하면 새 응답이 시작돼요. 대화가 계속되는 게 아니에요.

핵심 호출은 다음과 같아요.

# One socket, two independent conversations.
send_create(connection, "planner", "Draft a deployment plan.")
send_create(connection, "research", "List deployment risks.")

# Fork the planner response, then continue the original branch in parallel.
send_create(
    connection,
    "critic",
    "Find gaps in this plan.",
    previous_response_id=planner_response_id,
)
wait_for_in_progress(connection, "critic")
send_create(
    connection,
    "planner",
    "Add rollback steps.",
    previous_response_id=planner_response_id,
)

완전한 예시

병렬 대화 실행 후 하나 포크

import OpenAI from "openai";
import { ResponsesWS } from "openai/resources/responses/ws";

const client = new OpenAI();

const latestResponseIdByLane = new Map();

function sendCreate(
  ws,
  streamId,
  text,
  previousResponseId = latestResponseIdByLane.get(streamId)
) {
  ws.send({
    type: "response.create",
    stream_id: streamId,
    model: "gpt-6-astra",
    store: false,
    input: [
      {
        type: "message",
        role: "user",
        content: [{ type: "input_text", text }],
      },
    ],
    previous_response_id: previousResponseId,
  });
}

async function readMessage(events) {
  while (true) {
    const { value: event, done } = await events.next();
    if (done)
      throw new Error("Connection closed before all responses finished.");
    if (event.type === "error") throw event.error;
    if (event.type !== "message") continue;
    const message = event.message;
    if (
      message.type === "response.failed" ||
      message.type === "response.incomplete"
    ) {
      throw new Error(
        `Lane ${message.stream_id} failed: ${JSON.stringify(message)}`
      );
    }
    return message;
  }
}

async function drainUntilComplete(events, expectedStreamIds) {
  const remaining = new Set(expectedStreamIds);
  while (remaining.size > 0) {
    const message = await readMessage(events);
    const streamId = message.stream_id;
    if (!streamId || !remaining.has(streamId)) continue;
    if (message.type === "response.completed") {
      latestResponseIdByLane.set(streamId, message.response.id);
      remaining.delete(streamId);
    }
  }
}

async function waitForInProgress(events, streamId) {
  while (true) {
    const message = await readMessage(events);
    if (
      message.type === "response.in_progress" &&
      message.stream_id === streamId
    )
      return;
  }
}

const ws = new ResponsesWS(client);
// Keep one iterator so events stay queued while moving between phases.
const events = ws.stream();
try {
  // Run two independent conversations in parallel.
  sendCreate(
    ws,
    "planner",
    "Draft a deployment plan for a stateless API service."
  );
  sendCreate(
    ws,
    "research",
    "List common deployment risks for a stateless API service."
  );
  await drainUntilComplete(events, new Set(["planner", "research"]));

  // Fork the planner conversation and continue its original branch in parallel.
  const plannerResponseId = latestResponseIdByLane.get("planner");
  sendCreate(
    ws,
    "critic",
    "Find gaps in this deployment plan.",
    plannerResponseId
  );
  // Let the fork load its parent before advancing the original lane's cache.
  await waitForInProgress(events, "critic");
  sendCreate(
    ws,
    "planner",
    "Add rollback and monitoring steps to the plan.",
    plannerResponseId
  );
  await drainUntilComplete(events, new Set(["critic", "planner"]));
} finally {
  await events.return?.();
  ws.close();
}
from openai import OpenAI
from openai.resources.responses.responses import ResponsesConnection

client = OpenAI()
latest_response_id_by_lane: dict[str, str] = {}


def send_create(
    connection: ResponsesConnection,
    stream_id: str,
    text: str,
    previous_response_id: str | None = None,
):
    if previous_response_id is None:
        previous_response_id = latest_response_id_by_lane.get(stream_id)
    connection.response.create(
        stream_id=stream_id,
        model="gpt-6-astra",
        store=False,
        input=[
            {
                "type": "message",
                "role": "user",
                "content": [{"type": "input_text", "text": text}],
            }
        ],
        previous_response_id=previous_response_id,
    )


def drain_until_complete(
    connection: ResponsesConnection, expected_stream_ids: set[str]
):
    completed: set[str] = set()
    for event in connection:
        stream_id = event.stream_id
        if event.type == "error" and stream_id is None:
            raise RuntimeError(f"Connection error: {event.to_json()}")
        if stream_id is None or stream_id not in expected_stream_ids:
            continue

        if event.type == "response.completed":
            latest_response_id_by_lane[stream_id] = event.response.id
            completed.add(stream_id)
            if completed == expected_stream_ids:
                return
        elif event.type in {"response.failed", "response.incomplete", "error"}:
            raise RuntimeError(f"Lane {stream_id} failed: {event.to_json()}")
    raise RuntimeError("Connection closed before all responses finished.")


def wait_for_in_progress(connection: ResponsesConnection, expected_stream_id: str):
    for event in connection:
        if event.type == "error" and event.stream_id is None:
            raise RuntimeError(f"Connection error: {event.to_json()}")
        if event.stream_id != expected_stream_id:
            continue
        if event.type == "response.in_progress":
            return
        if event.type in {"response.failed", "response.incomplete", "error"}:
            raise RuntimeError(f"Lane {expected_stream_id} failed: {event.to_json()}")
    raise RuntimeError("Connection closed before the fork started.")


with client.responses.connect() as connection:
    # 1. Run two independent conversations in parallel.
    send_create(
        connection, "planner", "Draft a deployment plan for a stateless API service."
    )
    send_create(
        connection,
        "research",
        "List common deployment risks for a stateless API service.",
    )
    drain_until_complete(connection, {"planner", "research"})

    # 2. Fork the planner conversation and continue the original branch in parallel.
    planner_response_id = latest_response_id_by_lane["planner"]
    send_create(
        connection,
        "critic",
        "Find gaps in this deployment plan.",
        previous_response_id=planner_response_id,
    )
    # Let the fork bind its parent before advancing the original lane.
    wait_for_in_progress(connection, "critic")
    send_create(
        connection,
        "planner",
        "Add rollback and monitoring steps to the plan.",
        previous_response_id=planner_response_id,
    )
    drain_until_complete(connection, {"critic", "planner"})
require "async"
require "openai"
require "json"

def send_create(connection, stream_id, text, previous_response_id = nil)
  payload = {
    stream_id: stream_id,
    model: "gpt-6-astra",
    store: false,
    input: [
      {
        role: "user",
        content: text
      }
    ]
  }
  payload[:previous_response_id] = previous_response_id if previous_response_id
  connection.response.create(**payload)
end

def read_event(connection)
  event = connection.receive or raise "Connection closed before all responses finished"
  if ["response.failed", "response.incomplete", "error"].include?(event.type.to_s)
    raise "Response failed: #{event.to_json}"
  end

  event
end

def drain_responses(connection, lanes, latest_ids)
  remaining = lanes.dup
  until remaining.empty?
    event = read_event(connection)
    next unless event.type.to_s == "response.completed"

    lane = event.stream_id
    next unless remaining.include?(lane)

    latest_ids[lane] = event.response.id
    remaining.delete(lane)
  end
end

client = OpenAI::Client.new
Sync do |task|
  task.with_timeout(120) do
    client.responses.connect(request_options: { timeout: 10 }) do |connection|
      latest_ids = {}
      send_create(connection, "planner", "Draft a deployment plan for a stateless API service.")
      send_create(connection, "research", "List common deployment risks for a stateless API service.")
      drain_responses(connection, ["planner", "research"], latest_ids)
      parent_id = latest_ids.fetch("planner")
      send_create(connection, "critic", "Find gaps in this deployment plan.", parent_id)
      # Let the fork load its parent before advancing the original lane's cache.
      loop do
        event = read_event(connection)
        break if event.type.to_s == "response.in_progress" && event.stream_id == "critic"
      end
      send_create(connection, "planner", "Add rollback and monitoring steps.", parent_id)
      drain_responses(connection, ["critic", "planner"], latest_ids)
      puts(JSON.generate(latest_ids))
    end
  end
end

stream_id는 1~256자여야 하고 문자, 숫자, 밑줄(_), 하이픈(-), 마침표(.)만 포함할 수 있어요. WebSocket response.create 이벤트에서만 사용하세요. HTTP POST /v1/responses에는 포함하지 마세요.

이름 붙은 스트림에서는 서버 이벤트가 일치하는 stream_id를 포함해요(종료 이벤트와 요청 범위 오류 포함).

stream_id를 생략하면 요청은 암시적 기본 레인을 사용하고, 그 이벤트는 stream_id를 포함하지 않아요. 기본 레인은 그 외에는 이름 붙은 스트림과 같은 순서 및 동시성 규칙을 따라요. 빈 문자열은 유효한 stream_id가 아니에요. 필드를 생략해 기본 레인을 선택하세요.

연결 동작과 한계

  • 각 응답 내 이벤트는 기존 Responses 스트리밍 이벤트 모델을 따라요. 다른 레인의 이벤트는 서로 섞일 수 있어요.
  • 같은 stream_id를 가진 요청은 선입선출 순서로 실행되고 겹치지 않아요. 다른 레인의 요청은 동시에 실행될 수 있어요.
  • 연결은 최대 60분 동안 지속돼요. 한계에서 다시 연결하세요.

재연결과 복구

연결이 닫히면(또는 60분 한계에 도달하면) 연결 로컬 캐시는 모든 레인에서 사라져요. 새 WebSocket 연결을 열고 다음 패턴 중 하나로 각 레인을 복구하세요.

  1. 이전 응답을 저장했고(store=true) 유효한 응답 ID가 있다면 previous_response_id와 새 입력 항목으로 그 레인을 계속하세요.
  2. 레인을 계속할 수 없다면(예: store=false/ZDR 또는 previous_response_id가 previous_response_not_found), previous_response_id를 null로 설정하고 새 응답을 시작해 그 레인의 다음 턴 전체 입력 컨텍스트를 보내세요.
  3. /responses/compact로 컨텍스트를 컴팩트했다면 반환된 컴팩트 창을 그 새 응답의 기본 input으로 사용한 다음 최신 사용자/도구 항목을 추가하세요.

처리할 오류

서버가 오류를 이름 붙은 레인과 연관시킬 수 있으면 오류 이벤트는 stream_id를 포함해요. 다른 레인은 요청 범위 오류 후에도 계속할 수 있어요.

previous_response_not_found

{
  "type": "error",
  "status": 400,
  "stream_id": "main",
  "error": {
    "type": "invalid_request_error",
    "code": "previous_response_not_found",
    "message": "Previous response with id 'resp_abc' not found.",
    "param": "previous_response_id"
  }
}

invalid_stream_id

{
  "type": "error",
  "status": 400,
  "error": {
    "type": "invalid_request_error",
    "code": "invalid_stream_id",
    "message": "The 'stream_id' field must be a non-empty string with at most 256 characters and may only contain letters, numbers, underscores, hyphens, and periods.",
    "param": "stream_id"
  }
}

websocket_stream_limit_reached

{
  "type": "error",
  "status": 400,
  "stream_id": "agent_33",
  "error": {
    "type": "invalid_request_error",
    "code": "websocket_stream_limit_reached",
    "message": "This WebSocket connection has reached its maximum number of distinct stream IDs (32). Reuse an existing stream_id or open a new WebSocket connection.",
    "param": "stream_id"
  }
}

websocket_connection_limit_reached

{
  "type": "error",
  "error": {
    "type": "invalid_request_error",
    "code": "websocket_connection_limit_reached",
    "message": "Responses websocket connection limit reached (60 minutes). Create a new websocket connection to continue."
  },
  "status": 400
}

관련 가이드

더 알아보기 (Learn more)