Class: Wampproto::Dealer
- Inherits:
-
Object
- Object
- Wampproto::Dealer
- Defined in:
- lib/wampproto/dealer.rb
Overview
Wamprpoto Dealer handler
Defined Under Namespace
Classes: PendingInvocation
Instance Attribute Summary collapse
-
#id_gen ⇒ Object
readonly
Returns the value of attribute id_gen.
-
#pending_calls ⇒ Object
readonly
Returns the value of attribute pending_calls.
-
#registrations_by_procedure ⇒ Object
readonly
Returns the value of attribute registrations_by_procedure.
-
#registrations_by_session ⇒ Object
readonly
Returns the value of attribute registrations_by_session.
-
#sessions ⇒ Object
readonly
Returns the value of attribute sessions.
Instance Method Summary collapse
- #add_session(details) ⇒ Object
-
#handle_call(session_id, message) ⇒ Object
rubocop:disable Metrics/MethodLength, Metrics/AbcSize.
- #handle_register(session_id, message) ⇒ Object
-
#handle_unregister(session_id, message) ⇒ Object
rubocop:disable Metrics/AbcSize, Metrics/MethodLength).
-
#handle_yield(session_id, message) ⇒ Object
rubocop:disable Metrics/AbcSize.
-
#initialize(id_gen = IdGenerator.new) ⇒ Dealer
constructor
A new instance of Dealer.
- #invocation_details_for(session_id, message) ⇒ Object
- #receive_message(session_id, message) ⇒ Object
- #registration?(procedure) ⇒ Boolean
- #remove_session(session_id) ⇒ Object
- #result_details_for(session_id, message) ⇒ Object
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
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 = "cannot add session twice" raise KeyError, 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, ) # rubocop:disable Metrics/MethodLength, Metrics/AbcSize registrations = registrations_by_procedure.fetch(.procedure, {}) if registrations.empty? error = Message::Error.new(Message::Type::CALL, .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, ) pending_calls[[callee_id, request_id]] = PendingInvocation.new( session_id, callee_id, .request_id, request_id, details[:receive_progress] ) invocation = Message::Invocation.new( request_id, registration_id, details, *.args, **.kwargs ) MessageWithRecipient.new(invocation, callee_id) end |
#handle_register(session_id, message) ⇒ Object
123 124 125 126 127 128 129 130 131 132 133 |
# File 'lib/wampproto/dealer.rb', line 123 def handle_register(session_id, ) = "cannot register, session #{session_id} doesn't exist" raise ValueError, unless registrations_by_session.include?(session_id) registration_id = id_gen.next registrations_by_procedure[.procedure][registration_id] = session_id registrations_by_session[session_id][registration_id] = .procedure registered = Message::Registered.new(.request_id, registration_id) MessageWithRecipient.new(registered, session_id) end |
#handle_unregister(session_id, message) ⇒ Object
rubocop:disable Metrics/AbcSize, Metrics/MethodLength)
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, ) # rubocop:disable Metrics/AbcSize, Metrics/MethodLength) = "cannot unregister, session #{session_id} doesn't exist" raise ValueError, unless registrations_by_session.include?(session_id) registrations = registrations_by_session.fetch(session_id) unless registrations.include?(.registration_id) error = Message::Error.new(Message::Type::UNREGISTER, .request_id, {}, "wamp.error.no_such_registration") return MessageWithRecipient.new(error, session_id) end procedure = registrations.fetch(.registration_id) registrations_by_procedure[procedure].delete(.registration_id) registrations_by_session[session_id].delete(.registration_id) unregistered = Message::Unregistered.new(.request_id) MessageWithRecipient.new(unregistered, session_id) end |
#handle_yield(session_id, message) ⇒ Object
rubocop:disable Metrics/AbcSize
99 100 101 102 103 104 105 106 107 108 109 110 |
# File 'lib/wampproto/dealer.rb', line 99 def handle_yield(session_id, ) # rubocop:disable Metrics/AbcSize pending_invocation = pending_calls[[session_id, .request_id]] = "no pending calls for session #{session_id}" raise ValueError, if pending_invocation.nil? caller_id = pending_invocation.caller_id request_id = pending_invocation.call_id pending_calls.delete([session_id, .request_id]) unless .[:progress] result = Message::Result.new(request_id, result_details_for(session_id, ), *.args, **.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, ) = {} return if ..empty? receive_progress = .[:receive_progress] .merge!(receive_progress: true) if receive_progress return unless ..include?(:disclose_me) session = sessions[session_id] .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 (session_id, ) case when Wampproto::Message::Call then handle_call(session_id, ) when Message::Yield then handle_yield(session_id, ) when Message::Register then handle_register(session_id, ) when Message::Unregister then handle_unregister(session_id, ) else raise ValueError, "message type not supported" end end |
#registration?(procedure) ⇒ 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
29 30 31 32 33 34 35 36 37 38 |
# File 'lib/wampproto/dealer.rb', line 29 def remove_session(session_id) = "cannot remove non-existing session" raise KeyError, 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, ) = {} return if ..empty? pending_invocation = pending_calls[[session_id, .request_id]] progress = .[:progress] && pending_invocation.receive_progress .merge!(progress:) if progress end |