Mercurial > prosody-modules
comparison mod_mam_archive/mod_mam_archive.lua @ 1471:153df603f73d
mod_mam_archive: Initial commit
| author | syn@syn.im |
|---|---|
| date | Mon, 14 Jul 2014 23:00:13 +0200 |
| parents | |
| children | 08ca6dd36e39 |
comparison
equal
deleted
inserted
replaced
| 1470:b291a9423e0f | 1471:153df603f73d |
|---|---|
| 1 -- Prosody IM | |
| 2 -- | |
| 3 -- This project is MIT/X11 licensed. Please see the | |
| 4 -- COPYING file in the source package for more information. | |
| 5 -- | |
| 6 local get_prefs = module:require"mod_mam/mamprefs".get; | |
| 7 local set_prefs = module:require"mod_mam/mamprefs".set; | |
| 8 local rsm = module:require "mod_mam/rsm"; | |
| 9 local jid_bare = require "util.jid".bare; | |
| 10 local jid_prep = require "util.jid".prep; | |
| 11 local date_parse = require "util.datetime".parse; | |
| 12 | |
| 13 local st = require "util.stanza"; | |
| 14 local archive_store = "archive2"; | |
| 15 local archive = module:open_store(archive_store, "archive"); | |
| 16 local global_default_policy = module:get_option("default_archive_policy", false); | |
| 17 local default_max_items, max_max_items = 20, module:get_option_number("max_archive_query_results", 50); | |
| 18 local conversation_interval = module:get_option_number("archive_conversation_interval", 86400); | |
| 19 | |
| 20 -- Feature discovery | |
| 21 local xmlns_archive = "urn:xmpp:archive" | |
| 22 local feature_archive = st.stanza('feature', {xmlns=xmlns_archive}); | |
| 23 feature_archive:tag('optional'); | |
| 24 if(global_default_policy) then | |
| 25 feature_archive:tag('default'); | |
| 26 end | |
| 27 module:add_extension(feature_archive); | |
| 28 module:add_feature("urn:xmpp:archive:auto"); | |
| 29 module:add_feature("urn:xmpp:archive:manage"); | |
| 30 module:add_feature("urn:xmpp:archive:pref"); | |
| 31 module:add_feature("http://jabber.org/protocol/rsm"); | |
| 32 -- -------------------------------------------------- | |
| 33 | |
| 34 local function os_date() | |
| 35 return os.date("!*t"); | |
| 36 end | |
| 37 local function date_format(s) | |
| 38 return os.date("%Y-%m-%dT%H:%M:%SZ", s); | |
| 39 end | |
| 40 | |
| 41 local function prefs_to_stanza(prefs) | |
| 42 local prefstanza = st.stanza("pref", { xmlns='urn:xmpp:archive' }); | |
| 43 local default = prefs[false] ~= nil and prefs[false] or global_default_policy; | |
| 44 | |
| 45 prefstanza:tag('default', {otr='oppose', save=default and 'true' or 'false'}):up(); | |
| 46 prefstanza:tag('method', {type='auto', use='concede'}):up(); | |
| 47 prefstanza:tag('method', {type='local', use='concede'}):up(); | |
| 48 prefstanza:tag('method', {type='manual', use='concede'}):up(); | |
| 49 | |
| 50 for jid, choice in pairs(prefs) do | |
| 51 if jid then | |
| 52 prefstanza:tag('item', {jid=jid, otr='prefer', save=choice and 'message' or 'false' }):up() | |
| 53 end | |
| 54 end | |
| 55 | |
| 56 return prefstanza; | |
| 57 end | |
| 58 local function prefs_from_stanza(stanza) | |
| 59 local current_prefs = get_prefs(origin.username); | |
| 60 | |
| 61 -- "default" | "item" | "session" | "method" | |
| 62 for elem in stanza:children() do | |
| 63 if elem.name == "default" then | |
| 64 current_prefs[false] = elem.attr['save'] == 'true'; | |
| 65 elseif elem.name == "item" then | |
| 66 current_prefs[elem.attr['jid']] = not elem.attr['save'] == 'false'; | |
| 67 elseif elem.name == "session" then | |
| 68 module:log("info", "element is not supported: " .. tostring(elem)); | |
| 69 -- local found = false; | |
| 70 -- for child in data:children() do | |
| 71 -- if child.name == elem.name and child.attr["thread"] == elem.attr["thread"] then | |
| 72 -- for k, v in pairs(elem.attr) do | |
| 73 -- child.attr[k] = v; | |
| 74 -- end | |
| 75 -- found = true; | |
| 76 -- break; | |
| 77 -- end | |
| 78 -- end | |
| 79 -- if not found then | |
| 80 -- data:tag(elem.name, elem.attr):up(); | |
| 81 -- end | |
| 82 elseif elem.name == "method" then | |
| 83 module:log("info", "element is not supported: " .. tostring(elem)); | |
| 84 -- local newpref = stanza.tags[1]; -- iq:pref | |
| 85 -- for _, e in ipairs(newpref.tags) do | |
| 86 -- -- if e.name ~= "method" then continue end | |
| 87 -- local found = false; | |
| 88 -- for child in data:children() do | |
| 89 -- if child.name == "method" and child.attr["type"] == e.attr["type"] then | |
| 90 -- child.attr["use"] = e.attr["use"]; | |
| 91 -- found = true; | |
| 92 -- break; | |
| 93 -- end | |
| 94 -- end | |
| 95 -- if not found then | |
| 96 -- data:tag(e.name, e.attr):up(); | |
| 97 -- end | |
| 98 -- end | |
| 99 end | |
| 100 end | |
| 101 end | |
| 102 | |
| 103 ------------------------------------------------------------ | |
| 104 -- Preferences | |
| 105 ------------------------------------------------------------ | |
| 106 local function preferences_handler(event) | |
| 107 local origin, stanza = event.origin, event.stanza; | |
| 108 local user = origin.username; | |
| 109 local reply = st.reply(stanza); | |
| 110 | |
| 111 if stanza.attr.type == "get" then | |
| 112 reply:add_child(prefs_to_stanza(get_prefs(user))); | |
| 113 end | |
| 114 if stanza.attr.type == "set" then | |
| 115 local new_prefs = stanza:get_child("pref", xmlns_archive); | |
| 116 if not new_prefs then return false; end | |
| 117 | |
| 118 local prefs = prefs_from_stanza(stanza); | |
| 119 local ok, err = set_prefs(user, prefs); | |
| 120 | |
| 121 if not ok then | |
| 122 return origin.send(st.error_reply(stanza, "cancel", "internal-server-error", "Error storing preferences: "..tostring(err))); | |
| 123 end | |
| 124 end | |
| 125 return origin.send(reply); | |
| 126 end | |
| 127 local function auto_handler(event) | |
| 128 local origin, stanza = event.origin, event.stanza; | |
| 129 if not stanza.attr['type'] == 'set' then return false; end | |
| 130 | |
| 131 local user = origin.username; | |
| 132 local prefs = get_prefs(user); | |
| 133 local auto = stanza:get_child('auto', xmlns_archive); | |
| 134 | |
| 135 prefs[false] = auto.attr['save'] ~= nil and auto.attr['save'] == 'true' or false; | |
| 136 set_prefs(user, prefs); | |
| 137 | |
| 138 return origin.send(st.reply(stanza)); | |
| 139 end | |
| 140 | |
| 141 -- excerpt from mod_storage_sql2 | |
| 142 local function get_db() | |
| 143 local mod_sql = module:require("sql"); | |
| 144 local params = module:get_option("sql"); | |
| 145 local engine; | |
| 146 | |
| 147 params = params or { driver = "SQLite3" }; | |
| 148 if params.driver == "SQLite3" then | |
| 149 params.database = resolve_relative_path(prosody.paths.data or ".", params.database or "prosody.sqlite"); | |
| 150 end | |
| 151 | |
| 152 assert(params.driver and params.database, "Both the SQL driver and the database need to be specified"); | |
| 153 engine = mod_sql:create_engine(params); | |
| 154 engine:set_encoding(); | |
| 155 | |
| 156 return engine; | |
| 157 end | |
| 158 | |
| 159 ------------------------------------------------------------ | |
| 160 -- Collections. In our case there is one conversation with each contact for the whole day for simplicity | |
| 161 ------------------------------------------------------------ | |
| 162 local function list_stanza_to_query(origin, list_el) | |
| 163 local sql = "SELECT `with`, `when` / ".. conversation_interval .." as `day`, COUNT(0) FROM `prosodyarchive` WHERE `host`=? AND `user`=? AND `store`=? "; | |
| 164 local args = {origin.host, origin.username, archive_store}; | |
| 165 | |
| 166 local with = list_el.attr['with']; | |
| 167 if with ~= nil then | |
| 168 sql = sql .. "AND `with` = ? "; | |
| 169 table.insert(args, jid_bare(with)); | |
| 170 end | |
| 171 | |
| 172 local after = list_el.attr['start']; | |
| 173 if after ~= nil then | |
| 174 sql = sql .. "AND `when` >= ?"; | |
| 175 table.insert(args, date_parse(after)); | |
| 176 end | |
| 177 | |
| 178 local before = list_el.attr['end']; | |
| 179 if before ~= nil then | |
| 180 sql = sql .. "AND `when` <= ? "; | |
| 181 table.insert(args, date_parse(before)); | |
| 182 end | |
| 183 | |
| 184 sql = sql .. "GROUP BY `with`, `when` / ".. conversation_interval .." ORDER BY `when` / ".. conversation_interval .." ASC "; | |
| 185 | |
| 186 local qset = rsm.get(list_el); | |
| 187 local limit = math.min(qset and qset.max or default_max_items, max_max_items); | |
| 188 sql = sql..'LIMIT ?'; | |
| 189 table.insert(args, limit); | |
| 190 | |
| 191 table.insert(args, 1, sql); | |
| 192 return args; | |
| 193 end | |
| 194 local function list_handler(event) | |
| 195 local db = get_db(); | |
| 196 local origin, stanza = event.origin, event.stanza; | |
| 197 local reply = st.reply(stanza); | |
| 198 | |
| 199 local query = list_stanza_to_query(origin, stanza.tags[1]); | |
| 200 local list = reply:tag('list', {xmlns=xmlns_archive}); | |
| 201 | |
| 202 for row in db:select(unpack(query)) do | |
| 203 list:tag('chat', { | |
| 204 xmlns=xmlns_archive, | |
| 205 with=row[1], | |
| 206 start=date_format(row[2] * conversation_interval), | |
| 207 version=row[3] | |
| 208 }):up(); | |
| 209 end | |
| 210 | |
| 211 origin.send(reply); | |
| 212 return true; | |
| 213 end | |
| 214 | |
| 215 ------------------------------------------------------------ | |
| 216 -- Message archive retrieval | |
| 217 ------------------------------------------------------------ | |
| 218 | |
| 219 local function retrieve_handler(event) | |
| 220 local origin, stanza = event.origin, event.stanza; | |
| 221 local reply = st.reply(stanza); | |
| 222 | |
| 223 local retrieve = stanza:get_child('retrieve', xmlns_archive); | |
| 224 | |
| 225 local qwith = retrieve.attr['with']; | |
| 226 local qstart = retrieve.attr['start']; | |
| 227 | |
| 228 module:log("debug", "Archive query, with %s from %s)", | |
| 229 qwith or "anyone", qstart or "the dawn of time"); | |
| 230 | |
| 231 if qstart then -- Validate timestamps | |
| 232 local vstart = (qstart and date_parse(qstart)); | |
| 233 if (qstart and not vstart) then | |
| 234 origin.send(st.error_reply(stanza, "modify", "bad-request", "Invalid timestamp")) | |
| 235 return true | |
| 236 end | |
| 237 qstart = vstart; | |
| 238 end | |
| 239 | |
| 240 if qwith then -- Validate the 'with' jid | |
| 241 local pwith = qwith and jid_prep(qwith); | |
| 242 if pwith and not qwith then -- it failed prepping | |
| 243 origin.send(st.error_reply(stanza, "modify", "bad-request", "Invalid JID")) | |
| 244 return true | |
| 245 end | |
| 246 qwith = jid_bare(pwith); | |
| 247 end | |
| 248 | |
| 249 -- RSM stuff | |
| 250 local qset = rsm.get(retrieve); | |
| 251 local qmax = math.min(qset and qset.max or default_max_items, max_max_items); | |
| 252 local reverse = qset and qset.before or false; | |
| 253 local before, after = qset and qset.before, qset and qset.after; | |
| 254 if type(before) ~= "string" then before = nil; end | |
| 255 | |
| 256 -- Load all the data! | |
| 257 local data, err = archive:find(origin.username, { | |
| 258 start = qstart; ["end"] = qstart + conversation_interval; | |
| 259 with = qwith; | |
| 260 limit = qmax; | |
| 261 before = before; after = after; | |
| 262 reverse = reverse; | |
| 263 total = true; | |
| 264 }); | |
| 265 | |
| 266 if not data then | |
| 267 return origin.send(st.error_reply(stanza, "cancel", "internal-server-error", err)); | |
| 268 end | |
| 269 local count = err; | |
| 270 | |
| 271 local chat = reply:tag('chat', {xmlns=xmlns_archive, with=qwith, start=date_format(qstart), version=count}); | |
| 272 | |
| 273 module:log("debug", 'Count '..count); | |
| 274 for id, item, when in data do | |
| 275 local tag = jid_bare(item['attr']['from']) == jid_bare(origin.full_jid) and 'from' or 'to'; | |
| 276 tag = chat:tag(tag, {secs = when - qstart}); | |
| 277 tag:tag('body'):text(item[2][1]):up():up(); | |
| 278 end | |
| 279 | |
| 280 origin.send(reply); | |
| 281 return true; | |
| 282 end | |
| 283 | |
| 284 local function not_implemented(event) | |
| 285 local origin, stanza = event.origin, event.stanza; | |
| 286 local reply = st.reply(stanza):tag('error', {type='cancel'}); | |
| 287 reply:tag('feature-not-implemented', {xmlns='urn:ietf:params:xml:ns:xmpp-stanzas'}):up(); | |
| 288 origin.send(reply); | |
| 289 end | |
| 290 | |
| 291 -- Preferences | |
| 292 module:hook("iq/self/urn:xmpp:archive:pref", preferences_handler); | |
| 293 module:hook("iq/self/urn:xmpp:archive:auto", auto_handler); | |
| 294 module:hook("iq/self/urn:xmpp:archive:itemremove", not_implemented); | |
| 295 module:hook("iq/self/urn:xmpp:archive:sessionremove", not_implemented); | |
| 296 | |
| 297 -- Message Archive Management | |
| 298 module:hook("iq/self/urn:xmpp:archive:list", list_handler); | |
| 299 module:hook("iq/self/urn:xmpp:archive:retrieve", retrieve_handler); | |
| 300 module:hook("iq/self/urn:xmpp:archive:remove", not_implemented); | |
| 301 | |
| 302 -- manual archiving | |
| 303 module:hook("iq/self/urn:xmpp:archive:save", not_implemented); | |
| 304 -- replication | |
| 305 module:hook("iq/self/urn:xmpp:archive:modified", not_implemented); |
