Class: AggregateStreams::Handle

Inherits:
Object
  • Object
show all
Includes:
Log::Dependency, Messaging::Handle, Messaging::StreamName, Settings::Setting
Defined in:
lib/aggregate_streams/handle.rb

Constant Summary collapse

TransformError =
Class.new(RuntimeError)

Instance Method Summary collapse

Instance Method Details

#assure_message_data(message_data) ⇒ Object



105
106
107
108
109
# File 'lib/aggregate_streams/handle.rb', line 105

def assure_message_data(message_data)
  unless message_data.instance_of?(MessageStore::MessageData::Write)
    raise TransformError, "Not an instance of MessageData::Write"
  end
end

#configure(session: nil, settings:) ⇒ Object



20
21
22
23
24
25
26
27
28
29
# File 'lib/aggregate_streams/handle.rb', line 20

def configure(session: nil, settings:)
  settings.set(self)

  writer_session = self.writer_session
  writer_session ||= session

  Store.configure(self, category: category, session: writer_session, snapshot_interval: snapshot_interval)

  MessageStore::Postgres::Write.configure(self, session: writer_session)
end

#handle(message_data) ⇒ Object



31
32
33
34
35
36
37
38
39
40
41
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
# File 'lib/aggregate_streams/handle.rb', line 31

def handle(message_data)
  logger.trace { "Handling message (Stream: #{message_data.stream_name}, Global Position: #{message_data.global_position})" }

  Retry.(MessageStore::ExpectedVersion::Error, millisecond_intervals: [0, 10, 100, 1000]) do
    stream_id = Messaging::StreamName.get_id(message_data.stream_name)

    aggregation, version = store.fetch(stream_id, include: :version)

    if aggregation.processed?(message_data)
      logger.info(tag: :ignored) { "Message already handled (Stream: #{message_data.stream_name}, Global Position: #{message_data.global_position})" }
      return
    end

    raw_input_data = Messaging::Message::Transform::MessageData.read(message_data)
     = Messaging::Message::.build(raw_input_data[:metadata])

     = ()

    write_message_data = MessageStore::MessageData::Write.new

    SetAttributes.(write_message_data, message_data, copy: [:type, :data])

    write_message_data. = 

    input_category = Messaging::StreamName.get_category(message_data.stream_name)
    write_message_data = transform(write_message_data, input_category)

    if write_message_data
      assure_message_data(write_message_data)
    else
      logger.info(tag: :ignored) { "Message ignored (Stream: #{message_data.stream_name}, Global Position: #{message_data.global_position})" }
      return
    end

    stream_name = stream_name(stream_id)
    write.(write_message_data, stream_name, expected_version: version)

    logger.info do
      message_type = message_data.type
      unless write_message_data.type == message_type
        message_type = "#{write_message_data.type} ← #{message_type}"
      end

      "Message copied (Message Type: #{message_type}, Stream: #{message_data.stream_name}, Global Position: #{message_data.global_position})"
    end
  end
end

#raw_metadata(metadata) ⇒ Object



79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
# File 'lib/aggregate_streams/handle.rb', line 79

def ()
   = Messaging::Message::.build

  .follow()

   = .to_h

  .delete(:local_properties)

  if [:properties].empty?
    .delete(:properties)
  end

  .delete_if { |_, v| v.nil? }

  
end

#transform(write_message_data, stream_name) ⇒ Object



97
98
99
100
101
102
103
# File 'lib/aggregate_streams/handle.rb', line 97

def transform(write_message_data, stream_name)
  if transform_action.nil?
    write_message_data
  else
    transform_action.(write_message_data, stream_name)
  end
end