Fix streaming subscription race
This commit is contained in:
12
CHANGELOG.md
12
CHANGELOG.md
@@ -1,5 +1,17 @@
|
|||||||
# Changelog
|
# 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
|
## 0.0.1.alpha5 - 2026-07-15
|
||||||
|
|
||||||
### Added
|
### Added
|
||||||
|
|||||||
@@ -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
|
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)
|
### Synchronous send (no streaming)
|
||||||
|
|
||||||
```ruby
|
```ruby
|
||||||
@@ -187,7 +192,9 @@ bundle install
|
|||||||
bundle exec rake test
|
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
|
## License
|
||||||
|
|
||||||
|
|||||||
@@ -131,21 +131,35 @@ module Opencode
|
|||||||
on_activity_tick: nil,
|
on_activity_tick: nil,
|
||||||
&block
|
&block
|
||||||
)
|
)
|
||||||
send_message_async(
|
|
||||||
session_id, text,
|
|
||||||
model: model, agent: agent, system: system, message_id: message_id
|
|
||||||
)
|
|
||||||
|
|
||||||
reply = Opencode::Reply.new
|
reply = Opencode::Reply.new
|
||||||
reply.add_observer(StreamBlockObserver.new(&block)) if block_given?
|
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,
|
session_id: session_id,
|
||||||
timeout: stream_timeout,
|
timeout: stream_timeout,
|
||||||
first_event_timeout: first_event_timeout,
|
first_event_timeout: first_event_timeout,
|
||||||
idle_stream_timeout: idle_stream_timeout,
|
idle_stream_timeout: idle_stream_timeout,
|
||||||
reply: reply,
|
reply: reply,
|
||||||
on_activity_tick: on_activity_tick
|
on_activity_tick: on_activity_tick,
|
||||||
|
on_subscribed: on_subscribed
|
||||||
) do |event|
|
) do |event|
|
||||||
reply.apply(event)
|
reply.apply(event)
|
||||||
end
|
end
|
||||||
@@ -285,6 +299,20 @@ module Opencode
|
|||||||
def stream_events(session_id:, timeout: 600, first_event_timeout: 120,
|
def stream_events(session_id:, timeout: 600, first_event_timeout: 120,
|
||||||
idle_stream_timeout: nil,
|
idle_stream_timeout: nil,
|
||||||
reply: nil, on_activity_tick: nil, &block)
|
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")
|
uri = build_uri("/event")
|
||||||
deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout
|
deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout
|
||||||
first_event_deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + first_event_timeout
|
first_event_deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + first_event_timeout
|
||||||
@@ -319,6 +347,8 @@ module Opencode
|
|||||||
http.open_timeout = 10
|
http.open_timeout = 10
|
||||||
http.read_timeout = 30
|
http.read_timeout = 30
|
||||||
|
|
||||||
|
subscription_callback_error = nil
|
||||||
|
subscription_ready = on_subscribed.nil?
|
||||||
begin
|
begin
|
||||||
buffer = String.new
|
buffer = String.new
|
||||||
|
|
||||||
@@ -352,6 +382,36 @@ module Opencode
|
|||||||
event = parse_sse_event(raw_event, session_id)
|
event = parse_sse_event(raw_event, session_id)
|
||||||
next unless event
|
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.")
|
unless event[:type]&.start_with?("server.")
|
||||||
received_session_event = true
|
received_session_event = true
|
||||||
last_meaningful_event_at = Process.clock_gettime(Process::CLOCK_MONOTONIC)
|
last_meaningful_event_at = Process.clock_gettime(Process::CLOCK_MONOTONIC)
|
||||||
@@ -368,6 +428,8 @@ module Opencode
|
|||||||
end
|
end
|
||||||
end
|
end
|
||||||
rescue *TRANSIENT_SSE_ERRORS
|
rescue *TRANSIENT_SSE_ERRORS
|
||||||
|
raise if subscription_callback_error
|
||||||
|
|
||||||
# Treat transport-level SSE disconnects like clean EOF: reconnect
|
# Treat transport-level SSE disconnects like clean EOF: reconnect
|
||||||
# until an idle session event, the overall timeout, or first-event
|
# until an idle session event, the overall timeout, or first-event
|
||||||
# timeout.
|
# timeout.
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
# frozen_string_literal: true
|
# frozen_string_literal: true
|
||||||
|
|
||||||
module Opencode
|
module Opencode
|
||||||
VERSION = "0.0.1.alpha5"
|
VERSION = "0.0.1.alpha6"
|
||||||
end
|
end
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ class SmokeTest < Minitest::Test
|
|||||||
BASE = "http://opencode.test"
|
BASE = "http://opencode.test"
|
||||||
PASSWORD = "test-secret"
|
PASSWORD = "test-secret"
|
||||||
SESSION_ID = "ses_smoke_1"
|
SESSION_ID = "ses_smoke_1"
|
||||||
|
CONNECTED_EVENT = { type: "server.connected", properties: {} }.freeze
|
||||||
|
|
||||||
def setup
|
def setup
|
||||||
@client = Opencode::Client.new(
|
@client = Opencode::Client.new(
|
||||||
@@ -130,6 +131,7 @@ class SmokeTest < Minitest::Test
|
|||||||
.to_return(status: 204, body: "")
|
.to_return(status: 204, body: "")
|
||||||
|
|
||||||
sse = [
|
sse = [
|
||||||
|
CONNECTED_EVENT,
|
||||||
{ type: "message.part.delta",
|
{ type: "message.part.delta",
|
||||||
properties: { sessionID: SESSION_ID, partID: "p1", field: "text", delta: "hello " } },
|
properties: { sessionID: SESSION_ID, partID: "p1", field: "text", delta: "hello " } },
|
||||||
{ type: "message.part.delta",
|
{ type: "message.part.delta",
|
||||||
@@ -157,11 +159,164 @@ class SmokeTest < Minitest::Test
|
|||||||
refute_empty parts_yielded
|
refute_empty parts_yielded
|
||||||
end
|
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
|
def test_stream_block_is_optional
|
||||||
stub_request(:post, "#{BASE}/session/#{SESSION_ID}/prompt_async")
|
stub_request(:post, "#{BASE}/session/#{SESSION_ID}/prompt_async")
|
||||||
.to_return(status: 204, body: "")
|
.to_return(status: 204, body: "")
|
||||||
|
|
||||||
sse = [
|
sse = [
|
||||||
|
CONNECTED_EVENT,
|
||||||
{ type: "message.part.delta",
|
{ type: "message.part.delta",
|
||||||
properties: { sessionID: SESSION_ID, partID: "p1", field: "text", delta: "ack" } },
|
properties: { sessionID: SESSION_ID, partID: "p1", field: "text", delta: "ack" } },
|
||||||
{ type: "session.idle", properties: { sessionID: SESSION_ID } }
|
{ type: "session.idle", properties: { sessionID: SESSION_ID } }
|
||||||
@@ -199,6 +354,7 @@ class SmokeTest < Minitest::Test
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
sse = [
|
sse = [
|
||||||
|
CONNECTED_EVENT,
|
||||||
{ type: "todo.updated", properties: { sessionID: SESSION_ID, todos: [] } },
|
{ 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: skill_part } },
|
||||||
{ type: "message.part.updated", properties: { sessionID: SESSION_ID, part: task_part } },
|
{ type: "message.part.updated", properties: { sessionID: SESSION_ID, part: task_part } },
|
||||||
|
|||||||
Reference in New Issue
Block a user