Compare commits
9 Commits
v0.0.1.alp
...
v0.0.1.alp
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
78b6f9c9e9 | ||
| edad9c7018 | |||
|
|
a116b2708c | ||
| b16292723e | |||
| 8adc95985a | |||
| 5113a953db | |||
| de14d57634 | |||
| 2e866a618b | |||
| 8dbd3de81c |
23
.github/workflows/release.yml
vendored
Normal file
23
.github/workflows/release.yml
vendored
Normal file
@@ -0,0 +1,23 @@
|
|||||||
|
name: Push gem
|
||||||
|
|
||||||
|
on:
|
||||||
|
push:
|
||||||
|
tags:
|
||||||
|
- "v*"
|
||||||
|
|
||||||
|
jobs:
|
||||||
|
push:
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
permissions:
|
||||||
|
contents: write
|
||||||
|
id-token: write
|
||||||
|
environment: release
|
||||||
|
steps:
|
||||||
|
- uses: actions/checkout@v5
|
||||||
|
with:
|
||||||
|
persist-credentials: false
|
||||||
|
- uses: ruby/setup-ruby@v1
|
||||||
|
with:
|
||||||
|
ruby-version: ruby
|
||||||
|
bundler-cache: true
|
||||||
|
- uses: rubygems/release-gem@v1
|
||||||
4
.github/workflows/test.yml
vendored
4
.github/workflows/test.yml
vendored
@@ -14,7 +14,7 @@ jobs:
|
|||||||
matrix:
|
matrix:
|
||||||
ruby: ["3.2", "3.3", "3.4"]
|
ruby: ["3.2", "3.3", "3.4"]
|
||||||
steps:
|
steps:
|
||||||
- uses: actions/checkout@v4
|
- uses: actions/checkout@v5
|
||||||
|
|
||||||
- name: Set up Ruby ${{ matrix.ruby }}
|
- name: Set up Ruby ${{ matrix.ruby }}
|
||||||
uses: ruby/setup-ruby@v1
|
uses: ruby/setup-ruby@v1
|
||||||
@@ -30,5 +30,5 @@ jobs:
|
|||||||
|
|
||||||
- name: Verify gem loads after install
|
- name: Verify gem loads after install
|
||||||
run: |
|
run: |
|
||||||
gem install --local opencode-ruby-*.gem
|
gem install opencode-ruby-*.gem --no-document
|
||||||
ruby -ropencode-ruby -e 'puts Opencode::VERSION'
|
ruby -ropencode-ruby -e 'puts Opencode::VERSION'
|
||||||
|
|||||||
53
CHANGELOG.md
53
CHANGELOG.md
@@ -1,5 +1,58 @@
|
|||||||
# Changelog
|
# Changelog
|
||||||
|
|
||||||
|
## 0.0.1.alpha7 - 2026-07-18
|
||||||
|
|
||||||
|
### Fixed
|
||||||
|
|
||||||
|
- Add an at-most-once `on_subscribed` hook to the lower-level
|
||||||
|
`Opencode::Client#stream_events` API. Higher-level orchestrators can now
|
||||||
|
wait for `server.connected` before submitting `prompt_async` without
|
||||||
|
abandoning their own Reply observers, persistence, or recovery pipeline.
|
||||||
|
- Propagate subscription-hook failures directly and never invoke the hook
|
||||||
|
again on SSE reconnect. Ambiguous prompt transport failures therefore
|
||||||
|
cannot silently become duplicate turns.
|
||||||
|
|
||||||
|
## 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
|
||||||
|
|
||||||
|
- Extend `Opencode::Client#create_session` with OpenCode's native parent,
|
||||||
|
agent, model, metadata, and workspace fields while preserving the existing
|
||||||
|
title and permission call shape. Session model strings are encoded with the
|
||||||
|
session endpoint's `{ providerID, id }` shape rather than the message
|
||||||
|
endpoint's `{ providerID, modelID }` shape.
|
||||||
|
|
||||||
|
## 0.0.1.alpha4 - 2026-07-12
|
||||||
|
|
||||||
|
### Fixed
|
||||||
|
|
||||||
|
- End SSE streams on current OpenCode `session.status` idle events while
|
||||||
|
retaining compatibility with legacy `session.idle` events.
|
||||||
|
- Reconcile every assistant message in the current user turn after multi-step
|
||||||
|
tool loops, preserving stream-only parts without duplicating final text.
|
||||||
|
- Parse terminal tool parts in standalone Ruby clients without relying on the
|
||||||
|
Rails-loaded `Object#in?` extension.
|
||||||
|
|
||||||
|
## 0.0.1.alpha3 - 2026-07-10
|
||||||
|
|
||||||
|
### Added
|
||||||
|
|
||||||
|
- `Opencode::Client#update_session` for applying permission rules through
|
||||||
|
OpenCode's session PATCH endpoint. OpenCode appends these rules, so hosts
|
||||||
|
should fingerprint their ordered policy and call this only when it changes.
|
||||||
|
|
||||||
## 0.0.1.alpha2 — 2026-05-20
|
## 0.0.1.alpha2 — 2026-05-20
|
||||||
|
|
||||||
### Added
|
### Added
|
||||||
|
|||||||
67
README.md
67
README.md
@@ -48,6 +48,27 @@ Multi-tenant apps construct multiple clients with different `base_url`s — each
|
|||||||
|
|
||||||
## Core API
|
## Core API
|
||||||
|
|
||||||
|
### Configured and parent-linked sessions
|
||||||
|
|
||||||
|
OpenCode can create a session under an existing parent and select its agent,
|
||||||
|
model, metadata, workspace, and permission policy in the same request:
|
||||||
|
|
||||||
|
```ruby
|
||||||
|
child = client.create_session(
|
||||||
|
title: "Destination curator",
|
||||||
|
parent_id: parent_session_id,
|
||||||
|
agent: "destination-list-curator",
|
||||||
|
model: "openai/gpt-5.5",
|
||||||
|
metadata: { run_id: "9" },
|
||||||
|
workspace_id: workspace_id,
|
||||||
|
permissions: permission_rules
|
||||||
|
)
|
||||||
|
```
|
||||||
|
|
||||||
|
Model strings use OpenCode's `provider/model` form; a preformatted model hash
|
||||||
|
with `providerID` and `id` keys is also accepted. These configured-session
|
||||||
|
fields require OpenCode 1.16.1 or newer.
|
||||||
|
|
||||||
### Streaming (the headline)
|
### Streaming (the headline)
|
||||||
|
|
||||||
```ruby
|
```ruby
|
||||||
@@ -65,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
|
||||||
@@ -72,16 +98,44 @@ result = client.send_message(session_id, "Quick yes/no: is Ruby fun?")
|
|||||||
# result is the OpenCode response hash; see API docs for fields.
|
# result is the OpenCode response hash; see API docs for fields.
|
||||||
```
|
```
|
||||||
|
|
||||||
|
### Updating session permissions
|
||||||
|
|
||||||
|
```ruby
|
||||||
|
client.update_session(session_id, permissions: permission_rules)
|
||||||
|
```
|
||||||
|
|
||||||
|
OpenCode appends PATCHed permission rules and evaluates the last matching
|
||||||
|
rule. Hosts should send a complete ordered policy and fingerprint it so the
|
||||||
|
same policy is not appended on every turn. This endpoint requires OpenCode
|
||||||
|
1.16.1 or newer; the rest of the client remains compatible with 1.15.
|
||||||
|
|
||||||
### Lower-level event firehose
|
### Lower-level event firehose
|
||||||
|
|
||||||
If you need raw SSE events (every server tick, todo update, prompt asked/replied), use `stream_events` directly:
|
If you need raw SSE events (every server tick, todo update, prompt asked/replied), use `stream_events` directly:
|
||||||
|
|
||||||
```ruby
|
```ruby
|
||||||
client.stream_events(session_id: session_id) do |event|
|
client.stream_events(session_id: session_id) do |event|
|
||||||
puts event[:type] # "message.part.delta", "todo.updated", "session.idle", ...
|
puts event[:type] # "message.part.delta", "todo.updated", "session.status", ...
|
||||||
end
|
end
|
||||||
```
|
```
|
||||||
|
|
||||||
|
Orchestrators that submit their own async prompt must do so through
|
||||||
|
`on_subscribed`; the callback runs only after the first `server.connected`
|
||||||
|
frame and at most once across automatic reconnects:
|
||||||
|
|
||||||
|
```ruby
|
||||||
|
client.stream_events(
|
||||||
|
session_id: session_id,
|
||||||
|
on_subscribed: -> { client.send_message_async(session_id, prompt) }
|
||||||
|
) do |event|
|
||||||
|
reply.apply(event)
|
||||||
|
end
|
||||||
|
```
|
||||||
|
|
||||||
|
If the callback raises, `stream_events` propagates that error and does not
|
||||||
|
retry it. This is intentional: a timed-out prompt response is ambiguous, so
|
||||||
|
reposting could duplicate the model turn and its cost.
|
||||||
|
|
||||||
### Interactive prompts
|
### Interactive prompts
|
||||||
|
|
||||||
When the agent uses the `question` or `permission` tools, opencode emits `question.asked` / `permission.asked` events. Answer them via:
|
When the agent uses the `question` or `permission` tools, opencode emits `question.asked` / `permission.asked` events. Answer them via:
|
||||||
@@ -101,7 +155,7 @@ begin
|
|||||||
rescue Opencode::ConnectionError # server unreachable
|
rescue Opencode::ConnectionError # server unreachable
|
||||||
rescue Opencode::TimeoutError # client-side timeout
|
rescue Opencode::TimeoutError # client-side timeout
|
||||||
rescue Opencode::SessionNotFoundError # 404 on a session
|
rescue Opencode::SessionNotFoundError # 404 on a session
|
||||||
rescue Opencode::StaleSessionError # session.idle never arrived
|
rescue Opencode::StaleSessionError # no session event arrived after the prompt
|
||||||
rescue Opencode::IdleStreamError # mid-turn SSE wedge
|
rescue Opencode::IdleStreamError # mid-turn SSE wedge
|
||||||
rescue Opencode::ServerError # 5xx
|
rescue Opencode::ServerError # 5xx
|
||||||
rescue Opencode::BadRequestError # 4xx other than 404
|
rescue Opencode::BadRequestError # 4xx other than 404
|
||||||
@@ -155,7 +209,14 @@ bundle install
|
|||||||
bundle exec rake test
|
bundle exec rake test
|
||||||
```
|
```
|
||||||
|
|
||||||
12-test smoke 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.
|
||||||
|
|
||||||
|
Releases use RubyGems trusted publishing. After the repository's
|
||||||
|
`release.yml` workflow is registered as a trusted publisher with the `release`
|
||||||
|
environment, pushing a `v*` tag builds, attests, and publishes the gem without
|
||||||
|
a long-lived RubyGems API key.
|
||||||
|
|
||||||
## License
|
## License
|
||||||
|
|
||||||
|
|||||||
@@ -25,8 +25,24 @@ module Opencode
|
|||||||
@workspace = workspace
|
@workspace = workspace
|
||||||
end
|
end
|
||||||
|
|
||||||
def create_session(title: nil, permissions: nil)
|
def create_session(
|
||||||
body = { title: title, permission: permissions }.compact
|
title: nil,
|
||||||
|
permissions: nil,
|
||||||
|
parent_id: nil,
|
||||||
|
agent: nil,
|
||||||
|
model: nil,
|
||||||
|
metadata: nil,
|
||||||
|
workspace_id: nil
|
||||||
|
)
|
||||||
|
body = {
|
||||||
|
title: title,
|
||||||
|
permission: permissions,
|
||||||
|
parentID: parent_id,
|
||||||
|
agent: agent,
|
||||||
|
model: format_session_model(model),
|
||||||
|
metadata: metadata,
|
||||||
|
workspaceID: workspace_id
|
||||||
|
}.compact
|
||||||
post("/session", body)
|
post("/session", body)
|
||||||
end
|
end
|
||||||
|
|
||||||
@@ -115,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
|
||||||
@@ -144,6 +174,10 @@ module Opencode
|
|||||||
execute(request)
|
execute(request)
|
||||||
end
|
end
|
||||||
|
|
||||||
|
def update_session(session_id, permissions:)
|
||||||
|
patch("/session/#{session_id}", { permission: permissions })
|
||||||
|
end
|
||||||
|
|
||||||
def children(session_id)
|
def children(session_id)
|
||||||
uri = build_uri("/session/#{session_id}/children")
|
uri = build_uri("/session/#{session_id}/children")
|
||||||
request = Net::HTTP::Get.new(uri)
|
request = Net::HTTP::Get.new(uri)
|
||||||
@@ -233,8 +267,10 @@ module Opencode
|
|||||||
].freeze
|
].freeze
|
||||||
|
|
||||||
# Opens SSE connection to GET /event, yields parsed events filtered by session_id.
|
# Opens SSE connection to GET /event, yields parsed events filtered by session_id.
|
||||||
# Blocks until session goes idle or timeout, reconnecting across dropped
|
# Blocks until the session reports idle or timeout, reconnecting across
|
||||||
# event-stream connections.
|
# dropped event-stream connections. Current OpenCode emits
|
||||||
|
# `session.status` with `status.type == "idle"`; older versions emitted the
|
||||||
|
# standalone `session.idle` event, so both remain terminal.
|
||||||
#
|
#
|
||||||
# first_event_timeout: seconds to wait for a session-specific event before
|
# first_event_timeout: seconds to wait for a session-specific event before
|
||||||
# declaring the session stale. Server heartbeats don't count — they're global
|
# declaring the session stale. Server heartbeats don't count — they're global
|
||||||
@@ -249,6 +285,13 @@ module Opencode
|
|||||||
# without nuking real reasoning. Callers that know their agent is
|
# without nuking real reasoning. Callers that know their agent is
|
||||||
# short-prompt + fast can pass a lower value.
|
# short-prompt + fast can pass a lower value.
|
||||||
#
|
#
|
||||||
|
# on_subscribed: optional callable invoked at most once, after the first
|
||||||
|
# `server.connected` frame proves the SSE response body is flowing. This
|
||||||
|
# is the safe place for higher-level orchestrators to submit prompt_async;
|
||||||
|
# reconnects wait for their own connected frame but never invoke it again.
|
||||||
|
# A raised callback error is propagated directly and is never treated as
|
||||||
|
# a reconnectable SSE transport failure.
|
||||||
|
#
|
||||||
# idle_stream_timeout: seconds to wait BETWEEN meaningful events once
|
# idle_stream_timeout: seconds to wait BETWEEN meaningful events once
|
||||||
# the session has started producing them. Default nil = no check
|
# the session has started producing them. Default nil = no check
|
||||||
# (preserves the overall `timeout` ceiling behavior). Opt-in heartbeat
|
# (preserves the overall `timeout` ceiling behavior). Opt-in heartbeat
|
||||||
@@ -262,12 +305,28 @@ module Opencode
|
|||||||
# translate into a user-visible error / retry affordance.
|
# translate into a user-visible error / retry affordance.
|
||||||
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, on_subscribed: 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,
|
||||||
|
on_subscribed: on_subscribed,
|
||||||
|
&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
|
||||||
received_session_event = false
|
received_session_event = false
|
||||||
last_meaningful_event_at = Process.clock_gettime(Process::CLOCK_MONOTONIC)
|
last_meaningful_event_at = Process.clock_gettime(Process::CLOCK_MONOTONIC)
|
||||||
|
subscription_callback_attempted = on_subscribed.nil?
|
||||||
|
|
||||||
loop do
|
loop do
|
||||||
now = Process.clock_gettime(Process::CLOCK_MONOTONIC)
|
now = Process.clock_gettime(Process::CLOCK_MONOTONIC)
|
||||||
@@ -297,6 +356,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
|
||||||
|
|
||||||
@@ -330,6 +391,44 @@ 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"
|
||||||
|
|
||||||
|
turn_started = false
|
||||||
|
unless subscription_callback_attempted
|
||||||
|
# Mark the attempt before invoking the callback. A timeout
|
||||||
|
# can mean the prompt reached OpenCode even though its HTTP
|
||||||
|
# response did not reach us; retrying would duplicate the
|
||||||
|
# turn. The caller receives the original error instead.
|
||||||
|
subscription_callback_attempted = true
|
||||||
|
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
|
||||||
|
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)
|
||||||
@@ -341,13 +440,16 @@ module Opencode
|
|||||||
# the reaper doesn't kill it mid-wait.
|
# the reaper doesn't kill it mid-wait.
|
||||||
on_activity_tick&.call(event)
|
on_activity_tick&.call(event)
|
||||||
block.call(event)
|
block.call(event)
|
||||||
return if event[:type] == "session.idle"
|
return if terminal_session_event?(event)
|
||||||
end
|
end
|
||||||
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 session.idle, the overall timeout, or first-event timeout.
|
# until an idle session event, the overall timeout, or first-event
|
||||||
|
# timeout.
|
||||||
ensure
|
ensure
|
||||||
begin
|
begin
|
||||||
http&.finish if http&.started?
|
http&.finish if http&.started?
|
||||||
@@ -381,13 +483,16 @@ module Opencode
|
|||||||
# the caller's reply is still a usable Result either way.
|
# the caller's reply is still a usable Result either way.
|
||||||
def merge_final_exchange(session_id, reply)
|
def merge_final_exchange(session_id, reply)
|
||||||
exchange = get_messages(session_id)
|
exchange = get_messages(session_id)
|
||||||
last_assistant = Array(exchange).reverse_each.find do |message|
|
polled = current_turn_parts(exchange)
|
||||||
message.dig(:info, :role) == "assistant"
|
return if polled.empty?
|
||||||
end
|
|
||||||
return unless last_assistant
|
|
||||||
|
|
||||||
polled = Opencode::ResponseParser.extract_interleaved_parts(last_assistant)
|
merged = merge_stream_only_parts(reply.result.parts_json, polled)
|
||||||
reply.sync_recovered_parts(polled) if polled.any?
|
reply.sync_recovered_parts(merged)
|
||||||
|
# sync_recovered_parts intentionally never deletes live parts because it
|
||||||
|
# is also used during mid-stream recovery. This is the terminal poll, so
|
||||||
|
# the wire snapshot is authoritative: remove any replayed trailing wire
|
||||||
|
# part after observers have seen recovered additions/updates.
|
||||||
|
reply.replace_parts(merged) unless reply.result.parts_json == merged
|
||||||
rescue Opencode::Error
|
rescue Opencode::Error
|
||||||
# Stream's result is still complete; the merge was a polish, not a
|
# Stream's result is still complete; the merge was a polish, not a
|
||||||
# requirement.
|
# requirement.
|
||||||
@@ -404,6 +509,45 @@ module Opencode
|
|||||||
deadline
|
deadline
|
||||||
end
|
end
|
||||||
|
|
||||||
|
def terminal_session_event?(event)
|
||||||
|
return true if event[:type] == "session.idle"
|
||||||
|
return false unless event[:type] == "session.status"
|
||||||
|
|
||||||
|
status = event.dig(:properties, :status)
|
||||||
|
status = status[:type] || status["type"] if status.is_a?(Hash)
|
||||||
|
status == "idle"
|
||||||
|
end
|
||||||
|
|
||||||
|
# OpenCode persists one assistant message per model step. A tool loop can
|
||||||
|
# therefore produce several assistant messages for one user turn (for
|
||||||
|
# example skill -> task -> final text). Reconcile the complete current turn
|
||||||
|
# instead of aligning the live parts array with only the last assistant
|
||||||
|
# message, which corrupts tool parts and duplicates final text.
|
||||||
|
def current_turn_parts(exchange)
|
||||||
|
messages = Array(exchange)
|
||||||
|
last_user_index = messages.rindex { |message| message.dig(:info, :role) == "user" }
|
||||||
|
current_turn = last_user_index ? messages.drop(last_user_index + 1) : messages
|
||||||
|
|
||||||
|
current_turn
|
||||||
|
.select { |message| message.dig(:info, :role) == "assistant" }
|
||||||
|
.flat_map { |message| Opencode::ResponseParser.extract_interleaved_parts(message) }
|
||||||
|
end
|
||||||
|
|
||||||
|
def merge_stream_only_parts(stream_parts, wire_parts)
|
||||||
|
remaining_wire = Array(wire_parts).dup
|
||||||
|
merged = []
|
||||||
|
|
||||||
|
Array(stream_parts).each do |part|
|
||||||
|
if Opencode::PartSource.stream_only?(part)
|
||||||
|
merged << part
|
||||||
|
elsif remaining_wire.any?
|
||||||
|
merged << remaining_wire.shift
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
||||||
|
merged.concat(remaining_wire)
|
||||||
|
end
|
||||||
|
|
||||||
def prompt_payload(text, parts:, model:, agent:, system:, message_id:, no_reply:, tools:, format:, variant:)
|
def prompt_payload(text, parts:, model:, agent:, system:, message_id:, no_reply:, tools:, format:, variant:)
|
||||||
message_parts = parts || [ { type: "text", text: text } ]
|
message_parts = parts || [ { type: "text", text: text } ]
|
||||||
{
|
{
|
||||||
@@ -427,6 +571,14 @@ module Opencode
|
|||||||
{ providerID: provider, modelID: model_id }
|
{ providerID: provider, modelID: model_id }
|
||||||
end
|
end
|
||||||
|
|
||||||
|
def format_session_model(model)
|
||||||
|
return nil unless model
|
||||||
|
return model if model.is_a?(Hash)
|
||||||
|
|
||||||
|
provider, model_id = model.split("/", 2)
|
||||||
|
{ providerID: provider, id: model_id }
|
||||||
|
end
|
||||||
|
|
||||||
def post(path, body)
|
def post(path, body)
|
||||||
uri = build_uri(path)
|
uri = build_uri(path)
|
||||||
request = Net::HTTP::Post.new(uri)
|
request = Net::HTTP::Post.new(uri)
|
||||||
@@ -434,6 +586,13 @@ module Opencode
|
|||||||
execute(request)
|
execute(request)
|
||||||
end
|
end
|
||||||
|
|
||||||
|
def patch(path, body)
|
||||||
|
uri = build_uri(path)
|
||||||
|
request = Net::HTTP::Patch.new(uri)
|
||||||
|
request.body = body.to_json
|
||||||
|
execute(request)
|
||||||
|
end
|
||||||
|
|
||||||
def build_uri(path, scoped: true)
|
def build_uri(path, scoped: true)
|
||||||
uri = @uri.dup
|
uri = @uri.dup
|
||||||
uri.path = path
|
uri.path = path
|
||||||
|
|||||||
@@ -27,7 +27,7 @@ module Opencode
|
|||||||
def self.extract_tool_summary(response_body)
|
def self.extract_tool_summary(response_body)
|
||||||
parts = response_body[:parts] || []
|
parts = response_body[:parts] || []
|
||||||
parts
|
parts
|
||||||
.select { |p| p[:type] == "tool" && p.dig(:state, :status).in?(TERMINAL_STATUSES) }
|
.select { |p| p[:type] == "tool" && TERMINAL_STATUSES.include?(p.dig(:state, :status)) }
|
||||||
.map { |p| build_tool_summary(p) }
|
.map { |p| build_tool_summary(p) }
|
||||||
end
|
end
|
||||||
|
|
||||||
@@ -42,7 +42,7 @@ module Opencode
|
|||||||
{ "type" => "reasoning", "content" => part[:text] }
|
{ "type" => "reasoning", "content" => part[:text] }
|
||||||
when "tool"
|
when "tool"
|
||||||
status = part.dig(:state, :status)
|
status = part.dig(:state, :status)
|
||||||
next unless status.in?(TERMINAL_STATUSES)
|
next unless TERMINAL_STATUSES.include?(status)
|
||||||
|
|
||||||
build_tool_summary(part)
|
build_tool_summary(part)
|
||||||
else
|
else
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
# frozen_string_literal: true
|
# frozen_string_literal: true
|
||||||
|
|
||||||
module Opencode
|
module Opencode
|
||||||
VERSION = "0.0.1.alpha2"
|
VERSION = "0.0.1.alpha7"
|
||||||
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(
|
||||||
@@ -46,6 +47,7 @@ class SmokeTest < Minitest::Test
|
|||||||
|
|
||||||
def test_create_session_returns_session_id
|
def test_create_session_returns_session_id
|
||||||
stub_request(:post, "#{BASE}/session")
|
stub_request(:post, "#{BASE}/session")
|
||||||
|
.with(body: { title: "smoke", permission: [] }.to_json)
|
||||||
.to_return(status: 200, body: { id: SESSION_ID, title: "smoke" }.to_json,
|
.to_return(status: 200, body: { id: SESSION_ID, title: "smoke" }.to_json,
|
||||||
headers: { "Content-Type" => "application/json" })
|
headers: { "Content-Type" => "application/json" })
|
||||||
|
|
||||||
@@ -53,6 +55,69 @@ class SmokeTest < Minitest::Test
|
|||||||
assert_equal SESSION_ID, response[:id]
|
assert_equal SESSION_ID, response[:id]
|
||||||
end
|
end
|
||||||
|
|
||||||
|
def test_create_session_sends_native_child_and_configuration_fields
|
||||||
|
permissions = [ { permission: "skill", pattern: "*", action: "deny" } ]
|
||||||
|
expected_body = {
|
||||||
|
title: "curator",
|
||||||
|
permission: permissions,
|
||||||
|
parentID: "ses_parent",
|
||||||
|
agent: "destination-list-curator",
|
||||||
|
model: { providerID: "openrouter", id: "anthropic/claude-sonnet-4" },
|
||||||
|
metadata: { run: "9" },
|
||||||
|
workspaceID: "wrk_1"
|
||||||
|
}
|
||||||
|
|
||||||
|
stub_request(:post, "#{BASE}/session")
|
||||||
|
.with(body: expected_body.to_json)
|
||||||
|
.to_return(status: 200, body: { id: SESSION_ID }.to_json,
|
||||||
|
headers: { "Content-Type" => "application/json" })
|
||||||
|
|
||||||
|
response = @client.create_session(
|
||||||
|
title: "curator",
|
||||||
|
permissions: permissions,
|
||||||
|
parent_id: "ses_parent",
|
||||||
|
agent: "destination-list-curator",
|
||||||
|
model: "openrouter/anthropic/claude-sonnet-4",
|
||||||
|
metadata: { run: "9" },
|
||||||
|
workspace_id: "wrk_1"
|
||||||
|
)
|
||||||
|
|
||||||
|
assert_equal SESSION_ID, response[:id]
|
||||||
|
assert_requested :post, "#{BASE}/session", body: expected_body.to_json, times: 1
|
||||||
|
end
|
||||||
|
|
||||||
|
def test_create_session_preserves_a_preformatted_model
|
||||||
|
model = { providerID: "openai", id: "gpt-5.5", variant: "high" }
|
||||||
|
|
||||||
|
stub_request(:post, "#{BASE}/session")
|
||||||
|
.with(body: { model: model }.to_json)
|
||||||
|
.to_return(status: 200, body: { id: SESSION_ID }.to_json,
|
||||||
|
headers: { "Content-Type" => "application/json" })
|
||||||
|
|
||||||
|
response = @client.create_session(model: model)
|
||||||
|
|
||||||
|
assert_equal SESSION_ID, response[:id]
|
||||||
|
assert_requested :post, "#{BASE}/session", body: { model: model }.to_json, times: 1
|
||||||
|
end
|
||||||
|
|
||||||
|
def test_update_session_patches_permissions_and_returns_the_updated_session
|
||||||
|
permissions = [
|
||||||
|
{ permission: "skill", pattern: "*", action: "deny" },
|
||||||
|
{ permission: "skill", pattern: "core-details", action: "allow" }
|
||||||
|
]
|
||||||
|
|
||||||
|
stub_request(:patch, "#{BASE}/session/#{SESSION_ID}")
|
||||||
|
.with(body: { permission: permissions }.to_json)
|
||||||
|
.to_return(status: 200, body: { id: SESSION_ID, permission: permissions }.to_json,
|
||||||
|
headers: { "Content-Type" => "application/json" })
|
||||||
|
|
||||||
|
response = @client.update_session(SESSION_ID, permissions: permissions)
|
||||||
|
|
||||||
|
assert_equal SESSION_ID, response[:id]
|
||||||
|
assert_equal permissions, response[:permission]
|
||||||
|
assert_requested :patch, "#{BASE}/session/#{SESSION_ID}", times: 1
|
||||||
|
end
|
||||||
|
|
||||||
def test_send_message_async_returns_empty_body
|
def test_send_message_async_returns_empty_body
|
||||||
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: "")
|
||||||
@@ -66,11 +131,12 @@ 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",
|
||||||
properties: { sessionID: SESSION_ID, partID: "p1", field: "text", delta: "world" } },
|
properties: { sessionID: SESSION_ID, partID: "p1", field: "text", delta: "world" } },
|
||||||
{ type: "session.idle", properties: { sessionID: SESSION_ID } }
|
{ type: "session.status", properties: { sessionID: SESSION_ID, status: { type: "idle" } } }
|
||||||
].map { |e| "data: #{e.to_json}\n\n" }.join
|
].map { |e| "data: #{e.to_json}\n\n" }.join
|
||||||
|
|
||||||
stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z})
|
stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z})
|
||||||
@@ -93,11 +159,228 @@ 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_invokes_on_subscribed_once_after_connected_across_reconnects
|
||||||
|
order = []
|
||||||
|
connections = [
|
||||||
|
[ CONNECTED_EVENT, { type: "server.heartbeat", properties: {} } ],
|
||||||
|
[
|
||||||
|
CONNECTED_EVENT,
|
||||||
|
{
|
||||||
|
type: "session.status",
|
||||||
|
properties: { sessionID: SESSION_ID, status: { type: "idle" } }
|
||||||
|
}
|
||||||
|
]
|
||||||
|
]
|
||||||
|
|
||||||
|
event_stream = stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z})
|
||||||
|
.to_return do
|
||||||
|
events = connections.shift or raise "unexpected third SSE connection"
|
||||||
|
order << :sse_accepted
|
||||||
|
{
|
||||||
|
status: 200,
|
||||||
|
body: events.map { |event| "data: #{event.to_json}\n\n" }.join,
|
||||||
|
headers: { "Content-Type" => "text/event-stream" }
|
||||||
|
}
|
||||||
|
end
|
||||||
|
|
||||||
|
subscribed_calls = 0
|
||||||
|
@client.stream_events(
|
||||||
|
session_id: SESSION_ID,
|
||||||
|
timeout: 1,
|
||||||
|
first_event_timeout: 1,
|
||||||
|
on_subscribed: -> {
|
||||||
|
subscribed_calls += 1
|
||||||
|
order << :prompt
|
||||||
|
true
|
||||||
|
}
|
||||||
|
) { |_event| }
|
||||||
|
|
||||||
|
assert_equal 1, subscribed_calls
|
||||||
|
assert_equal [ :sse_accepted, :prompt, :sse_accepted ], order
|
||||||
|
assert_requested event_stream, times: 2
|
||||||
|
end
|
||||||
|
|
||||||
|
def test_stream_events_surfaces_on_subscribed_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" })
|
||||||
|
|
||||||
|
calls = 0
|
||||||
|
error = assert_raises(Net::ReadTimeout) do
|
||||||
|
@client.stream_events(
|
||||||
|
session_id: SESSION_ID,
|
||||||
|
timeout: 1,
|
||||||
|
first_event_timeout: 1,
|
||||||
|
on_subscribed: -> {
|
||||||
|
calls += 1
|
||||||
|
raise Net::ReadTimeout, "ambiguous prompt response"
|
||||||
|
}
|
||||||
|
) { |_event| }
|
||||||
|
end
|
||||||
|
|
||||||
|
assert_equal "Net::ReadTimeout with \"ambiguous prompt response\"", error.message
|
||||||
|
assert_equal 1, calls
|
||||||
|
assert_requested event_stream, 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 } }
|
||||||
@@ -115,6 +398,60 @@ class SmokeTest < Minitest::Test
|
|||||||
assert_equal "ack", reply.full_text
|
assert_equal "ack", reply.full_text
|
||||||
end
|
end
|
||||||
|
|
||||||
|
def test_stream_merges_a_multi_assistant_tool_loop_without_duplicate_text
|
||||||
|
stub_request(:post, "#{BASE}/session/#{SESSION_ID}/prompt_async")
|
||||||
|
.to_return(status: 204, body: "")
|
||||||
|
|
||||||
|
skill_part = {
|
||||||
|
id: "p_skill", sessionID: SESSION_ID, messageID: "m_skill",
|
||||||
|
type: "tool", tool: "skill", callID: "call_skill",
|
||||||
|
state: { status: "completed", input: { name: "travelwolf-itinerary" }, output: "loaded" }
|
||||||
|
}
|
||||||
|
task_part = {
|
||||||
|
id: "p_task", sessionID: SESSION_ID, messageID: "m_task",
|
||||||
|
type: "tool", tool: "task", callID: "call_task",
|
||||||
|
state: {
|
||||||
|
status: "completed",
|
||||||
|
input: { subagent_type: "itinerary-planner" },
|
||||||
|
output: "{\"days\":[]}",
|
||||||
|
metadata: { sessionId: "ses_child" }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
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 } },
|
||||||
|
{ type: "message.part.delta",
|
||||||
|
properties: { sessionID: SESSION_ID, partID: "p_text", field: "text", delta: "SUBAGENT_OK" } },
|
||||||
|
{ type: "message.part.delta",
|
||||||
|
properties: { sessionID: SESSION_ID, partID: "p_text_replay", field: "text", delta: "SUBAGENT_OK" } },
|
||||||
|
{ type: "session.status", properties: { sessionID: SESSION_ID, status: { type: "idle" } } }
|
||||||
|
].map { |event| "data: #{event.to_json}\n\n" }.join
|
||||||
|
|
||||||
|
stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z})
|
||||||
|
.to_return(status: 200, body: sse,
|
||||||
|
headers: { "Content-Type" => "text/event-stream" })
|
||||||
|
|
||||||
|
exchange = [
|
||||||
|
{ info: { role: "user" }, parts: [ { type: "text", text: "previous" } ] },
|
||||||
|
{ info: { role: "assistant" }, parts: [ { type: "text", text: "previous answer" } ] },
|
||||||
|
{ info: { role: "user" }, parts: [ { type: "text", text: "plan" } ] },
|
||||||
|
{ info: { role: "assistant" }, parts: [ skill_part ] },
|
||||||
|
{ info: { role: "assistant" }, parts: [ task_part ] },
|
||||||
|
{ info: { role: "assistant" }, parts: [ { type: "text", text: "SUBAGENT_OK" } ] }
|
||||||
|
]
|
||||||
|
stub_request(:get, "#{BASE}/session/#{SESSION_ID}/message")
|
||||||
|
.to_return(status: 200, body: exchange.to_json,
|
||||||
|
headers: { "Content-Type" => "application/json" })
|
||||||
|
|
||||||
|
reply = @client.stream(SESSION_ID, "plan", stream_timeout: 1)
|
||||||
|
|
||||||
|
assert_equal "SUBAGENT_OK", reply.full_text
|
||||||
|
assert_equal %w[todowrite skill task], reply.tool_parts.map { |part| part.fetch("tool") }
|
||||||
|
assert_equal "ses_child", reply.tool_parts.last.dig("metadata", "sessionId")
|
||||||
|
end
|
||||||
|
|
||||||
def test_connection_refused_raises_ConnectionError
|
def test_connection_refused_raises_ConnectionError
|
||||||
stub_request(:get, "http://opencode.dead/global/health")
|
stub_request(:get, "http://opencode.dead/global/health")
|
||||||
.to_raise(Errno::ECONNREFUSED)
|
.to_raise(Errno::ECONNREFUSED)
|
||||||
|
|||||||
Reference in New Issue
Block a user