Merge pull request #326: Reuse push connections pinned to the address the guard just approved

Reviewed and merged by GPT on behalf of DHH.
This commit is contained in:
GPT on behalf of DHH
2026-10-07 10:33:57 +02:00
8 changed files with 489 additions and 38 deletions
-1
View File
@@ -43,7 +43,6 @@ gem "rqrcode"
gem "rails_autolink"
gem "geared_pagination"
gem "jbuilder"
gem "net-http-persistent"
gem "surfguard", github: "basecamp/surfguard" # The SSRF address policy behind RestrictedHTTP
gem "kredis"
gem "platform_agent"
-3
View File
@@ -241,8 +241,6 @@ GEM
multi_json (1.17.0)
mustermann (3.0.4)
ruby2_keywords (~> 0.0.1)
net-http-persistent (4.0.6)
connection_pool (~> 2.2, >= 2.2.4)
net-imap (0.6.7)
date
net-protocol
@@ -451,7 +449,6 @@ DEPENDENCIES
kredis
lexxy (~> 0.9.24)
mocha
net-http-persistent
ostruct
platform_agent
propshaft!
+22 -33
View File
@@ -1,4 +1,5 @@
require "web-push"
require "web_push/connections"
require "web_push/pool"
require "web_push/notification"
@@ -17,43 +18,31 @@ end
module WebPush::PersistentRequest
def perform
if endpoint_ip = @options[:endpoint_ip]
# Pin the connection to the public IP resolved (and guarded) by
# Push::Subscription so delivery can't be rebound to a private address
# between resolution and connect. Bypasses the shared persistent pool,
# which would re-resolve the host itself.
#
# The explicit nil proxy address disables proxy discovery from
# http_proxy/https_proxy. An egress proxy would open the TCP connection
# itself and re-resolve the endpoint host, so http.ipaddr would no longer
# pin the destination and the DNS-rebinding guarantee would be lost. This
# path is already committed to a direct connection (it bypasses the pool);
# push delivery to public vendor endpoints goes direct.
http = Net::HTTP.new(uri.host, uri.port, nil)
http.ipaddr = endpoint_ip
http.use_ssl = true
http.ssl_timeout = @options[:ssl_timeout] unless @options[:ssl_timeout].nil?
http.open_timeout = @options[:open_timeout] unless @options[:open_timeout].nil?
http.read_timeout = @options[:read_timeout] unless @options[:read_timeout].nil?
elsif @options[:connection]
http = @options[:connection]
else
http = Net::HTTP.new(uri.host, uri.port, *proxy_options)
http.use_ssl = true
http.ssl_timeout = @options[:ssl_timeout] unless @options[:ssl_timeout].nil?
http.open_timeout = @options[:open_timeout] unless @options[:open_timeout].nil?
http.read_timeout = @options[:read_timeout] unless @options[:read_timeout].nil?
end
# Pin the connection to the public IP resolved (and guarded) by
# Push::Subscription so delivery can't be rebound to a private address
# between resolution and connect. There is no unpinned path: a delivery
# without a resolved IP is not sent.
endpoint_ip = @options[:endpoint_ip] or raise ArgumentError, "Push deliveries must be pinned to a resolved endpoint IP"
# The explicit nil proxy address disables proxy discovery from
# http_proxy/https_proxy. An egress proxy would open the TCP connection
# itself and re-resolve the endpoint host, so http.ipaddr would no longer
# pin the destination and the DNS-rebinding guarantee would be lost.
# Push delivery to public vendor endpoints goes direct.
http = Net::HTTP.new(uri.host, uri.port, nil)
http.ipaddr = endpoint_ip
http.use_ssl = true
http.ssl_timeout = @options[:ssl_timeout] unless @options[:ssl_timeout].nil?
http.open_timeout = @options[:open_timeout] unless @options[:open_timeout].nil?
http.read_timeout = @options[:read_timeout] unless @options[:read_timeout].nil?
req = Net::HTTP::Post.new(uri.request_uri, headers)
req.body = body
if http.is_a?(Net::HTTP::Persistent)
response = http.request uri, req
else
resp = http.request(req)
verify_response(resp)
end
# WebPush::Connections reuses an open connection only for this same host
# and pinned IP, so the guarantee holds for every request it sends.
resp = @options[:connection] ? @options[:connection].request(http, req) : http.request(req)
verify_response(resp)
resp
end
+119
View File
@@ -0,0 +1,119 @@
# This is in lib so we can use it in a thread pool without the Rails executor
#
# Keeps TLS connections to push services open between deliveries, so a push doesn't pay for a new TCP and TLS
# handshake every time. Each connection is pinned to the address Push::Subscription resolved and guarded for the
# delivery that opened it, and is handed out again only to a delivery whose own resolution, just made, returned that
# same address for the same host. Net::HTTP reconnects to its pinned address, never to a new lookup of the host.
class WebPush::Connections
class ConnectionLost < StandardError; end
class StaleConnection < StandardError; end
# Where a request failed. Before writing it, Net::HTTP checks an idle connection (:checking) and connects again
# if the push service closed it (:connecting); then it writes (:sent). Only a failed check means a dead idle
# connection that a new one can replace without sending the push twice.
module Stages
attr_reader :stage
private
def begin_transport(...)
@stage = :checking
super.tap { @stage = :sent }
end
def connect(...)
@stage = :connecting if @stage == :checking
super
end
end
def initialize(keep_alive_timeout: 30)
@keep_alive_timeout = keep_alive_timeout
@idle = Hash.new { |idle, address| idle[address] = [] }
@mutex = Mutex.new
@pid = Process.pid
end
def request(http, request)
unless http.ipaddr && !http.proxy? && http.use_ssl? && !http.started?
raise ArgumentError, "Only new, direct TLS connections pinned to an address are pooled"
end
address = [ http.address, http.port, http.ipaddr ]
if idle = checkout(address)
begin
return send_over(idle, request, reused: true).tap { checkin(address, idle) }
rescue StaleConnection
# The push service had closed it: a new connection takes the push
end
end
http.extend Stages
http.keep_alive_timeout = @keep_alive_timeout
http.start
send_over(http, request, reused: false).tap { checkin(address, http) }
end
def shutdown
@mutex.synchronize do
forget_after_fork
@shut_down = true
@idle.each_value { |connections| connections.each { |http, _| close(http) } }
@idle.clear
end
end
private
def send_over(http, request, reused:)
http.request(request)
rescue IOError, SystemCallError, OpenSSL::SSL::SSLError => error
close(http)
# Only a connection that was idle can have been dead already: on a new one the error is the push service's
raise StaleConnection if reused && http.stage == :checking
# The push service may have it: don't send it twice, and don't report a dropped connection as a TLS failure
raise ConnectionLost, "#{error.class}: #{error.message}" if http.stage == :sent
raise
end
def checkout(address)
@mutex.synchronize do
forget_after_fork
close_expired
@idle[address].pop&.first
end
end
def checkin(address, http)
@mutex.synchronize do
close_expired
if @shut_down
close(http)
else
@idle[address].push [ http, now ]
end
end
end
# A forked process must not write on its parent's TLS sessions. They're left to the parent, not closed.
def forget_after_fork
unless @pid == Process.pid
@idle = Hash.new { |idle, address| idle[address] = [] }
@pid = Process.pid
end
end
def close_expired
@idle.each_value do |connections|
connections.reject! { |http, idle_since| (now - idle_since > @keep_alive_timeout).tap { |expired| close(http) if expired } }
end
@idle.delete_if { |_, connections| connections.empty? }
end
def close(http)
http.finish if http.started?
rescue IOError, SystemCallError, OpenSSL::SSL::SSLError
end
def now
Process.clock_gettime(Process::CLOCK_MONOTONIC)
end
end
+1 -1
View File
@@ -5,7 +5,7 @@ class WebPush::Pool
def initialize(invalid_subscription_handler:)
@delivery_pool = Concurrent::ThreadPoolExecutor.new(max_threads: 50, max_queue: 10000)
@invalidation_pool = Concurrent::FixedThreadPool.new(1)
@connection = Net::HTTP::Persistent.new(name: "web_push", pool_size: 150)
@connection = WebPush::Connections.new
@invalid_subscription_handler = invalid_subscription_handler
end
+156
View File
@@ -0,0 +1,156 @@
require "test_helper"
class WebPush::ConnectionsTest < ActiveSupport::TestCase
include PushServiceTestHelper
setup { @connections = WebPush::Connections.new }
teardown { @connections.shutdown }
test "a delivery to the same host and address reuses the open connection" do
with_push_service do |server|
2.times { |i| assert_kind_of Net::HTTPCreated, @connections.request(pinned_connection(server), push_request("/push/#{i}")) }
assert_equal 1, server.connections
assert_equal %w[ /push/0 /push/1 ], server.requests
end
end
test "a delivery to another address never reuses a connection" do
with_push_service do |server|
@connections.request(pinned_connection(server), push_request)
other = pinned_connection(server, "127.0.0.2")
other.expects(:start)
other.expects(:request).returns(:sent_to_the_other_address)
assert_equal :sent_to_the_other_address, @connections.request(other, push_request)
assert_equal 1, server.requests.size
end
end
test "a connection the push service closed while idle is replaced, and the push sent once" do
[ :fin, :close_notify ].each do |hang_up|
with_push_service(hang_up_after_response: hang_up) do |server|
@connections.request(pinned_connection(server), push_request("/push/0"))
assert server.hung_up?
assert_kind_of Net::HTTPCreated, @connections.request(pinned_connection(server), push_request("/push/1"))
assert_equal 2, server.connections
assert_equal %w[ /push/0 /push/1 ], server.requests
end
end
end
test "a certificate for another name is refused, and nothing is sent" do
with_push_service(certificate: ->(_) { "other.test" }) do |server|
assert_raises(OpenSSL::SSL::SSLError) { @connections.request(pinned_connection(server), push_request) }
assert_empty server.requests
end
end
test "a certificate for another name is refused when a dead connection is replaced too" do
with_push_service(hang_up_after_response: :fin, certificate: ->(connection) { connection == 1 ? HOST : "other.test" }) do |server|
@connections.request(pinned_connection(server), push_request)
assert server.hung_up?
assert_raises(OpenSSL::SSL::SSLError) { @connections.request(pinned_connection(server), push_request) }
assert_equal 1, server.requests.size
end
end
test "a push the service may have received is not sent again" do
with_push_service(drop_request: ->(number) { number == 2 }) do |server|
@connections.request(pinned_connection(server), push_request("/push/0"))
error = assert_raises(WebPush::Connections::ConnectionLost) { @connections.request(pinned_connection(server), push_request("/push/1")) }
assert_not_kind_of OpenSSL::OpenSSLError, error
assert_equal %w[ /push/0 /push/1 ], server.requests
assert_equal 1, server.connections
end
end
test "a push the service may have received on a new connection isn't reported as a TLS failure either" do
with_push_service(drop_request: ->(number) { number == 1 }) do |server|
error = assert_raises(WebPush::Connections::ConnectionLost) { @connections.request(pinned_connection(server), push_request) }
assert_not_kind_of OpenSSL::OpenSSLError, error
assert_equal 1, server.requests.size
end
end
test "a TLS failure on a new connection before the push is written is reported as it is" do
failing_check = Module.new { private def begin_transport(*) = raise(OpenSSL::SSL::SSLError, "alert right after the handshake") }
with_push_service do |server|
assert_raises(OpenSSL::SSL::SSLError) { @connections.request(pinned_connection(server).extend(failing_check), push_request) }
assert_empty server.requests
end
end
test "a certificate for another name when reconnecting a cleanly closed connection is reported, even if a new one would work" do
certificates = { 2 => "other.test" }
with_push_service(hang_up_after_response: :close_notify, certificate: ->(connection) { certificates.fetch(connection, HOST) }) do |server|
@connections.request(pinned_connection(server), push_request)
assert server.hung_up?
assert_raises(OpenSSL::SSL::SSLError) { @connections.request(pinned_connection(server), push_request) }
assert_equal 1, server.requests.size
assert_equal 2, server.connections
end
end
test "only new, direct TLS connections pinned to an address are pooled" do
with_push_service do |server|
unpinned = Net::HTTP.new(HOST, server.port, nil).tap { it.use_ssl = true }
proxied = Net::HTTP.new(HOST, server.port, "127.0.0.1", 3128).tap { it.ipaddr = IP; it.use_ssl = true }
plain = pinned_connection(server).tap { it.use_ssl = false }
started = pinned_connection(server).tap(&:start)
[ unpinned, proxied, plain, started ].each do |http|
assert_raises(ArgumentError) { @connections.request(http, push_request) }
end
assert_empty server.requests
ensure
started&.finish
end
end
test "a forked process opens its own connections" do
with_push_service do |server|
@connections.request(pinned_connection(server), push_request)
child = Process.pid + 1
Process.stubs(:pid).returns(child)
@connections.request(pinned_connection(server), push_request)
assert_equal 2, server.connections
end
end
test "a forked process shutting down leaves its parent's connections open" do
with_push_service do |server|
@connections.request(pinned_connection(server), push_request)
child = Process.pid + 1
Process.stubs(:pid).returns(child)
@connections.shutdown
assert_not server.hung_up?(within: 0.5)
end
end
test "idle connections are closed after the keep-alive timeout and on shutdown" do
connections = WebPush::Connections.new(keep_alive_timeout: 0)
with_push_service do |server|
other = Server.new
connections.request(pinned_connection(server), push_request)
connections.request(pinned_connection(other), push_request)
assert server.hung_up?
connections.shutdown
assert other.hung_up?
connections.request(pinned_connection(server), push_request)
assert server.hung_up?
ensure
other&.stop
end
end
end
@@ -1,6 +1,8 @@
require "test_helper"
class WebPush::PersistentRequestTest < ActiveSupport::TestCase
include PushServiceTestHelper
ENDPOINT = "https://fcm.googleapis.com/fcm/send/test123"
# The delivery must connect to the public IP resolved and guarded by
@@ -50,4 +52,72 @@ class WebPush::PersistentRequestTest < ActiveSupport::TestCase
%w[ http_proxy https_proxy HTTP_PROXY HTTPS_PROXY ].each { |k| ENV.delete(k) }
saved.each { |k, v| ENV[k] = v }
end
test "deliveries through the pool to the same host and address share one connection" do
pool = WebPush::Pool.new(invalid_subscription_handler: nil)
with_push_service do |server|
2.times { deliver_to server, connection: pool.connection }
assert_equal 2, server.requests.size
assert_equal 1, server.connections
end
ensure
pool.shutdown
end
test "a pooled delivery checks the response like any other" do
pool = WebPush::Pool.new(invalid_subscription_handler: nil)
with_push_service(status: "410 Gone") do |server|
assert_raises(WebPush::ExpiredSubscription) { deliver_to server, connection: pool.connection }
end
ensure
pool.shutdown
end
test "pooled deliveries are pinned and ignore proxy env too" do
host = URI(ENDPOINT).host
saved = ENV.slice("http_proxy", "https_proxy", "HTTP_PROXY", "HTTPS_PROXY")
%w[ http_proxy https_proxy HTTP_PROXY HTTPS_PROXY ].each { |k| ENV[k] = "http://proxy.internal:3128" }
WebMock.disable_net_connect! allow: [ host ]
TCPSocket.expects(:open).with { |*args, **| args.first == host || args.first == "proxy.internal" }.never
TCPSocket.expects(:open).with { |*args, **| args.first == DnsTestHelper::WEB_PUSH_PUBLIC_TEST_IP && args[1] == 443 }.throws(:pinned_to_ip)
assert_throws :pinned_to_ip do
WebPush.payload_send \
message: "",
endpoint: ENDPOINT,
endpoint_ip: DnsTestHelper::WEB_PUSH_PUBLIC_TEST_IP,
p256dh: "", auth: "", vapid: {},
connection: WebPush::Connections.new,
urgency: "high"
end
ensure
%w[ http_proxy https_proxy HTTP_PROXY HTTPS_PROXY ].each { |k| ENV.delete(k) }
saved.each { |k, v| ENV[k] = v }
end
test "nothing is sent without a resolved endpoint IP" do
WebMock.disable_net_connect! allow: [ URI(ENDPOINT).host ]
TCPSocket.expects(:open).never
error = assert_raises(ArgumentError) do
WebPush.payload_send message: "", endpoint: ENDPOINT, p256dh: "", auth: "", vapid: {}, urgency: "high"
end
assert_equal "Push deliveries must be pinned to a resolved endpoint IP", error.message
end
private
def deliver_to(server, connection:)
WebPush.payload_send \
message: "",
endpoint: "https://#{PushServiceTestHelper::HOST}:#{server.port}/fcm/send/test123",
endpoint_ip: PushServiceTestHelper::IP,
p256dh: "", auth: "", vapid: {},
connection: connection,
urgency: "high"
end
end
@@ -0,0 +1,121 @@
module PushServiceTestHelper
# A push service on 127.0.0.1 for delivery tests: TLS with a certificate for HOST from a test CA that the
# process trusts, HTTP/1.1 keep-alive. HOST doesn't resolve, so a delivery that looked it up instead of
# using its pinned address would fail.
HOST = "push.test"
IP = "127.0.0.1"
class Server
attr_reader :port
# status: the answer to every request.
# hang_up_after_response: :fin or :close_notify to hang up after each answer, as an idle timeout would.
# drop_request: ->(number) whether to hang up after reading that request, without answering it.
# certificate: ->(connection) the name its certificate is for, HOST unless told otherwise.
def initialize(status: "201 Created", hang_up_after_response: nil, drop_request: ->(_) { false }, certificate: ->(_) { HOST })
@status, @hang_up_after_response, @drop_request, @certificate = status, hang_up_after_response, drop_request, certificate
@connections, @requests, @closed = 0, [], Queue.new
@mutex = Mutex.new
@listener = TCPServer.new(IP, 0)
@port = @listener.addr[1]
@thread = Thread.new { loop { serve(@listener.accept) } rescue IOError }
end
def connections = @mutex.synchronize { @connections }
def requests = @mutex.synchronize { @requests.dup }
# Whether a connection to it closed, waiting for that up to the given seconds.
def hung_up?(within: 5)
!!@closed.pop(timeout: within)
end
def stop
@listener.close
@thread.join
end
private
def serve(socket)
connection = @mutex.synchronize { @connections += 1 }
Thread.new do
tls = OpenSSL::SSL::SSLSocket.new(socket, PushServiceTestHelper.context_for(@certificate.(connection))).tap(&:accept)
while (request_line = tls.gets("\r\n"))
length = 0
while (header = tls.gets("\r\n")) != "\r\n"
length = header.split(":", 2).last.to_i if header.downcase.start_with?("content-length:")
end
tls.read(length)
number = @mutex.synchronize { @requests << request_line.split[1]; @requests.size }
break if @drop_request.(number)
tls.write "HTTP/1.1 #{@status}\r\nContent-Length: 0\r\n\r\n"
if @hang_up_after_response
tls.close if @hang_up_after_response == :close_notify
break
end
end
rescue IOError, SystemCallError, OpenSSL::SSL::SSLError
ensure
socket.close
@closed << true
end
end
end
class << self
def context_for(name)
@contexts ||= {}
@contexts[name] ||= OpenSSL::SSL::SSLContext.new.tap do |context|
context.key = OpenSSL::PKey::EC.generate("prime256v1")
context.cert = sign(subject: "CN=#{name}", key: context.key, ca: ca, ca_key: ca_key) do |extensions|
[ extensions.create_extension("subjectAltName", "DNS:#{name}") ]
end
end
end
private
def ca_key
@ca_key ||= OpenSSL::PKey::EC.generate("prime256v1")
end
def ca
@ca ||= sign(subject: "CN=Push service test CA", key: ca_key, ca_key: ca_key) do |extensions|
[ extensions.create_extension("basicConstraints", "CA:TRUE", true), extensions.create_extension("keyUsage", "keyCertSign", true) ]
end.tap { OpenSSL::SSL::SSLContext::DEFAULT_CERT_STORE.add_cert(it) }
end
def sign(subject:, key:, ca_key:, ca: nil)
OpenSSL::X509::Certificate.new.tap do |certificate|
certificate.version, certificate.serial = 2, SecureRandom.random_number(1 << 64)
certificate.subject = OpenSSL::X509::Name.parse(subject)
certificate.issuer = ca ? ca.subject : certificate.subject
certificate.public_key = key
certificate.not_before, certificate.not_after = 1.minute.ago, 1.hour.from_now
extensions = OpenSSL::X509::ExtensionFactory.new(ca || certificate, certificate)
yield(extensions).each { certificate.add_extension(it) }
certificate.sign(ca_key, "SHA256")
end
end
end
private
def with_push_service(**options)
WebMock.disable!
server = Server.new(**options)
yield server
ensure
server&.stop
WebMock.enable!
end
# Built as WebPush::PersistentRequest builds it for a delivery.
def pinned_connection(server, ip = IP)
Net::HTTP.new(HOST, server.port, nil).tap do |http|
http.ipaddr = ip
http.use_ssl = true
end
end
def push_request(path = "/push")
Net::HTTP::Post.new(path).tap { it.body = "payload" }
end
end