Module: Cosmo::Stream

Defined in:
lib/cosmo/stream.rb,
lib/cosmo/stream/data.rb,
lib/cosmo/stream/message.rb,
lib/cosmo/stream/processor.rb,
lib/cosmo/stream/serializer.rb,
sig/cosmo/stream.rbs,
sig/cosmo/stream/data.rbs,
sig/cosmo/stream/message.rbs,
sig/cosmo/stream/processor.rbs,
sig/cosmo/stream/serializer.rbs

Defined Under Namespace

Modules: ClassMethods, Serializer Classes: Data, Message, Processor

Class Method Summary collapse

Instance Method Summary collapse

Class Method Details

.included(base) ⇒ void

This method returns an undefined value.

Parameters:

  • base (Class)


10
11
12
# File 'lib/cosmo/stream.rb', line 10

def self.included(base)
  base.extend(ClassMethods)
end

Instance Method Details

#loggerObject

Returns:

  • (Object)


65
66
67
# File 'lib/cosmo/stream.rb', line 65

def logger
  Logger.instance
end

#messageMessage?

Returns:



69
70
71
# File 'lib/cosmo/stream.rb', line 69

def message
  Thread.current[:cosmo_message]
end

#process(messages) ⇒ void Also known as: process_many, process_batch

This method returns an undefined value.

Parameters:



50
51
52
53
54
55
56
57
# File 'lib/cosmo/stream.rb', line 50

def process(messages)
  messages.each do |message|
    Thread.current[:cosmo_message] = message
    process_one
  ensure
    Thread.current[:cosmo_message] = nil
  end
end

#process_onevoid

This method returns an undefined value.



61
62
63
# File 'lib/cosmo/stream.rb', line 61

def process_one
  raise NotImplementedError, "#{self.class}#process_one must be implemented"
end

#publish(data, subject, **options) ⇒ Boolean

Parameters:

  • data (Object)
  • subject (::String)
  • options (Object)

Returns:

  • (Boolean)


73
74
75
# File 'lib/cosmo/stream.rb', line 73

def publish(data, subject, **options)
  self.class.publish(data, subject:, **options)
end