Class: Karafka::Pro::RecurringTasks::Serializer

Inherits:
Object
  • Object
show all
Defined in:
lib/karafka/pro/recurring_tasks/serializer.rb

Overview

Converts schedule command and log details into data we can dispatch to Kafka.

Constant Summary collapse

SCHEMA_VERSION =

Current recurring tasks related schema structure

'1.0'

Instance Method Summary collapse

Instance Method Details

#command(command_name, task_id) ⇒ String

Returns serialized and compressed command data.

Parameters:

  • command_name (String)

    command name

  • task_id (String)

    task id or ‘*’ to match all.

Returns:

  • (String)

    serialized and compressed command data



47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
# File 'lib/karafka/pro/recurring_tasks/serializer.rb', line 47

def command(command_name, task_id)
  data = {
    schema_version: SCHEMA_VERSION,
    schedule_version: Karafka::Pro::RecurringTasks.schedule.version,
    dispatched_at: Time.now.to_f,
    type: 'command',
    command: {
      name: command_name
    },
    task: {
      id: task_id
    }
  }

  compress(
    serialize(data)
  )
end

#log(event) ⇒ String

Returns serialized and compressed event log data.

Parameters:

  • event (Karafka::Core::Monitoring::Event)

    recurring task dispatch event

Returns:

  • (String)

    serialized and compressed event log data



68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
# File 'lib/karafka/pro/recurring_tasks/serializer.rb', line 68

def log(event)
  task = event[:task]

  data = {
    schema_version: SCHEMA_VERSION,
    schedule_version: Karafka::Pro::RecurringTasks.schedule.version,
    dispatched_at: Time.now.to_f,
    type: 'log',
    task: {
      id: task.id,
      time_taken: event.payload[:time] || -1,
      result: event.payload.key?(:error) ? 'failure' : 'success'
    }
  }

  compress(
    serialize(data)
  )
end

#schedule(schedule) ⇒ String

Serializes and compresses the schedule with all its tasks and their execution state

Parameters:

Returns:

  • (String)

    serialized and compressed current schedule data with its tasks and their current state.



18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
# File 'lib/karafka/pro/recurring_tasks/serializer.rb', line 18

def schedule(schedule)
  tasks = {}

  schedule.each do |task|
    tasks[task.id] = {
      id: task.id,
      cron: task.cron.original,
      previous_time: task.previous_time.to_i,
      next_time: task.next_time.to_i,
      enabled: task.enabled?
    }
  end

  data = {
    schema_version: SCHEMA_VERSION,
    schedule_version: schedule.version,
    dispatched_at: Time.now.to_f,
    type: 'schedule',
    tasks: tasks
  }

  compress(
    serialize(data)
  )
end