diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index db94508..752cac7 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -29,4 +29,6 @@ jobs: 'expected = ENV.fetch("RELEASE_TAG").delete_prefix("v"); abort "tag #{expected.inspect} does not match gem #{Opencode::VERSION.inspect}" unless Opencode::VERSION == expected' + - name: Run tests + run: bundle exec rake test - uses: rubygems/release-gem@052cc82692552de3ef2b81fd670e41d13cba8092 # v1.4.0 diff --git a/CHANGELOG.md b/CHANGELOG.md index 699e8b8..0643a60 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,15 @@ with or without the optional leading space; join multiple `data:` fields; ignore the stream's optional leading UTF-8 BOM and comment fields; and decode correctly when a transport chunk splits any byte of the event framing. +- Keep request instrumentation free of query credentials and private workspace + paths. +- Load non-empty todo events and written-file artifacts in standalone Ruby + clients. +- Bound reconnects after an interactive prompt loses its event stream, handle + DNS failures as transient stream errors, and map non-404 4xx responses to + `BadRequestError` consistently. +- Retry the Rails recovery recipe with its replacement session and keep + exception details out of user-facing message content. ### Changed diff --git a/examples/conversation_recipe.rb b/examples/conversation_recipe.rb index d3ea436..a19315e 100644 --- a/examples/conversation_recipe.rb +++ b/examples/conversation_recipe.rb @@ -128,12 +128,13 @@ class GenerateAssistantReplyJob < ApplicationJob end rescue Opencode::SessionNotFoundError, Opencode::StaleSessionError raise if attempted_recreate - message.conversation.recreate_opencode_session!(client) + session_id = message.conversation.recreate_opencode_session!(client) attempted_recreate = true retry end rescue StandardError => e - message&.update!(status: :errored, content: "An error occurred: #{e.message.truncate(200)}") + Rails.logger.error(e.full_message) + message&.update!(status: :errored, content: "An error occurred. Please try again.") end private diff --git a/lib/opencode-ruby.rb b/lib/opencode-ruby.rb index c4cc0d4..48ba990 100644 --- a/lib/opencode-ruby.rb +++ b/lib/opencode-ruby.rb @@ -6,8 +6,10 @@ # non-Rails apps. require "active_support/core_ext/object/blank" # provides blank?, present?, presence require "active_support/core_ext/object/duplicable" +require "active_support/core_ext/hash/keys" # provides Hash#deep_stringify_keys require "active_support/core_ext/string/filters" # provides String#truncate require "active_support/core_ext/numeric/bytes" # provides Integer#megabytes +require "marcel" require_relative "opencode/version" require_relative "opencode/error" diff --git a/lib/opencode/client.rb b/lib/opencode/client.rb index 624cd1b..074f1be 100644 --- a/lib/opencode/client.rb +++ b/lib/opencode/client.rb @@ -230,18 +230,12 @@ module Opencode request = Net::HTTP::Get.new(uri) add_auth_header(request) - response = Opencode::Instrumentation.instrument("opencode.request", method: request.method, path: request.path) do + response = Opencode::Instrumentation.instrument("opencode.request", method: request.method, path: uri.path) do http_client.request(request) end - unless response.code.to_i.between?(200, 299) - raise ServerError, "list_questions failed: HTTP #{response.code} — #{response.body.to_s[0, 200]}" - end - - return [] if response.body.blank? - JSON.parse(response.body, symbolize_names: true) - rescue JSON::ParserError => e - raise ServerError, "list_questions returned invalid JSON: #{e.message}" + body = handle_response(response) + body.is_a?(Array) ? body : [] rescue Net::OpenTimeout, Net::ReadTimeout, Net::WriteTimeout => e raise TimeoutError, "OpenCode timeout after #{@timeout}s: #{e.message}" rescue Errno::ECONNREFUSED, SocketError => e @@ -267,7 +261,8 @@ module Opencode Net::ReadTimeout, Errno::ECONNREFUSED, Errno::ECONNRESET, - Errno::EPIPE + Errno::EPIPE, + SocketError ].freeze # Opens SSE connection to GET /event, yields parsed events filtered by session_id. @@ -334,7 +329,7 @@ module Opencode loop do now = Process.clock_gettime(Process::CLOCK_MONOTONIC) - deadline = check_deadline_or_suspend(now, deadline, timeout, reply) + deadline = check_deadline(now, deadline, timeout) # NOTE: first_event_deadline is *not* suspension-eligible. If the agent # never gets started we want to fail fast — a session that's blocked on @@ -361,6 +356,7 @@ module Opencode http.read_timeout = 30 subscription_callback_error = nil + event_callback_error = nil subscription_ready = on_subscribed.nil? begin buffer = String.new @@ -448,14 +444,19 @@ module Opencode # that's the whole point: a healthy long wait (user thinking # for 30 minutes) keeps the container warm via heartbeats so # the reaper doesn't kill it mid-wait. - on_activity_tick&.call(event) - block.call(event) + begin + on_activity_tick&.call(event) + block.call(event) + rescue StandardError => error + event_callback_error = error + raise + end return if terminal_session_event?(event) end end end rescue *TRANSIENT_SSE_ERRORS - raise if subscription_callback_error + raise if subscription_callback_error || event_callback_error # Treat transport-level SSE disconnects like clean EOF: reconnect # until an idle session event, the overall timeout, or first-event @@ -514,6 +515,11 @@ module Opencode # the prompt resolves. Otherwise apply the normal deadline check. def check_deadline_or_suspend(now, deadline, timeout, reply) return now + timeout if reply&.prompt_blocked? + + check_deadline(now, deadline, timeout) + end + + def check_deadline(now, deadline, timeout) raise TimeoutError, "SSE stream timed out after #{timeout}s" if now > deadline deadline @@ -628,7 +634,7 @@ module Opencode add_auth_header(request) response = nil - result = Opencode::Instrumentation.instrument("opencode.request", method: request.method, path: request.path) do + result = Opencode::Instrumentation.instrument("opencode.request", method: request.method, path: request.uri.path) do response = http_client.request(request) handle_response(response) end @@ -680,23 +686,30 @@ module Opencode end def handle_response(response) - return {} if response.code.to_i == 204 + status = response.code.to_i + return {} if status == 204 - body = if response.body.present? - JSON.parse(response.body, symbolize_names: true) - else - {} - end + body = parse_response_body(response, status) - case response.code.to_i + case status when 200..299 then body - when 400 then raise BadRequestError.new(error_message(body, "Bad request"), response: body) when 404 then raise SessionNotFoundError.new(error_message(body, "Session not found"), response: body) + when 400..499 then raise BadRequestError.new(error_message(body, "Bad request"), response: body) when 500..599 then raise ServerError.new(error_message(body, "Server error"), response: body) else raise Error.new("Unexpected response: #{response.code}", response: body) end + end + + def parse_response_body(response, status) + return {} if response.body.blank? + + JSON.parse(response.body, symbolize_names: true) rescue JSON::ParserError - raise ServerError.new("Invalid JSON from OpenCode (HTTP #{response.code}): #{response.body&.truncate(200)}") + if status.between?(200, 299) + raise ServerError.new("Invalid JSON from OpenCode (HTTP #{response.code}): #{response.body&.truncate(200)}") + end + + { message: response.body.to_s.truncate(200) } end # OpenCode HTTP error bodies use a wrapped shape: { name:, data: { message:, kind?: } }. @@ -704,6 +717,9 @@ module Opencode # `body[:message]` is no longer populated for errors — only `body[:data][:message]`. # We read both to keep older mock servers working in tests. def error_message(body, fallback) + return body.truncate(200) if body.is_a?(String) && body.present? + return fallback unless body.is_a?(Hash) + body.dig(:data, :message) || body[:message] || fallback end end diff --git a/opencode-ruby.gemspec b/opencode-ruby.gemspec index cfca6d1..31076fd 100644 --- a/opencode-ruby.gemspec +++ b/opencode-ruby.gemspec @@ -30,14 +30,10 @@ Gem::Specification.new do |spec| %w[README.md LICENSE CHANGELOG.md opencode-ruby.gemspec] spec.require_paths = ["lib"] - # The only runtime dependency is ActiveSupport (NOT Rails). ActiveSupport - # is a standalone gem providing the `present?`/`blank?`/`presence`/ - # `truncate`/`duplicable?` helpers used in this gem's code. It does NOT - # pull in ActiveRecord, ActionView, ActionController, Turbo, or any other - # Rails-only piece. Most Ruby apps in the wild already have ActiveSupport - # transitively via another gem; in the rare case yours doesn't, ~250 LOC - # of core_ext is added when this gem installs. + # ActiveSupport supplies the small set of core extensions used by the + # client without pulling in Rails. Marcel identifies artifact content types. spec.add_runtime_dependency "activesupport", ">= 6.1", "< 9.0" + spec.add_runtime_dependency "marcel", "~> 1.0" spec.add_development_dependency "minitest", "~> 5.20" spec.add_development_dependency "rake", "~> 13.0" diff --git a/test/conversation_recipe_test.rb b/test/conversation_recipe_test.rb new file mode 100644 index 0000000..4171281 --- /dev/null +++ b/test/conversation_recipe_test.rb @@ -0,0 +1,22 @@ +# frozen_string_literal: true + +require "test_helper" + +class ConversationRecipeTest < Minitest::Test + RECIPE = File.read(File.expand_path("../examples/conversation_recipe.rb", __dir__)) + + def test_recovery_retries_with_the_recreated_session + assert_includes RECIPE, + "session_id = message.conversation.recreate_opencode_session!(client)" + end + + def test_user_facing_error_does_not_include_exception_details + assert_includes RECIPE, 'content: "An error occurred. Please try again."' + refute_includes RECIPE, 'content: "An error occurred: #{e.message.truncate(200)}"' + end + + def test_error_reporting_works_on_the_supported_rails_6_1_baseline + assert_includes RECIPE, "Rails.logger.error(e.full_message)" + refute_includes RECIPE, "Rails.error.report" + end +end diff --git a/test/opencode/smoke_test.rb b/test/opencode/smoke_test.rb index 209b1b5..5a36643 100644 --- a/test/opencode/smoke_test.rb +++ b/test/opencode/smoke_test.rb @@ -1,6 +1,7 @@ # frozen_string_literal: true require "test_helper" +require "timeout" # End-to-end smoke test of the gem's public surface. Validates that the # headline `client.stream(...)` API + Reply::Result + error model + the @@ -499,6 +500,70 @@ class SmokeTest < Minitest::Test assert_equal [ true, false, true, false, false ], wait_states end + def test_prompt_wait_does_not_suspend_timeout_after_disconnect + question = { + type: "question.asked", + properties: { id: "que_1", sessionID: SESSION_ID, questions: [] } + } + stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z}) + .to_return( + status: 200, + body: "data: #{question.to_json}\n\n", + headers: { "Content-Type" => "text/event-stream" } + ).then.to_raise(Errno::ECONNREFUSED) + + reply = Opencode::Reply.new + assert_raises(Opencode::TimeoutError) do + Timeout.timeout(0.5) do + @client.stream_events( + session_id: SESSION_ID, + timeout: 0.01, + first_event_timeout: 1, + reply: reply + ) { |event| reply.apply(event) } + end + end + end + + def test_connected_prompt_wait_suspends_timeout_while_events_keep_arriving + question = { + type: "question.asked", + properties: { id: "que_1", sessionID: SESSION_ID, questions: [] } + } + replied = { + type: "question.replied", + properties: { requestID: "que_1", sessionID: SESSION_ID, answers: [ [ "yes" ] ] } + } + idle = { + type: "session.status", + properties: { sessionID: SESSION_ID, status: { type: "idle" } } + } + + response = Net::HTTPOK.new("1.1", "200", "OK") + response.define_singleton_method(:read_body) do |&callback| + callback.call("data: #{question.to_json}\n\n") + sleep 0.02 + callback.call("data: #{replied.to_json}\n\n") + callback.call("data: #{idle.to_json}\n\n") + end + http = Object.new + http.define_singleton_method(:use_ssl=) { |_value| } + http.define_singleton_method(:open_timeout=) { |_value| } + http.define_singleton_method(:read_timeout=) { |_value| } + http.define_singleton_method(:request) { |_request, &callback| callback.call(response) } + http.define_singleton_method(:started?) { false } + + reply = Opencode::Reply.new + Net::HTTP.stub(:new, ->(_host, _port) { http }) do + @client.stream_events( + session_id: SESSION_ID, + timeout: 0.01, + first_event_timeout: 1, + reply: reply + ) { |event| reply.apply(event) } + end + end + def test_stream_block_is_optional stub_request(:post, "#{BASE}/session/#{SESSION_ID}/prompt_async") .to_return(status: 204, body: "") @@ -543,7 +608,13 @@ class SmokeTest < Minitest::Test } sse = [ CONNECTED_EVENT, - { type: "todo.updated", properties: { sessionID: SESSION_ID, todos: [] } }, + { + type: "todo.updated", + properties: { + sessionID: SESSION_ID, + todos: [ { content: "Plan", status: "in-progress", priority: "high" } ] + } + }, { 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", @@ -573,9 +644,33 @@ class SmokeTest < Minitest::Test assert_equal "SUBAGENT_OK", reply.full_text assert_equal %w[todowrite skill task], reply.tool_parts.map { |part| part.fetch("tool") } + assert_equal( + [ { "content" => "Plan", "status" => "in_progress", "priority" => "high" } ], + reply.tool_parts.first.dig("input", "todos") + ) assert_equal "ses_child", reply.tool_parts.last.dig("metadata", "sessionId") end + def test_write_artifacts_work_in_a_standalone_client + response = { + parts: [ + { + type: "tool", + tool: "write", + state: { + status: "completed", + input: { filePath: "/tmp/report.md", content: "# Report" } + } + } + ] + } + + assert_equal( + [ { filename: "report.md", content: "# Report", content_type: "text/markdown" } ], + Opencode::ResponseParser.extract_artifact_files(response) + ) + end + def test_connection_refused_raises_ConnectionError stub_request(:get, "http://opencode.dead/global/health") .to_raise(Errno::ECONNREFUSED) @@ -594,6 +689,85 @@ class SmokeTest < Minitest::Test end end + def test_plain_text_404_preserves_SessionNotFoundError + stub_request(:get, "#{BASE}/session/missing/message") + .to_return(status: 404, body: "not found", headers: { "Content-Type" => "text/plain" }) + + assert_raises(Opencode::SessionNotFoundError) do + @client.get_messages("missing") + end + end + + def test_non_404_client_errors_raise_BadRequestError + stub_request(:get, "#{BASE}/session") + .to_return(status: 401, body: { error: "unauthorized" }.to_json, + headers: { "Content-Type" => "application/json" }) + + assert_raises(Opencode::BadRequestError) { @client.list_sessions } + end + + def test_plain_text_client_errors_raise_BadRequestError + stub_request(:get, "#{BASE}/session") + .to_return(status: 401, body: "unauthorized", headers: { "Content-Type" => "text/plain" }) + + assert_raises(Opencode::BadRequestError) { @client.list_sessions } + end + + def test_non_object_json_client_errors_raise_BadRequestError + [ '"unauthorized"', "[]", "null" ].each do |body| + stub_request(:get, "#{BASE}/session") + .to_return(status: 401, body: body, headers: { "Content-Type" => "application/json" }) + + assert_raises(Opencode::BadRequestError) { @client.list_sessions } + end + end + + def test_list_questions_uses_the_public_error_hierarchy + stub_request(:get, "#{BASE}/question") + .to_return(status: 401, body: { error: "unauthorized" }.to_json, + headers: { "Content-Type" => "application/json" }) + + assert_raises(Opencode::BadRequestError) { @client.list_questions } + end + + def test_stream_retries_socket_errors_until_its_timeout + stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z}) + .to_raise(SocketError.new("name resolution failed")) + + assert_raises(Opencode::StaleSessionError) do + @client.stream_events( + session_id: SESSION_ID, + timeout: 1, + first_event_timeout: 0.01 + ) { } + end + end + + def test_stream_does_not_swallow_socket_errors_from_the_caller_block + event = { + type: "message.part.updated", + properties: { + sessionID: SESSION_ID, + part: { id: "part_1", type: "text", text: "hello" } + } + } + stub_request(:get, %r{#{Regexp.escape(BASE)}/event(\?.*)?\z}) + .to_return( + status: 200, + body: "data: #{event.to_json}\n\n", + headers: { "Content-Type" => "text/event-stream" } + ) + + error = assert_raises(SocketError) do + Timeout.timeout(0.5) do + @client.stream_events(session_id: SESSION_ID, timeout: 1, first_event_timeout: 1) do + raise SocketError, "caller lookup failed" + end + end + end + assert_equal "caller lookup failed", error.message + end + def test_instrumentation_adapter_receives_request_events events = [] Opencode::Instrumentation.adapter = ->(name, payload, &block) { @@ -610,6 +784,47 @@ class SmokeTest < Minitest::Test "instrumentation adapter must receive opencode.request events" end + def test_instrumentation_does_not_expose_scoped_query_values + events = [] + Opencode::Instrumentation.adapter = ->(name, payload, &block) { + events << [ name, payload ] + block.call + } + client = Opencode::Client.new( + base_url: "#{BASE}?token=base-secret", + directory: "/private/workspace", + workspace: "workspace-secret" + ) + stub_request(:get, %r{#{Regexp.escape(BASE)}/session\?.*}) + .to_return(status: 200, body: "[]", + headers: { "Content-Type" => "application/json" }) + + client.list_sessions + + payload = events.find { |name, _| name == "opencode.request" }.last + assert_equal "/session", payload.fetch(:path) + end + + def test_list_questions_instrumentation_does_not_expose_scoped_query_values + events = [] + Opencode::Instrumentation.adapter = ->(name, payload, &block) { + events << [ name, payload ] + block.call + } + client = Opencode::Client.new( + base_url: "#{BASE}?token=base-secret", + directory: "/private/workspace", + workspace: "workspace-secret" + ) + stub_request(:get, %r{#{Regexp.escape(BASE)}/question\?.*}) + .to_return(status: 200, body: "[]", headers: { "Content-Type" => "application/json" }) + + client.list_questions + + payload = events.find { |name, _| name == "opencode.request" }.last + assert_equal "/question", payload.fetch(:path) + end + def test_Reply_distill_returns_typed_Result parts = [ { "type" => "text", "content" => "hi" }, diff --git a/test/release_workflow_test.rb b/test/release_workflow_test.rb index ecbc78d..723c13b 100644 --- a/test/release_workflow_test.rb +++ b/test/release_workflow_test.rb @@ -50,4 +50,15 @@ class ReleaseWorkflowTest < Minitest::Test assert_equal "${{ github.ref_name }}", preflight.dig("env", "RELEASE_TAG") assert_includes preflight.fetch("run"), "unless Opencode::VERSION == expected" end + + def test_release_runs_the_test_suite_before_publish + steps = push_job.fetch("steps") + test_index = steps.index { |step| step["name"] == "Run tests" } + publish_index = steps.index { |step| step["uses"] == RELEASE_GEM_ACTION } + + refute_nil test_index + refute_nil publish_index + assert_operator test_index, :<, publish_index + assert_equal "bundle exec rake test", steps.fetch(test_index).fetch("run") + end end