| 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 |