plugins/smacks.lua

1
local verse = require "verse";
2
local now = require"socket".gettime;
3
 
4
local xmlns_sm = "urn:xmpp:sm:3";
5
 
6
function verse.plugins.smacks(stream)
7
	-- State for outgoing stanzas
8
	local outgoing_queue = nil;
9
	local last_ack = nil;
10
	local last_stanza_time = nil;
11
	local timer_active;
12
 
13
	-- State for incoming stanzas
14
	local handled_stanza_count = nil;
15
 
16
	-- Catch incoming stanzas
17
	local function incoming_stanza(stanza)
18
		if handled_stanza_count and (stanza.attr.xmlns == "jabber:client" or not stanza.attr.xmlns) then
19
			handled_stanza_count = handled_stanza_count + 1;
20
			stream:debug("Increasing handled stanzas to %d for %s", handled_stanza_count, stanza:top_tag());
21
		end
22
	end
23
 
24
	-- Catch outgoing stanzas
25
	local function outgoing_stanza(stanza)
26
		-- NOTE: This will not behave nice if stanzas are serialized before this point
27
		if outgoing_queue and (stanza.name and not stanza.attr.xmlns) then
28
			-- serialize stanzas in order to bypass this on resumption
29
			outgoing_queue[#outgoing_queue+1] = tostring(stanza);
30
			last_stanza_time = now();
31
			if not timer_active then
32
				timer_active = true;
33
				stream:debug("Waiting to send ack request...");
34
				verse.add_task(1, function()
35
					if #outgoing_queue == 0 then
36
						timer_active = false;
37
						return;
38
					end
39
					local time_since_last_stanza = now() - last_stanza_time;
40
					if time_since_last_stanza < 1 and #outgoing_queue < 10 then
41
						return 1 - time_since_last_stanza;
42
					end
43
					stream:debug("Time up, sending <r>...");
44
					timer_active = false;
45
					stream:send(verse.stanza("r", { xmlns = xmlns_sm }));
46
				end);
47
			end
48
		end
49
	end
50
 
51
	local function on_disconnect()
52
		stream:debug("smacks: connection lost");
53
		stream.stream_management_supported = nil;
54
		if stream.resumption_token then
55
			stream:debug("smacks: have resumption token, reconnecting in 1s...");
56
			stream.authenticated = nil;
57
			verse.add_task(1, function ()
58
				stream:connect(stream.connect_host or stream.host, stream.connect_port or 5222);
59
			end);
60
			return true;
61
		end
62
	end
63
 
64
	-- Graceful shutdown
65
	local function on_close()
66
		stream.resumption_token = nil;
67
	end
68
 
69
	local function handle_sm_command(stanza)
70
		if stanza.name == "r" then -- Request for acks for stanzas we received
71
			stream:debug("Ack requested... acking %d handled stanzas", handled_stanza_count);
72
			stream:send(verse.stanza("a", { xmlns = xmlns_sm, h = tostring(handled_stanza_count) }));
73
		elseif stanza.name == "a" then -- Ack for stanzas we sent
74
			local new_ack = tonumber(stanza.attr.h);
75
			if new_ack > last_ack then
76
				local old_unacked = #outgoing_queue;
77
				for i=last_ack+1,new_ack do
78
					table.remove(outgoing_queue, 1);
79
				end
80
				stream:debug("Received ack: New ack: "..new_ack.." Last ack: "..last_ack.." Unacked stanzas now: "..#outgoing_queue.." (was "..old_unacked..")");
81
				last_ack = new_ack;
82
			elseif new_ack < last_ack then
83
				stream:warn("Received bad ack for "..new_ack.." when last ack was "..last_ack);
84
			end
85
		elseif stanza.name == "enabled" then
86
			handled_stanza_count = 0;
87
			stream.pre_smacks_features = nil;
88
 
89
			if stanza.attr.id then
90
				stream.resumption_token = stanza.attr.id;
91
			end
92
		elseif stanza.name == "resumed" then
93
			stream.pre_smacks_features = nil;
94
			local new_ack = tonumber(stanza.attr.h);
95
			if new_ack > last_ack then
96
				local old_unacked = #outgoing_queue;
97
				for i=last_ack+1,new_ack do
98
					table.remove(outgoing_queue, 1);
99
				end
100
				stream:debug("Received ack: New ack: "..new_ack.." Last ack: "..last_ack.." Unacked stanzas now: "..#outgoing_queue.." (was "..old_unacked..")");
101
				last_ack = new_ack;
102
			end
103
			for i=1,#outgoing_queue do
104
				stream:send(outgoing_queue[i]);
105
			end
106
			outgoing_queue = {};
107
			stream:debug("Resumed successfully");
108
			stream:event("resumed");
109
		elseif stanza.name == "failed" then
110
			stream.bound = nil
111
			stream.smacks = nil
112
			last_ack = nil
113
			handled_stanza_count = nil
114
 
115
			-- TODO ack using final h value from <failed/> if present
116
			outgoing_queue = {}; -- TODO fire some delivery failures
117
 
118
			local features = stream.pre_smacks_features;
119
			stream.pre_smacks_features = nil;
120
 
121
			-- should trigger a bind and then a new smacks session
122
			stream:event("stream-features", features);
123
		else
124
			stream:warn("Don't know how to handle "..xmlns_sm.."/"..stanza.name);
125
		end
126
	end
127
 
128
	local function on_bind_success()
129
		if stream.stream_management_supported and not stream.smacks then
130
			--stream:unhook("bind-success", on_bind_success);
131
			stream:debug("smacks: sending enable");
132
			outgoing_queue = {};
133
			last_ack = 0;
134
			last_stanza_time = now();
135
			stream:send(verse.stanza("enable", { xmlns = xmlns_sm, resume = "true" }));
136
			stream.smacks = true;
137
		end
138
	end
139
 
140
	local function on_features(features)
141
		if features:get_child("sm", xmlns_sm) then
142
			stream.pre_smacks_features = features;
143
			stream.stream_management_supported = true;
144
			if stream.smacks and stream.bound then -- Already enabled in a previous session - resume
145
				stream:debug("Resuming stream with %d handled stanzas", handled_stanza_count);
146
				stream:send(verse.stanza("resume", { xmlns = xmlns_sm,
147
					h = tostring(handled_stanza_count), previd = stream.resumption_token }));
148
				return true;
149
			else
150
			end
151
		end
152
	end
153
 
154
	stream:hook("stream-features", on_features, 250);
155
	stream:hook("stream/"..xmlns_sm, handle_sm_command);
156
	stream:hook("bind-success", on_bind_success, 1);
157
 
158
	-- Catch stanzas
159
	stream:hook("stanza", incoming_stanza);
160
	stream:hook("outgoing", outgoing_stanza);
161
 
162
	stream:hook("closed", on_close, 100);
163
	stream:hook("disconnected", on_disconnect, 100);
164
 
165
	--stream:hook("ready", on_stream_ready, 500);
166
end