| 1 | local new_xmpp_stream = require "prosody.util.xmppstream".new; |
| 2 | local st = require "prosody.util.stanza"; |
| 3 | local http = require "prosody.net.http"; |
| 4 | |
| 5 | local stream_mt = setmetatable({}, { __index = verse.stream_mt }); |
| 6 | stream_mt.__index = stream_mt; |
| 7 | |
| 8 | local xmlns_stream = "http://etherx.jabber.org/streams"; |
| 9 | local xmlns_bosh = "http://jabber.org/protocol/httpbind"; |
| 10 | |
| 11 | local reconnect_timeout = 5; |
| 12 | |
| 13 | function verse.new_bosh(logger, url) |
| 14 | local stream = { |
| 15 | bosh_conn_pool = {}; |
| 16 | bosh_waiting_requests = {}; |
| 17 | bosh_rid = math.random(1,999999); |
| 18 | bosh_outgoing_buffer = {}; |
| 19 | bosh_url = url; |
| 20 | conn = {}; |
| 21 | }; |
| 22 | function stream:reopen() |
| 23 | self.bosh_need_restart = true; |
| 24 | self:flush(); |
| 25 | end |
| 26 | local conn = verse.new(logger, stream); |
| 27 | return setmetatable(conn, stream_mt); |
| 28 | end |
| 29 | |
| 30 | function stream_mt:connect() |
| 31 | self:_send_session_request(); |
| 32 | end |
| 33 | |
| 34 | function stream_mt:send(data) |
| 35 | self:debug("Putting into BOSH send buffer: %s", tostring(data)); |
| 36 | self.bosh_outgoing_buffer[#self.bosh_outgoing_buffer+1] = st.clone(data); |
| 37 | self:flush(); --TODO: Optimize by doing this on next tick (give a chance for data to buffer) |
| 38 | end |
| 39 | |
| 40 | function stream_mt:flush() |
| 41 | if self.connected |
| 42 | and #self.bosh_waiting_requests < self.bosh_max_requests |
| 43 | and (#self.bosh_waiting_requests == 0 |
| 44 | or #self.bosh_outgoing_buffer > 0 |
| 45 | or self.bosh_need_restart) then |
| 46 | self:debug("Flushing..."); |
| 47 | local payload = self:_make_body(); |
| 48 | local buffer = self.bosh_outgoing_buffer; |
| 49 | for i, stanza in ipairs(buffer) do |
| 50 | payload:add_child(stanza); |
| 51 | buffer[i] = nil; |
| 52 | end |
| 53 | self:_make_request(payload); |
| 54 | else |
| 55 | self:debug("Decided not to flush."); |
| 56 | end |
| 57 | end |
| 58 | |
| 59 | function stream_mt:_make_request(payload) |
| 60 | local request, err = http.request(self.bosh_url, { body = tostring(payload) }, function (response, code, request) |
| 61 | if code ~= 0 then |
| 62 | self.inactive_since = nil; |
| 63 | return self:_handle_response(response, code, request); |
| 64 | end |
| 65 | |
| 66 | -- Connection issues, we need to retry this request |
| 67 | local time = os.time(); |
| 68 | if not self.inactive_since then |
| 69 | self.inactive_since = time; -- So we know when it is time to give up |
| 70 | elseif time - self.inactive_since > self.bosh_max_inactivity then |
| 71 | return self:_disconnected(); |
| 72 | else |
| 73 | self:debug("%d seconds left to reconnect, retrying in %d seconds...", |
| 74 | self.bosh_max_inactivity - (time - self.inactive_since), reconnect_timeout); |
| 75 | end |
| 76 | |
| 77 | -- Set up reconnect timer |
| 78 | timer.add_task(reconnect_timeout, function () |
| 79 | self:debug("Retrying request..."); |
| 80 | -- Remove old request |
| 81 | for i, waiting_request in ipairs(self.bosh_waiting_requests) do |
| 82 | if waiting_request == request then |
| 83 | table.remove(self.bosh_waiting_requests, i); |
| 84 | break; |
| 85 | end |
| 86 | end |
| 87 | self:_make_request(payload); |
| 88 | end); |
| 89 | end); |
| 90 | if request then |
| 91 | table.insert(self.bosh_waiting_requests, request); |
| 92 | else |
| 93 | self:warn("Request failed instantly: %s", err); |
| 94 | end |
| 95 | end |
| 96 | |
| 97 | function stream_mt:_disconnected() |
| 98 | self.connected = nil; |
| 99 | self:event("disconnected"); |
| 100 | end |
| 101 | |
| 102 | function stream_mt:_send_session_request() |
| 103 | local body = self:_make_body(); |
| 104 | |
| 105 | -- XEP-0124 |
| 106 | body.attr.hold = "1"; |
| 107 | body.attr.wait = "60"; |
| 108 | body.attr["xml:lang"] = "en"; |
| 109 | body.attr.ver = "1.6"; |
| 110 | |
| 111 | -- XEP-0206 |
| 112 | body.attr.from = self.jid; |
| 113 | body.attr.to = self.host; |
| 114 | body.attr.secure = 'true'; |
| 115 | |
| 116 | http.request(self.bosh_url, { body = tostring(body) }, function (response, code) |
| 117 | if code == 0 then |
| 118 | -- Failed to connect |
| 119 | return self:_disconnected(); |
| 120 | end |
| 121 | -- Handle session creation response |
| 122 | local payload = self:_parse_response(response) |
| 123 | if not payload then |
| 124 | self:warn("Invalid session creation response"); |
| 125 | self:_disconnected(); |
| 126 | return; |
| 127 | end |
| 128 | self.bosh_sid = payload.attr.sid; -- Session id |
| 129 | self.bosh_wait = tonumber(payload.attr.wait); -- How long the server may hold connections for |
| 130 | self.bosh_hold = tonumber(payload.attr.hold); -- How many connections the server may hold |
| 131 | self.bosh_max_inactivity = tonumber(payload.attr.inactivity); -- Max amount of time with no connections |
| 132 | self.bosh_max_requests = tonumber(payload.attr.requests) or self.bosh_hold; -- Max simultaneous requests we can make |
| 133 | self.connected = true; |
| 134 | self:event("connected"); |
| 135 | self:_handle_response_payload(payload); |
| 136 | end); |
| 137 | end |
| 138 | |
| 139 | function stream_mt:_handle_response(response, code, request) |
| 140 | if self.bosh_waiting_requests[1] ~= request then |
| 141 | self:warn("Server replied to request that wasn't the oldest"); |
| 142 | for i, waiting_request in ipairs(self.bosh_waiting_requests) do |
| 143 | if waiting_request == request then |
| 144 | self.bosh_waiting_requests[i] = nil; |
| 145 | break; |
| 146 | end |
| 147 | end |
| 148 | else |
| 149 | table.remove(self.bosh_waiting_requests, 1); |
| 150 | end |
| 151 | local payload = self:_parse_response(response); |
| 152 | if payload then |
| 153 | self:_handle_response_payload(payload); |
| 154 | end |
| 155 | self:flush(); |
| 156 | end |
| 157 | |
| 158 | function stream_mt:_handle_response_payload(payload) |
| 159 | local stanzas = payload.tags; |
| 160 | for i = 1, #stanzas do |
| 161 | local stanza = stanzas[i]; |
| 162 | if stanza.attr.xmlns == xmlns_stream then |
| 163 | self:event("stream-"..stanza.name, stanza); |
| 164 | elseif stanza.attr.xmlns then |
| 165 | self:event("stream/"..stanza.attr.xmlns, stanza); |
| 166 | else |
| 167 | self:event("stanza", stanza); |
| 168 | end |
| 169 | end |
| 170 | if payload.attr.type == "terminate" then |
| 171 | self:_disconnected({reason = payload.attr.condition}); |
| 172 | end |
| 173 | end |
| 174 | |
| 175 | local stream_callbacks = { |
| 176 | stream_ns = "http://jabber.org/protocol/httpbind", stream_tag = "body", |
| 177 | default_ns = "jabber:client", |
| 178 | streamopened = function (session, attr) session.notopen = nil; session.payload = verse.stanza("body", attr); return true; end; |
| 179 | handlestanza = function (session, stanza) session.payload:add_child(stanza); end; |
| 180 | }; |
| 181 | function stream_mt:_parse_response(response) |
| 182 | self:debug("Parsing response: %s", response); |
| 183 | if response == nil then |
| 184 | self:debug("%s", debug.traceback()); |
| 185 | self:_disconnected(); |
| 186 | return; |
| 187 | end |
| 188 | local session = { notopen = true, stream = self }; |
| 189 | local stream = new_xmpp_stream(session, stream_callbacks); |
| 190 | stream:feed(response); |
| 191 | return session.payload; |
| 192 | end |
| 193 | |
| 194 | function stream_mt:_make_body() |
| 195 | self.bosh_rid = self.bosh_rid + 1; |
| 196 | local body = verse.stanza("body", { |
| 197 | xmlns = xmlns_bosh; |
| 198 | content = "text/xml; charset=utf-8"; |
| 199 | sid = self.bosh_sid; |
| 200 | rid = self.bosh_rid; |
| 201 | }); |
| 202 | if self.bosh_need_restart then |
| 203 | self.bosh_need_restart = nil; |
| 204 | body.attr.restart = 'true'; |
| 205 | end |
| 206 | return body; |
| 207 | end |