bosh.lua

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