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



54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
# File 'lib/karafka/pro/recurring_tasks/serializer.rb', line 54

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



75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
# File 'lib/karafka/pro/recurring_tasks/serializer.rb', line 75

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

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

Parameters:

Returns:

  • (String)

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



25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
# File 'lib/karafka/pro/recurring_tasks/serializer.rb', line 25

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