Class: Wampproto::Dealer

Inherits:
Object
  • Object
show all
Defined in:
lib/wampproto/dealer.rb

Overview

Wamprpoto Dealer handler

Defined Under Namespace

Classes: PendingInvocation

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(id_gen = IdGenerator.new) ⇒ Dealer

Returns a new instance of Dealer.



11
12
13
14
15
16
17
# File 'lib/wampproto/dealer.rb', line 11

def initialize(id_gen = IdGenerator.new)
  @registrations_by_session = {}
  @registrations_by_procedure = Hash.new { |h, k| h[k] = {} }
  @pending_calls = {}
  @id_gen = id_gen
  @sessions = {}
end

Instance Attribute Details

#id_gen ⇒ Object (readonly)

Returns the value of attribute id_gen.



8
9
10
# File 'lib/wampproto/dealer.rb', line 8

def id_gen
  @id_gen
end

#pending_calls ⇒ Object (readonly)

Returns the value of attribute pending_calls.



8
9
10
# File 'lib/wampproto/dealer.rb', line 8

def pending_calls
  @pending_calls
end

#registrations_by_procedure ⇒ Object (readonly)

Returns the value of attribute registrations_by_procedure.



8
9
10
# File 'lib/wampproto/dealer.rb', line 8

def registrations_by_procedure
  @registrations_by_procedure
end

#registrations_by_session ⇒ Object (readonly)

Returns the value of attribute registrations_by_session.



8
9
10
# File 'lib/wampproto/dealer.rb', line 8

def registrations_by_session
  @registrations_by_session
end

#sessions ⇒ Object (readonly)

Returns the value of attribute sessions.



8
9
10
# File 'lib/wampproto/dealer.rb', line 8

def sessions
  @sessions
end

Instance Method Details

#add_session(details) ⇒ Object

Raises:

  • (KeyError)


19
20
21
22
23
24
25
26
27
# File 'lib/wampproto/dealer.rb', line 19

def add_session(details)
  session_id = details.session_id

  error_message = "cannot add session twice"
  raise KeyError, error_message if registrations_by_session.include?(session_id)

  registrations_by_session[session_id] = {}
  sessions[session_id] = details
end

#handle_call(session_id, message) ⇒ Object

rubocop:disable Metrics/MethodLength, Metrics/AbcSize



58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
# File 'lib/wampproto/dealer.rb', line 58

def handle_call(session_id, message) # rubocop:disable Metrics/MethodLength, Metrics/AbcSize
  registrations = registrations_by_procedure.fetch(message.procedure, {})
  if registrations.empty?
    error = Message::Error.new(Message::Type::CALL, message.request_id, {}, "wamp.error.no_such_procedure")
    return MessageWithRecipient.new(error, session_id)
  end

  registration_id, callee_id = registrations.first

  request_id = id_gen.next

  details = invocation_details_for(session_id, message)

  pending_calls[[callee_id, request_id]] = PendingInvocation.new(
    session_id, callee_id, message.request_id, request_id, details[:receive_progress]
  )

  invocation = Message::Invocation.new(
    request_id,
    registration_id,
    details,
    *message.args,
    **message.kwargs
  )

  MessageWithRecipient.new(invocation, callee_id)
end

#handle_register(session_id, message) ⇒ Object

Raises:



123
124
125
126
127
128
129
130
131
132
133
# File 'lib/wampproto/dealer.rb', line 123

def handle_register(session_id, message)
  error_message = "cannot register, session #{session_id} doesn't exist"
  raise ValueError, error_message unless registrations_by_session.include?(session_id)

  registration_id = id_gen.next
  registrations_by_procedure[message.procedure][registration_id] = session_id
  registrations_by_session[session_id][registration_id] = message.procedure

  registered = Message::Registered.new(message.request_id, registration_id)
  MessageWithRecipient.new(registered, session_id)
end

#handle_unregister(session_id, message) ⇒ Object

rubocop:disable Metrics/AbcSize, Metrics/MethodLength)

Raises:



135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
# File 'lib/wampproto/dealer.rb', line 135

def handle_unregister(session_id, message) # rubocop:disable Metrics/AbcSize, Metrics/MethodLength)
  error_message = "cannot unregister, session #{session_id} doesn't exist"
  raise ValueError, error_message unless registrations_by_session.include?(session_id)

  registrations = registrations_by_session.fetch(session_id)

  unless registrations.include?(message.registration_id)
    error = Message::Error.new(Message::Type::UNREGISTER, message.request_id, {}, "wamp.error.no_such_registration")
    return MessageWithRecipient.new(error, session_id)
  end

  procedure = registrations.fetch(message.registration_id)

  registrations_by_procedure[procedure].delete(message.registration_id)
  registrations_by_session[session_id].delete(message.registration_id)

  unregistered = Message::Unregistered.new(message.request_id)
  MessageWithRecipient.new(unregistered, session_id)
end

#handle_yield(session_id, message) ⇒ Object

rubocop:disable Metrics/AbcSize

Raises:



99
100
101
102
103
104
105
106
107
108
109
110
# File 'lib/wampproto/dealer.rb', line 99

def handle_yield(session_id, message) # rubocop:disable Metrics/AbcSize
  pending_invocation = pending_calls[[session_id, message.request_id]]
  error_message = "no pending calls for session #{session_id}"
  raise ValueError, error_message if pending_invocation.nil?

  caller_id   = pending_invocation.caller_id
  request_id  = pending_invocation.call_id
  pending_calls.delete([session_id, message.request_id]) unless message.options[:progress]

  result = Message::Result.new(request_id, result_details_for(session_id, message), *message.args, **message.kwargs)
  MessageWithRecipient.new(result, caller_id)
end

#invocation_details_for(session_id, message) ⇒ Object



86
87
88
89
90
91
92
93
94
95
96
97
# File 'lib/wampproto/dealer.rb', line 86

def invocation_details_for(session_id, message)
  options = {}
  return options if message.options.empty?

  receive_progress = message.options[:receive_progress]
  options.merge!(receive_progress: true) if receive_progress

  return options unless message.options.include?(:disclose_me)

  session = sessions[session_id]
  options.merge({ caller: session_id, caller_authid: session.authid, caller_authrole: session.authrole })
end

#receive_message(session_id, message) ⇒ Object



47
48
49
50
51
52
53
54
55
56
# File 'lib/wampproto/dealer.rb', line 47

def receive_message(session_id, message)
  case message
  when Wampproto::Message::Call then handle_call(session_id, message)
  when Message::Yield then handle_yield(session_id, message)
  when Message::Register then handle_register(session_id, message)
  when Message::Unregister then handle_unregister(session_id, message)
  else
    raise ValueError, "message type not supported"
  end
end

#registration?(procedure) ⇒ Boolean

Returns:

  • (Boolean)


40
41
42
43
44
45
# File 'lib/wampproto/dealer.rb', line 40

def registration?(procedure)
  registrations = registrations_by_procedure[procedure]
  return false unless registrations

  registrations.any?
end

#remove_session(session_id) ⇒ Object

Raises:

  • (KeyError)


29
30
31
32
33
34
35
36
37
38
# File 'lib/wampproto/dealer.rb', line 29

def remove_session(session_id)
  error_message = "cannot remove non-existing session"
  raise KeyError, error_message unless registrations_by_session.include?(session_id)

  registrations = registrations_by_session.delete(session_id) || {}
  registrations.each do |registration_id, procedure|
    registrations_by_procedure[procedure].delete(registration_id)
  end
  sessions.delete(session_id)
end

#result_details_for(session_id, message) ⇒ Object



112
113
114
115
116
117
118
119
120
121
# File 'lib/wampproto/dealer.rb', line 112

def result_details_for(session_id, message)
  options = {}
  return options if message.options.empty?

  pending_invocation = pending_calls[[session_id, message.request_id]]

  progress = message.options[:progress] && pending_invocation.receive_progress
  options.merge!(progress:) if progress
  options
end