sync.c 64 KB

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