Skip to content

Commit aec05bc

Browse files
cardmagicclaude
andcommitted
fix: prove the wake-up choice before claiming it
Two claims in selection were not supported by what it measured. `configured` overwrote the capability that an adapter reports about itself. A configured `SolidObjects::WakeUp` signals in one process only and says so, but selection relabelled it `crosses_processes: true`, so the doctor reported PASS and the process registry dropped the warning that this exact topology needs. Selection now keeps the capability that an adapter provides and labels only an adapter that provides none. The PostgreSQL probe set `application_name` and read it back on the next statement. A transaction pooler can hand the same backend to two consecutive autocommit statements, so the read-back succeeds while `LISTEN` still has no session affinity, and selection memoised notification support that never fires. The probe now listens, sends one `NOTIFY` from a second connection, and waits up to two seconds for it to arrive. Nothing but a delivered notification selects notifications, and a probe that cannot run falls back to polling rather than assume. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent 09e7064 commit aec05bc

6 files changed

Lines changed: 155 additions & 47 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -12,10 +12,14 @@
1212
cross-process wake-up, a connection per waiting thread outside the pool, and
1313
one `NOTIFY` per enqueue after the commit. Set
1414
`config.wake_up_adapter = :in_process` to keep polling.
15-
- Probe the PostgreSQL session before selecting notifications, because `LISTEN`
16-
does not survive a transaction pooler such as PgBouncer. A session that does
17-
not outlive a statement falls back to polling and warns once. A probe that
18-
cannot run is not treated as a pooler.
15+
- Prove the PostgreSQL notification path before selecting it, because `LISTEN`
16+
does not survive a transaction pooler such as PgBouncer. Selection listens,
17+
sends one `NOTIFY` from a second connection, and waits up to two seconds for
18+
it to arrive. A probe that does not deliver falls back to polling and warns
19+
once.
20+
- Keep the capability that a configured adapter reports about itself. A
21+
configured `SolidObjects::WakeUp` now reports `:in_process` and warns, rather
22+
than claim that it crosses processes.
1923
- Report the resolved choice. `SolidObjects.wake_up.capability` names the
2024
adapter, whether it crosses processes, its measured floor, and why it was
2125
chosen. The doctor reports it, and the polling-only warning now fires on what

‎docs/roadmap.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -131,7 +131,9 @@
131131
adapter opens a connection per waiting thread outside the pool and adds one
132132
`NOTIFY` per enqueue, and Redis is not a dependency of this gem.
133133
`LISTEN` does not survive a transaction-pooling proxy such as PgBouncer, so
134-
the session is probed and a pooled one falls back to polling and warns once.
134+
selection listens, sends one `NOTIFY` from a second connection, and waits for
135+
it to arrive. A probe that does not deliver falls back to polling and warns
136+
once.
135137
`SolidObjects.wake_up.capability` reports the adapter, whether it crosses
136138
processes, its floor, and why, and the doctor shows the same record.
137139
MySQL still polls. It has no notification channel, and no MySQL notifier has

‎lib/solid_objects/wake_up_adapters.rb‎

Lines changed: 37 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ module WakeUpAdapters
55
POSTGRESQL_FLOOR_MS = 2.9
66
REDIS_FLOOR_MS = 5.7
77
REDIS_URL_VARIABLE = "SOLID_OBJECTS_REDIS_URL"
8+
PROBE_TIMEOUT_SECONDS = 2.0
89

910
@pooled_warning_mutex = Thread::Mutex.new
1011
@pooled_warning_emitted = false
@@ -40,6 +41,7 @@ def named(name)
4041
# @rbs (untyped) -> untyped
4142
def configured(adapter)
4243
return adapter unless adapter.respond_to?(:capability=)
44+
return adapter if adapter.respond_to?(:default_capability)
4345

4446
labelled(adapter, :configured, true, nil, "an adapter was configured, so selection did not run")
4547
end
@@ -61,21 +63,37 @@ def select(connection = Record.connection)
6163
return redis_selection(url) if url
6264

6365
family = DatabaseAdapter.family(connection)
64-
return postgresql_selection(connection) if family == :postgresql
66+
return postgresql_selection if family == :postgresql
6567

6668
polling_selection(family)
6769
end
6870

69-
# @rbs (untyped) -> bool?
70-
def session_survives_transactions?(connection)
71-
previous = connection.select_value("SELECT current_setting('application_name')")
72-
token = SecureRandom.hex(8)
73-
connection.execute("SET application_name = #{connection.quote(token)}")
74-
connection.select_value("SELECT current_setting('application_name')") == token
71+
# @rbs (untyped) -> bool
72+
def notifications_deliver?(adapter)
73+
return false unless adapter.listen
74+
return false unless notify_probe_channel(adapter.channel)
75+
76+
adapter.wait(timeout: PROBE_TIMEOUT_SECONDS)
7577
rescue
76-
nil
78+
false
79+
ensure
80+
adapter.stop
81+
end
82+
83+
# @rbs (String) -> bool
84+
def notify_probe_channel(channel)
85+
connection = Record.connection_pool.send(:new_connection)
86+
connection.execute("NOTIFY #{connection.quote_table_name(channel)}")
87+
true
7788
ensure
78-
restore_application_name(connection, previous)
89+
disconnect_probe(connection)
90+
end
91+
92+
# @rbs (untyped) -> void
93+
def disconnect_probe(connection)
94+
connection&.disconnect!
95+
rescue
96+
nil
7997
end
8098

8199
# @rbs () -> void
@@ -97,24 +115,23 @@ def redis_selection(url)
97115
)
98116
end
99117

100-
# @rbs (untyped) -> untyped
101-
def postgresql_selection(connection)
102-
survives = session_survives_transactions?(connection)
103-
return pooled_selection if survives == false
118+
# @rbs () -> untyped
119+
def postgresql_selection
120+
adapter = Postgresql.new
121+
return pooled_selection unless notifications_deliver?(adapter)
104122

105123
labelled(
106-
Postgresql.new, :postgresql_notify, true, POSTGRESQL_FLOOR_MS,
107-
survives ? "PostgreSQL LISTEN is available and the session outlives a transaction"
108-
: "PostgreSQL LISTEN was selected without a session probe"
124+
adapter, :postgresql_notify, true, POSTGRESQL_FLOOR_MS,
125+
"a probe notification arrived, so PostgreSQL LISTEN carries the signal between processes"
109126
)
110127
end
111128

112129
# @rbs () -> untyped
113130
def pooled_selection
114131
warn_pooled_session_once
115132
polling_adapter(
116-
"the PostgreSQL session does not outlive a transaction, which a transaction " \
117-
"pooler such as PgBouncer causes, so LISTEN would never fire"
133+
"a probe notification did not arrive, so LISTEN cannot carry the signal " \
134+
"between processes; a transaction pooler such as PgBouncer is the usual cause"
118135
)
119136
end
120137

@@ -141,20 +158,11 @@ def warn_pooled_session_once
141158

142159
SolidObjects.configuration.logger.warn(
143160
event: "solid_objects.wake_up.pooled_session",
144-
reason: "PostgreSQL notifications were not selected because the session " \
145-
"does not outlive a transaction"
161+
reason: "PostgreSQL notifications were not selected because a probe " \
162+
"notification did not arrive"
146163
)
147164
@pooled_warning_emitted = true
148165
end
149166
end
150-
151-
# @rbs (untyped, untyped) -> void
152-
def restore_application_name(connection, previous)
153-
return if previous.nil?
154-
155-
connection.execute("SET application_name = #{connection.quote(previous)}")
156-
rescue
157-
nil
158-
end
159167
end
160168
end

‎sig/generated/lib/solid_objects/wake_up_adapters.rbs‎

Lines changed: 12 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,8 @@ module SolidObjects
88

99
REDIS_URL_VARIABLE: ::String
1010

11+
PROBE_TIMEOUT_SECONDS: ::Float
12+
1113
NAMES: untyped
1214

1315
# @rbs (?untyped) -> untyped
@@ -28,8 +30,14 @@ module SolidObjects
2830
# @rbs (?untyped) -> untyped
2931
def self?.select: (?untyped) -> untyped
3032

31-
# @rbs (untyped) -> bool?
32-
def self?.session_survives_transactions?: (untyped) -> bool?
33+
# @rbs (untyped) -> bool
34+
def self?.notifications_deliver?: (untyped) -> bool
35+
36+
# @rbs (String) -> bool
37+
def self?.notify_probe_channel: (String) -> bool
38+
39+
# @rbs (untyped) -> void
40+
def self?.disconnect_probe: (untyped) -> void
3341

3442
# @rbs () -> void
3543
def self?.reset_pooled_warning!: () -> void
@@ -40,8 +48,8 @@ module SolidObjects
4048
# @rbs (String) -> untyped
4149
def self?.redis_selection: (String) -> untyped
4250

43-
# @rbs (untyped) -> untyped
44-
def self?.postgresql_selection: (untyped) -> untyped
51+
# @rbs () -> untyped
52+
def self?.postgresql_selection: () -> untyped
4553

4654
# @rbs () -> untyped
4755
def self?.pooled_selection: () -> untyped
@@ -54,8 +62,5 @@ module SolidObjects
5462

5563
# @rbs () -> void
5664
def self?.warn_pooled_session_once: () -> void
57-
58-
# @rbs (untyped, untyped) -> void
59-
def self?.restore_application_name: (untyped, untyped) -> void
6065
end
6166
end

‎test/integration/wake_up_selection_test.rb‎

Lines changed: 62 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,14 +4,26 @@
44
require "solid_objects/doctor"
55

66
class WakeUpSelectionTest < ActiveSupport::TestCase
7+
class CustomAdapter
8+
attr_accessor :capability
9+
10+
def signal = true
11+
12+
def watch = self
13+
14+
def wait(timeout:) = false
15+
end
16+
717
setup do
818
SolidObjects.reset_wake_up!
919
@redis_url = ENV.delete("SOLID_OBJECTS_REDIS_URL")
20+
@configured_adapter = SolidObjects.configuration.wake_up_adapter
1021
end
1122

1223
teardown do
1324
ENV.delete("SOLID_OBJECTS_REDIS_URL")
1425
ENV["SOLID_OBJECTS_REDIS_URL"] = @redis_url if @redis_url
26+
SolidObjects.configuration.wake_up_adapter = @configured_adapter
1527
SolidObjects.reset_wake_up!
1628
end
1729

@@ -20,7 +32,32 @@ class WakeUpSelectionTest < ActiveSupport::TestCase
2032
SolidObjects.configuration.wake_up_adapter = explicit
2133

2234
assert_same explicit, SolidObjects.wake_up
23-
assert_equal :configured, SolidObjects.wake_up.capability.adapter
35+
end
36+
37+
test "a configured adapter keeps the capability it reports about itself" do
38+
SolidObjects.configuration.wake_up_adapter = SolidObjects::WakeUp.new
39+
40+
capability = SolidObjects.wake_up.capability
41+
assert_equal :in_process, capability.adapter
42+
assert_not capability.crosses_processes
43+
assert_match(/another process/i, capability.reason)
44+
end
45+
46+
test "the doctor warns about a configured in-process adapter" do
47+
SolidObjects.configuration.authorize_administration = ->(**) { true }
48+
SolidObjects.configuration.wake_up_adapter = SolidObjects::WakeUp.new
49+
50+
check = SolidObjects::Doctor.new.call.check(:wake_up)
51+
assert_equal :warn, check.status
52+
assert_match(/cannot wake another/i, check.message)
53+
end
54+
55+
test "a configured adapter that reports no capability is recorded as configured" do
56+
SolidObjects.configuration.wake_up_adapter = CustomAdapter.new
57+
58+
capability = SolidObjects.wake_up.capability
59+
assert_equal :configured, capability.adapter
60+
assert capability.crosses_processes
2461
end
2562

2663
test "in_process opts out of selection" do
@@ -56,13 +93,26 @@ class WakeUpSelectionTest < ActiveSupport::TestCase
5693
assert_match(/redis/i, capability.reason)
5794
end
5895

59-
test "postgresql selects notifications when the session survives transactions" do
96+
test "postgresql selects notifications when a probe notification arrives" do
6097
skip unless database_family == :postgresql
6198

6299
capability = SolidObjects.wake_up.capability
63100
assert_equal :postgresql_notify, capability.adapter
64101
assert capability.crosses_processes
65102
assert_operator capability.measured_floor_ms, :<, 100
103+
assert_match(/probe notification arrived/i, capability.reason)
104+
end
105+
106+
test "postgresql polls when a probe notification does not arrive" do
107+
skip unless database_family == :postgresql
108+
109+
with_undelivered_notifications do
110+
capability = SolidObjects.wake_up.capability
111+
112+
assert_equal :polling, capability.adapter
113+
assert_not capability.crosses_processes
114+
assert_match(/pooler/i, capability.reason)
115+
end
66116
end
67117

68118
test "a database without a channel polls and reports its floor" do
@@ -134,7 +184,16 @@ def with_module_method(name, replacement)
134184
end
135185

136186
def with_pooled_session(&block)
137-
with_module_method(:session_survives_transactions?, ->(_connection) { false }, &block)
187+
with_undelivered_notifications(&block)
188+
end
189+
190+
def with_undelivered_notifications
191+
adapter_class = SolidObjects::WakeUpAdapters::Postgresql
192+
original = adapter_class.instance_method(:wait)
193+
adapter_class.define_method(:wait) { |timeout:| false }
194+
yield
195+
ensure
196+
adapter_class.define_method(:wait, original)
138197
end
139198

140199
def with_unreachable_database(&block)

‎test/unit/wake_up_adapters_test.rb‎

Lines changed: 33 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,10 +5,21 @@
55
class WakeUpAdaptersTest < ActiveSupport::TestCase
66
Connection = Struct.new(:adapter_name)
77

8-
test "selects PostgreSQL notifications for a PostgreSQL connection" do
9-
adapter = SolidObjects::WakeUpAdapters.for(Connection.new("PostgreSQL"))
8+
test "selects PostgreSQL notifications when a probe notification arrives" do
9+
with_delivered_notifications do
10+
adapter = SolidObjects::WakeUpAdapters.for(Connection.new("PostgreSQL"))
1011

11-
assert_instance_of SolidObjects::WakeUpAdapters::Postgresql, adapter
12+
assert_instance_of SolidObjects::WakeUpAdapters::Postgresql, adapter
13+
end
14+
end
15+
16+
test "polls when a PostgreSQL probe notification does not arrive" do
17+
with_undelivered_notifications do
18+
adapter = SolidObjects::WakeUpAdapters.for(Connection.new("PostgreSQL"))
19+
20+
assert_instance_of SolidObjects::WakeUp, adapter
21+
assert_equal :polling, adapter.capability.adapter
22+
end
1223
end
1324

1425
test "falls back to the in-process wake-up for MySQL" do
@@ -44,4 +55,23 @@ class WakeUpAdaptersTest < ActiveSupport::TestCase
4455

4556
assert_equal true, watch.wait(timeout: 1.0)
4657
end
58+
59+
private
60+
61+
def with_delivered_notifications(&block)
62+
with_probe(->(_adapter) { true }, &block)
63+
end
64+
65+
def with_undelivered_notifications(&block)
66+
with_probe(->(_adapter) { false }, &block)
67+
end
68+
69+
def with_probe(replacement)
70+
adapters = SolidObjects::WakeUpAdapters
71+
original = adapters.method(:notifications_deliver?)
72+
adapters.define_singleton_method(:notifications_deliver?, replacement)
73+
yield
74+
ensure
75+
adapters.define_singleton_method(:notifications_deliver?, original)
76+
end
4777
end

0 commit comments

Comments
 (0)