diff --git a/CHANGELOG.md b/CHANGELOG.md index 31e19f3..c789f62 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,17 @@ # Changelog +## 0.0.1.alpha6 - 2026-07-18 + +### Fixed + +- Make `Opencode::Client#stream` wait for OpenCode's initial + `server.connected` SSE readiness frame before submitting `prompt_async`, + closing the fast-response window where a turn could emit events before the + client was listening. +- Keep prompt submission at-most-once across automatic SSE reconnects. A + reconnect now reopens only the event stream; it never posts the user prompt + again, and prompt transport failures remain visible to the caller. + ## 0.0.1.alpha5 - 2026-07-15 ### Added diff --git a/README.md b/README.md index 15ff018..789fff9 100644 --- a/README.md +++ b/README.md @@ -86,6 +86,11 @@ reply.reasoning_text # => the model's hidden reasoning, if any reply.parts_json # => the full ordered parts array, ready for persistence ``` +`stream` waits for OpenCode's initial `server.connected` SSE readiness frame +before it submits the asynchronous prompt. If the event connection drops +afterward, the client reconnects only the subscription; it never reposts the +prompt. This prevents both missed fast responses and duplicate turns. + ### Synchronous send (no streaming) ```ruby @@ -187,7 +192,9 @@ bundle install bundle exec rake test ``` -The smoke suite covers Client end-to-end against WebMock-stubbed OpenCode endpoints. +The smoke suite covers Client end-to-end against WebMock-stubbed OpenCode +endpoints, including subscription-before-prompt ordering and +reconnect-without-repost. ## License diff --git a/lib/opencode/client.rb b/lib/opencode/client.rb index 75e8179..33780c6 100644 --- a/lib/opencode/client.rb +++ b/lib/opencode/client.rb @@ -131,21 +131,35 @@ module Opencode on_activity_tick: nil, &block ) - send_message_async( - session_id, text, - model: model, agent: agent, system: system, message_id: message_id - ) - reply = Opencode::Reply.new reply.add_observer(StreamBlockObserver.new(&block)) if block_given? - stream_events( + # Opening the event stream after prompt_async leaves a race where a fast + # turn can emit (and finish) before the client is subscribed. Wait for + # OpenCode's initial server.connected SSE frame, then submit the prompt + # exactly once. Reconnects invoke on_subscribed again, so mark the attempt + # before the POST: an ambiguous prompt response must never cause the same + # turn to be submitted twice. + prompt_attempted = false + on_subscribed = lambda do + next false if prompt_attempted + + prompt_attempted = true + send_message_async( + session_id, text, + model: model, agent: agent, system: system, message_id: message_id + ) + true + end + + consume_event_stream( session_id: session_id, timeout: stream_timeout, first_event_timeout: first_event_timeout, idle_stream_timeout: idle_stream_timeout, reply: reply, - on_activity_tick: on_activity_tick + on_activity_tick: on_activity_tick, + on_subscribed: on_subscribed ) do |event| reply.apply(event) end @@ -285,6 +299,20 @@ module Opencode def stream_events(session_id:, timeout: 600, first_event_timeout: 120, idle_stream_timeout: nil, reply: nil, on_activity_tick: nil, &block) + consume_event_stream( + session_id: session_id, + timeout: timeout, + first_event_timeout: first_event_timeout, + idle_stream_timeout: idle_stream_timeout, + reply: reply, + on_activity_tick: on_activity_tick, + &block + ) + end + + private def consume_event_stream(session_id:, timeout:, first_event_timeout:, + idle_stream_timeout:, reply:, on_activity_tick:, + on_subscribed: nil, &block) uri = build_uri("/event") deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout first_event_deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + first_event_timeout @@ -319,6 +347,8 @@ module Opencode http.open_timeout = 10 http.read_timeout = 30 + subscription_callback_error = nil + subscription_ready = on_subscribed.nil? begin buffer = String.new @@ -352,6 +382,36 @@ module Opencode event = parse_sse_event(raw_event, session_id) next unless event + unless subscription_ready + # Every supported OpenCode server starts /event with this + # frame. Receiving it proves the stream body is flowing; on + # current servers the bus listener is registered eagerly, + # and on older lazy-stream servers it is the strongest + # available readiness handshake before prompting. + next unless event[:type] == "server.connected" + + begin + turn_started = on_subscribed.call + rescue StandardError => error + # Prompt submission happens inside the open SSE response. + # Do not mistake its transport failure for an SSE disconnect + # and hide it behind a reconnect/first-event timeout. + subscription_callback_error = error + raise + end + if turn_started + # Before this fix stream_events began only after the prompt + # POST returned. Preserve those timeout semantics: the turn + # and first-session-event windows begin after prompt_async + # succeeds, not while establishing the readiness handshake. + started_at = Process.clock_gettime(Process::CLOCK_MONOTONIC) + deadline = started_at + timeout + first_event_deadline = started_at + first_event_timeout + last_meaningful_event_at = started_at + end + subscription_ready = true + end + unless event[:type]&.start_with?("server.") received_session_event = true last_meaningful_event_at = Process.clock_gettime(Process::CLOCK_MONOTONIC) @@ -368,6 +428,8 @@ module Opencode end end rescue *TRANSIENT_SSE_ERRORS + raise if subscription_callback_error + # Treat transport-level SSE disconnects like clean EOF: reconnect # until an idle session event, the overall timeout, or first-event # timeout. diff --git a/lib/opencode/version.rb b/lib/opencode/version.rb index 7ce8ac9..e309a3c 100644 --- a/lib/opencode/version.rb +++ b/lib/opencode/version.rb @@ -1,5 +1,5 @@ # frozen_string_literal: true module Opencode - VERSION = "0.0.1.alpha5" + VERSION = "0.0.1.alpha6" end diff --git a/test/opencode/smoke_test.rb b/test/opencode/smoke_test.rb index 14b8413..1890385 100644 --- a/test/opencode/smoke_test.rb +++ b/test/opencode/smoke_test.rb @@ -11,6 +11,7 @@ class SmokeTest < Minitest::Test BASE = "http://opencode.test" PASSWORD = "test-secret" SESSION_ID = "ses_smoke_1" + CONNECTED_EVENT = { type: "server.connected", properties: {} }.freeze def setup @client = Opencode::Client.new( @@ -130,6 +131,7 @@ class SmokeTest < Minitest::Test .to_return(status: 204, body: "") sse = [ + CONNECTED_EVENT, { type: "message.part.delta", properties: { sessionID: SESSION_ID, partID: "p1", field: "text", delta: "hello " } }, { type: "message.part.delta", @@ -157,11 +159,164 @@ class SmokeTest < Minitest::Test refute_empty parts_yielded end + def test_stream_waits_for_server_connected_before_posting_the_prompt + request_order = [] + connection_count = 0 + terminal_event = { + type: "session.status", + properties: { sessionID: SESSION_ID, status: { type: "idle" } } + } + + stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z}) + .to_return do + connection_count += 1 + request_order << :sse_accepted + events = if connection_count == 1 + [ { type: "server.heartbeat", properties: {} } ] + else + [ CONNECTED_EVENT, terminal_event ] + end + { + status: 200, + body: events.map { |event| "data: #{event.to_json}\n\n" }.join, + headers: { "Content-Type" => "text/event-stream" } + } + end + + stub_request(:post, "#{BASE}/session/#{SESSION_ID}/prompt_async") + .to_return do + request_order << :prompt + { status: 204, body: "" } + end + + stub_request(:get, "#{BASE}/session/#{SESSION_ID}/message") + .to_return(status: 200, body: [].to_json, + headers: { "Content-Type" => "application/json" }) + + @client.stream(SESSION_ID, "ping", stream_timeout: 1, first_event_timeout: 1) + + assert_equal [ :sse_accepted, :sse_accepted, :prompt ], request_order + end + + def test_stream_does_not_post_when_sse_subscription_is_rejected + stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z}) + .to_return(status: 503, body: "unavailable") + prompt = stub_request(:post, "#{BASE}/session/#{SESSION_ID}/prompt_async") + .to_return(status: 204, body: "") + + error = assert_raises(Opencode::ServerError) do + @client.stream(SESSION_ID, "ping", stream_timeout: 1, first_event_timeout: 1) + end + + assert_match "SSE connection failed: HTTP 503", error.message + assert_not_requested prompt + end + + def test_stream_surfaces_prompt_timeout_without_reconnecting + event_stream = stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z}) + .to_return(status: 200, body: "data: #{CONNECTED_EVENT.to_json}\n\n", + headers: { "Content-Type" => "text/event-stream" }) + prompt = stub_request(:post, "#{BASE}/session/#{SESSION_ID}/prompt_async") + .to_raise(Net::ReadTimeout.new("prompt timed out")) + + error = assert_raises(Opencode::TimeoutError) do + @client.stream(SESSION_ID, "ping", stream_timeout: 1, first_event_timeout: 1) + end + + assert_match "OpenCode timeout after 5s", error.message + assert_requested event_stream, times: 1 + assert_requested prompt, times: 1 + end + + def test_stream_reconnects_without_reposting_the_prompt + first_connection = [ + CONNECTED_EVENT, + { type: "server.heartbeat", properties: {} } + ] + second_connection = [ + CONNECTED_EVENT, + { + type: "message.part.delta", + properties: { sessionID: SESSION_ID, partID: "p1", field: "text", delta: "once" } + }, + { + type: "session.status", + properties: { sessionID: SESSION_ID, status: { type: "idle" } } + } + ] + + event_stream = stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z}) + .to_return( + { + status: 200, + body: first_connection.map { |event| "data: #{event.to_json}\n\n" }.join, + headers: { "Content-Type" => "text/event-stream" } + }, + { + status: 200, + body: second_connection.map { |event| "data: #{event.to_json}\n\n" }.join, + headers: { "Content-Type" => "text/event-stream" } + } + ) + prompt = stub_request(:post, "#{BASE}/session/#{SESSION_ID}/prompt_async") + .to_return(status: 204, body: "") + stub_request(:get, "#{BASE}/session/#{SESSION_ID}/message") + .to_return(status: 200, body: [].to_json, + headers: { "Content-Type" => "application/json" }) + + reply = @client.stream(SESSION_ID, "ping", stream_timeout: 1, first_event_timeout: 1) + + assert_equal "once", reply.full_text + assert_requested event_stream, times: 2 + assert_requested prompt, times: 1 + end + + def test_stream_events_preserves_question_and_permission_wait_state + events = [ + { + type: "question.asked", + properties: { id: "que_1", sessionID: SESSION_ID, questions: [] } + }, + { + type: "question.replied", + properties: { requestID: "que_1", sessionID: SESSION_ID, answers: [ [ "yes" ] ] } + }, + { + type: "permission.asked", + properties: { id: "per_1", sessionID: SESSION_ID, permission: "bash" } + }, + { + type: "permission.replied", + properties: { requestID: "per_1", sessionID: SESSION_ID, reply: "once" } + }, + { + type: "session.status", + properties: { sessionID: SESSION_ID, status: { type: "idle" } } + } + ] + stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z}) + .to_return( + status: 200, + body: events.map { |event| "data: #{event.to_json}\n\n" }.join, + headers: { "Content-Type" => "text/event-stream" } + ) + + reply = Opencode::Reply.new + wait_states = [] + @client.stream_events(session_id: SESSION_ID, reply: reply) do |event| + reply.apply(event) + wait_states << reply.prompt_blocked? + end + + assert_equal [ true, false, true, false, false ], wait_states + end + def test_stream_block_is_optional stub_request(:post, "#{BASE}/session/#{SESSION_ID}/prompt_async") .to_return(status: 204, body: "") sse = [ + CONNECTED_EVENT, { type: "message.part.delta", properties: { sessionID: SESSION_ID, partID: "p1", field: "text", delta: "ack" } }, { type: "session.idle", properties: { sessionID: SESSION_ID } } @@ -199,6 +354,7 @@ class SmokeTest < Minitest::Test } } sse = [ + CONNECTED_EVENT, { type: "todo.updated", properties: { sessionID: SESSION_ID, todos: [] } }, { type: "message.part.updated", properties: { sessionID: SESSION_ID, part: skill_part } }, { type: "message.part.updated", properties: { sessionID: SESSION_ID, part: task_part } },