diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 1c6ead0..e7da970 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -102,7 +102,7 @@ jobs: options: --health-cmd "redis-cli ping" --health-interval 10s --health-timeout 5s --health-retries 5 steps: - name: Install packages - run: sudo apt-get update && sudo apt-get install --no-install-recommends -y libsqlite3-0 libvips curl ffmpeg + run: sudo apt-get update && sudo apt-get install --no-install-recommends -y libsqlite3-0 sqlite3 libvips curl ffmpeg - name: Checkout code uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 @@ -120,6 +120,9 @@ jobs: - name: Run tests run: bin/rails db:setup test + - name: Test benchmark response validation + run: bundle exec ruby bench/test_http_client.rb + test_system: runs-on: ubuntu-latest permissions: @@ -132,7 +135,7 @@ jobs: options: --health-cmd "redis-cli ping" --health-interval 10s --health-timeout 5s --health-retries 5 steps: - name: Install packages - run: sudo apt-get update && sudo apt-get install --no-install-recommends -y libsqlite3-0 libvips curl ffmpeg + run: sudo apt-get update && sudo apt-get install --no-install-recommends -y libsqlite3-0 sqlite3 libvips curl ffmpeg - name: Checkout code uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 diff --git a/app/controllers/searches_controller.rb b/app/controllers/searches_controller.rb index 735da0c..9bfae79 100644 --- a/app/controllers/searches_controller.rb +++ b/app/controllers/searches_controller.rb @@ -20,7 +20,7 @@ class SearchesController < ApplicationController private def set_messages if query.present? - @messages = Current.user.reachable_messages.search(query).with_presentation.last_page_of_matches(100) + @messages = Message.search_reachable(Current.user, query) else @messages = Message.none end diff --git a/app/models/message/searchable.rb b/app/models/message/searchable.rb index 43b35d9..f7b4962 100644 --- a/app/models/message/searchable.rb +++ b/app/models/message/searchable.rb @@ -19,6 +19,24 @@ module Message::Searchable end class_methods do + def search_reachable(user, query, size: 100) + relation = user.reachable_messages.search(query).with_presentation + probe = connection.select_rows(sanitize_sql([ <<~SQL, user.id, match_terms(query) ])) + SELECT m.id, mm.user_id IS NOT NULL FROM message_search_index idx + JOIN messages m ON m.id = idx.rowid + LEFT JOIN memberships mm ON mm.room_id = m.room_id AND mm.user_id = ? + WHERE idx.body MATCH ? ORDER BY idx.rowid DESC LIMIT 1000 + SQL + ids = probe.filter_map { |id, reachable| id if reachable == 1 }.first(size) + + # A scoped fallback preserves older reachable hits behind a large private history. + if ids.size < size && probe.size == 1000 + relation.last_page_of_matches(size) + else + Message::Pagination::Page.load(relation.where(id: ids).reorder(:id), :last, size) + end + end + # Orders by the index's rowid, which is the message id, so SQLite walks the full-text index # newest first and stops at the page. Ordering by created_at sorted every match before paging. def last_page_of_matches(size) diff --git a/bench/README.md b/bench/README.md index ca5cd5e..8204952 100644 --- a/bench/README.md +++ b/bench/README.md @@ -22,5 +22,6 @@ CSRF disabling. Unread fanout excludes adapter I/O. The HTTP driver uses production Puma/Redis with one worker and five threads. Ruby threads each maintain a keep-alive connection, request uncompressed responses, and consume the whole body. Login uses normal CSRF protection; all warmup and measured responses must be -HTTP 200 without transport errors. Measurements exclude Thruster, TLS and gzip. Client +HTTP 200 without transport errors, complete HTML, exact independently seeded message windows +and the visible sidebar rooms; this applies to every warmup and measured response. Measurements exclude Thruster, TLS and gzip. Client CPU, JIT warmup and GC can affect throughput; repeat runs and check client saturation. diff --git a/bench/compare_http.rb b/bench/compare_http.rb index 58eb421..519b016 100644 --- a/bench/compare_http.rb +++ b/bench/compare_http.rb @@ -3,6 +3,7 @@ require "socket" require_relative "support" require_relative "http_client" +require_relative "response_contract" include BenchmarkSupport options = parse_options("Compare HTTP throughput with Ruby keep-alive clients; every response must be HTTP 200.", @@ -58,9 +59,10 @@ begin cookie = client.login(labels) results = {} paths.each do |name, path| - client.measure(path, cookie, concurrency: 1, duration: 3) + contract = BenchmarkResponseContract.new(name, File.join(data, "storage/db/production.sqlite3"), labels) + client.measure(path, cookie, concurrency: 1, duration: 3, contract: contract) concurrencies.each do |concurrency| - results["#{name}_#{concurrency}"] = client.measure(path, cookie, concurrency: concurrency, duration: options[:duration]) + results["#{name}_#{concurrency}"] = client.measure(path, cookie, concurrency: concurrency, duration: options[:duration], contract: contract) end end write_json(File.join(options[:output], "#{side}-#{iteration + 1}.json"), results) diff --git a/bench/http_client.rb b/bench/http_client.rb index 574f3b7..d126862 100644 --- a/bench/http_client.rb +++ b/bench/http_client.rb @@ -33,12 +33,12 @@ class BenchmarkHTTPClient cookie_header(cookies) end - def measure(path, cookie, concurrency:, duration:) + def measure(path, cookie, concurrency:, duration:, contract:) start = clock deadline = start + duration workers = Array.new(concurrency) do Thread.new do - result = { latencies: [], statuses: Hash.new(0), bytes: 0, errors: 0 } + result = { latencies: [], statuses: Hash.new(0), bytes: 0, errors: 0, invalid_responses: 0 } while clock < deadline begin connection.start do |http| @@ -48,6 +48,7 @@ class BenchmarkHTTPClient result[:latencies] << (clock - requested) * 1000 result[:statuses][response.code] += 1 result[:bytes] += response.body.bytesize + result[:invalid_responses] += 1 unless contract.valid?(response) end end rescue IOError, SystemCallError, Timeout::Error, SocketError, Net::HTTPBadResponse @@ -64,9 +65,10 @@ class BenchmarkHTTPClient statuses = Hash.new(0) samples.each { |sample| sample[:statuses].each { |status, count| statuses[status] += count } } errors = samples.sum { |sample| sample[:errors] } - raise "#{path}: HTTP statuses #{statuses}, #{errors} transport errors" unless errors.zero? && statuses.keys == [ "200" ] + invalid = samples.sum { |sample| sample[:invalid_responses] } + raise "#{path}: HTTP statuses #{statuses}, #{errors} transport errors, #{invalid} invalid bodies" unless errors.zero? && invalid.zero? && statuses.keys == [ "200" ] { path: path, conc: concurrency, gzip: false, secs: elapsed, rps: latencies.size / elapsed, - ok: latencies.size, statuses: statuses, errors: errors, + ok: latencies.size, statuses: statuses, errors: errors, invalid_responses: invalid, validation: "route-contract-v1", avg_bytes: samples.sum { |sample| sample[:bytes] } / latencies.size, latency_ms: { p50: percentile(latencies, 0.50), p95: percentile(latencies, 0.95), p99: percentile(latencies, 0.99) } } end diff --git a/bench/response_contract.rb b/bench/response_contract.rb new file mode 100644 index 0000000..fc12bb3 --- /dev/null +++ b/bench/response_contract.rb @@ -0,0 +1,43 @@ +require "json" +require "open3" +require "cgi" + +# Independent seed SQL determines the required result window before any response is sampled. +class BenchmarkResponseContract + def initialize(kind, database, labels) + @kind = kind + room = Integer(labels.fetch("rooms.watercooler")) + user = labels.fetch("emails.david").gsub("'", "''") + @ids = case kind + when "room" + query(database, "SELECT id FROM messages WHERE room_id=#{room} ORDER BY created_at DESC LIMIT 40").reverse.map { |row| row.fetch("id") } + when "messages" + before = Integer(labels.fetch("messages.busy_060")) + query(database, "SELECT id FROM messages WHERE room_id=#{room} AND created_at<(SELECT created_at FROM messages WHERE id=#{before}) ORDER BY created_at DESC LIMIT 40").reverse.map { |row| row.fetch("id") } + when "search" + query(database, "SELECT m.id FROM message_search_index idx JOIN messages m ON m.id=idx.rowid JOIN memberships mm ON mm.room_id=m.room_id JOIN users u ON u.id=mm.user_id WHERE u.email_address='#{user}' AND idx.body MATCH 'coffee' ORDER BY m.id DESC LIMIT 100").reverse.map { |row| row.fetch("id") } + when "sidebar" + @names = query(database, "SELECT r.name FROM rooms r JOIN memberships mm ON mm.room_id=r.id JOIN users u ON u.id=mm.user_id WHERE u.email_address='#{user}' AND r.type='Rooms::Open' AND mm.involvement<>'invisible'").map { |row| CGI.escapeHTML(row.fetch("name")) } + nil + else + raise "Unknown route contract #{kind}" + end + end + + def valid?(response) + body = response.body.to_s + return false unless response.code == "200" && response["content-type"].to_s.split(";").first == "text/html" && + [ nil, "identity" ].include?(response["content-encoding"]) && !body.empty? && body.dup.force_encoding("UTF-8").valid_encoding? + return false if @kind != "messages" && !(body.lstrip.start_with?("") && body.rstrip.end_with?("")) + return false if @ids && body.scan(/data-message-id="(\d+)"/).flatten.map(&:to_i) != @ids + return false if @names && (!body.include?("shared_rooms") || @names.any? { |name| !body.include?(name) }) + true + end + + private + def query(database, sql) + output, errors, status = Open3.capture3("sqlite3", "-readonly", "-json", database, sql) + raise "Cannot read benchmark seed: #{errors}" unless status.success? + JSON.parse(output) + end +end diff --git a/bench/test_http_client.rb b/bench/test_http_client.rb new file mode 100644 index 0000000..3c67845 --- /dev/null +++ b/bench/test_http_client.rb @@ -0,0 +1,71 @@ +require "minitest/autorun" +require "fileutils" +require "socket" +require "sqlite3" +require "tmpdir" +require_relative "http_client" +require_relative "response_contract" + +class BenchmarkHTTPClientTest < Minitest::Test + def setup + @directory = Dir.mktmpdir("campfire-benchmark-test") + database = File.join(@directory, "seed.sqlite3") + SQLite3::Database.new(database).tap do |db| + db.execute("CREATE TABLE messages (id INTEGER PRIMARY KEY, room_id INTEGER, created_at TEXT)") + db.execute("INSERT INTO messages VALUES (42, 1, '2026-10-01'), (43, 1, '2026-10-02')") + db.close + end + @contract = BenchmarkResponseContract.new("room", database, "rooms.watercooler" => 1, "emails.david" => "david@example.com") + @body = '
coffee
meeting
' + end + + def teardown + FileUtils.remove_entry(@directory) + end + + def test_seeded_contract_rejects_error_pages_truncation_wrong_windows_and_encoding + assert @contract.valid?(response(@body)) + refute @contract.valid?(response("Something went wrong")) + refute @contract.valid?(response(@body.delete_suffix(""))) + refute @contract.valid?(response(@body.sub('data-message-id="42"', 'data-message-id="99"'))) + refute @contract.valid?(response(@body, "text/plain")) + refute @contract.valid?(response(@body + "\xFF".b)) + refute @contract.valid?(response(@body, "text/html", "500")) + end + + def test_intermittent_200_error_page_rejects_the_loaded_run + server = TCPServer.new("127.0.0.1", 0) + worker = Thread.new do + loop do + socket = server.accept + count = 0 + while socket.gets + while (line = socket.gets) && line != "\r\n" + end + body = count % 8 == 0 ? "Something went wrong" : @body + count += 1 + socket.write("HTTP/1.1 200 OK\r\nContent-Type: text/html\r\nContent-Length: #{body.bytesize}\r\n\r\n#{body}") + end + socket.close + end + rescue IOError, SystemCallError + nil + end + client = BenchmarkHTTPClient.new("http://127.0.0.1:#{server.addr[1]}") + error = assert_raises(RuntimeError) { client.measure("/rooms/1", "", concurrency: 1, duration: 0.1, contract: @contract) } + assert_match(/HTTP statuses \{"200"\s*=>\s*\d+\}, 0 transport errors, [1-9]\d* invalid bodies/, error.message) + ensure + server&.close + worker&.kill + worker&.join + end + + private + def response(body, type = "text/html", status = "200") + Net::HTTPResponse.new("1.1", status, "OK").tap do |result| + result["content-type"] = type + result.instance_variable_set(:@body, body) + result.instance_variable_set(:@read, true) + end + end +end diff --git a/test/models/message/searchable_test.rb b/test/models/message/searchable_test.rb index 713b370..f856ecc 100644 --- a/test/models/message/searchable_test.rb +++ b/test/models/message/searchable_test.rb @@ -68,6 +68,22 @@ class Message::SearchableTest < ActiveSupport::TestCase assert_no_match(/TEMP B-TREE/, Message.connection.select_rows("EXPLAIN QUERY PLAN #{queries.sole}").map(&:last).join(" | ")) end + test "bounded search keeps older reachable hits behind inaccessible matches" do + accessible = rooms(:designers).messages.create! body: "bounded sparse needle", creator: users(:david) + private_room = Room.create! name: "Private search history", type: "Rooms::Closed", creator: users(:jason) + now = Time.current + records = 1100.times.map do |i| + { room_id: private_room.id, creator_id: users(:jason).id, client_message_id: "private-#{i}", created_at: now, updated_at: now } + end + ids = Message.insert_all!(records, returning: %w[id]).rows.flatten + ids.each do |id| + Message.connection.execute Message.sanitize_sql([ "INSERT INTO message_search_index(rowid,body) VALUES (?,?)", id, "bounded sparse needle" ]) + end + assert_equal [ accessible.id ], Message.search_reachable(users(:david), "bounded sparse").map(&:id) + rooms(:designers).memberships.where(user: users(:david)).delete_all + assert_empty Message.search_reachable(users(:david), "bounded sparse") + end + test "rich text body is converted to plain text for indexing" do message = rooms(:designers).messages.create! body: "My hovercraft is full of eels", client_message_id: "earth", creator: users(:david)