Class: Kaal::Dispatch::RedisEngine

Inherits:
Registry
  • Object
show all
Defined in:
lib/kaal/dispatch/redis_engine.rb,
sig/kaal/dispatch/redis_engine.rbs

Overview

Redis-backed dispatch registry.

Stores dispatch records in Redis as JSON-serialized values. Keys are automatically expired based on TTL to prevent unbounded growth.

Examples:

Usage

redis = Redis.new(url: ENV['REDIS_URL'])
registry = Kaal::Dispatch::RedisEngine.new(redis, namespace: 'myapp')
registry.log_dispatch('daily_report', Time.now.utc, 'node-1')

Constant Summary collapse

DEFAULT_TTL =

Default TTL for dispatch records (7 days in seconds)

Returns:

  • (Kaal::rbs_any)
7 * 24 * 60 * 60

Instance Method Summary collapse

Methods inherited from Registry

#dispatched?

Constructor Details

#initialize(redis, namespace: 'kaal', ttl: DEFAULT_TTL) ⇒ RedisEngine

Initialize a new Redis-backed registry.

Parameters:

  • redis (Redis) —

    Redis client instance

  • namespace (String) (defaults to: 'kaal') —

    namespace prefix for Redis keys

  • ttl (Integer) (defaults to: DEFAULT_TTL) —

    TTL in seconds for dispatch records

  • namespace: (::String) (defaults to: 'kaal')
  • ttl: (Kaal::rbs_any) (defaults to: DEFAULT_TTL)


32
33
34
35
36
37
# File 'lib/kaal/dispatch/redis_engine.rb', line 32

def initialize(redis, namespace: 'kaal', ttl: DEFAULT_TTL)
  super()
  @redis = redis
  @namespace = namespace
  @ttl = ttl
end

Instance Method Details

#build_redis_key(key, fire_time) ⇒ String

Build a Redis key from job key and fire time.

Format: "namespace:cron_dispatch:key:fire_time"

Parameters:

  • key (String) —

    the cron job key

  • fire_time (Time) —

    the fire time

Returns:

  • (String) —

    Redis key



86
87
88
# File 'lib/kaal/dispatch/redis_engine.rb', line 86

def build_redis_key(key, fire_time)
  "#{@namespace}:cron_dispatch:#{key}:#{fire_time.to_i}"
end

#convert_timestamps(record) ⇒ Hash

Convert Unix timestamps in record to Time objects.

Parameters:

  • record (Hash) —

    the dispatch record with Unix timestamps

Returns:

  • (Hash) —

    the record with Time objects



95
96
97
98
99
# File 'lib/kaal/dispatch/redis_engine.rb', line 95

def convert_timestamps(record)
  record[:fire_time] = Time.at(record[:fire_time])
  record[:dispatched_at] = Time.at(record[:dispatched_at])
  record
end

#find_dispatch(key, fire_time) ⇒ Hash?

Find a dispatch record for a specific job and fire time.

Parameters:

  • key (String) —

    the cron job key

  • fire_time (Time) —

    when the job was scheduled to fire

Returns:

  • (Hash, nil) —

    dispatch record or nil if not found



67
68
69
70
71
72
73
74
# File 'lib/kaal/dispatch/redis_engine.rb', line 67

def find_dispatch(key, fire_time)
  redis_key = build_redis_key(key, fire_time)
  value = @redis.get(redis_key)
  return nil unless value

  record = JSON.parse(value, symbolize_names: true)
  convert_timestamps(record)
end

#log_dispatch(key, fire_time, node_id, status = 'dispatched') ⇒ Hash

Log a dispatch attempt in Redis.

Parameters:

  • key (String) —

    the cron job key

  • fire_time (Time) —

    when the job was scheduled to fire

  • node_id (String) —

    identifier for the dispatching node

  • status (String) (defaults to: 'dispatched') —

    dispatch status ('dispatched', 'failed', etc.)

Returns:

  • (Hash) —

    the stored dispatch record



47
48
49
50
51
52
53
54
55
56
57
58
59
# File 'lib/kaal/dispatch/redis_engine.rb', line 47

def log_dispatch(key, fire_time, node_id, status = 'dispatched')
  redis_key = build_redis_key(key, fire_time)
  record = {
    key: key,
    fire_time: fire_time.to_i,
    dispatched_at: Time.now.utc.to_i,
    node_id: node_id,
    status: status
  }

  @redis.setex(redis_key, @ttl, JSON.generate(record))
  record
end