|
2 | 2 |
|
3 | 3 | module SolidObjects |
4 | 4 | module WakeUpAdapters |
| 5 | + POSTGRESQL_FLOOR_MS = 2.9 |
| 6 | + REDIS_FLOOR_MS = 5.7 |
| 7 | + REDIS_URL_VARIABLE = "SOLID_OBJECTS_REDIS_URL" |
| 8 | + |
| 9 | + @pooled_warning_mutex = Thread::Mutex.new |
| 10 | + @pooled_warning_emitted = false |
| 11 | + |
5 | 12 | module_function |
6 | 13 |
|
7 | | - # Returns the best wake-up strategy for a connection: cross-process |
8 | | - # notifications where the database provides them, and the in-process |
9 | | - # default everywhere else. |
10 | | - # |
11 | | - # This is deliberately not the default. A notification adapter opens a |
12 | | - # connection per waiting thread outside the pool, and `LISTEN` does not |
13 | | - # survive a transaction-pooling proxy such as PgBouncer, so adopting it is |
14 | | - # a deployment decision rather than an upgrade side effect. |
15 | | - # |
| 14 | + NAMES = %i[automatic in_process postgresql redis].freeze |
| 15 | + |
16 | 16 | # @rbs (?untyped) -> untyped |
17 | 17 | def for(connection = Record.connection) |
18 | | - return Postgresql.new if DatabaseAdapter.family(connection) == :postgresql |
| 18 | + select(connection) |
| 19 | + end |
| 20 | + |
| 21 | + # @rbs (untyped) -> untyped |
| 22 | + def build(setting) |
| 23 | + return select if setting.nil? || setting == :automatic |
| 24 | + return named(setting) if setting.is_a?(Symbol) |
| 25 | + |
| 26 | + configured(setting) |
| 27 | + end |
| 28 | + |
| 29 | + # @rbs (Symbol) -> untyped |
| 30 | + def named(name) |
| 31 | + case name |
| 32 | + when :in_process then labelled(WakeUp.new, :in_process, false, nil, "in-process signalling was requested") |
| 33 | + when :postgresql then labelled(Postgresql.new, :postgresql_notify, true, POSTGRESQL_FLOOR_MS, "PostgreSQL LISTEN was requested") |
| 34 | + when :redis then labelled(Redis.new(url: redis_url), :redis, true, REDIS_FLOOR_MS, "Redis was requested") |
| 35 | + else |
| 36 | + raise ArgumentError, "unknown wake_up_adapter #{name.inspect}, expected one of #{NAMES.join(", ")} or an adapter" |
| 37 | + end |
| 38 | + end |
| 39 | + |
| 40 | + # @rbs (untyped) -> untyped |
| 41 | + def configured(adapter) |
| 42 | + return adapter unless adapter.respond_to?(:capability=) |
| 43 | + |
| 44 | + labelled(adapter, :configured, true, nil, "an adapter was configured, so selection did not run") |
| 45 | + end |
| 46 | + |
| 47 | + # @rbs (untyped, Symbol, bool, Numeric?, String) -> untyped |
| 48 | + def labelled(adapter, name, crosses_processes, floor, reason) |
| 49 | + adapter.capability = WakeUpCapability.new( |
| 50 | + adapter: name, |
| 51 | + crosses_processes:, |
| 52 | + measured_floor_ms: floor, |
| 53 | + reason: |
| 54 | + ) |
| 55 | + adapter |
| 56 | + end |
| 57 | + |
| 58 | + # @rbs (?untyped) -> untyped |
| 59 | + def select(connection = Record.connection) |
| 60 | + url = redis_url |
| 61 | + return redis_selection(url) if url |
| 62 | + |
| 63 | + family = DatabaseAdapter.family(connection) |
| 64 | + return postgresql_selection(connection) if family == :postgresql |
| 65 | + |
| 66 | + polling_selection(family) |
| 67 | + end |
| 68 | + |
| 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 |
| 75 | + rescue |
| 76 | + nil |
| 77 | + ensure |
| 78 | + restore_application_name(connection, previous) |
| 79 | + end |
| 80 | + |
| 81 | + # @rbs () -> void |
| 82 | + def reset_pooled_warning! |
| 83 | + @pooled_warning_mutex.synchronize { @pooled_warning_emitted = false } |
| 84 | + end |
| 85 | + |
| 86 | + # @rbs () -> String? |
| 87 | + def redis_url |
| 88 | + value = ENV[REDIS_URL_VARIABLE].to_s |
| 89 | + value.empty? ? nil : value |
| 90 | + end |
| 91 | + |
| 92 | + # @rbs (String) -> untyped |
| 93 | + def redis_selection(url) |
| 94 | + labelled( |
| 95 | + Redis.new(url:), :redis, true, REDIS_FLOOR_MS, |
| 96 | + "#{REDIS_URL_VARIABLE} is set, so Redis carries the signal between processes" |
| 97 | + ) |
| 98 | + end |
| 99 | + |
| 100 | + # @rbs (untyped) -> untyped |
| 101 | + def postgresql_selection(connection) |
| 102 | + survives = session_survives_transactions?(connection) |
| 103 | + return pooled_selection if survives == false |
| 104 | + |
| 105 | + 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" |
| 109 | + ) |
| 110 | + end |
| 111 | + |
| 112 | + # @rbs () -> untyped |
| 113 | + def pooled_selection |
| 114 | + warn_pooled_session_once |
| 115 | + 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" |
| 118 | + ) |
| 119 | + end |
| 120 | + |
| 121 | + # @rbs (Symbol?) -> untyped |
| 122 | + def polling_selection(family) |
| 123 | + polling_adapter( |
| 124 | + "#{family || "this database"} has no notification channel and " \ |
| 125 | + "#{REDIS_URL_VARIABLE} is not set" |
| 126 | + ) |
| 127 | + end |
| 128 | + |
| 129 | + # @rbs (String) -> untyped |
| 130 | + def polling_adapter(reason) |
| 131 | + labelled( |
| 132 | + WakeUp.new, :polling, false, |
| 133 | + SolidObjects.configuration.idle_polling_interval * 1_000, reason |
| 134 | + ) |
| 135 | + end |
| 136 | + |
| 137 | + # @rbs () -> void |
| 138 | + def warn_pooled_session_once |
| 139 | + @pooled_warning_mutex.synchronize do |
| 140 | + return if @pooled_warning_emitted |
| 141 | + |
| 142 | + SolidObjects.configuration.logger.warn( |
| 143 | + event: "solid_objects.wake_up.pooled_session", |
| 144 | + reason: "PostgreSQL notifications were not selected because the session " \ |
| 145 | + "does not outlive a transaction" |
| 146 | + ) |
| 147 | + @pooled_warning_emitted = true |
| 148 | + end |
| 149 | + end |
| 150 | + |
| 151 | + # @rbs (untyped, untyped) -> void |
| 152 | + def restore_application_name(connection, previous) |
| 153 | + return if previous.nil? |
19 | 154 |
|
20 | | - WakeUp.new |
| 155 | + connection.execute("SET application_name = #{connection.quote(previous)}") |
| 156 | + rescue |
| 157 | + nil |
21 | 158 | end |
22 | 159 | end |
23 | 160 | end |
0 commit comments