Class: Mammoth::DeadLetterStore
- Inherits:
-
Object
- Object
- Mammoth::DeadLetterStore
- Defined in:
- lib/mammoth/dead_letter_store.rb
Overview
Persists failed deliveries in Mammoth's SQLite dead letter queue.
Instance Attribute Summary collapse
-
#sqlite_store ⇒ Object
readonly
Returns the value of attribute sqlite_store.
Instance Method Summary collapse
-
#count(status: nil, destination: nil) ⇒ Integer
Count dead letters by status.
-
#counts_by_destination(status: nil) ⇒ Array<Hash>
Count dead letters grouped by destination.
-
#fetch(id) ⇒ Hash?
Fetch one dead letter by id.
-
#ignore(id) ⇒ void
Mark a dead letter as ignored.
-
#initialize(sqlite_store) ⇒ DeadLetterStore
constructor
A new instance of DeadLetterStore.
-
#pending(limit: 100, destination: nil, failed_after: nil, failed_before: nil) ⇒ Array<Hash>
Fetch pending dead letters.
-
#resolve(id) ⇒ void
Mark a dead letter as resolved.
-
#rows(status: nil, limit: 100, destination: nil, failed_after: nil, failed_before: nil) ⇒ Array<Hash>
Fetch dead letters, optionally filtered by status.
-
#write(event:, destination_name:, error: nil, retry_count: 0, serializer: EventSerializer) ⇒ Integer
Store a failed delivery.
-
#write_payload(payload:, destination_name:, error: nil, retry_count: 0) ⇒ Integer
Store the exact payload prepared for one destination.
Constructor Details
#initialize(sqlite_store) ⇒ DeadLetterStore
Returns a new instance of DeadLetterStore.
12 13 14 |
# File 'lib/mammoth/dead_letter_store.rb', line 12 def initialize(sqlite_store) @sqlite_store = sqlite_store end |
Instance Attribute Details
#sqlite_store ⇒ Object (readonly)
Returns the value of attribute sqlite_store.
9 10 11 |
# File 'lib/mammoth/dead_letter_store.rb', line 9 def sqlite_store @sqlite_store end |
Instance Method Details
#count(status: nil, destination: nil) ⇒ Integer
Count dead letters by status.
116 117 118 119 |
# File 'lib/mammoth/dead_letter_store.rb', line 116 def count(status: nil, destination: nil) where, values = row_filters(status:, destination:) database.get_first_value("SELECT COUNT(*) FROM dead_letters#{where}", values) end |
#counts_by_destination(status: nil) ⇒ Array<Hash>
Count dead letters grouped by destination.
125 126 127 128 129 130 131 |
# File 'lib/mammoth/dead_letter_store.rb', line 125 def counts_by_destination(status: nil) where, values = row_filters(status:) database.execute( "SELECT destination_name, COUNT(*) AS count FROM dead_letters#{where} GROUP BY destination_name", values ) end |
#fetch(id) ⇒ Hash?
Fetch one dead letter by id.
108 109 110 |
# File 'lib/mammoth/dead_letter_store.rb', line 108 def fetch(id) database.get_first_row("SELECT * FROM dead_letters WHERE id = ?", [id]) end |
#ignore(id) ⇒ void
This method returns an undefined value.
Mark a dead letter as ignored.
145 146 147 |
# File 'lib/mammoth/dead_letter_store.rb', line 145 def ignore(id) update_status(id, "ignored") end |
#pending(limit: 100, destination: nil, failed_after: nil, failed_before: nil) ⇒ Array<Hash>
Fetch pending dead letters.
86 87 88 89 |
# File 'lib/mammoth/dead_letter_store.rb', line 86 def pending(limit: 100, destination: nil, failed_after: nil, failed_before: nil) rows(status: "pending", limit: limit, destination: destination, failed_after: failed_after, failed_before: failed_before) end |
#resolve(id) ⇒ void
This method returns an undefined value.
Mark a dead letter as resolved.
137 138 139 |
# File 'lib/mammoth/dead_letter_store.rb', line 137 def resolve(id) update_status(id, "resolved") end |
#rows(status: nil, limit: 100, destination: nil, failed_after: nil, failed_before: nil) ⇒ Array<Hash>
Fetch dead letters, optionally filtered by status.
96 97 98 99 100 101 102 |
# File 'lib/mammoth/dead_letter_store.rb', line 96 def rows(status: nil, limit: 100, destination: nil, failed_after: nil, failed_before: nil) where, values = row_filters(status:, destination:, failed_after:, failed_before:) database.execute( "SELECT * FROM dead_letters#{where} ORDER BY failed_at ASC LIMIT ?", values + [limit] ) end |
#write(event:, destination_name:, error: nil, retry_count: 0, serializer: EventSerializer) ⇒ Integer
Store a failed delivery.
23 24 25 26 27 28 29 30 |
# File 'lib/mammoth/dead_letter_store.rb', line 23 def write(event:, destination_name:, error: nil, retry_count: 0, serializer: EventSerializer) write_payload( payload: serializer.call(event), destination_name: destination_name, error: error, retry_count: retry_count ) end |
#write_payload(payload:, destination_name:, error: nil, retry_count: 0) ⇒ Integer
Store the exact payload prepared for one destination.
Persisting the prepared payload prevents retries and replay from restoring fields removed by a destination payload policy.
42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 |
# File 'lib/mammoth/dead_letter_store.rb', line 42 def write_payload(payload:, destination_name:, error: nil, retry_count: 0) # rubocop:disable Metrics/MethodLength now = Time.now.utc.iso8601 database.execute( <<~SQL, INSERT INTO dead_letters( event_id, source_name, destination_name, operation, namespace, entity, source_position, payload_json, error_class, error_message, retry_count, status, failed_at, updated_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?, ?) SQL [ payload.fetch("event_id"), payload.fetch("source"), destination_name, payload["operation"] || payload["type"], payload["namespace"], payload["entity"], payload["source_position"], JSON.generate(payload), error&.class&.name, error&., retry_count, now, now ] ) database.last_insert_row_id end |