Class: Mxrb::Runtime::SQLiteSharedStore

Inherits:
Object
  • Object
show all
Includes:
SchedulerCoordinator, SessionStore
Defined in:
lib/mxrb/runtime/shared_store.rb

Overview

SQLite implementation safe for multiple Ruby processes sharing a file. BEGIN IMMEDIATE serializes claim decisions; primary keys make a schedule slot idempotent even after a process restart.

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(path, busy_timeout: 5_000) ⇒ SQLiteSharedStore

Returns a new instance of SQLiteSharedStore.

Raises:

  • (ArgumentError)


135
136
137
138
139
140
141
142
143
144
145
146
147
148
# File 'lib/mxrb/runtime/shared_store.rb', line 135

def initialize(path, busy_timeout: 5_000)
  raw_path = path.to_s
  raise ArgumentError, 'shared store path must not be empty' if raw_path.empty?

  expanded = raw_path == ':memory:' ? raw_path : File.expand_path(raw_path)

  FileUtils.mkdir_p(File.dirname(expanded)) unless expanded == ':memory:'
  @database = SQLite3::Database.new(expanded)
  @database.results_as_hash = true
  @database.busy_timeout = Integer(busy_timeout)
  @mutex = Mutex.new
  configure!
  migrate!
end

Instance Attribute Details

#databaseObject (readonly)

Returns the value of attribute database.



133
134
135
# File 'lib/mxrb/runtime/shared_store.rb', line 133

def database
  @database
end

Instance Method Details

#claim_scheduled_event(event:, slot:, owner:, now:, lease_until:, skip_overlap:) ⇒ Object



188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
# File 'lib/mxrb/runtime/shared_store.rb', line 188

def claim_scheduled_event(event:, slot:, owner:, now:, lease_until:, skip_overlap:)
  transaction do
    claim = database.get_first_row(
      <<~SQL, [event.to_s, slot.to_s]
        SELECT completed_at, lease_expires_at
        FROM mxrb_runtime_scheduled_claims
        WHERE event_key = ? AND slot = ?
      SQL
    )
    next false if claim && (claim['completed_at'] || claim.fetch('lease_expires_at').to_f > timestamp(now))

    if skip_overlap
      active = database.get_first_value(
        'SELECT 1 FROM mxrb_runtime_scheduler_leases WHERE event_key = ? AND expires_at > ?',
        [event.to_s, timestamp(now)]
      )
      next false if active
    end

    database.execute(
      <<~SQL, [event.to_s, slot.to_s, owner.to_s, timestamp(now), timestamp(lease_until)]
        INSERT INTO mxrb_runtime_scheduled_claims
          (event_key, slot, owner, claimed_at, lease_expires_at, completed_at)
        VALUES (?, ?, ?, ?, ?, NULL)
        ON CONFLICT(event_key, slot) DO UPDATE SET
          owner = excluded.owner,
          claimed_at = excluded.claimed_at,
          lease_expires_at = excluded.lease_expires_at,
          completed_at = NULL
      SQL
    )
    if skip_overlap
      database.execute(
        <<~SQL, [event.to_s, slot.to_s, owner.to_s, timestamp(lease_until)]
          INSERT INTO mxrb_runtime_scheduler_leases (event_key, slot, owner, expires_at)
          VALUES (?, ?, ?, ?)
          ON CONFLICT(event_key) DO UPDATE SET
            slot = excluded.slot,
            owner = excluded.owner,
            expires_at = excluded.expires_at
        SQL
      )
    end
    true
  end
end

#closeObject



276
277
278
# File 'lib/mxrb/runtime/shared_store.rb', line 276

def close
  @mutex.synchronize { database.close unless database.closed? }
end

#complete_scheduled_event(event:, slot:, owner:, now: Time.now.utc) ⇒ Object



235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
# File 'lib/mxrb/runtime/shared_store.rb', line 235

def complete_scheduled_event(event:, slot:, owner:, now: Time.now.utc)
  transaction do
    database.execute(
      <<~SQL, [timestamp(now), event.to_s, slot.to_s, owner.to_s]
        UPDATE mxrb_runtime_scheduled_claims
        SET completed_at = ?
        WHERE event_key = ? AND slot = ? AND owner = ?
      SQL
    )
    database.execute(
      <<~SQL, [event.to_s, slot.to_s, owner.to_s]
        DELETE FROM mxrb_runtime_scheduler_leases
        WHERE event_key = ? AND slot = ? AND owner = ?
      SQL
    )
  end
  true
end

#delete_session(token) ⇒ Object



181
182
183
184
185
186
# File 'lib/mxrb/runtime/shared_store.rb', line 181

def delete_session(token)
  synchronize do
    database.execute('DELETE FROM mxrb_runtime_sessions WHERE token = ?', [session_key(token)])
    database.changes.positive?
  end
end

#read_session(token, now: Time.now.utc) ⇒ Object



165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
# File 'lib/mxrb/runtime/shared_store.rb', line 165

def read_session(token, now: Time.now.utc)
  row = synchronize do
    database.execute('DELETE FROM mxrb_runtime_sessions WHERE expires_at <= ?', [timestamp(now)])
    database.get_first_row(
      'SELECT token, identity_json, expires_at FROM mxrb_runtime_sessions WHERE token = ?',
      [session_key(token)]
    )
  end
  return unless row

  SessionRecord.new(
    token.to_s, JSON.parse(row.fetch('identity_json')),
    Time.at(row.fetch('expires_at').to_f).utc
  )
end

#renew_scheduled_event(event:, slot:, owner:, lease_until:) ⇒ Object



254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
# File 'lib/mxrb/runtime/shared_store.rb', line 254

def renew_scheduled_event(event:, slot:, owner:, lease_until:)
  transaction do
    database.execute(
      <<~SQL, [timestamp(lease_until), event.to_s, slot.to_s, owner.to_s]
        UPDATE mxrb_runtime_scheduled_claims
        SET lease_expires_at = ?
        WHERE event_key = ? AND slot = ? AND owner = ? AND completed_at IS NULL
      SQL
    )
    next false unless database.changes.positive?

    database.execute(
      <<~SQL, [timestamp(lease_until), event.to_s, slot.to_s, owner.to_s]
        UPDATE mxrb_runtime_scheduler_leases
        SET expires_at = ?
        WHERE event_key = ? AND slot = ? AND owner = ?
      SQL
    )
    true
  end
end

#write_session(token:, identity:, expires_at:) ⇒ Object



150
151
152
153
154
155
156
157
158
159
160
161
162
163
# File 'lib/mxrb/runtime/shared_store.rb', line 150

def write_session(token:, identity:, expires_at:)
  synchronize do
    database.execute(
      <<~SQL, [session_key(token), JSON.generate(identity), timestamp(expires_at)]
        INSERT INTO mxrb_runtime_sessions (token, identity_json, expires_at)
        VALUES (?, ?, ?)
        ON CONFLICT(token) DO UPDATE SET
          identity_json = excluded.identity_json,
          expires_at = excluded.expires_at
      SQL
    )
  end
  token
end