mirror of
https://github.com/basecamp/once-campfire.git
synced 2026-10-08 15:50:08 +09:00
Bound accessible search probes and reject invalid benchmark responses
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
+2
-1
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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?("<!DOCTYPE html>") && body.rstrip.end_with?("</html>"))
|
||||
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
|
||||
@@ -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 = '<!DOCTYPE html><html><article data-message-id="42">coffee</article><article data-message-id="43">meeting</article></html>'
|
||||
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("</html>")))
|
||||
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
|
||||
@@ -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: "<span>My hovercraft is full of eels</span>", client_message_id: "earth", creator: users(:david)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user