sync.c 56 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924
  1. /*
  2. * mbsync - mailbox synchronizer
  3. * Copyright (C) 2000-2002 Michael R. Elkins <me@mutt.org>
  4. * Copyright (C) 2002-2006,2010-2013 Oswald Buddenhagen <ossi@users.sf.net>
  5. *
  6. * This program is free software; you can redistribute it and/or modify
  7. * it under the terms of the GNU General Public License as published by
  8. * the Free Software Foundation; either version 2 of the License, or
  9. * (at your option) any later version.
  10. *
  11. * This program is distributed in the hope that it will be useful,
  12. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  13. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  14. * GNU General Public License for more details.
  15. *
  16. * You should have received a copy of the GNU General Public License
  17. * along with this program. If not, see <http://www.gnu.org/licenses/>.
  18. *
  19. * As a special exception, mbsync may be linked with the OpenSSL library,
  20. * despite that library's more restrictive license.
  21. */
  22. #include "isync.h"
  23. #include <assert.h>
  24. #include <stdio.h>
  25. #include <limits.h>
  26. #include <stdlib.h>
  27. #include <stddef.h>
  28. #include <unistd.h>
  29. #include <time.h>
  30. #include <fcntl.h>
  31. #include <ctype.h>
  32. #include <string.h>
  33. #include <errno.h>
  34. #include <sys/stat.h>
  35. #ifndef _POSIX_SYNCHRONIZED_IO
  36. # define fdatasync fsync
  37. #endif
  38. const char *str_ms[] = { "master", "slave" }, *str_hl[] = { "push", "pull" };
  39. void
  40. Fclose( FILE *f, int safe )
  41. {
  42. if ((safe && (fflush( f ) || (UseFSync && fdatasync( fileno( f ) )))) || fclose( f ) == EOF) {
  43. sys_error( "Error: cannot close file. Disk full?" );
  44. exit( 1 );
  45. }
  46. }
  47. void
  48. Fprintf( FILE *f, const char *msg, ... )
  49. {
  50. int r;
  51. va_list va;
  52. va_start( va, msg );
  53. r = vfprintf( f, msg, va );
  54. va_end( va );
  55. if (r < 0) {
  56. sys_error( "Error: cannot write file. Disk full?" );
  57. exit( 1 );
  58. }
  59. }
  60. static const char Flags[] = { 'D', 'F', 'R', 'S', 'T' };
  61. static int
  62. parse_flags( const char *buf )
  63. {
  64. unsigned flags, i, d;
  65. for (flags = i = d = 0; i < as(Flags); i++)
  66. if (buf[d] == Flags[i]) {
  67. flags |= (1 << i);
  68. d++;
  69. }
  70. return flags;
  71. }
  72. static int
  73. make_flags( int flags, char *buf )
  74. {
  75. unsigned i, d;
  76. for (i = d = 0; i < as(Flags); i++)
  77. if (flags & (1 << i))
  78. buf[d++] = Flags[i];
  79. buf[d] = 0;
  80. return d;
  81. }
  82. #define S_DEAD (1<<0) /* ephemeral: the entry was killed and should be ignored */
  83. #define S_DEL(ms) (1<<(2+(ms))) /* ephemeral: m/s message would be subject to expunge */
  84. #define S_EXPIRED (1<<4) /* the entry is expired (slave message removal confirmed) */
  85. #define S_EXPIRE (1<<5) /* the entry is being expired (slave message removal scheduled) */
  86. #define S_NEXPIRE (1<<6) /* temporary: new expiration state */
  87. #define S_DELETE (1<<7) /* ephemeral: flags propagation is a deletion */
  88. #define mvBit(in,ib,ob) ((unsigned char)(((unsigned)in) * (ob) / (ib)))
  89. typedef struct sync_rec {
  90. struct sync_rec *next;
  91. /* string_list_t *keywords; */
  92. int uid[2]; /* -2 = pending (use tuid), -1 = skipped (too big), 0 = expired */
  93. message_t *msg[2];
  94. unsigned char status, flags, aflags[2], dflags[2];
  95. char tuid[TUIDL];
  96. } sync_rec_t;
  97. /* cases:
  98. a) both non-null
  99. b) only master null
  100. b.1) uid[M] 0
  101. b.2) uid[M] -1
  102. b.3) master not scanned
  103. b.4) master gone
  104. c) only slave null
  105. c.1) uid[S] 0
  106. c.2) uid[S] -1
  107. c.3) slave not scanned
  108. c.4) slave gone
  109. d) both null
  110. d.1) both gone
  111. d.2) uid[M] 0, slave not scanned
  112. d.3) uid[M] -1, slave not scanned
  113. d.4) master gone, slave not scanned
  114. d.5) uid[M] 0, slave gone
  115. d.6) uid[M] -1, slave gone
  116. d.7) uid[S] 0, master not scanned
  117. d.8) uid[S] -1, master not scanned
  118. d.9) slave gone, master not scanned
  119. d.10) uid[S] 0, master gone
  120. d.11) uid[S] -1, master gone
  121. impossible cases: both uid[M] & uid[S] 0 or -1, both not scanned
  122. */
  123. typedef struct {
  124. int t[2];
  125. void (*cb)( int sts, void *aux ), *aux;
  126. char *dname, *jname, *nname, *lname;
  127. FILE *jfp, *nfp;
  128. sync_rec_t *srecs, **srecadd;
  129. channel_conf_t *chan;
  130. store_t *ctx[2];
  131. driver_t *drv[2];
  132. int state[2], ref_count, nsrecs, ret, lfd;
  133. int new_total[2], new_done[2];
  134. int flags_total[2], flags_done[2];
  135. int trash_total[2], trash_done[2];
  136. int maxuid[2]; /* highest UID that was already propagated */
  137. int newmaxuid[2]; /* highest UID that is currently being propagated */
  138. int uidval[2]; /* UID validity value */
  139. int newuid[2]; /* TUID lookup makes sense only for UIDs >= this */
  140. int mmaxxuid; /* highest expired UID on master during new message propagation */
  141. int smaxxuid; /* highest expired UID on slave */
  142. } sync_vars_t;
  143. static void sync_ref( sync_vars_t *svars ) { ++svars->ref_count; }
  144. static int sync_deref( sync_vars_t *svars );
  145. static int deref_check_cancel( sync_vars_t *svars );
  146. static int check_cancel( sync_vars_t *svars );
  147. #define DRIVER_CALL_RET(call) \
  148. do { \
  149. sync_ref( svars ); \
  150. svars->drv[t]->call; \
  151. return deref_check_cancel( svars ); \
  152. } while (0)
  153. #define DRIVER_CALL(call) \
  154. do { \
  155. sync_ref( svars ); \
  156. svars->drv[t]->call; \
  157. if (deref_check_cancel( svars )) \
  158. return; \
  159. } while (0)
  160. #define AUX &svars->t[t]
  161. #define INV_AUX &svars->t[1-t]
  162. #define DECL_SVARS \
  163. int t; \
  164. sync_vars_t *svars
  165. #define INIT_SVARS(aux) \
  166. t = *(int *)aux; \
  167. svars = (sync_vars_t *)(((char *)(&((int *)aux)[-t])) - offsetof(sync_vars_t, t))
  168. #define DECL_INIT_SVARS(aux) \
  169. int t = *(int *)aux; \
  170. sync_vars_t *svars = (sync_vars_t *)(((char *)(&((int *)aux)[-t])) - offsetof(sync_vars_t, t))
  171. /* operation dependencies:
  172. select(S): -
  173. select(M): select(S) | -
  174. new(M), new(S), flags(M): select(M) & select(S)
  175. flags(S): count(new(S))
  176. find_new(x): new(x)
  177. trash(x): flags(x)
  178. close(x): trash(x) & find_new(x) // with expunge
  179. cleanup: close(M) & close(S)
  180. */
  181. #define ST_LOADED (1<<0)
  182. #define ST_FIND_OLD (1<<1)
  183. #define ST_SENT_NEW (1<<2)
  184. #define ST_FIND_NEW (1<<3)
  185. #define ST_FOUND_NEW (1<<4)
  186. #define ST_SENT_FLAGS (1<<5)
  187. #define ST_SENT_TRASH (1<<6)
  188. #define ST_CLOSED (1<<7)
  189. #define ST_SENT_CANCEL (1<<8)
  190. #define ST_CANCELED (1<<9)
  191. #define ST_SELECTED (1<<10)
  192. #define ST_DID_EXPUNGE (1<<11)
  193. static void
  194. match_tuids( sync_vars_t *svars, int t )
  195. {
  196. sync_rec_t *srec;
  197. message_t *tmsg, *ntmsg = 0;
  198. const char *diag;
  199. int num_lost = 0;
  200. for (srec = svars->srecs; srec; srec = srec->next) {
  201. if (srec->status & S_DEAD)
  202. continue;
  203. if (srec->uid[t] == -2 && srec->tuid[0]) {
  204. debug( " pair(%d,%d): lookup %s, TUID %." stringify(TUIDL) "s\n", srec->uid[M], srec->uid[S], str_ms[t], srec->tuid );
  205. for (tmsg = ntmsg; tmsg; tmsg = tmsg->next) {
  206. if (tmsg->status & M_DEAD)
  207. continue;
  208. if (tmsg->tuid[0] && !memcmp( tmsg->tuid, srec->tuid, TUIDL )) {
  209. diag = (tmsg == ntmsg) ? "adjacently" : "after gap";
  210. goto mfound;
  211. }
  212. }
  213. for (tmsg = svars->ctx[t]->msgs; tmsg != ntmsg; tmsg = tmsg->next) {
  214. if (tmsg->status & M_DEAD)
  215. continue;
  216. if (tmsg->tuid[0] && !memcmp( tmsg->tuid, srec->tuid, TUIDL )) {
  217. diag = "after reset";
  218. goto mfound;
  219. }
  220. }
  221. debug( " -> TUID lost\n" );
  222. Fprintf( svars->jfp, "& %d %d\n", srec->uid[M], srec->uid[S] );
  223. srec->flags = 0;
  224. srec->tuid[0] = 0;
  225. num_lost++;
  226. continue;
  227. mfound:
  228. debug( " -> new UID %d %s\n", tmsg->uid, diag );
  229. Fprintf( svars->jfp, "%c %d %d %d\n", "<>"[t], srec->uid[M], srec->uid[S], tmsg->uid );
  230. tmsg->srec = srec;
  231. srec->msg[t] = tmsg;
  232. ntmsg = tmsg->next;
  233. srec->uid[t] = tmsg->uid;
  234. srec->tuid[0] = 0;
  235. }
  236. }
  237. if (num_lost)
  238. warn( "Warning: lost track of %d %sed message(s)\n", num_lost, str_hl[t] );
  239. }
  240. typedef struct copy_vars {
  241. void (*cb)( int sts, int uid, struct copy_vars *vars );
  242. void *aux;
  243. sync_rec_t *srec; /* also ->tuid */
  244. message_t *msg;
  245. msg_data_t data;
  246. } copy_vars_t;
  247. static void msg_fetched( int sts, void *aux );
  248. static int
  249. copy_msg( copy_vars_t *vars )
  250. {
  251. DECL_INIT_SVARS(vars->aux);
  252. t ^= 1;
  253. vars->data.flags = vars->msg->flags;
  254. vars->data.date = svars->chan->use_internal_date ? -1 : 0;
  255. DRIVER_CALL_RET(fetch_msg( svars->ctx[t], vars->msg, &vars->data, msg_fetched, vars ));
  256. }
  257. static void msg_stored( int sts, int uid, void *aux );
  258. static void
  259. msg_fetched( int sts, void *aux )
  260. {
  261. copy_vars_t *vars = (copy_vars_t *)aux;
  262. DECL_SVARS;
  263. char *fmap, *buf;
  264. int i, len, extra, scr, tcr, lcrs, hcrs, bcrs, lines;
  265. int start, sbreak = 0, ebreak = 0;
  266. char c;
  267. switch (sts) {
  268. case DRV_OK:
  269. INIT_SVARS(vars->aux);
  270. if (check_cancel( svars )) {
  271. free( vars->data.data );
  272. vars->cb( SYNC_CANCELED, 0, vars );
  273. return;
  274. }
  275. vars->msg->flags = vars->data.flags;
  276. scr = (svars->drv[1-t]->flags / DRV_CRLF) & 1;
  277. tcr = (svars->drv[t]->flags / DRV_CRLF) & 1;
  278. if (vars->srec || scr != tcr) {
  279. fmap = vars->data.data;
  280. len = vars->data.len;
  281. extra = lines = hcrs = bcrs = i = 0;
  282. if (vars->srec) {
  283. nloop:
  284. start = i;
  285. lcrs = 0;
  286. while (i < len) {
  287. c = fmap[i++];
  288. if (c == '\r')
  289. lcrs++;
  290. else if (c == '\n') {
  291. if (!memcmp( fmap + start, "X-TUID: ", 8 )) {
  292. extra = (sbreak = start) - (ebreak = i);
  293. goto oke;
  294. }
  295. lines++;
  296. hcrs += lcrs;
  297. if (i - lcrs - 1 == start) {
  298. sbreak = ebreak = start;
  299. goto oke;
  300. }
  301. goto nloop;
  302. }
  303. }
  304. /* invalid message */
  305. warn( "Warning: message %d from %s has incomplete header.\n",
  306. vars->msg->uid, str_ms[1-t] );
  307. free( fmap );
  308. vars->cb( SYNC_NOGOOD, 0, vars );
  309. return;
  310. oke:
  311. extra += 8 + TUIDL + 1 + (tcr && (!scr || hcrs));
  312. }
  313. if (tcr != scr) {
  314. for (; i < len; i++) {
  315. c = fmap[i];
  316. if (c == '\r')
  317. bcrs++;
  318. else if (c == '\n')
  319. lines++;
  320. }
  321. extra -= hcrs + bcrs;
  322. if (tcr)
  323. extra += lines;
  324. }
  325. vars->data.len = len + extra;
  326. buf = vars->data.data = nfmalloc( vars->data.len );
  327. i = 0;
  328. if (vars->srec) {
  329. if (tcr != scr) {
  330. if (tcr) {
  331. for (; i < sbreak; i++)
  332. if ((c = fmap[i]) != '\r') {
  333. if (c == '\n')
  334. *buf++ = '\r';
  335. *buf++ = c;
  336. }
  337. } else {
  338. for (; i < sbreak; i++)
  339. if ((c = fmap[i]) != '\r')
  340. *buf++ = c;
  341. }
  342. } else {
  343. memcpy( buf, fmap, sbreak );
  344. buf += sbreak;
  345. }
  346. memcpy( buf, "X-TUID: ", 8 );
  347. buf += 8;
  348. memcpy( buf, vars->srec->tuid, TUIDL );
  349. buf += TUIDL;
  350. if (tcr && (!scr || hcrs))
  351. *buf++ = '\r';
  352. *buf++ = '\n';
  353. i = ebreak;
  354. }
  355. if (tcr != scr) {
  356. if (tcr) {
  357. for (; i < len; i++)
  358. if ((c = fmap[i]) != '\r') {
  359. if (c == '\n')
  360. *buf++ = '\r';
  361. *buf++ = c;
  362. }
  363. } else {
  364. for (; i < len; i++)
  365. if ((c = fmap[i]) != '\r')
  366. *buf++ = c;
  367. }
  368. } else
  369. memcpy( buf, fmap + i, len - i );
  370. free( fmap );
  371. }
  372. svars->drv[t]->store_msg( svars->ctx[t], &vars->data, !vars->srec, msg_stored, vars );
  373. break;
  374. case DRV_CANCELED:
  375. vars->cb( SYNC_CANCELED, 0, vars );
  376. break;
  377. case DRV_MSG_BAD:
  378. vars->cb( SYNC_NOGOOD, 0, vars );
  379. break;
  380. default:
  381. vars->cb( SYNC_FAIL, 0, vars );
  382. break;
  383. }
  384. }
  385. static void
  386. msg_stored( int sts, int uid, void *aux )
  387. {
  388. copy_vars_t *vars = (copy_vars_t *)aux;
  389. DECL_SVARS;
  390. switch (sts) {
  391. case DRV_OK:
  392. vars->cb( SYNC_OK, uid, vars );
  393. break;
  394. case DRV_CANCELED:
  395. vars->cb( SYNC_CANCELED, 0, vars );
  396. break;
  397. case DRV_MSG_BAD:
  398. INIT_SVARS(vars->aux);
  399. (void)svars;
  400. warn( "Warning: %s refuses to store message %d from %s.\n",
  401. str_ms[t], vars->msg->uid, str_ms[1-t] );
  402. vars->cb( SYNC_NOGOOD, 0, vars );
  403. break;
  404. default:
  405. vars->cb( SYNC_FAIL, 0, vars );
  406. break;
  407. }
  408. }
  409. static void
  410. stats( sync_vars_t *svars )
  411. {
  412. char buf[2][64];
  413. char *cs;
  414. int t, l;
  415. static int cols = -1;
  416. if (cols < 0 && (!(cs = getenv( "COLUMNS" )) || !(cols = atoi( cs ) / 2)))
  417. cols = 36;
  418. if (!(DFlags & QUIET)) {
  419. for (t = 0; t < 2; t++) {
  420. l = sprintf( buf[t], "+%d/%d *%d/%d #%d/%d",
  421. svars->new_done[t], svars->new_total[t],
  422. svars->flags_done[t], svars->flags_total[t],
  423. svars->trash_done[t], svars->trash_total[t] );
  424. if (l > cols)
  425. buf[t][cols - 1] = '~';
  426. }
  427. infon( "\v\rM: %.*s S: %.*s", cols, buf[0], cols, buf[1] );
  428. }
  429. }
  430. static void sync_bail( sync_vars_t *svars );
  431. static void sync_bail1( sync_vars_t *svars );
  432. static void sync_bail2( sync_vars_t *svars );
  433. static void sync_bail3( sync_vars_t *svars );
  434. static void cancel_done( void *aux );
  435. static void
  436. cancel_sync( sync_vars_t *svars )
  437. {
  438. int t;
  439. for (t = 0; t < 2; t++) {
  440. int other_state = svars->state[1-t];
  441. if (svars->ret & SYNC_BAD(t)) {
  442. cancel_done( AUX );
  443. } else if (!(svars->state[t] & ST_SENT_CANCEL)) {
  444. /* ignore subsequent failures from in-flight commands */
  445. svars->state[t] |= ST_SENT_CANCEL;
  446. svars->drv[t]->cancel( svars->ctx[t], cancel_done, AUX );
  447. }
  448. if (other_state & ST_CANCELED)
  449. break;
  450. }
  451. }
  452. static void
  453. cancel_done( void *aux )
  454. {
  455. DECL_INIT_SVARS(aux);
  456. svars->state[t] |= ST_CANCELED;
  457. if (svars->state[1-t] & ST_CANCELED) {
  458. if (svars->lfd) {
  459. Fclose( svars->nfp, 0 );
  460. Fclose( svars->jfp, 0 );
  461. sync_bail( svars );
  462. } else {
  463. sync_bail2( svars );
  464. }
  465. }
  466. }
  467. static void
  468. store_bad( void *aux )
  469. {
  470. DECL_INIT_SVARS(aux);
  471. svars->drv[t]->cancel_store( svars->ctx[t] );
  472. svars->ret |= SYNC_BAD(t);
  473. cancel_sync( svars );
  474. }
  475. static int
  476. deref_check_cancel( sync_vars_t *svars )
  477. {
  478. if (sync_deref( svars ))
  479. return -1;
  480. return check_cancel( svars );
  481. }
  482. static int
  483. check_cancel( sync_vars_t *svars )
  484. {
  485. return (svars->state[M] | svars->state[S]) & (ST_SENT_CANCEL | ST_CANCELED);
  486. }
  487. static int
  488. check_ret( int sts, void *aux )
  489. {
  490. DECL_SVARS;
  491. if (sts == DRV_CANCELED)
  492. return 1;
  493. INIT_SVARS(aux);
  494. if (sts == DRV_BOX_BAD) {
  495. svars->ret |= SYNC_FAIL;
  496. cancel_sync( svars );
  497. return 1;
  498. }
  499. return check_cancel( svars );
  500. }
  501. #define SVARS_CHECK_RET \
  502. DECL_SVARS; \
  503. if (check_ret( sts, aux )) \
  504. return; \
  505. INIT_SVARS(aux)
  506. #define SVARS_CHECK_RET_VARS(type) \
  507. type *vars = (type *)aux; \
  508. DECL_SVARS; \
  509. if (check_ret( sts, vars->aux )) { \
  510. free( vars ); \
  511. return; \
  512. } \
  513. INIT_SVARS(vars->aux)
  514. #define SVARS_CHECK_CANCEL_RET \
  515. DECL_SVARS; \
  516. if (sts == SYNC_CANCELED) { \
  517. free( vars ); \
  518. return; \
  519. } \
  520. INIT_SVARS(vars->aux)
  521. static char *
  522. clean_strdup( const char *s )
  523. {
  524. char *cs;
  525. int i;
  526. cs = nfstrdup( s );
  527. for (i = 0; cs[i]; i++)
  528. if (cs[i] == '/')
  529. cs[i] = '!';
  530. return cs;
  531. }
  532. #define JOURNAL_VERSION "2"
  533. static void box_selected( int sts, void *aux );
  534. void
  535. sync_boxes( store_t *ctx[], const char *names[], channel_conf_t *chan,
  536. void (*cb)( int sts, void *aux ), void *aux )
  537. {
  538. sync_vars_t *svars;
  539. int t;
  540. svars = nfcalloc( sizeof(*svars) );
  541. svars->t[1] = 1;
  542. svars->ref_count = 1;
  543. svars->cb = cb;
  544. svars->aux = aux;
  545. svars->ctx[0] = ctx[0];
  546. svars->ctx[1] = ctx[1];
  547. svars->chan = chan;
  548. svars->uidval[0] = svars->uidval[1] = -1;
  549. svars->srecadd = &svars->srecs;
  550. for (t = 0; t < 2; t++) {
  551. ctx[t]->orig_name =
  552. (!names[t] || (ctx[t]->conf->map_inbox && !strcmp( ctx[t]->conf->map_inbox, names[t] ))) ?
  553. "INBOX" : names[t];
  554. if (!ctx[t]->conf->flat_delim) {
  555. ctx[t]->name = nfstrdup( ctx[t]->orig_name );
  556. } else if (map_name( ctx[t]->orig_name, &ctx[t]->name, 0, "/", ctx[t]->conf->flat_delim ) < 0) {
  557. error( "Error: canonical mailbox name '%s' contains flattened hierarchy delimiter\n", ctx[t]->name );
  558. svars->ret = SYNC_FAIL;
  559. sync_bail3( svars );
  560. return;
  561. }
  562. ctx[t]->uidvalidity = -1;
  563. set_bad_callback( ctx[t], store_bad, AUX );
  564. svars->drv[t] = ctx[t]->conf->driver;
  565. }
  566. /* Both boxes must be fully set up at this point, so that error exit paths
  567. * don't run into uninitialized variables. */
  568. for (t = 0; t < 2; t++) {
  569. info( "Selecting %s %s...\n", str_ms[t], ctx[t]->orig_name );
  570. DRIVER_CALL(select( ctx[t], (chan->ops[t] & OP_CREATE) != 0, box_selected, AUX ));
  571. }
  572. }
  573. static int load_box( sync_vars_t *svars, int t, int minwuid, int *mexcs, int nmexcs );
  574. static void
  575. box_selected( int sts, void *aux )
  576. {
  577. DECL_SVARS;
  578. sync_rec_t *srec, *nsrec;
  579. char *s, *cmname, *csname;
  580. store_t *ctx[2];
  581. channel_conf_t *chan;
  582. FILE *jfp;
  583. int opts[2], line, t1, t2, t3;
  584. int *mexcs, nmexcs, rmexcs, minwuid;
  585. struct stat st;
  586. struct flock lck;
  587. char fbuf[16]; /* enlarge when support for keywords is added */
  588. char buf[128], buf1[64], buf2[64];
  589. if (check_ret( sts, aux ))
  590. return;
  591. INIT_SVARS(aux);
  592. ctx[0] = svars->ctx[0];
  593. ctx[1] = svars->ctx[1];
  594. svars->state[t] |= ST_SELECTED;
  595. if (!(svars->state[1-t] & ST_SELECTED))
  596. return;
  597. chan = svars->chan;
  598. if (!strcmp( chan->sync_state ? chan->sync_state : global_conf.sync_state, "*" )) {
  599. if (!ctx[S]->path) {
  600. error( "Error: store '%s' does not support in-box sync state\n", chan->stores[S]->name );
  601. sbail:
  602. svars->ret = SYNC_FAIL;
  603. sync_bail2( svars );
  604. return;
  605. }
  606. nfasprintf( &svars->dname, "%s/." EXE "state", ctx[S]->path );
  607. } else {
  608. csname = clean_strdup( ctx[S]->name );
  609. if (chan->sync_state)
  610. nfasprintf( &svars->dname, "%s%s", chan->sync_state, csname );
  611. else {
  612. cmname = clean_strdup( ctx[M]->name );
  613. nfasprintf( &svars->dname, "%s:%s:%s_:%s:%s", global_conf.sync_state,
  614. chan->stores[M]->name, cmname, chan->stores[S]->name, csname );
  615. free( cmname );
  616. }
  617. free( csname );
  618. if (!(s = strrchr( svars->dname, '/' ))) {
  619. error( "Error: invalid SyncState location '%s'\n", svars->dname );
  620. goto sbail;
  621. }
  622. *s = 0;
  623. if (mkdir( svars->dname, 0700 ) && errno != EEXIST) {
  624. sys_error( "Error: cannot create SyncState directory '%s'", svars->dname );
  625. goto sbail;
  626. }
  627. *s = '/';
  628. }
  629. nfasprintf( &svars->jname, "%s.journal", svars->dname );
  630. nfasprintf( &svars->nname, "%s.new", svars->dname );
  631. nfasprintf( &svars->lname, "%s.lock", svars->dname );
  632. memset( &lck, 0, sizeof(lck) );
  633. #if SEEK_SET != 0
  634. lck.l_whence = SEEK_SET;
  635. #endif
  636. #if F_WRLCK != 0
  637. lck.l_type = F_WRLCK;
  638. #endif
  639. if ((svars->lfd = open( svars->lname, O_WRONLY|O_CREAT, 0666 )) < 0) {
  640. sys_error( "Error: cannot create lock file %s", svars->lname );
  641. svars->ret = SYNC_FAIL;
  642. sync_bail2( svars );
  643. return;
  644. }
  645. if (fcntl( svars->lfd, F_SETLK, &lck )) {
  646. error( "Error: channel :%s:%s-:%s:%s is locked\n",
  647. chan->stores[M]->name, ctx[M]->orig_name, chan->stores[S]->name, ctx[S]->orig_name );
  648. svars->ret = SYNC_FAIL;
  649. sync_bail1( svars );
  650. return;
  651. }
  652. if ((jfp = fopen( svars->dname, "r" ))) {
  653. debug( "reading sync state %s ...\n", svars->dname );
  654. line = 0;
  655. while (fgets( buf, sizeof(buf), jfp )) {
  656. line++;
  657. if (!(t = strlen( buf )) || buf[t - 1] != '\n') {
  658. error( "Error: incomplete sync state header entry at %s:%d\n", svars->dname, line );
  659. jbail:
  660. fclose( jfp );
  661. bail:
  662. svars->ret = SYNC_FAIL;
  663. sync_bail( svars );
  664. return;
  665. }
  666. if (t == 1)
  667. goto gothdr;
  668. if (line == 1 && isdigit( buf[0] )) {
  669. if (sscanf( buf, "%63s %63s", buf1, buf2 ) != 2 ||
  670. sscanf( buf1, "%d:%d", &svars->uidval[M], &svars->maxuid[M] ) < 2 ||
  671. sscanf( buf2, "%d:%d:%d", &svars->uidval[S], &svars->smaxxuid, &svars->maxuid[S] ) < 3) {
  672. error( "Error: invalid sync state header in %s\n", svars->dname );
  673. goto jbail;
  674. }
  675. goto gothdr;
  676. }
  677. if (sscanf( buf, "%63s %d", buf1, &t1 ) != 2) {
  678. error( "Error: malformed sync state header entry at %s:%d\n", svars->dname, line );
  679. goto jbail;
  680. }
  681. if (!strcmp( buf1, "MasterUidValidity" ))
  682. svars->uidval[M] = t1;
  683. else if (!strcmp( buf1, "SlaveUidValidity" ))
  684. svars->uidval[S] = t1;
  685. else if (!strcmp( buf1, "MaxPulledUid" ))
  686. svars->maxuid[M] = t1;
  687. else if (!strcmp( buf1, "MaxPushedUid" ))
  688. svars->maxuid[S] = t1;
  689. else if (!strcmp( buf1, "MaxExpiredSlaveUid" ))
  690. svars->smaxxuid = t1;
  691. else {
  692. error( "Error: unrecognized sync state header entry at %s:%d\n", svars->dname, line );
  693. goto jbail;
  694. }
  695. }
  696. error( "Error: unterminated sync state header in %s\n", svars->dname );
  697. goto jbail;
  698. gothdr:
  699. while (fgets( buf, sizeof(buf), jfp )) {
  700. line++;
  701. if (!(t = strlen( buf )) || buf[t - 1] != '\n') {
  702. error( "Error: incomplete sync state entry at %s:%d\n", svars->dname, line );
  703. goto jbail;
  704. }
  705. fbuf[0] = 0;
  706. if (sscanf( buf, "%d %d %15s", &t1, &t2, fbuf ) < 2) {
  707. error( "Error: invalid sync state entry at %s:%d\n", svars->dname, line );
  708. goto jbail;
  709. }
  710. srec = nfmalloc( sizeof(*srec) );
  711. srec->uid[M] = t1;
  712. srec->uid[S] = t2;
  713. s = fbuf;
  714. if (*s == 'X') {
  715. s++;
  716. srec->status = S_EXPIRE | S_EXPIRED;
  717. } else
  718. srec->status = 0;
  719. srec->flags = parse_flags( s );
  720. debug( " entry (%d,%d,%u,%s)\n", srec->uid[M], srec->uid[S], srec->flags, srec->status & S_EXPIRED ? "X" : "" );
  721. srec->msg[M] = srec->msg[S] = 0;
  722. srec->tuid[0] = 0;
  723. srec->next = 0;
  724. *svars->srecadd = srec;
  725. svars->srecadd = &srec->next;
  726. svars->nsrecs++;
  727. }
  728. fclose( jfp );
  729. } else {
  730. if (errno != ENOENT) {
  731. error( "Error: cannot read sync state %s\n", svars->dname );
  732. goto bail;
  733. }
  734. }
  735. svars->newmaxuid[M] = svars->maxuid[M];
  736. svars->newmaxuid[S] = svars->maxuid[S];
  737. svars->mmaxxuid = INT_MAX;
  738. line = 0;
  739. if ((jfp = fopen( svars->jname, "r" ))) {
  740. if (!stat( svars->nname, &st ) && fgets( buf, sizeof(buf), jfp )) {
  741. debug( "recovering journal ...\n" );
  742. if (!(t = strlen( buf )) || buf[t - 1] != '\n') {
  743. error( "Error: incomplete journal header in %s\n", svars->jname );
  744. goto jbail;
  745. }
  746. if (memcmp( buf, JOURNAL_VERSION "\n", strlen(JOURNAL_VERSION) + 1 )) {
  747. error( "Error: incompatible journal version "
  748. "(got %.*s, expected " JOURNAL_VERSION ")\n", t - 1, buf );
  749. goto jbail;
  750. }
  751. srec = 0;
  752. line = 1;
  753. while (fgets( buf, sizeof(buf), jfp )) {
  754. line++;
  755. if (!(t = strlen( buf )) || buf[t - 1] != '\n') {
  756. error( "Error: incomplete journal entry at %s:%d\n", svars->jname, line );
  757. goto jbail;
  758. }
  759. if (buf[0] == '#' ?
  760. (t3 = 0, (sscanf( buf + 2, "%d %d %n", &t1, &t2, &t3 ) < 2) || !t3 || (t - t3 != TUIDL + 3)) :
  761. buf[0] == '(' || buf[0] == ')' || buf[0] == '{' || buf[0] == '}' || buf[0] == '!' ?
  762. (sscanf( buf + 2, "%d", &t1 ) != 1) :
  763. buf[0] == '+' || buf[0] == '&' || buf[0] == '-' || buf[0] == '|' || buf[0] == '/' || buf[0] == '\\' ?
  764. (sscanf( buf + 2, "%d %d", &t1, &t2 ) != 2) :
  765. (sscanf( buf + 2, "%d %d %d", &t1, &t2, &t3 ) != 3))
  766. {
  767. error( "Error: malformed journal entry at %s:%d\n", svars->jname, line );
  768. goto jbail;
  769. }
  770. if (buf[0] == '(')
  771. svars->maxuid[M] = t1;
  772. else if (buf[0] == ')')
  773. svars->maxuid[S] = t1;
  774. else if (buf[0] == '{')
  775. svars->newuid[M] = t1;
  776. else if (buf[0] == '}')
  777. svars->newuid[S] = t1;
  778. else if (buf[0] == '!')
  779. svars->smaxxuid = t1;
  780. else if (buf[0] == '|') {
  781. svars->uidval[M] = t1;
  782. svars->uidval[S] = t2;
  783. } else if (buf[0] == '+') {
  784. srec = nfmalloc( sizeof(*srec) );
  785. srec->uid[M] = t1;
  786. srec->uid[S] = t2;
  787. if (svars->newmaxuid[M] < t1)
  788. svars->newmaxuid[M] = t1;
  789. if (svars->newmaxuid[S] < t2)
  790. svars->newmaxuid[S] = t2;
  791. debug( " new entry(%d,%d)\n", t1, t2 );
  792. srec->msg[M] = srec->msg[S] = 0;
  793. srec->status = 0;
  794. srec->flags = 0;
  795. srec->tuid[0] = 0;
  796. srec->next = 0;
  797. *svars->srecadd = srec;
  798. svars->srecadd = &srec->next;
  799. svars->nsrecs++;
  800. } else {
  801. for (nsrec = srec; srec; srec = srec->next)
  802. if (srec->uid[M] == t1 && srec->uid[S] == t2)
  803. goto syncfnd;
  804. for (srec = svars->srecs; srec != nsrec; srec = srec->next)
  805. if (srec->uid[M] == t1 && srec->uid[S] == t2)
  806. goto syncfnd;
  807. error( "Error: journal entry at %s:%d refers to non-existing sync state entry\n", svars->jname, line );
  808. goto jbail;
  809. syncfnd:
  810. debugn( " entry(%d,%d,%u) ", srec->uid[M], srec->uid[S], srec->flags );
  811. switch (buf[0]) {
  812. case '-':
  813. debug( "killed\n" );
  814. if (srec->msg[M])
  815. srec->msg[M]->srec = 0;
  816. srec->status = S_DEAD;
  817. break;
  818. case '#':
  819. debug( "TUID now %." stringify(TUIDL) "s\n", buf + t3 + 2 );
  820. memcpy( srec->tuid, buf + t3 + 2, TUIDL );
  821. break;
  822. case '&':
  823. debug( "TUID %." stringify(TUIDL) "s lost\n", srec->tuid );
  824. srec->flags = 0;
  825. srec->tuid[0] = 0;
  826. break;
  827. case '<':
  828. debug( "master now %d\n", t3 );
  829. srec->uid[M] = t3;
  830. srec->tuid[0] = 0;
  831. break;
  832. case '>':
  833. debug( "slave now %d\n", t3 );
  834. srec->uid[S] = t3;
  835. srec->tuid[0] = 0;
  836. break;
  837. case '*':
  838. debug( "flags now %d\n", t3 );
  839. srec->flags = t3;
  840. break;
  841. case '~':
  842. debug( "expire now %d\n", t3 );
  843. if (t3)
  844. srec->status |= S_EXPIRE;
  845. else
  846. srec->status &= ~S_EXPIRE;
  847. break;
  848. case '\\':
  849. t3 = (srec->status & S_EXPIRED);
  850. debug( "expire back to %d\n", t3 / S_EXPIRED );
  851. if (t3)
  852. srec->status |= S_EXPIRE;
  853. else
  854. srec->status &= ~S_EXPIRE;
  855. break;
  856. case '/':
  857. t3 = (srec->status & S_EXPIRE);
  858. debug( "expired now %d\n", t3 / S_EXPIRE );
  859. if (t3) {
  860. if (svars->smaxxuid < srec->uid[S])
  861. svars->smaxxuid = srec->uid[S];
  862. srec->status |= S_EXPIRED;
  863. } else
  864. srec->status &= ~S_EXPIRED;
  865. break;
  866. default:
  867. error( "Error: unrecognized journal entry at %s:%d\n", svars->jname, line );
  868. goto jbail;
  869. }
  870. }
  871. }
  872. }
  873. fclose( jfp );
  874. } else {
  875. if (errno != ENOENT) {
  876. error( "Error: cannot read journal %s\n", svars->jname );
  877. goto bail;
  878. }
  879. }
  880. t1 = 0;
  881. for (t = 0; t < 2; t++)
  882. if (svars->uidval[t] >= 0 && svars->uidval[t] != ctx[t]->uidvalidity) {
  883. error( "Error: UIDVALIDITY of %s changed (got %d, expected %d)\n",
  884. str_ms[t], ctx[t]->uidvalidity, svars->uidval[t] );
  885. t1++;
  886. }
  887. if (t1)
  888. goto bail;
  889. if (!(svars->nfp = fopen( svars->nname, "w" ))) {
  890. error( "Error: cannot write new sync state %s\n", svars->nname );
  891. goto bail;
  892. }
  893. if (!(svars->jfp = fopen( svars->jname, "a" ))) {
  894. error( "Error: cannot write journal %s\n", svars->jname );
  895. fclose( svars->nfp );
  896. goto bail;
  897. }
  898. setlinebuf( svars->jfp );
  899. if (!line)
  900. Fprintf( svars->jfp, JOURNAL_VERSION "\n" );
  901. opts[M] = opts[S] = 0;
  902. for (t = 0; t < 2; t++) {
  903. if (chan->ops[t] & (OP_DELETE|OP_FLAGS)) {
  904. opts[t] |= OPEN_SETFLAGS;
  905. opts[1-t] |= OPEN_OLD;
  906. if (chan->ops[t] & OP_FLAGS)
  907. opts[1-t] |= OPEN_FLAGS;
  908. }
  909. if (chan->ops[t] & (OP_NEW|OP_RENEW)) {
  910. opts[t] |= OPEN_APPEND;
  911. if (chan->ops[t] & OP_RENEW)
  912. opts[1-t] |= OPEN_OLD;
  913. if (chan->ops[t] & OP_NEW)
  914. opts[1-t] |= OPEN_NEW;
  915. if (chan->ops[t] & OP_EXPUNGE)
  916. opts[1-t] |= OPEN_FLAGS;
  917. if (chan->stores[t]->max_size != INT_MAX)
  918. opts[1-t] |= OPEN_SIZE;
  919. }
  920. if (chan->ops[t] & OP_EXPUNGE) {
  921. opts[t] |= OPEN_EXPUNGE;
  922. if (chan->stores[t]->trash) {
  923. if (!chan->stores[t]->trash_only_new)
  924. opts[t] |= OPEN_OLD;
  925. opts[t] |= OPEN_NEW|OPEN_FLAGS;
  926. } else if (chan->stores[1-t]->trash && chan->stores[1-t]->trash_remote_new)
  927. opts[t] |= OPEN_NEW|OPEN_FLAGS;
  928. }
  929. }
  930. if ((chan->ops[S] & (OP_NEW|OP_RENEW|OP_FLAGS)) && chan->max_messages)
  931. opts[S] |= OPEN_OLD|OPEN_NEW|OPEN_FLAGS;
  932. if (line)
  933. for (srec = svars->srecs; srec; srec = srec->next) {
  934. if (srec->status & S_DEAD)
  935. continue;
  936. if (srec->tuid[0]) {
  937. if (srec->uid[M] == -2)
  938. opts[M] |= OPEN_NEW|OPEN_FIND, svars->state[M] |= ST_FIND_OLD;
  939. else if (srec->uid[S] == -2)
  940. opts[S] |= OPEN_NEW|OPEN_FIND, svars->state[S] |= ST_FIND_OLD;
  941. else
  942. assert( !"sync record with stray TUID" );
  943. }
  944. }
  945. svars->drv[M]->prepare_opts( ctx[M], opts[M] );
  946. svars->drv[S]->prepare_opts( ctx[S], opts[S] );
  947. mexcs = 0;
  948. nmexcs = rmexcs = 0;
  949. if (svars->ctx[M]->opts & OPEN_OLD) {
  950. if (chan->max_messages) {
  951. /* When messages have been expired on the slave, the master fetch is split into
  952. * two ranges: The bulk fetch which corresponds with the most recent messages, and an
  953. * exception list of messages which would have been expired if they weren't important. */
  954. debug( "preparing master selection - max expired slave uid is %d\n", svars->smaxxuid );
  955. /* First, find out the lower bound for the bulk fetch. */
  956. minwuid = INT_MAX;
  957. for (srec = svars->srecs; srec; srec = srec->next) {
  958. if (srec->status & S_DEAD)
  959. continue;
  960. if (srec->status & S_EXPIRED) {
  961. if (!srec->uid[S]) {
  962. /* The expired message was already gone. */
  963. continue;
  964. }
  965. /* The expired message was not expunged yet, so re-examine it.
  966. * This will happen en masse, so just extend the bulk fetch. */
  967. } else {
  968. if (svars->smaxxuid >= srec->uid[S]) {
  969. /* The non-expired message is in the generally expired range, so don't
  970. * make it contribute to the bulk fetch. */
  971. continue;
  972. }
  973. /* Usual non-expired message. */
  974. }
  975. if (minwuid > srec->uid[M])
  976. minwuid = srec->uid[M];
  977. }
  978. debug( " min non-orphaned master uid is %d\n", minwuid );
  979. /* Next, calculate the exception fetch. */
  980. for (srec = svars->srecs; srec; srec = srec->next) {
  981. if (srec->status & S_DEAD)
  982. continue;
  983. if (srec->uid[M] > 0 && srec->uid[S] > 0 && minwuid > srec->uid[M] &&
  984. (!(svars->ctx[M]->opts & OPEN_NEW) || svars->maxuid[M] >= srec->uid[M])) {
  985. /* The pair is alive, but outside the bulk range. */
  986. if (nmexcs == rmexcs) {
  987. rmexcs = rmexcs * 2 + 100;
  988. mexcs = nfrealloc( mexcs, rmexcs * sizeof(int) );
  989. }
  990. mexcs[nmexcs++] = srec->uid[M];
  991. }
  992. }
  993. debugn( " exception list is:" );
  994. for (t = 0; t < nmexcs; t++)
  995. debugn( " %d", mexcs[t] );
  996. debug( "\n" );
  997. } else {
  998. minwuid = 1;
  999. }
  1000. } else {
  1001. minwuid = INT_MAX;
  1002. }
  1003. if (load_box( svars, M, minwuid, mexcs, nmexcs ))
  1004. return;
  1005. load_box( svars, S, (ctx[S]->opts & OPEN_OLD) ? 1 : INT_MAX, 0, 0 );
  1006. }
  1007. static void box_loaded( int sts, void *aux );
  1008. static int
  1009. load_box( sync_vars_t *svars, int t, int minwuid, int *mexcs, int nmexcs )
  1010. {
  1011. sync_rec_t *srec;
  1012. int maxwuid;
  1013. if (svars->ctx[t]->opts & OPEN_NEW) {
  1014. if (minwuid > svars->maxuid[t] + 1)
  1015. minwuid = svars->maxuid[t] + 1;
  1016. maxwuid = INT_MAX;
  1017. } else if (svars->ctx[t]->opts & OPEN_OLD) {
  1018. maxwuid = 0;
  1019. for (srec = svars->srecs; srec; srec = srec->next)
  1020. if (!(srec->status & S_DEAD) && srec->uid[t] > maxwuid)
  1021. maxwuid = srec->uid[t];
  1022. } else
  1023. maxwuid = 0;
  1024. info( "Loading %s...\n", str_ms[t] );
  1025. debug( maxwuid == INT_MAX ? "loading %s [%d,inf]\n" : "loading %s [%d,%d]\n", str_ms[t], minwuid, maxwuid );
  1026. DRIVER_CALL_RET(load( svars->ctx[t], minwuid, maxwuid, svars->newuid[t], mexcs, nmexcs, box_loaded, AUX ));
  1027. }
  1028. typedef struct {
  1029. void *aux;
  1030. sync_rec_t *srec;
  1031. int aflags, dflags;
  1032. } flag_vars_t;
  1033. typedef struct {
  1034. int uid;
  1035. sync_rec_t *srec;
  1036. } sync_rec_map_t;
  1037. static void flags_set( int sts, void *aux );
  1038. static void flags_set_p2( sync_vars_t *svars, sync_rec_t *srec, int t );
  1039. static int msgs_flags_set( sync_vars_t *svars, int t );
  1040. static void msg_copied( int sts, int uid, copy_vars_t *vars );
  1041. static void msg_copied_p2( sync_vars_t *svars, sync_rec_t *srec, int t, int uid );
  1042. static void msgs_copied( sync_vars_t *svars, int t );
  1043. static void
  1044. box_loaded( int sts, void *aux )
  1045. {
  1046. DECL_SVARS;
  1047. sync_rec_t *srec;
  1048. sync_rec_map_t *srecmap;
  1049. message_t *tmsg;
  1050. copy_vars_t *cv;
  1051. flag_vars_t *fv;
  1052. int uid, no[2], del[2], alive, todel, t1, t2;
  1053. int sflags, nflags, aflags, dflags, nex;
  1054. unsigned hashsz, idx;
  1055. char fbuf[16]; /* enlarge when support for keywords is added */
  1056. if (check_ret( sts, aux ))
  1057. return;
  1058. INIT_SVARS(aux);
  1059. svars->state[t] |= ST_LOADED;
  1060. info( "%s: %d messages, %d recent\n", str_ms[t], svars->ctx[t]->count, svars->ctx[t]->recent );
  1061. if (svars->state[t] & ST_FIND_OLD) {
  1062. debug( "matching previously copied messages on %s\n", str_ms[t] );
  1063. match_tuids( svars, t );
  1064. }
  1065. debug( "matching messages on %s against sync records\n", str_ms[t] );
  1066. hashsz = bucketsForSize( svars->nsrecs * 3 );
  1067. srecmap = nfcalloc( hashsz * sizeof(*srecmap) );
  1068. for (srec = svars->srecs; srec; srec = srec->next) {
  1069. if (srec->status & S_DEAD)
  1070. continue;
  1071. uid = srec->uid[t];
  1072. idx = (unsigned)((unsigned)uid * 1103515245U) % hashsz;
  1073. while (srecmap[idx].uid)
  1074. if (++idx == hashsz)
  1075. idx = 0;
  1076. srecmap[idx].uid = uid;
  1077. srecmap[idx].srec = srec;
  1078. }
  1079. for (tmsg = svars->ctx[t]->msgs; tmsg; tmsg = tmsg->next) {
  1080. if (tmsg->srec) /* found by TUID */
  1081. continue;
  1082. uid = tmsg->uid;
  1083. if (DFlags & DEBUG) {
  1084. make_flags( tmsg->flags, fbuf );
  1085. printf( svars->ctx[t]->opts & OPEN_SIZE ? " message %5d, %-4s, %6lu: " : " message %5d, %-4s: ", uid, fbuf, tmsg->size );
  1086. }
  1087. idx = (unsigned)((unsigned)uid * 1103515245U) % hashsz;
  1088. while (srecmap[idx].uid) {
  1089. if (srecmap[idx].uid == uid) {
  1090. srec = srecmap[idx].srec;
  1091. goto found;
  1092. }
  1093. if (++idx == hashsz)
  1094. idx = 0;
  1095. }
  1096. debug( "new\n" );
  1097. continue;
  1098. found:
  1099. tmsg->srec = srec;
  1100. srec->msg[t] = tmsg;
  1101. debug( "pairs %5d\n", srec->uid[1-t] );
  1102. }
  1103. free( srecmap );
  1104. if (!(svars->state[1-t] & ST_LOADED))
  1105. return;
  1106. if (svars->uidval[M] < 0 || svars->uidval[S] < 0) {
  1107. svars->uidval[M] = svars->ctx[M]->uidvalidity;
  1108. svars->uidval[S] = svars->ctx[S]->uidvalidity;
  1109. Fprintf( svars->jfp, "| %d %d\n", svars->uidval[M], svars->uidval[S] );
  1110. }
  1111. info( "Synchronizing...\n" );
  1112. debug( "synchronizing old entries\n" );
  1113. for (srec = svars->srecs; srec; srec = srec->next) {
  1114. if (srec->status & S_DEAD)
  1115. continue;
  1116. debug( "pair (%d,%d)\n", srec->uid[M], srec->uid[S] );
  1117. no[M] = !srec->msg[M] && (svars->ctx[M]->opts & OPEN_OLD);
  1118. no[S] = !srec->msg[S] && (svars->ctx[S]->opts & OPEN_OLD);
  1119. if (no[M] && no[S]) {
  1120. debug( " vanished\n" );
  1121. /* d.1) d.5) d.6) d.10) d.11) */
  1122. srec->status = S_DEAD;
  1123. Fprintf( svars->jfp, "- %d %d\n", srec->uid[M], srec->uid[S] );
  1124. } else {
  1125. del[M] = no[M] && (srec->uid[M] > 0);
  1126. del[S] = no[S] && (srec->uid[S] > 0);
  1127. for (t = 0; t < 2; t++) {
  1128. srec->aflags[t] = srec->dflags[t] = 0;
  1129. if (srec->msg[t] && (srec->msg[t]->flags & F_DELETED))
  1130. srec->status |= S_DEL(t);
  1131. /* excludes (push) c.3) d.2) d.3) d.4) / (pull) b.3) d.7) d.8) d.9) */
  1132. if (!srec->uid[t]) {
  1133. /* b.1) / c.1) */
  1134. debug( " no more %s\n", str_ms[t] );
  1135. } else if (del[1-t]) {
  1136. /* c.4) d.9) / b.4) d.4) */
  1137. if ((t == M) && (srec->status & (S_EXPIRE|S_EXPIRED))) {
  1138. /* Don't propagate deletion resulting from expiration. */
  1139. debug( " slave expired, orphaning master\n" );
  1140. Fprintf( svars->jfp, "> %d %d 0\n", srec->uid[M], srec->uid[S] );
  1141. srec->uid[S] = 0;
  1142. } else {
  1143. if (srec->msg[t] && (srec->msg[t]->status & M_FLAGS) && srec->msg[t]->flags != srec->flags)
  1144. info( "Info: conflicting changes in (%d,%d)\n", srec->uid[M], srec->uid[S] );
  1145. if (svars->chan->ops[t] & OP_DELETE) {
  1146. debug( " %sing delete\n", str_hl[t] );
  1147. srec->aflags[t] = F_DELETED;
  1148. srec->status |= S_DELETE;
  1149. } else {
  1150. debug( " not %sing delete\n", str_hl[t] );
  1151. }
  1152. }
  1153. } else if (!srec->msg[1-t])
  1154. /* c.1) c.2) d.7) d.8) / b.1) b.2) d.2) d.3) */
  1155. ;
  1156. else if (srec->uid[t] < 0)
  1157. /* b.2) / c.2) */
  1158. ; /* handled as new messages (sort of) */
  1159. else if (!del[t]) {
  1160. /* a) & b.3) / c.3) */
  1161. if (svars->chan->ops[t] & OP_FLAGS) {
  1162. sflags = srec->msg[1-t]->flags;
  1163. if ((t == M) && (srec->status & (S_EXPIRE|S_EXPIRED))) {
  1164. /* Don't propagate deletion resulting from expiration. */
  1165. debug( " slave expiring\n" );
  1166. sflags &= ~F_DELETED;
  1167. }
  1168. srec->aflags[t] = sflags & ~srec->flags;
  1169. srec->dflags[t] = ~sflags & srec->flags;
  1170. if (DFlags & DEBUG) {
  1171. char afbuf[16], dfbuf[16]; /* enlarge when support for keywords is added */
  1172. make_flags( srec->aflags[t], afbuf );
  1173. make_flags( srec->dflags[t], dfbuf );
  1174. debug( " %sing flags: +%s -%s\n", str_hl[t], afbuf, dfbuf );
  1175. }
  1176. } else
  1177. debug( " not %sing flags\n", str_hl[t] );
  1178. } /* else b.4) / c.4) */
  1179. }
  1180. }
  1181. }
  1182. debug( "synchronizing new entries\n" );
  1183. for (t = 0; t < 2; t++) {
  1184. for (tmsg = svars->ctx[1-t]->msgs; tmsg; tmsg = tmsg->next) {
  1185. /* If we have a srec:
  1186. * - message is old (> 0) or expired (0) => ignore
  1187. * - message was skipped (-1) => ReNew
  1188. * - message was attempted, but failed (-2) => New
  1189. * If new have no srec, the message is always New. If messages were previously ignored
  1190. * due to being excessive, they would now appear to be newer than the messages that
  1191. * got actually synced, so make sure to look only at the newest ones. As some messages
  1192. * may be already propagated before an interruption, and maxuid logging is delayed,
  1193. * we need to track the newmaxuid separately. */
  1194. srec = tmsg->srec;
  1195. if (srec ? srec->uid[t] < 0 && (svars->chan->ops[t] & (srec->uid[t] == -1 ? OP_RENEW : OP_NEW))
  1196. : svars->newmaxuid[1-t] < tmsg->uid && (svars->chan->ops[t] & OP_NEW)) {
  1197. debug( "new message %d on %s\n", tmsg->uid, str_ms[1-t] );
  1198. if ((svars->chan->ops[t] & OP_EXPUNGE) && (tmsg->flags & F_DELETED)) {
  1199. debug( " -> not %sing - would be expunged anyway\n", str_hl[t] );
  1200. } else {
  1201. if (srec) {
  1202. debug( " -> pair(%d,%d) exists\n", srec->uid[M], srec->uid[S] );
  1203. } else {
  1204. srec = nfmalloc( sizeof(*srec) );
  1205. srec->next = 0;
  1206. *svars->srecadd = srec;
  1207. svars->srecadd = &srec->next;
  1208. svars->nsrecs++;
  1209. srec->status = 0;
  1210. srec->flags = 0;
  1211. srec->tuid[0] = 0;
  1212. srec->uid[1-t] = tmsg->uid;
  1213. srec->uid[t] = -2;
  1214. srec->msg[1-t] = tmsg;
  1215. srec->msg[t] = 0;
  1216. tmsg->srec = srec;
  1217. if (svars->newmaxuid[1-t] < tmsg->uid)
  1218. svars->newmaxuid[1-t] = tmsg->uid;
  1219. Fprintf( svars->jfp, "+ %d %d\n", srec->uid[M], srec->uid[S] );
  1220. debug( " -> pair(%d,%d) created\n", srec->uid[M], srec->uid[S] );
  1221. }
  1222. if (svars->maxuid[1-t] < tmsg->uid) {
  1223. /* We do this here for simplicity. However, logging must be delayed until
  1224. * all messages were propagated, as skipped messages could otherwise be
  1225. * logged before the propagation of messages with lower UIDs completes. */
  1226. svars->maxuid[1-t] = tmsg->uid;
  1227. }
  1228. if ((tmsg->flags & F_FLAGGED) || tmsg->size <= svars->chan->stores[t]->max_size) {
  1229. if (tmsg->flags) {
  1230. srec->flags = tmsg->flags;
  1231. Fprintf( svars->jfp, "* %d %d %u\n", srec->uid[M], srec->uid[S], srec->flags );
  1232. debug( " -> updated flags to %u\n", tmsg->flags );
  1233. }
  1234. for (t1 = 0; t1 < TUIDL; t1++) {
  1235. t2 = arc4_getbyte() & 0x3f;
  1236. srec->tuid[t1] = t2 < 26 ? t2 + 'A' : t2 < 52 ? t2 + 'a' - 26 : t2 < 62 ? t2 + '0' - 52 : t2 == 62 ? '+' : '/';
  1237. }
  1238. Fprintf( svars->jfp, "# %d %d %." stringify(TUIDL) "s\n", srec->uid[M], srec->uid[S], srec->tuid );
  1239. debug( " -> %sing message, TUID %." stringify(TUIDL) "s\n", str_hl[t], srec->tuid );
  1240. } else {
  1241. if (srec->uid[t] == -1) {
  1242. debug( " -> not %sing - still too big\n", str_hl[t] );
  1243. } else {
  1244. debug( " -> not %sing - too big\n", str_hl[t] );
  1245. msg_copied_p2( svars, srec, t, -1 );
  1246. }
  1247. }
  1248. }
  1249. }
  1250. }
  1251. }
  1252. if ((svars->chan->ops[S] & (OP_NEW|OP_RENEW|OP_FLAGS)) && svars->chan->max_messages) {
  1253. /* Note: When this branch is entered, we have loaded all slave messages. */
  1254. /* Expire excess messages. Important (flagged, unread, or unpropagated) messages
  1255. * older than the first not expired message are not counted towards the total. */
  1256. debug( "preparing message expiration\n" );
  1257. alive = 0;
  1258. for (tmsg = svars->ctx[S]->msgs; tmsg; tmsg = tmsg->next) {
  1259. if (tmsg->status & M_DEAD)
  1260. continue;
  1261. if ((srec = tmsg->srec) && srec->uid[M] > 0 &&
  1262. ((tmsg->flags | srec->aflags[S]) & ~srec->dflags[S] & F_DELETED) &&
  1263. !(srec->status & (S_EXPIRE|S_EXPIRED))) {
  1264. /* Message was not propagated yet, or is deleted. */
  1265. } else {
  1266. alive++;
  1267. }
  1268. }
  1269. for (tmsg = svars->ctx[M]->msgs; tmsg; tmsg = tmsg->next) {
  1270. if ((srec = tmsg->srec) && srec->tuid[0] && !(tmsg->flags & F_DELETED))
  1271. alive++;
  1272. }
  1273. todel = alive - svars->chan->max_messages;
  1274. debug( "%d alive messages, %d excess - expiring\n", alive, todel );
  1275. alive = 0;
  1276. for (tmsg = svars->ctx[S]->msgs; tmsg; tmsg = tmsg->next) {
  1277. if (tmsg->status & M_DEAD)
  1278. continue;
  1279. if (!(srec = tmsg->srec) || srec->uid[M] <= 0) {
  1280. /* We did not push the message, so it must be kept. */
  1281. debug( " old pair(%d,%d) unpropagated\n", srec->uid[M], srec->uid[S] );
  1282. todel--;
  1283. } else {
  1284. nflags = (tmsg->flags | srec->aflags[S]) & ~srec->dflags[S];
  1285. if (!(nflags & F_DELETED) || (srec->status & (S_EXPIRE|S_EXPIRED))) {
  1286. /* The message is not deleted, or is already (being) expired. */
  1287. if ((nflags & F_FLAGGED) || !((nflags & F_SEEN) || ((void)(todel > 0 && alive++), svars->chan->expire_unread > 0))) {
  1288. /* Important messages are always kept. */
  1289. debug( " old pair(%d,%d) important\n", srec->uid[M], srec->uid[S] );
  1290. todel--;
  1291. } else if (todel > 0 ||
  1292. ((srec->status & (S_EXPIRE|S_EXPIRED)) == (S_EXPIRE|S_EXPIRED)) ||
  1293. ((srec->status & (S_EXPIRE|S_EXPIRED)) && (tmsg->flags & F_DELETED))) {
  1294. /* The message is excess or was already (being) expired. */
  1295. srec->status |= S_NEXPIRE;
  1296. debug( " old pair(%d,%d) expired\n", srec->uid[M], srec->uid[S] );
  1297. todel--;
  1298. }
  1299. }
  1300. }
  1301. }
  1302. for (tmsg = svars->ctx[M]->msgs; tmsg; tmsg = tmsg->next) {
  1303. if ((srec = tmsg->srec) && srec->tuid[0]) {
  1304. nflags = tmsg->flags;
  1305. if (!(nflags & F_DELETED)) {
  1306. if ((nflags & F_FLAGGED) || !((nflags & F_SEEN) || ((void)(todel > 0 && alive++), svars->chan->expire_unread > 0))) {
  1307. /* Important messages are always fetched. */
  1308. debug( " new pair(%d,%d) important\n", srec->uid[M], srec->uid[S] );
  1309. todel--;
  1310. } else if (todel > 0) {
  1311. /* The message is excess. */
  1312. srec->status |= S_NEXPIRE;
  1313. debug( " new pair(%d,%d) expired\n", srec->uid[M], srec->uid[S] );
  1314. svars->mmaxxuid = srec->uid[M];
  1315. todel--;
  1316. }
  1317. }
  1318. }
  1319. }
  1320. debug( "%d excess messages remain\n", todel );
  1321. if (svars->chan->expire_unread < 0 && (unsigned)alive * 2 > svars->chan->max_messages) {
  1322. error( "%s: %d unread messages in excess of MaxMessages (%d).\n"
  1323. "Please set ExpireUnread to decide outcome. Skipping mailbox.\n",
  1324. svars->ctx[S]->orig_name, alive, svars->chan->max_messages );
  1325. svars->ret |= SYNC_FAIL;
  1326. cancel_sync( svars );
  1327. return;
  1328. }
  1329. for (srec = svars->srecs; srec; srec = srec->next) {
  1330. if (srec->status & S_DEAD)
  1331. continue;
  1332. if (!srec->tuid[0]) {
  1333. if (!srec->msg[S])
  1334. continue;
  1335. nex = (srec->status / S_NEXPIRE) & 1;
  1336. if (nex != ((srec->status / S_EXPIRED) & 1)) {
  1337. /* The record needs a state change ... */
  1338. if (nex != ((srec->status / S_EXPIRE) & 1)) {
  1339. /* ... and we need to start a transaction. */
  1340. Fprintf( svars->jfp, "~ %d %d %d\n", srec->uid[M], srec->uid[S], nex );
  1341. debug( " pair(%d,%d): %d (pre)\n", srec->uid[M], srec->uid[S], nex );
  1342. srec->status = (srec->status & ~S_EXPIRE) | (nex * S_EXPIRE);
  1343. } else {
  1344. /* ... but the "right" transaction is already pending. */
  1345. debug( " pair(%d,%d): %d (pending)\n", srec->uid[M], srec->uid[S], nex );
  1346. }
  1347. } else {
  1348. /* Note: the "wrong" transaction may be pending here,
  1349. * e.g.: S_NEXPIRE = 0, S_EXPIRE = 1, S_EXPIRED = 0. */
  1350. }
  1351. } else {
  1352. if (srec->status & S_NEXPIRE) {
  1353. Fprintf( svars->jfp, "- %d %d\n", srec->uid[M], srec->uid[S] );
  1354. debug( " pair(%d,%d): 1 (abort)\n", srec->uid[M], srec->uid[S] );
  1355. srec->msg[M]->srec = 0;
  1356. srec->status = S_DEAD;
  1357. }
  1358. }
  1359. }
  1360. }
  1361. debug( "synchronizing flags\n" );
  1362. for (srec = svars->srecs; srec; srec = srec->next) {
  1363. if ((srec->status & S_DEAD) || srec->uid[M] <= 0 || srec->uid[S] <= 0)
  1364. continue;
  1365. for (t = 0; t < 2; t++) {
  1366. aflags = srec->aflags[t];
  1367. dflags = srec->dflags[t];
  1368. if (srec->status & S_DELETE) {
  1369. if (!aflags) {
  1370. /* This deletion propagation goes the other way round. */
  1371. continue;
  1372. }
  1373. } else {
  1374. /* The trigger is an expiration transaction being ongoing ... */
  1375. if ((t == S) && ((mvBit(srec->status, S_EXPIRE, S_EXPIRED) ^ srec->status) & S_EXPIRED)) {
  1376. /* ... but the actual action derives from the wanted state. */
  1377. if (srec->status & S_NEXPIRE)
  1378. aflags |= F_DELETED;
  1379. else
  1380. dflags |= F_DELETED;
  1381. }
  1382. }
  1383. if ((svars->chan->ops[t] & OP_EXPUNGE) && (((srec->msg[t] ? srec->msg[t]->flags : 0) | aflags) & ~dflags & F_DELETED) &&
  1384. (!svars->ctx[t]->conf->trash || svars->ctx[t]->conf->trash_only_new))
  1385. {
  1386. /* If the message is going to be expunged, don't propagate anything but the deletion. */
  1387. srec->aflags[t] &= F_DELETED;
  1388. aflags &= F_DELETED;
  1389. srec->dflags[t] = dflags = 0;
  1390. }
  1391. if (srec->msg[t] && (srec->msg[t]->status & M_FLAGS)) {
  1392. /* If we know the target message's state, optimize away non-changes. */
  1393. aflags &= ~srec->msg[t]->flags;
  1394. dflags &= srec->msg[t]->flags;
  1395. }
  1396. if (aflags | dflags) {
  1397. svars->flags_total[t]++;
  1398. stats( svars );
  1399. fv = nfmalloc( sizeof(*fv) );
  1400. fv->aux = AUX;
  1401. fv->srec = srec;
  1402. fv->aflags = aflags;
  1403. fv->dflags = dflags;
  1404. DRIVER_CALL(set_flags( svars->ctx[t], srec->msg[t], srec->uid[t], aflags, dflags, flags_set, fv ));
  1405. } else
  1406. flags_set_p2( svars, srec, t );
  1407. }
  1408. }
  1409. for (t = 0; t < 2; t++) {
  1410. svars->drv[t]->commit( svars->ctx[t] );
  1411. svars->state[t] |= ST_SENT_FLAGS;
  1412. if (msgs_flags_set( svars, t ))
  1413. return;
  1414. }
  1415. debug( "propagating new messages\n" );
  1416. if (UseFSync)
  1417. fdatasync( fileno( svars->jfp ) );
  1418. for (t = 0; t < 2; t++) {
  1419. Fprintf( svars->jfp, "%c %d\n", "{}"[t], svars->ctx[t]->uidnext );
  1420. for (tmsg = svars->ctx[1-t]->msgs; tmsg; tmsg = tmsg->next) {
  1421. if ((srec = tmsg->srec) && srec->tuid[0]) {
  1422. svars->new_total[t]++;
  1423. stats( svars );
  1424. cv = nfmalloc( sizeof(*cv) );
  1425. cv->cb = msg_copied;
  1426. cv->aux = AUX;
  1427. cv->srec = srec;
  1428. cv->msg = tmsg;
  1429. if (copy_msg( cv ))
  1430. return;
  1431. }
  1432. }
  1433. svars->state[t] |= ST_SENT_NEW;
  1434. msgs_copied( svars, t );
  1435. }
  1436. }
  1437. static void
  1438. msg_copied( int sts, int uid, copy_vars_t *vars )
  1439. {
  1440. SVARS_CHECK_CANCEL_RET;
  1441. switch (sts) {
  1442. case SYNC_OK:
  1443. if (uid < 0)
  1444. svars->state[t] |= ST_FIND_NEW;
  1445. msg_copied_p2( svars, vars->srec, t, uid );
  1446. break;
  1447. case SYNC_NOGOOD:
  1448. debug( " -> killing (%d,%d)\n", vars->srec->uid[M], vars->srec->uid[S] );
  1449. vars->srec->status = S_DEAD;
  1450. Fprintf( svars->jfp, "- %d %d\n", vars->srec->uid[M], vars->srec->uid[S] );
  1451. break;
  1452. default:
  1453. cancel_sync( svars );
  1454. free( vars );
  1455. return;
  1456. }
  1457. free( vars );
  1458. svars->new_done[t]++;
  1459. stats( svars );
  1460. msgs_copied( svars, t );
  1461. }
  1462. static void
  1463. msg_copied_p2( sync_vars_t *svars, sync_rec_t *srec, int t, int uid )
  1464. {
  1465. /* Possible previous UIDs:
  1466. * - -2 when the entry is new
  1467. * - -1 when re-newing an entry
  1468. * Possible new UIDs:
  1469. * - a real UID when storing a message to a UIDPLUS mailbox
  1470. * - -2 when storing a message to a dumb mailbox
  1471. * - -1 when not actually storing a message */
  1472. if (srec->uid[t] != uid) {
  1473. debug( " -> new UID %d on %s\n", uid, str_ms[t] );
  1474. Fprintf( svars->jfp, "%c %d %d %d\n", "<>"[t], srec->uid[M], srec->uid[S], uid );
  1475. srec->uid[t] = uid;
  1476. srec->tuid[0] = 0;
  1477. }
  1478. if (t == S && svars->mmaxxuid < srec->uid[M]) {
  1479. /* If we have so many new messages that some of them are instantly expired,
  1480. * but some are still propagated because they are important, we need to
  1481. * ensure explicitly that the bulk fetch limit is upped. */
  1482. svars->mmaxxuid = INT_MAX;
  1483. if (svars->smaxxuid < srec->uid[S] - 1) {
  1484. svars->smaxxuid = srec->uid[S] - 1;
  1485. Fprintf( svars->jfp, "! %d\n", svars->smaxxuid );
  1486. }
  1487. }
  1488. }
  1489. static void msgs_found_new( int sts, void *aux );
  1490. static void msgs_new_done( sync_vars_t *svars, int t );
  1491. static void sync_close( sync_vars_t *svars, int t );
  1492. static void
  1493. msgs_copied( sync_vars_t *svars, int t )
  1494. {
  1495. if (!(svars->state[t] & ST_SENT_NEW) || svars->new_done[t] < svars->new_total[t])
  1496. return;
  1497. Fprintf( svars->jfp, "%c %d\n", ")("[t], svars->maxuid[1-t] );
  1498. if (svars->state[t] & ST_FIND_NEW) {
  1499. debug( "finding just copied messages on %s\n", str_ms[t] );
  1500. svars->drv[t]->find_new_msgs( svars->ctx[t], msgs_found_new, AUX );
  1501. } else {
  1502. msgs_new_done( svars, t );
  1503. }
  1504. }
  1505. static void
  1506. msgs_found_new( int sts, void *aux )
  1507. {
  1508. SVARS_CHECK_RET;
  1509. switch (sts) {
  1510. case DRV_OK:
  1511. debug( "matching just copied messages on %s\n", str_ms[t] );
  1512. break;
  1513. default:
  1514. warn( "Warning: cannot find newly stored messages on %s.\n", str_ms[t] );
  1515. break;
  1516. }
  1517. match_tuids( svars, t );
  1518. msgs_new_done( svars, t );
  1519. }
  1520. static void
  1521. msgs_new_done( sync_vars_t *svars, int t )
  1522. {
  1523. svars->state[t] |= ST_FOUND_NEW;
  1524. sync_close( svars, t );
  1525. }
  1526. static void
  1527. flags_set( int sts, void *aux )
  1528. {
  1529. SVARS_CHECK_RET_VARS(flag_vars_t);
  1530. switch (sts) {
  1531. case DRV_OK:
  1532. if (vars->aflags & F_DELETED)
  1533. vars->srec->status |= S_DEL(t);
  1534. else if (vars->dflags & F_DELETED)
  1535. vars->srec->status &= ~S_DEL(t);
  1536. flags_set_p2( svars, vars->srec, t );
  1537. break;
  1538. }
  1539. free( vars );
  1540. svars->flags_done[t]++;
  1541. stats( svars );
  1542. msgs_flags_set( svars, t );
  1543. }
  1544. static void
  1545. flags_set_p2( sync_vars_t *svars, sync_rec_t *srec, int t )
  1546. {
  1547. if (srec->status & S_DELETE) {
  1548. debug( " pair(%d,%d): resetting %s UID\n", srec->uid[M], srec->uid[S], str_ms[1-t] );
  1549. Fprintf( svars->jfp, "%c %d %d 0\n", "><"[t], srec->uid[M], srec->uid[S] );
  1550. srec->uid[1-t] = 0;
  1551. } else {
  1552. int nflags = (srec->flags | srec->aflags[t]) & ~srec->dflags[t];
  1553. if (srec->flags != nflags) {
  1554. debug( " pair(%d,%d): updating flags (%u -> %u; %sed)\n", srec->uid[M], srec->uid[S], srec->flags, nflags, str_hl[t] );
  1555. srec->flags = nflags;
  1556. Fprintf( svars->jfp, "* %d %d %u\n", srec->uid[M], srec->uid[S], nflags );
  1557. }
  1558. if (t == S) {
  1559. int nex = (srec->status / S_NEXPIRE) & 1;
  1560. if (nex != ((srec->status / S_EXPIRED) & 1)) {
  1561. if (nex && (svars->smaxxuid < srec->uid[S]))
  1562. svars->smaxxuid = srec->uid[S];
  1563. Fprintf( svars->jfp, "/ %d %d\n", srec->uid[M], srec->uid[S] );
  1564. debug( " pair(%d,%d): expired %d (commit)\n", srec->uid[M], srec->uid[S], nex );
  1565. srec->status = (srec->status & ~S_EXPIRED) | (nex * S_EXPIRED);
  1566. } else if (nex != ((srec->status / S_EXPIRE) & 1)) {
  1567. Fprintf( svars->jfp, "\\ %d %d\n", srec->uid[M], srec->uid[S] );
  1568. debug( " pair(%d,%d): expire %d (cancel)\n", srec->uid[M], srec->uid[S], nex );
  1569. srec->status = (srec->status & ~S_EXPIRE) | (nex * S_EXPIRE);
  1570. }
  1571. }
  1572. }
  1573. }
  1574. static void msg_trashed( int sts, void *aux );
  1575. static void msg_rtrashed( int sts, int uid, copy_vars_t *vars );
  1576. static int
  1577. msgs_flags_set( sync_vars_t *svars, int t )
  1578. {
  1579. message_t *tmsg;
  1580. copy_vars_t *cv;
  1581. if (!(svars->state[t] & ST_SENT_FLAGS) || svars->flags_done[t] < svars->flags_total[t])
  1582. return 0;
  1583. if ((svars->chan->ops[t] & OP_EXPUNGE) &&
  1584. (svars->ctx[t]->conf->trash || (svars->ctx[1-t]->conf->trash && svars->ctx[1-t]->conf->trash_remote_new))) {
  1585. debug( "trashing in %s\n", str_ms[t] );
  1586. for (tmsg = svars->ctx[t]->msgs; tmsg; tmsg = tmsg->next)
  1587. if ((tmsg->flags & F_DELETED) && (t == M || !tmsg->srec || !(tmsg->srec->status & (S_EXPIRE|S_EXPIRED)))) {
  1588. if (svars->ctx[t]->conf->trash) {
  1589. if (!svars->ctx[t]->conf->trash_only_new || !tmsg->srec || tmsg->srec->uid[1-t] < 0) {
  1590. debug( "%s: trashing message %d\n", str_ms[t], tmsg->uid );
  1591. svars->trash_total[t]++;
  1592. stats( svars );
  1593. sync_ref( svars );
  1594. svars->drv[t]->trash_msg( svars->ctx[t], tmsg, msg_trashed, AUX );
  1595. if (deref_check_cancel( svars ))
  1596. return -1;
  1597. } else
  1598. debug( "%s: not trashing message %d - not new\n", str_ms[t], tmsg->uid );
  1599. } else {
  1600. if (!tmsg->srec || tmsg->srec->uid[1-t] < 0) {
  1601. if (tmsg->size <= svars->ctx[1-t]->conf->max_size) {
  1602. debug( "%s: remote trashing message %d\n", str_ms[t], tmsg->uid );
  1603. svars->trash_total[t]++;
  1604. stats( svars );
  1605. cv = nfmalloc( sizeof(*cv) );
  1606. cv->cb = msg_rtrashed;
  1607. cv->aux = INV_AUX;
  1608. cv->srec = 0;
  1609. cv->msg = tmsg;
  1610. if (copy_msg( cv ))
  1611. return -1;
  1612. } else
  1613. debug( "%s: not remote trashing message %d - too big\n", str_ms[t], tmsg->uid );
  1614. } else
  1615. debug( "%s: not remote trashing message %d - not new\n", str_ms[t], tmsg->uid );
  1616. }
  1617. }
  1618. }
  1619. svars->state[t] |= ST_SENT_TRASH;
  1620. sync_close( svars, t );
  1621. return 0;
  1622. }
  1623. static void
  1624. msg_trashed( int sts, void *aux )
  1625. {
  1626. DECL_SVARS;
  1627. if (sts == DRV_MSG_BAD)
  1628. sts = DRV_BOX_BAD;
  1629. if (check_ret( sts, aux ))
  1630. return;
  1631. INIT_SVARS(aux);
  1632. svars->trash_done[t]++;
  1633. stats( svars );
  1634. sync_close( svars, t );
  1635. }
  1636. static void
  1637. msg_rtrashed( int sts, int uid ATTR_UNUSED, copy_vars_t *vars )
  1638. {
  1639. SVARS_CHECK_CANCEL_RET;
  1640. switch (sts) {
  1641. case SYNC_OK:
  1642. case SYNC_NOGOOD: /* the message is gone or heavily busted */
  1643. break;
  1644. default:
  1645. cancel_sync( svars );
  1646. free( vars );
  1647. return;
  1648. }
  1649. free( vars );
  1650. t ^= 1;
  1651. svars->trash_done[t]++;
  1652. stats( svars );
  1653. sync_close( svars, t );
  1654. }
  1655. static void box_closed( int sts, void *aux );
  1656. static void box_closed_p2( sync_vars_t *svars, int t );
  1657. static void
  1658. sync_close( sync_vars_t *svars, int t )
  1659. {
  1660. if ((~svars->state[t] & (ST_FOUND_NEW|ST_SENT_TRASH)) ||
  1661. svars->trash_done[t] < svars->trash_total[t])
  1662. return;
  1663. if ((svars->chan->ops[t] & OP_EXPUNGE) /*&& !(svars->state[t] & ST_TRASH_BAD)*/) {
  1664. debug( "expunging %s\n", str_ms[t] );
  1665. svars->drv[t]->close( svars->ctx[t], box_closed, AUX );
  1666. } else {
  1667. box_closed_p2( svars, t );
  1668. }
  1669. }
  1670. static void
  1671. box_closed( int sts, void *aux )
  1672. {
  1673. SVARS_CHECK_RET;
  1674. svars->state[t] |= ST_DID_EXPUNGE;
  1675. box_closed_p2( svars, t );
  1676. }
  1677. static void
  1678. box_closed_p2( sync_vars_t *svars, int t )
  1679. {
  1680. sync_rec_t *srec;
  1681. int minwuid;
  1682. char fbuf[16]; /* enlarge when support for keywords is added */
  1683. svars->state[t] |= ST_CLOSED;
  1684. if (!(svars->state[1-t] & ST_CLOSED))
  1685. return;
  1686. if (((svars->state[M] | svars->state[S]) & ST_DID_EXPUNGE) || svars->chan->max_messages) {
  1687. /* This cleanup is not strictly necessary, as the next full sync
  1688. would throw out the dead entries anyway. But ... */
  1689. debug( "purging obsolete entries\n" );
  1690. minwuid = INT_MAX;
  1691. if (svars->chan->max_messages) {
  1692. debug( " max expired slave uid is %d\n", svars->smaxxuid );
  1693. for (srec = svars->srecs; srec; srec = srec->next) {
  1694. if (srec->status & S_DEAD)
  1695. continue;
  1696. if (!((srec->uid[S] <= 0 || ((srec->status & S_DEL(S)) && (svars->state[S] & ST_DID_EXPUNGE))) &&
  1697. (srec->uid[M] <= 0 || ((srec->status & S_DEL(M)) && (svars->state[M] & ST_DID_EXPUNGE)) || (srec->status & S_EXPIRED))) &&
  1698. svars->smaxxuid < srec->uid[S] && minwuid > srec->uid[M])
  1699. minwuid = srec->uid[M];
  1700. }
  1701. debug( " min non-orphaned master uid is %d\n", minwuid );
  1702. }
  1703. for (srec = svars->srecs; srec; srec = srec->next) {
  1704. if (srec->status & S_DEAD)
  1705. continue;
  1706. if (srec->uid[S] <= 0 || ((srec->status & S_DEL(S)) && (svars->state[S] & ST_DID_EXPUNGE))) {
  1707. if (srec->uid[M] <= 0 || ((srec->status & S_DEL(M)) && (svars->state[M] & ST_DID_EXPUNGE)) ||
  1708. ((srec->status & S_EXPIRED) && svars->maxuid[M] >= srec->uid[M] && minwuid > srec->uid[M])) {
  1709. debug( " -> killing (%d,%d)\n", srec->uid[M], srec->uid[S] );
  1710. srec->status = S_DEAD;
  1711. Fprintf( svars->jfp, "- %d %d\n", srec->uid[M], srec->uid[S] );
  1712. } else if (srec->uid[S] > 0) {
  1713. debug( " -> orphaning (%d,[%d])\n", srec->uid[M], srec->uid[S] );
  1714. Fprintf( svars->jfp, "> %d %d 0\n", srec->uid[M], srec->uid[S] );
  1715. srec->uid[S] = 0;
  1716. }
  1717. } else if (srec->uid[M] > 0 && ((srec->status & S_DEL(M)) && (svars->state[M] & ST_DID_EXPUNGE))) {
  1718. debug( " -> orphaning ([%d],%d)\n", srec->uid[M], srec->uid[S] );
  1719. Fprintf( svars->jfp, "< %d %d 0\n", srec->uid[M], srec->uid[S] );
  1720. srec->uid[M] = 0;
  1721. }
  1722. }
  1723. }
  1724. Fprintf( svars->nfp,
  1725. "MasterUidValidity %d\nSlaveUidValidity %d\nMaxPulledUid %d\nMaxPushedUid %d\n",
  1726. svars->uidval[M], svars->uidval[S], svars->maxuid[M], svars->maxuid[S] );
  1727. if (svars->smaxxuid)
  1728. Fprintf( svars->nfp, "MaxExpiredSlaveUid %d\n", svars->smaxxuid );
  1729. Fprintf( svars->nfp, "\n" );
  1730. for (srec = svars->srecs; srec; srec = srec->next) {
  1731. if (srec->status & S_DEAD)
  1732. continue;
  1733. make_flags( srec->flags, fbuf );
  1734. Fprintf( svars->nfp, "%d %d %s%s\n", srec->uid[M], srec->uid[S],
  1735. srec->status & S_EXPIRED ? "X" : "", fbuf );
  1736. }
  1737. Fclose( svars->nfp, 1 );
  1738. Fclose( svars->jfp, 0 );
  1739. if (!(DFlags & KEEPJOURNAL)) {
  1740. /* order is important! */
  1741. rename( svars->nname, svars->dname );
  1742. unlink( svars->jname );
  1743. }
  1744. sync_bail( svars );
  1745. }
  1746. static void
  1747. sync_bail( sync_vars_t *svars )
  1748. {
  1749. sync_rec_t *srec, *nsrec;
  1750. for (srec = svars->srecs; srec; srec = nsrec) {
  1751. nsrec = srec->next;
  1752. free( srec );
  1753. }
  1754. unlink( svars->lname );
  1755. sync_bail1( svars );
  1756. }
  1757. static void
  1758. sync_bail1( sync_vars_t *svars )
  1759. {
  1760. close( svars->lfd );
  1761. sync_bail2( svars );
  1762. }
  1763. static void
  1764. sync_bail2( sync_vars_t *svars )
  1765. {
  1766. free( svars->lname );
  1767. free( svars->nname );
  1768. free( svars->jname );
  1769. free( svars->dname );
  1770. flushn();
  1771. sync_bail3( svars );
  1772. }
  1773. static void
  1774. sync_bail3( sync_vars_t *svars )
  1775. {
  1776. free( svars->ctx[M]->name );
  1777. free( svars->ctx[S]->name );
  1778. sync_deref( svars );
  1779. }
  1780. static int sync_deref( sync_vars_t *svars )
  1781. {
  1782. if (!--svars->ref_count) {
  1783. void (*cb)( int sts, void *aux ) = svars->cb;
  1784. void *aux = svars->aux;
  1785. int ret = svars->ret;
  1786. free( svars );
  1787. cb( ret, aux );
  1788. return -1;
  1789. }
  1790. return 0;
  1791. }