likeopera-backend/operetta.js

246 lines
10 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters!

This file contains ambiguous Unicode characters that may be confused with others in your current locale. If your use case is intentional and legitimate, you can safely ignore this warning. Use the Escape button to highlight these characters.

process.env.NODE_TLS_REJECT_UNAUTHORIZED = "0";
var cfg = require('./cfg.json');
require('heapdump');
var gen = require('gen-thread');
var Imap = require('imap');
var inspect = require('util').inspect;
//var pg;
//try { require('pg-native'); pg = require('pg').native; }
//catch(e) { pg = require('pg'); }
//var pg_pool = new pg.Pool(cfg.pg);
var bricks = require('pg-bricks');
var pg = bricks.configure('postgresql://'+cfg.pg.user+':'+cfg.pg.password+'@'+(cfg.pg.host||'')+':'+cfg.pg.port+'/'+cfg.pg.database);
var flagNum = {
'\\recent': 1,
'\\flagged': 2,
'\\answered': 4,
'\\seen': 8,
};
function splitEmails(s)
{
var re = /^[\s,]*(?:(?:["'](.*?)["']|([^<]+))\s*<([^>]+)>|<?([^<>]+)>?)/; // '
var m, r = [];
while (m = re.exec(s))
{
s = s.substr(m[0].length);
r.push({ name: (m[1]||m[2]||'').trim(), email: (m[3]||m[4]||'').trim() });
}
return r;
}
function toPgArray(a)
{
a = JSON.stringify(a);
return '{'+a.substring(1, a.length-1)+'}';
}
function* main(NEXT, account)
{
var accountId;
var [ rows ] = yield pg.select('id').from('accounts').where({ email: account.email }).rows(NEXT.ef());
if (rows[0] && rows[0].id)
accountId = rows[0].id;
else
{
var [ row ] = pg.insert('accounts', {
name: account.name,
email: account.email,
settings: {
imap: account.imap
}
}).returning('id').row(NEXT.ef());
accountId = row.id;
}
var srv = new Imap(account.imap);
srv.once('ready', NEXT.cb());
yield srv.connect();
var [ boxes ] = yield srv.getBoxes(NEXT.ef());
for (var k in boxes)
{
var [ box ] = yield srv.openBox(k, true, NEXT.ef());
var boxId;
var [ rows ] = yield pg.update('folders', { uidvalidity: box.uidvalidity, unread_count: box.messages.new })
.where({ account_id: accountId, name: box.name }).returning('id').rows(NEXT.ef());
if (rows[0] && rows[0].id)
{
// IMAP sync: http://tools.ietf.org/html/rfc4549
// TODO: check old uidvalidity
boxId = rows[0].id;
}
else
{
var [ row ] = yield pg.insert('folders', {
name: box.name,
uidvalidity: box.uidvalidity,
account_id: accountId,
unread_count: box.messages.new,
// total_count: box.messages.count
}).returning('id').row(NEXT.ef());
boxId = row.id;
}
var f = srv.fetch('1:*', {
size: true,
bodies: 'HEADER'
});
var parsed = 0, paused = false;
f.on('message', function(msg, seqnum)
{
gen.run(function*(NEXT)
{
var msgrow = {};
var attrs;
msg.on('body', function(stream, info)
{
var buffer = '';
stream.on('data', function(chunk)
{
buffer += chunk.toString('utf8');
});
stream.once('end', function()
{
msgrow.body = '';
msgrow.headers = buffer;
});
});
msg.once('attributes', function(a) {
attrs = a;
});
yield msg.once('end', NEXT.cb());
// Workaround memory leak in node-imap
// TODO: send pull request
if (srv._curReq.fetchCache)
delete srv._curReq.fetchCache[seqnum];
parsed++;
if (!paused && parsed > 20)
{
// ГОРШОЧЕК, НЕ ВАРИ!!! И так уже кучу сообщений прочитал из сокета, хорош!
srv._parser._ignoreReadable = true;
paused = true;
}
var pgtx, end_transaction;
try
{
[ pgtx, end_transaction ] = yield pg.transaction(NEXT.cb(), function(e) { if (e) throw e; });
msgrow.uid = attrs.uid;
msgrow.folder_id = boxId;
msgrow.flags = 0;
for (var i = 0; i < attrs.flags.length; i++)
msgrow.flags = msgrow.flags || flagNum[attrs.flags[i].toLowerCase()];
msgrow.flags = (msgrow.flags & ~8) | (msgrow.flags & 8 ? 0 : 8); // invert "\seen" (unread) flag
var [ exists ] = yield pgtx.select('id').from('messages').where({ folder_id: msgrow.folder_id, uid: msgrow.uid }).rows(NEXT.ef());
if (exists.length)
{
process.stderr.write('\rsynchronizing '+msgrow.uid+'...');
yield pgtx.update('messages', { flags: msgrow.flags }).where({ folder_id: msgrow.folder_id, uid: msgrow.uid }).run(NEXT.ef());
}
else
{
var header = Imap.parseHeader(msgrow.headers);
for (var i in header)
for (var k = 0; k < header[i].length; k++)
header[i][k] = header[i][k].replace(/\x00/g, '');
header.from = header.from && splitEmails(header.from[0])[0];
header.replyto = header['reply-to'] && splitEmails(header['reply-to'][0])[0];
var re = /(<[^>]*>)/;
header.references = (header.references && header.references[0] || '').split(re).filter(a => a.match(re));
if (header.references.length)
{
if (header.references.length > 10)
header.references = [ header.references[0] ].concat(header.references.slice(header.references.length-9));
if (!header['in-reply-to'] || !header['in-reply-to'][0])
header['in-reply-to'] = [ header.references[header.references.length-1] ];
else if (header.references[header.references.length-1] != header['in-reply-to'][0])
header.references.push(header['in-reply-to'][0]);
}
if (header.date)
{
var t = Date.parse(header.date[0]);
if (!isNaN(t))
header.date = new Date(t);
}
if (!header.date)
header.date = new Date(attrs.date);
msgrow.from_email = header.from && header.from.email || '';
msgrow.from_name = header.from && header.from.name || '';
msgrow.replyto_email = header.replyto && header.replyto.email || '';
msgrow.replyto_name = header.replyto && header.replyto.name || '';
msgrow.to_list = header.to && header.to[0] || '';
msgrow.cc_list = header.cc && header.cc[0] || '';
msgrow.bcc_list = header.bcc && header.bcc[0] || '';
msgrow.subject = header.subject && header.subject[0] || '';
msgrow.messageid = header['message-id'] && header['message-id'][0] || '';
msgrow.inreplyto = header['in-reply-to'] && header['in-reply-to'][0] || '';
msgrow.inreplyto = msgrow.inreplyto.replace(/^[\s\S]*(<[^>]*>)[\s\S]*$/, '$1');
msgrow.time = header.date;
msgrow.refs = toPgArray(header.references);
if (header.references.length)
{
var [ threadId ] = yield pgtx.select('MAX(thread_id)').from('messages')
.where(pg.sql.in('messageid', header.references)).val(NEXT.ef());
if (!threadId)
{
[ threadId ] = yield pgtx.select('MAX(thread_id)').from('messages')
.where(new pg.sql.Binary('@>', 'refs', toPgArray([msgrow.messageid]))).val(NEXT.ef());
}
if (threadId)
{
try
{
yield pgtx.update('threads', { msg_count: pg.sql('msg_count+1') })
.where({ id: threadId }).run(NEXT.ef());
}
catch (e)
{
throw new Error(''+e);
}
}
msgrow.thread_id = threadId;
}
console.log(msgrow.time+' '+msgrow.from_email+' '+msgrow.subject);
[ msgrow.id ] = yield pgtx.insert('messages', msgrow).returning('id').val(NEXT.ef());
if (!msgrow.thread_id)
{
[ msgrow.thread_id ] = yield pgtx.insert('threads', {
first_msg: msgrow.id,
msg_count: 1
}).returning('id').val(NEXT.ef());
yield pgtx.update('messages', { thread_id: msgrow.thread_id }).where({ id: msgrow.id }).run(NEXT.ef());
}
}
end_transaction();
}
catch (e0)
{
if (end_transaction)
end_transaction();
throw e0;
}
parsed--;
if (paused && parsed <= 10)
{
paused = false;
srv._parser._ignoreReadable = false;
process.nextTick(srv._parser._cbReadable);
}
});
});
yield f.once('end', NEXT.cb());
yield srv.closeBox(NEXT.cb());
}
srv.end();
}
gen.run(main, cfg.accounts[0], function() { process.exit() });