io.c 44 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931
  1. /*
  2. * Socket and pipe I/O utilities used in rsync.
  3. *
  4. * Copyright (C) 1996-2001 Andrew Tridgell
  5. * Copyright (C) 1996 Paul Mackerras
  6. * Copyright (C) 2001, 2002 Martin Pool <mbp@samba.org>
  7. * Copyright (C) 2003-2008 Wayne Davison
  8. *
  9. * This program is free software; you can redistribute it and/or modify
  10. * it under the terms of the GNU General Public License as published by
  11. * the Free Software Foundation; either version 3 of the License, or
  12. * (at your option) any later version.
  13. *
  14. * This program is distributed in the hope that it will be useful,
  15. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  16. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  17. * GNU General Public License for more details.
  18. *
  19. * You should have received a copy of the GNU General Public License along
  20. * with this program; if not, visit the http://fsf.org website.
  21. */
  22. /* Rsync provides its own multiplexing system, which is used to send
  23. * stderr and stdout over a single socket.
  24. *
  25. * For historical reasons this is off during the start of the
  26. * connection, but it's switched on quite early using
  27. * io_start_multiplex_out() and io_start_multiplex_in(). */
  28. #include "rsync.h"
  29. #include "ifuncs.h"
  30. /** If no timeout is specified then use a 60 second select timeout */
  31. #define SELECT_TIMEOUT 60
  32. extern int bwlimit;
  33. extern size_t bwlimit_writemax;
  34. extern int io_timeout;
  35. extern int allowed_lull;
  36. extern int am_server;
  37. extern int am_daemon;
  38. extern int am_sender;
  39. extern int am_generator;
  40. extern int inc_recurse;
  41. extern int io_error;
  42. extern int eol_nulls;
  43. extern int flist_eof;
  44. extern int list_only;
  45. extern int read_batch;
  46. extern int csum_length;
  47. extern int protect_args;
  48. extern int checksum_seed;
  49. extern int protocol_version;
  50. extern int remove_source_files;
  51. extern int preserve_hard_links;
  52. extern struct stats stats;
  53. extern struct file_list *cur_flist;
  54. #ifdef ICONV_OPTION
  55. extern int filesfrom_convert;
  56. extern iconv_t ic_send, ic_recv;
  57. #endif
  58. const char phase_unknown[] = "unknown";
  59. int ignore_timeout = 0;
  60. int batch_fd = -1;
  61. int msgdone_cnt = 0;
  62. /* Ignore an EOF error if non-zero. See whine_about_eof(). */
  63. int kluge_around_eof = 0;
  64. int msg_fd_in = -1;
  65. int msg_fd_out = -1;
  66. int sock_f_in = -1;
  67. int sock_f_out = -1;
  68. static int iobuf_f_in = -1;
  69. static char *iobuf_in;
  70. static size_t iobuf_in_siz;
  71. static size_t iobuf_in_ndx;
  72. static size_t iobuf_in_remaining;
  73. static int iobuf_f_out = -1;
  74. static char *iobuf_out;
  75. static int iobuf_out_cnt;
  76. int flist_forward_from = -1;
  77. static int io_multiplexing_out;
  78. static int io_multiplexing_in;
  79. static time_t last_io_in;
  80. static time_t last_io_out;
  81. static int no_flush;
  82. static int write_batch_monitor_in = -1;
  83. static int write_batch_monitor_out = -1;
  84. static int io_filesfrom_f_in = -1;
  85. static int io_filesfrom_f_out = -1;
  86. static xbuf ff_buf = EMPTY_XBUF;
  87. static char ff_lastchar;
  88. #ifdef ICONV_OPTION
  89. static xbuf iconv_buf = EMPTY_XBUF;
  90. #endif
  91. static int defer_forwarding_messages = 0, keep_defer_forwarding = 0;
  92. static int select_timeout = SELECT_TIMEOUT;
  93. static int active_filecnt = 0;
  94. static OFF_T active_bytecnt = 0;
  95. static int first_message = 1;
  96. static char int_byte_extra[64] = {
  97. 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, /* (00 - 3F)/4 */
  98. 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, /* (40 - 7F)/4 */
  99. 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, /* (80 - BF)/4 */
  100. 2, 2, 2, 2, 2, 2, 2, 2, 3, 3, 3, 3, 4, 4, 5, 6, /* (C0 - FF)/4 */
  101. };
  102. #define REMOTE_OPTION_ERROR "rsync: on remote machine: -"
  103. #define REMOTE_OPTION_ERROR2 ": unknown option"
  104. enum festatus { FES_SUCCESS, FES_REDO, FES_NO_SEND };
  105. static void readfd(int fd, char *buffer, size_t N);
  106. static void writefd(int fd, const char *buf, size_t len);
  107. static void writefd_unbuffered(int fd, const char *buf, size_t len);
  108. static void mplex_write(int fd, enum msgcode code, const char *buf, size_t len, int convert);
  109. struct flist_ndx_item {
  110. struct flist_ndx_item *next;
  111. int ndx;
  112. };
  113. struct flist_ndx_list {
  114. struct flist_ndx_item *head, *tail;
  115. };
  116. static struct flist_ndx_list redo_list, hlink_list;
  117. struct msg_list_item {
  118. struct msg_list_item *next;
  119. char convert;
  120. char buf[1];
  121. };
  122. struct msg_list {
  123. struct msg_list_item *head, *tail;
  124. };
  125. static struct msg_list msg_queue;
  126. static void flist_ndx_push(struct flist_ndx_list *lp, int ndx)
  127. {
  128. struct flist_ndx_item *item;
  129. if (!(item = new(struct flist_ndx_item)))
  130. out_of_memory("flist_ndx_push");
  131. item->next = NULL;
  132. item->ndx = ndx;
  133. if (lp->tail)
  134. lp->tail->next = item;
  135. else
  136. lp->head = item;
  137. lp->tail = item;
  138. }
  139. static int flist_ndx_pop(struct flist_ndx_list *lp)
  140. {
  141. struct flist_ndx_item *next;
  142. int ndx;
  143. if (!lp->head)
  144. return -1;
  145. ndx = lp->head->ndx;
  146. next = lp->head->next;
  147. free(lp->head);
  148. lp->head = next;
  149. if (!next)
  150. lp->tail = NULL;
  151. return ndx;
  152. }
  153. static void got_flist_entry_status(enum festatus status, const char *buf)
  154. {
  155. int ndx = IVAL(buf, 0);
  156. struct file_list *flist = flist_for_ndx(ndx, "got_flist_entry_status");
  157. if (remove_source_files) {
  158. active_filecnt--;
  159. active_bytecnt -= F_LENGTH(flist->files[ndx - flist->ndx_start]);
  160. }
  161. if (inc_recurse)
  162. flist->in_progress--;
  163. switch (status) {
  164. case FES_SUCCESS:
  165. if (remove_source_files)
  166. send_msg(MSG_SUCCESS, buf, 4, 0);
  167. if (preserve_hard_links) {
  168. struct file_struct *file = flist->files[ndx - flist->ndx_start];
  169. if (F_IS_HLINKED(file)) {
  170. flist_ndx_push(&hlink_list, ndx);
  171. flist->in_progress++;
  172. }
  173. }
  174. break;
  175. case FES_REDO:
  176. if (inc_recurse)
  177. flist->to_redo++;
  178. flist_ndx_push(&redo_list, ndx);
  179. break;
  180. case FES_NO_SEND:
  181. break;
  182. }
  183. }
  184. static void check_timeout(void)
  185. {
  186. time_t t;
  187. if (!io_timeout || ignore_timeout)
  188. return;
  189. if (!last_io_in) {
  190. last_io_in = time(NULL);
  191. return;
  192. }
  193. t = time(NULL);
  194. if (t - last_io_in >= io_timeout) {
  195. if (!am_server && !am_daemon) {
  196. rprintf(FERROR, "io timeout after %d seconds -- exiting\n",
  197. (int)(t-last_io_in));
  198. }
  199. exit_cleanup(RERR_TIMEOUT);
  200. }
  201. }
  202. /* Note the fds used for the main socket (which might really be a pipe
  203. * for a local transfer, but we can ignore that). */
  204. void io_set_sock_fds(int f_in, int f_out)
  205. {
  206. sock_f_in = f_in;
  207. sock_f_out = f_out;
  208. }
  209. void set_io_timeout(int secs)
  210. {
  211. io_timeout = secs;
  212. if (!io_timeout || io_timeout > SELECT_TIMEOUT)
  213. select_timeout = SELECT_TIMEOUT;
  214. else
  215. select_timeout = io_timeout;
  216. allowed_lull = read_batch ? 0 : (io_timeout + 1) / 2;
  217. }
  218. /* Setup the fd used to receive MSG_* messages. Only needed during the
  219. * early stages of being a local sender (up through the sending of the
  220. * file list) or when we're the generator (to fetch the messages from
  221. * the receiver). */
  222. void set_msg_fd_in(int fd)
  223. {
  224. msg_fd_in = fd;
  225. }
  226. /* Setup the fd used to send our MSG_* messages. Only needed when
  227. * we're the receiver (to send our messages to the generator). */
  228. void set_msg_fd_out(int fd)
  229. {
  230. msg_fd_out = fd;
  231. set_nonblocking(msg_fd_out);
  232. }
  233. /* Add a message to the pending MSG_* list. */
  234. static void msg_list_add(struct msg_list *lst, int code, const char *buf, int len, int convert)
  235. {
  236. struct msg_list_item *m;
  237. int sz = len + 4 + sizeof m[0] - 1;
  238. if (!(m = (struct msg_list_item *)new_array(char, sz)))
  239. out_of_memory("msg_list_add");
  240. m->next = NULL;
  241. m->convert = convert;
  242. SIVAL(m->buf, 0, ((code+MPLEX_BASE)<<24) | len);
  243. memcpy(m->buf + 4, buf, len);
  244. if (lst->tail)
  245. lst->tail->next = m;
  246. else
  247. lst->head = m;
  248. lst->tail = m;
  249. }
  250. static inline int flush_a_msg(int fd)
  251. {
  252. struct msg_list_item *m = msg_queue.head;
  253. int len = IVAL(m->buf, 0) & 0xFFFFFF;
  254. int tag = *((uchar*)m->buf+3) - MPLEX_BASE;
  255. if (!(msg_queue.head = m->next))
  256. msg_queue.tail = NULL;
  257. defer_forwarding_messages++;
  258. mplex_write(fd, tag, m->buf + 4, len, m->convert);
  259. defer_forwarding_messages--;
  260. free(m);
  261. return len;
  262. }
  263. static void msg_flush(void)
  264. {
  265. if (am_generator) {
  266. while (msg_queue.head && io_multiplexing_out)
  267. stats.total_written += flush_a_msg(sock_f_out) + 4;
  268. } else {
  269. while (msg_queue.head)
  270. (void)flush_a_msg(msg_fd_out);
  271. }
  272. }
  273. static void check_for_d_option_error(const char *msg)
  274. {
  275. static char rsync263_opts[] = "BCDHIKLPRSTWabceghlnopqrtuvxz";
  276. char *colon;
  277. int saw_d = 0;
  278. if (*msg != 'r'
  279. || strncmp(msg, REMOTE_OPTION_ERROR, sizeof REMOTE_OPTION_ERROR - 1) != 0)
  280. return;
  281. msg += sizeof REMOTE_OPTION_ERROR - 1;
  282. if (*msg == '-' || (colon = strchr(msg, ':')) == NULL
  283. || strncmp(colon, REMOTE_OPTION_ERROR2, sizeof REMOTE_OPTION_ERROR2 - 1) != 0)
  284. return;
  285. for ( ; *msg != ':'; msg++) {
  286. if (*msg == 'd')
  287. saw_d = 1;
  288. else if (*msg == 'e')
  289. break;
  290. else if (strchr(rsync263_opts, *msg) == NULL)
  291. return;
  292. }
  293. if (saw_d) {
  294. rprintf(FWARNING,
  295. "*** Try using \"--old-d\" if remote rsync is <= 2.6.3 ***\n");
  296. }
  297. }
  298. /* Read a message from the MSG_* fd and handle it. This is called either
  299. * during the early stages of being a local sender (up through the sending
  300. * of the file list) or when we're the generator (to fetch the messages
  301. * from the receiver). */
  302. static void read_msg_fd(void)
  303. {
  304. char buf[2048];
  305. size_t n;
  306. struct file_list *flist;
  307. int fd = msg_fd_in;
  308. int tag, len;
  309. /* Temporarily disable msg_fd_in. This is needed to avoid looping back
  310. * to this routine from writefd_unbuffered(). */
  311. no_flush++;
  312. msg_fd_in = -1;
  313. defer_forwarding_messages++;
  314. readfd(fd, buf, 4);
  315. tag = IVAL(buf, 0);
  316. len = tag & 0xFFFFFF;
  317. tag = (tag >> 24) - MPLEX_BASE;
  318. switch (tag) {
  319. case MSG_DONE:
  320. if (len < 0 || len > 1 || !am_generator) {
  321. invalid_msg:
  322. rprintf(FERROR, "invalid message %d:%d [%s%s]\n",
  323. tag, len, who_am_i(),
  324. inc_recurse ? "/inc" : "");
  325. exit_cleanup(RERR_STREAMIO);
  326. }
  327. if (len) {
  328. readfd(fd, buf, len);
  329. stats.total_read = read_varlong(fd, 3);
  330. }
  331. msgdone_cnt++;
  332. break;
  333. case MSG_REDO:
  334. if (len != 4 || !am_generator)
  335. goto invalid_msg;
  336. readfd(fd, buf, 4);
  337. got_flist_entry_status(FES_REDO, buf);
  338. break;
  339. case MSG_FLIST:
  340. if (len != 4 || !am_generator || !inc_recurse)
  341. goto invalid_msg;
  342. readfd(fd, buf, 4);
  343. /* Read extra file list from receiver. */
  344. assert(iobuf_in != NULL);
  345. assert(iobuf_f_in == fd);
  346. if (verbose > 3) {
  347. rprintf(FINFO, "[%s] receiving flist for dir %d\n",
  348. who_am_i(), IVAL(buf,0));
  349. }
  350. flist = recv_file_list(fd);
  351. flist->parent_ndx = IVAL(buf,0);
  352. #ifdef SUPPORT_HARD_LINKS
  353. if (preserve_hard_links)
  354. match_hard_links(flist);
  355. #endif
  356. break;
  357. case MSG_FLIST_EOF:
  358. if (len != 0 || !am_generator || !inc_recurse)
  359. goto invalid_msg;
  360. flist_eof = 1;
  361. break;
  362. case MSG_IO_ERROR:
  363. if (len != 4)
  364. goto invalid_msg;
  365. readfd(fd, buf, len);
  366. io_error |= IVAL(buf, 0);
  367. break;
  368. case MSG_DELETED:
  369. if (len >= (int)sizeof buf || !am_generator)
  370. goto invalid_msg;
  371. readfd(fd, buf, len);
  372. send_msg(MSG_DELETED, buf, len, 1);
  373. break;
  374. case MSG_SUCCESS:
  375. if (len != 4 || !am_generator)
  376. goto invalid_msg;
  377. readfd(fd, buf, 4);
  378. got_flist_entry_status(FES_SUCCESS, buf);
  379. break;
  380. case MSG_NO_SEND:
  381. if (len != 4 || !am_generator)
  382. goto invalid_msg;
  383. readfd(fd, buf, 4);
  384. got_flist_entry_status(FES_NO_SEND, buf);
  385. break;
  386. case MSG_ERROR_SOCKET:
  387. case MSG_ERROR_UTF8:
  388. case MSG_CLIENT:
  389. if (!am_generator)
  390. goto invalid_msg;
  391. if (tag == MSG_ERROR_SOCKET)
  392. io_end_multiplex_out();
  393. /* FALL THROUGH */
  394. case MSG_INFO:
  395. case MSG_ERROR:
  396. case MSG_ERROR_XFER:
  397. case MSG_WARNING:
  398. case MSG_LOG:
  399. while (len) {
  400. n = len;
  401. if (n >= sizeof buf)
  402. n = sizeof buf - 1;
  403. readfd(fd, buf, n);
  404. rwrite((enum logcode)tag, buf, n, !am_generator);
  405. len -= n;
  406. }
  407. break;
  408. default:
  409. rprintf(FERROR, "unknown message %d:%d [%s]\n",
  410. tag, len, who_am_i());
  411. exit_cleanup(RERR_STREAMIO);
  412. }
  413. no_flush--;
  414. msg_fd_in = fd;
  415. if (!--defer_forwarding_messages && !no_flush)
  416. msg_flush();
  417. }
  418. /* This is used by the generator to limit how many file transfers can
  419. * be active at once when --remove-source-files is specified. Without
  420. * this, sender-side deletions were mostly happening at the end. */
  421. void increment_active_files(int ndx, int itemizing, enum logcode code)
  422. {
  423. /* TODO: tune these limits? */
  424. while (active_filecnt >= (active_bytecnt >= 128*1024 ? 10 : 50)) {
  425. check_for_finished_files(itemizing, code, 0);
  426. if (iobuf_out_cnt)
  427. io_flush(NORMAL_FLUSH);
  428. else
  429. read_msg_fd();
  430. }
  431. active_filecnt++;
  432. active_bytecnt += F_LENGTH(cur_flist->files[ndx - cur_flist->ndx_start]);
  433. }
  434. /* Write an message to a multiplexed stream. If this fails, rsync exits. */
  435. static void mplex_write(int fd, enum msgcode code, const char *buf, size_t len, int convert)
  436. {
  437. char buffer[BIGPATHBUFLEN]; /* Oversized for use by iconv code. */
  438. size_t n = len;
  439. #ifdef ICONV_OPTION
  440. /* We need to convert buf before doing anything else so that we
  441. * can include the (converted) byte length in the message header. */
  442. if (convert && ic_send != (iconv_t)-1) {
  443. xbuf outbuf, inbuf;
  444. INIT_XBUF(outbuf, buffer + 4, 0, sizeof buffer - 4);
  445. INIT_XBUF(inbuf, (char*)buf, len, -1);
  446. iconvbufs(ic_send, &inbuf, &outbuf,
  447. ICB_INCLUDE_BAD | ICB_INCLUDE_INCOMPLETE);
  448. if (inbuf.len > 0) {
  449. rprintf(FERROR, "overflowed conversion buffer in mplex_write");
  450. exit_cleanup(RERR_UNSUPPORTED);
  451. }
  452. n = len = outbuf.len;
  453. } else
  454. #endif
  455. if (n > 1024 - 4) /* BIGPATHBUFLEN can handle 1024 bytes */
  456. n = 0; /* We'd rather do 2 writes than too much memcpy(). */
  457. else
  458. memcpy(buffer + 4, buf, n);
  459. SIVAL(buffer, 0, ((MPLEX_BASE + (int)code)<<24) + len);
  460. keep_defer_forwarding++; /* defer_forwarding_messages++ on return */
  461. writefd_unbuffered(fd, buffer, n+4);
  462. keep_defer_forwarding--;
  463. if (len > n)
  464. writefd_unbuffered(fd, buf+n, len-n);
  465. if (!--defer_forwarding_messages && !no_flush)
  466. msg_flush();
  467. }
  468. int send_msg(enum msgcode code, const char *buf, int len, int convert)
  469. {
  470. if (msg_fd_out < 0) {
  471. if (!defer_forwarding_messages)
  472. return io_multiplex_write(code, buf, len, convert);
  473. if (!io_multiplexing_out)
  474. return 0;
  475. msg_list_add(&msg_queue, code, buf, len, convert);
  476. return 1;
  477. }
  478. if (flist_forward_from >= 0)
  479. msg_list_add(&msg_queue, code, buf, len, convert);
  480. else
  481. mplex_write(msg_fd_out, code, buf, len, convert);
  482. return 1;
  483. }
  484. void send_msg_int(enum msgcode code, int num)
  485. {
  486. char numbuf[4];
  487. SIVAL(numbuf, 0, num);
  488. send_msg(code, numbuf, 4, 0);
  489. }
  490. void wait_for_receiver(void)
  491. {
  492. if (io_flush(NORMAL_FLUSH))
  493. return;
  494. read_msg_fd();
  495. }
  496. int get_redo_num(void)
  497. {
  498. return flist_ndx_pop(&redo_list);
  499. }
  500. int get_hlink_num(void)
  501. {
  502. return flist_ndx_pop(&hlink_list);
  503. }
  504. /**
  505. * When we're the receiver and we have a local --files-from list of names
  506. * that needs to be sent over the socket to the sender, we have to do two
  507. * things at the same time: send the sender a list of what files we're
  508. * processing and read the incoming file+info list from the sender. We do
  509. * this by augmenting the read_timeout() function to copy this data. It
  510. * uses ff_buf to read a block of data from f_in (when it is ready, since
  511. * it might be a pipe) and then blast it out f_out (when it is ready to
  512. * receive more data).
  513. */
  514. void io_set_filesfrom_fds(int f_in, int f_out)
  515. {
  516. io_filesfrom_f_in = f_in;
  517. io_filesfrom_f_out = f_out;
  518. alloc_xbuf(&ff_buf, 2048);
  519. #ifdef ICONV_OPTION
  520. if (protect_args)
  521. alloc_xbuf(&iconv_buf, 1024);
  522. #endif
  523. }
  524. /* It's almost always an error to get an EOF when we're trying to read from the
  525. * network, because the protocol is (for the most part) self-terminating.
  526. *
  527. * There is one case for the receiver when it is at the end of the transfer
  528. * (hanging around reading any keep-alive packets that might come its way): if
  529. * the sender dies before the generator's kill-signal comes through, we can end
  530. * up here needing to loop until the kill-signal arrives. In this situation,
  531. * kluge_around_eof will be < 0.
  532. *
  533. * There is another case for older protocol versions (< 24) where the module
  534. * listing was not terminated, so we must ignore an EOF error in that case and
  535. * exit. In this situation, kluge_around_eof will be > 0. */
  536. static void whine_about_eof(int fd)
  537. {
  538. if (kluge_around_eof && fd == sock_f_in) {
  539. int i;
  540. if (kluge_around_eof > 0)
  541. exit_cleanup(0);
  542. /* If we're still here after 10 seconds, exit with an error. */
  543. for (i = 10*1000/20; i--; )
  544. msleep(20);
  545. }
  546. rprintf(FERROR, RSYNC_NAME ": connection unexpectedly closed "
  547. "(%.0f bytes received so far) [%s]\n",
  548. (double)stats.total_read, who_am_i());
  549. exit_cleanup(RERR_STREAMIO);
  550. }
  551. /**
  552. * Read from a socket with I/O timeout. return the number of bytes
  553. * read. If no bytes can be read then exit, never return a number <= 0.
  554. *
  555. * TODO: If the remote shell connection fails, then current versions
  556. * actually report an "unexpected EOF" error here. Since it's a
  557. * fairly common mistake to try to use rsh when ssh is required, we
  558. * should trap that: if we fail to read any data at all, we should
  559. * give a better explanation. We can tell whether the connection has
  560. * started by looking e.g. at whether the remote version is known yet.
  561. */
  562. static int read_timeout(int fd, char *buf, size_t len)
  563. {
  564. int n, cnt = 0;
  565. io_flush(FULL_FLUSH);
  566. while (cnt == 0) {
  567. /* until we manage to read *something* */
  568. fd_set r_fds, w_fds;
  569. struct timeval tv;
  570. int maxfd = fd;
  571. int count;
  572. FD_ZERO(&r_fds);
  573. FD_ZERO(&w_fds);
  574. FD_SET(fd, &r_fds);
  575. if (io_filesfrom_f_out >= 0) {
  576. int new_fd;
  577. if (ff_buf.len == 0) {
  578. if (io_filesfrom_f_in >= 0) {
  579. FD_SET(io_filesfrom_f_in, &r_fds);
  580. new_fd = io_filesfrom_f_in;
  581. } else {
  582. io_filesfrom_f_out = -1;
  583. new_fd = -1;
  584. }
  585. } else {
  586. FD_SET(io_filesfrom_f_out, &w_fds);
  587. new_fd = io_filesfrom_f_out;
  588. }
  589. if (new_fd > maxfd)
  590. maxfd = new_fd;
  591. }
  592. tv.tv_sec = select_timeout;
  593. tv.tv_usec = 0;
  594. errno = 0;
  595. count = select(maxfd + 1, &r_fds, &w_fds, NULL, &tv);
  596. if (count <= 0) {
  597. if (errno == EBADF) {
  598. defer_forwarding_messages = 0;
  599. exit_cleanup(RERR_SOCKETIO);
  600. }
  601. check_timeout();
  602. continue;
  603. }
  604. if (io_filesfrom_f_out >= 0) {
  605. if (ff_buf.len) {
  606. if (FD_ISSET(io_filesfrom_f_out, &w_fds)) {
  607. int l = write(io_filesfrom_f_out,
  608. ff_buf.buf + ff_buf.pos,
  609. ff_buf.len);
  610. if (l > 0) {
  611. if (!(ff_buf.len -= l))
  612. ff_buf.pos = 0;
  613. else
  614. ff_buf.pos += l;
  615. } else if (errno != EINTR) {
  616. /* XXX should we complain? */
  617. io_filesfrom_f_out = -1;
  618. }
  619. }
  620. } else if (io_filesfrom_f_in >= 0) {
  621. if (FD_ISSET(io_filesfrom_f_in, &r_fds)) {
  622. #ifdef ICONV_OPTION
  623. xbuf *ibuf = filesfrom_convert ? &iconv_buf : &ff_buf;
  624. #else
  625. xbuf *ibuf = &ff_buf;
  626. #endif
  627. int l = read(io_filesfrom_f_in, ibuf->buf, ibuf->size);
  628. if (l <= 0) {
  629. if (l == 0 || errno != EINTR) {
  630. /* Send end-of-file marker */
  631. memcpy(ff_buf.buf, "\0\0", 2);
  632. ff_buf.len = ff_lastchar? 2 : 1;
  633. ff_buf.pos = 0;
  634. io_filesfrom_f_in = -1;
  635. }
  636. } else {
  637. #ifdef ICONV_OPTION
  638. if (filesfrom_convert) {
  639. iconv_buf.pos = 0;
  640. iconv_buf.len = l;
  641. iconvbufs(ic_send, &iconv_buf, &ff_buf,
  642. ICB_EXPAND_OUT|ICB_INCLUDE_BAD|ICB_INCLUDE_INCOMPLETE);
  643. l = ff_buf.len;
  644. }
  645. #endif
  646. if (!eol_nulls) {
  647. char *s = ff_buf.buf + l;
  648. /* Transform CR and/or LF into '\0' */
  649. while (s-- > ff_buf.buf) {
  650. if (*s == '\n' || *s == '\r')
  651. *s = '\0';
  652. }
  653. }
  654. if (!ff_lastchar) {
  655. /* Last buf ended with a '\0', so don't
  656. * let this buf start with one. */
  657. while (l && ff_buf.buf[ff_buf.pos] == '\0')
  658. ff_buf.pos++, l--;
  659. }
  660. if (!l)
  661. ff_buf.pos = 0;
  662. else {
  663. char *f = ff_buf.buf + ff_buf.pos;
  664. char *t = f;
  665. char *eob = f + l;
  666. /* Eliminate any multi-'\0' runs. */
  667. while (f != eob) {
  668. if (!(*t++ = *f++)) {
  669. while (f != eob && !*f)
  670. f++, l--;
  671. }
  672. }
  673. ff_lastchar = f[-1];
  674. }
  675. ff_buf.len = l;
  676. }
  677. }
  678. }
  679. }
  680. if (!FD_ISSET(fd, &r_fds))
  681. continue;
  682. n = read(fd, buf, len);
  683. if (n <= 0) {
  684. if (n == 0)
  685. whine_about_eof(fd); /* Doesn't return. */
  686. if (errno == EINTR || errno == EWOULDBLOCK
  687. || errno == EAGAIN)
  688. continue;
  689. /* Don't write errors on a dead socket. */
  690. if (fd == sock_f_in) {
  691. io_end_multiplex_out();
  692. rsyserr(FERROR_SOCKET, errno, "read error");
  693. } else
  694. rsyserr(FERROR, errno, "read error");
  695. exit_cleanup(RERR_STREAMIO);
  696. }
  697. buf += n;
  698. len -= n;
  699. cnt += n;
  700. if (fd == sock_f_in && io_timeout)
  701. last_io_in = time(NULL);
  702. }
  703. return cnt;
  704. }
  705. /* Read a line into the "buf" buffer. */
  706. int read_line(int fd, char *buf, size_t bufsiz, int flags)
  707. {
  708. char ch, *s, *eob;
  709. int cnt;
  710. #ifdef ICONV_OPTION
  711. if (flags & RL_CONVERT && iconv_buf.size < bufsiz)
  712. realloc_xbuf(&iconv_buf, bufsiz + 1024);
  713. #endif
  714. start:
  715. #ifdef ICONV_OPTION
  716. s = flags & RL_CONVERT ? iconv_buf.buf : buf;
  717. #else
  718. s = buf;
  719. #endif
  720. eob = s + bufsiz - 1;
  721. while (1) {
  722. cnt = read(fd, &ch, 1);
  723. if (cnt < 0 && (errno == EWOULDBLOCK
  724. || errno == EINTR || errno == EAGAIN)) {
  725. struct timeval tv;
  726. fd_set r_fds, e_fds;
  727. FD_ZERO(&r_fds);
  728. FD_SET(fd, &r_fds);
  729. FD_ZERO(&e_fds);
  730. FD_SET(fd, &e_fds);
  731. tv.tv_sec = select_timeout;
  732. tv.tv_usec = 0;
  733. if (!select(fd+1, &r_fds, NULL, &e_fds, &tv))
  734. check_timeout();
  735. /*if (FD_ISSET(fd, &e_fds))
  736. rprintf(FINFO, "select exception on fd %d\n", fd); */
  737. continue;
  738. }
  739. if (cnt != 1)
  740. break;
  741. if (flags & RL_EOL_NULLS ? ch == '\0' : (ch == '\r' || ch == '\n')) {
  742. /* Skip empty lines if dumping comments. */
  743. if (flags & RL_DUMP_COMMENTS && s == buf)
  744. continue;
  745. break;
  746. }
  747. if (s < eob)
  748. *s++ = ch;
  749. }
  750. *s = '\0';
  751. if (flags & RL_DUMP_COMMENTS && (*buf == '#' || *buf == ';'))
  752. goto start;
  753. #ifdef ICONV_OPTION
  754. if (flags & RL_CONVERT) {
  755. xbuf outbuf;
  756. INIT_XBUF(outbuf, buf, 0, bufsiz);
  757. iconv_buf.pos = 0;
  758. iconv_buf.len = s - iconv_buf.buf;
  759. iconvbufs(ic_recv, &iconv_buf, &outbuf,
  760. ICB_INCLUDE_BAD | ICB_INCLUDE_INCOMPLETE);
  761. outbuf.buf[outbuf.len] = '\0';
  762. return outbuf.len;
  763. }
  764. #endif
  765. return s - buf;
  766. }
  767. void read_args(int f_in, char *mod_name, char *buf, size_t bufsiz, int rl_nulls,
  768. char ***argv_p, int *argc_p, char **request_p)
  769. {
  770. int maxargs = MAX_ARGS;
  771. int dot_pos = 0;
  772. int argc = 0;
  773. char **argv, *p;
  774. int rl_flags = (rl_nulls ? RL_EOL_NULLS : 0);
  775. #ifdef ICONV_OPTION
  776. rl_flags |= (protect_args && ic_recv != (iconv_t)-1 ? RL_CONVERT : 0);
  777. #endif
  778. if (!(argv = new_array(char *, maxargs)))
  779. out_of_memory("read_args");
  780. if (mod_name && !protect_args)
  781. argv[argc++] = "rsyncd";
  782. while (1) {
  783. if (read_line(f_in, buf, bufsiz, rl_flags) == 0)
  784. break;
  785. if (argc == maxargs-1) {
  786. maxargs += MAX_ARGS;
  787. if (!(argv = realloc_array(argv, char *, maxargs)))
  788. out_of_memory("read_args");
  789. }
  790. if (dot_pos) {
  791. if (request_p) {
  792. *request_p = strdup(buf);
  793. request_p = NULL;
  794. }
  795. if (mod_name)
  796. glob_expand_module(mod_name, buf, &argv, &argc, &maxargs);
  797. else
  798. glob_expand(buf, &argv, &argc, &maxargs);
  799. } else {
  800. if (!(p = strdup(buf)))
  801. out_of_memory("read_args");
  802. argv[argc++] = p;
  803. if (*p == '.' && p[1] == '\0')
  804. dot_pos = argc;
  805. }
  806. }
  807. argv[argc] = NULL;
  808. glob_expand(NULL, NULL, NULL, NULL);
  809. *argc_p = argc;
  810. *argv_p = argv;
  811. }
  812. int io_start_buffering_out(int f_out)
  813. {
  814. if (iobuf_out) {
  815. assert(f_out == iobuf_f_out);
  816. return 0;
  817. }
  818. if (!(iobuf_out = new_array(char, IO_BUFFER_SIZE)))
  819. out_of_memory("io_start_buffering_out");
  820. iobuf_out_cnt = 0;
  821. iobuf_f_out = f_out;
  822. return 1;
  823. }
  824. int io_start_buffering_in(int f_in)
  825. {
  826. if (iobuf_in) {
  827. assert(f_in == iobuf_f_in);
  828. return 0;
  829. }
  830. iobuf_in_siz = 2 * IO_BUFFER_SIZE;
  831. if (!(iobuf_in = new_array(char, iobuf_in_siz)))
  832. out_of_memory("io_start_buffering_in");
  833. iobuf_f_in = f_in;
  834. return 1;
  835. }
  836. void io_end_buffering_in(void)
  837. {
  838. if (!iobuf_in)
  839. return;
  840. free(iobuf_in);
  841. iobuf_in = NULL;
  842. iobuf_in_ndx = 0;
  843. iobuf_in_remaining = 0;
  844. iobuf_f_in = -1;
  845. }
  846. void io_end_buffering_out(void)
  847. {
  848. if (!iobuf_out)
  849. return;
  850. io_flush(FULL_FLUSH);
  851. free(iobuf_out);
  852. iobuf_out = NULL;
  853. iobuf_f_out = -1;
  854. }
  855. void maybe_flush_socket(int important)
  856. {
  857. if (iobuf_out && iobuf_out_cnt
  858. && (important || time(NULL) - last_io_out >= 5))
  859. io_flush(NORMAL_FLUSH);
  860. }
  861. void maybe_send_keepalive(void)
  862. {
  863. if (time(NULL) - last_io_out >= allowed_lull) {
  864. if (!iobuf_out || !iobuf_out_cnt) {
  865. if (protocol_version < 29)
  866. return; /* there's nothing we can do */
  867. if (protocol_version >= 30)
  868. send_msg(MSG_NOOP, "", 0, 0);
  869. else {
  870. write_int(sock_f_out, cur_flist->used);
  871. write_shortint(sock_f_out, ITEM_IS_NEW);
  872. }
  873. }
  874. if (iobuf_out)
  875. io_flush(NORMAL_FLUSH);
  876. }
  877. }
  878. void start_flist_forward(int f_in)
  879. {
  880. assert(iobuf_out != NULL);
  881. assert(iobuf_f_out == msg_fd_out);
  882. flist_forward_from = f_in;
  883. }
  884. void stop_flist_forward()
  885. {
  886. flist_forward_from = -1;
  887. io_flush(FULL_FLUSH);
  888. }
  889. /**
  890. * Continue trying to read len bytes - don't return until len has been
  891. * read.
  892. **/
  893. static void read_loop(int fd, char *buf, size_t len)
  894. {
  895. while (len) {
  896. int n = read_timeout(fd, buf, len);
  897. buf += n;
  898. len -= n;
  899. }
  900. }
  901. /**
  902. * Read from the file descriptor handling multiplexing - return number
  903. * of bytes read.
  904. *
  905. * Never returns <= 0.
  906. */
  907. static int readfd_unbuffered(int fd, char *buf, size_t len)
  908. {
  909. size_t msg_bytes;
  910. int tag, cnt = 0;
  911. char line[BIGPATHBUFLEN];
  912. if (!iobuf_in || fd != iobuf_f_in)
  913. return read_timeout(fd, buf, len);
  914. if (!io_multiplexing_in && iobuf_in_remaining == 0) {
  915. iobuf_in_remaining = read_timeout(fd, iobuf_in, iobuf_in_siz);
  916. iobuf_in_ndx = 0;
  917. }
  918. while (cnt == 0) {
  919. if (iobuf_in_remaining) {
  920. len = MIN(len, iobuf_in_remaining);
  921. memcpy(buf, iobuf_in + iobuf_in_ndx, len);
  922. iobuf_in_ndx += len;
  923. iobuf_in_remaining -= len;
  924. cnt = len;
  925. break;
  926. }
  927. read_loop(fd, line, 4);
  928. tag = IVAL(line, 0);
  929. msg_bytes = tag & 0xFFFFFF;
  930. tag = (tag >> 24) - MPLEX_BASE;
  931. switch (tag) {
  932. case MSG_DATA:
  933. if (msg_bytes > iobuf_in_siz) {
  934. if (!(iobuf_in = realloc_array(iobuf_in, char,
  935. msg_bytes)))
  936. out_of_memory("readfd_unbuffered");
  937. iobuf_in_siz = msg_bytes;
  938. }
  939. read_loop(fd, iobuf_in, msg_bytes);
  940. iobuf_in_remaining = msg_bytes;
  941. iobuf_in_ndx = 0;
  942. break;
  943. case MSG_NOOP:
  944. if (am_sender)
  945. maybe_send_keepalive();
  946. break;
  947. case MSG_IO_ERROR:
  948. if (msg_bytes != 4)
  949. goto invalid_msg;
  950. read_loop(fd, line, msg_bytes);
  951. send_msg_int(MSG_IO_ERROR, IVAL(line, 0));
  952. io_error |= IVAL(line, 0);
  953. break;
  954. case MSG_DELETED:
  955. if (msg_bytes >= sizeof line)
  956. goto overflow;
  957. #ifdef ICONV_OPTION
  958. if (ic_recv != (iconv_t)-1) {
  959. xbuf outbuf, inbuf;
  960. char ibuf[512];
  961. int add_null = 0;
  962. int pos = 0;
  963. INIT_CONST_XBUF(outbuf, line);
  964. INIT_XBUF(inbuf, ibuf, 0, -1);
  965. while (msg_bytes) {
  966. inbuf.len = msg_bytes > sizeof ibuf
  967. ? sizeof ibuf : msg_bytes;
  968. read_loop(fd, inbuf.buf, inbuf.len);
  969. if (!(msg_bytes -= inbuf.len)
  970. && !ibuf[inbuf.len-1])
  971. inbuf.len--, add_null = 1;
  972. if (iconvbufs(ic_send, &inbuf, &outbuf,
  973. ICB_INCLUDE_BAD | ICB_INCLUDE_INCOMPLETE) < 0)
  974. goto overflow;
  975. pos = -1;
  976. }
  977. if (add_null) {
  978. if (outbuf.len == outbuf.size)
  979. goto overflow;
  980. outbuf.buf[outbuf.len++] = '\0';
  981. }
  982. msg_bytes = outbuf.len;
  983. } else
  984. #endif
  985. read_loop(fd, line, msg_bytes);
  986. /* A directory name was sent with the trailing null */
  987. if (msg_bytes > 0 && !line[msg_bytes-1])
  988. log_delete(line, S_IFDIR);
  989. else {
  990. line[msg_bytes] = '\0';
  991. log_delete(line, S_IFREG);
  992. }
  993. break;
  994. case MSG_SUCCESS:
  995. if (msg_bytes != 4) {
  996. invalid_msg:
  997. rprintf(FERROR, "invalid multi-message %d:%ld [%s]\n",
  998. tag, (long)msg_bytes, who_am_i());
  999. exit_cleanup(RERR_STREAMIO);
  1000. }
  1001. read_loop(fd, line, msg_bytes);
  1002. successful_send(IVAL(line, 0));
  1003. break;
  1004. case MSG_NO_SEND:
  1005. if (msg_bytes != 4)
  1006. goto invalid_msg;
  1007. read_loop(fd, line, msg_bytes);
  1008. send_msg_int(MSG_NO_SEND, IVAL(line, 0));
  1009. break;
  1010. case MSG_INFO:
  1011. case MSG_ERROR:
  1012. case MSG_ERROR_XFER:
  1013. case MSG_WARNING:
  1014. if (msg_bytes >= sizeof line) {
  1015. overflow:
  1016. rprintf(FERROR,
  1017. "multiplexing overflow %d:%ld [%s]\n",
  1018. tag, (long)msg_bytes, who_am_i());
  1019. exit_cleanup(RERR_STREAMIO);
  1020. }
  1021. read_loop(fd, line, msg_bytes);
  1022. rwrite((enum logcode)tag, line, msg_bytes, 1);
  1023. if (first_message) {
  1024. if (list_only && !am_sender && tag == 1) {
  1025. line[msg_bytes] = '\0';
  1026. check_for_d_option_error(line);
  1027. }
  1028. first_message = 0;
  1029. }
  1030. break;
  1031. default:
  1032. rprintf(FERROR, "unexpected tag %d [%s]\n",
  1033. tag, who_am_i());
  1034. exit_cleanup(RERR_STREAMIO);
  1035. }
  1036. }
  1037. if (iobuf_in_remaining == 0)
  1038. io_flush(NORMAL_FLUSH);
  1039. return cnt;
  1040. }
  1041. /* Do a buffered read from fd. Don't return until all N bytes have
  1042. * been read. If all N can't be read then exit with an error. */
  1043. static void readfd(int fd, char *buffer, size_t N)
  1044. {
  1045. int cnt;
  1046. size_t total = 0;
  1047. while (total < N) {
  1048. cnt = readfd_unbuffered(fd, buffer + total, N-total);
  1049. total += cnt;
  1050. }
  1051. if (fd == write_batch_monitor_in) {
  1052. if ((size_t)write(batch_fd, buffer, total) != total)
  1053. exit_cleanup(RERR_FILEIO);
  1054. }
  1055. if (fd == flist_forward_from)
  1056. writefd(iobuf_f_out, buffer, total);
  1057. if (fd == sock_f_in)
  1058. stats.total_read += total;
  1059. }
  1060. unsigned short read_shortint(int f)
  1061. {
  1062. char b[2];
  1063. readfd(f, b, 2);
  1064. return (UVAL(b, 1) << 8) + UVAL(b, 0);
  1065. }
  1066. int32 read_int(int f)
  1067. {
  1068. char b[4];
  1069. int32 num;
  1070. readfd(f, b, 4);
  1071. num = IVAL(b, 0);
  1072. #if SIZEOF_INT32 > 4
  1073. if (num & (int32)0x80000000)
  1074. num |= ~(int32)0xffffffff;
  1075. #endif
  1076. return num;
  1077. }
  1078. int32 read_varint(int f)
  1079. {
  1080. union {
  1081. char b[5];
  1082. int32 x;
  1083. } u;
  1084. uchar ch;
  1085. int extra;
  1086. u.x = 0;
  1087. readfd(f, (char*)&ch, 1);
  1088. extra = int_byte_extra[ch / 4];
  1089. if (extra) {
  1090. uchar bit = ((uchar)1<<(8-extra));
  1091. if (extra >= (int)sizeof u.b) {
  1092. rprintf(FERROR, "Overflow in read_varint()\n");
  1093. exit_cleanup(RERR_STREAMIO);
  1094. }
  1095. readfd(f, u.b, extra);
  1096. u.b[extra] = ch & (bit-1);
  1097. } else
  1098. u.b[0] = ch;
  1099. #if CAREFUL_ALIGNMENT
  1100. u.x = IVAL(u.b,0);
  1101. #endif
  1102. #if SIZEOF_INT32 > 4
  1103. if (u.x & (int32)0x80000000)
  1104. u.x |= ~(int32)0xffffffff;
  1105. #endif
  1106. return u.x;
  1107. }
  1108. int64 read_varlong(int f, uchar min_bytes)
  1109. {
  1110. union {
  1111. char b[9];
  1112. int64 x;
  1113. } u;
  1114. char b2[8];
  1115. int extra;
  1116. #if SIZEOF_INT64 < 8
  1117. memset(u.b, 0, 8);
  1118. #else
  1119. u.x = 0;
  1120. #endif
  1121. readfd(f, b2, min_bytes);
  1122. memcpy(u.b, b2+1, min_bytes-1);
  1123. extra = int_byte_extra[CVAL(b2, 0) / 4];
  1124. if (extra) {
  1125. uchar bit = ((uchar)1<<(8-extra));
  1126. if (min_bytes + extra > (int)sizeof u.b) {
  1127. rprintf(FERROR, "Overflow in read_varlong()\n");
  1128. exit_cleanup(RERR_STREAMIO);
  1129. }
  1130. readfd(f, u.b + min_bytes - 1, extra);
  1131. u.b[min_bytes + extra - 1] = CVAL(b2, 0) & (bit-1);
  1132. #if SIZEOF_INT64 < 8
  1133. if (min_bytes + extra > 5 || u.b[4] || CVAL(u.b,3) & 0x80) {
  1134. rprintf(FERROR, "Integer overflow: attempted 64-bit offset\n");
  1135. exit_cleanup(RERR_UNSUPPORTED);
  1136. }
  1137. #endif
  1138. } else
  1139. u.b[min_bytes + extra - 1] = CVAL(b2, 0);
  1140. #if SIZEOF_INT64 < 8
  1141. u.x = IVAL(u.b,0);
  1142. #elif CAREFUL_ALIGNMENT
  1143. u.x = IVAL(u.b,0) | (((int64)IVAL(u.b,4))<<32);
  1144. #endif
  1145. return u.x;
  1146. }
  1147. int64 read_longint(int f)
  1148. {
  1149. #if SIZEOF_INT64 >= 8
  1150. char b[9];
  1151. #endif
  1152. int32 num = read_int(f);
  1153. if (num != (int32)0xffffffff)
  1154. return num;
  1155. #if SIZEOF_INT64 < 8
  1156. rprintf(FERROR, "Integer overflow: attempted 64-bit offset\n");
  1157. exit_cleanup(RERR_UNSUPPORTED);
  1158. #else
  1159. readfd(f, b, 8);
  1160. return IVAL(b,0) | (((int64)IVAL(b,4))<<32);
  1161. #endif
  1162. }
  1163. void read_buf(int f, char *buf, size_t len)
  1164. {
  1165. readfd(f,buf,len);
  1166. }
  1167. void read_sbuf(int f, char *buf, size_t len)
  1168. {
  1169. readfd(f, buf, len);
  1170. buf[len] = '\0';
  1171. }
  1172. uchar read_byte(int f)
  1173. {
  1174. uchar c;
  1175. readfd(f, (char *)&c, 1);
  1176. return c;
  1177. }
  1178. int read_vstring(int f, char *buf, int bufsize)
  1179. {
  1180. int len = read_byte(f);
  1181. if (len & 0x80)
  1182. len = (len & ~0x80) * 0x100 + read_byte(f);
  1183. if (len >= bufsize) {
  1184. rprintf(FERROR, "over-long vstring received (%d > %d)\n",
  1185. len, bufsize - 1);
  1186. return -1;
  1187. }
  1188. if (len)
  1189. readfd(f, buf, len);
  1190. buf[len] = '\0';
  1191. return len;
  1192. }
  1193. /* Populate a sum_struct with values from the socket. This is
  1194. * called by both the sender and the receiver. */
  1195. void read_sum_head(int f, struct sum_struct *sum)
  1196. {
  1197. int32 max_blength = protocol_version < 30 ? OLD_MAX_BLOCK_SIZE : MAX_BLOCK_SIZE;
  1198. sum->count = read_int(f);
  1199. if (sum->count < 0) {
  1200. rprintf(FERROR, "Invalid checksum count %ld [%s]\n",
  1201. (long)sum->count, who_am_i());
  1202. exit_cleanup(RERR_PROTOCOL);
  1203. }
  1204. sum->blength = read_int(f);
  1205. if (sum->blength < 0 || sum->blength > max_blength) {
  1206. rprintf(FERROR, "Invalid block length %ld [%s]\n",
  1207. (long)sum->blength, who_am_i());
  1208. exit_cleanup(RERR_PROTOCOL);
  1209. }
  1210. sum->s2length = protocol_version < 27 ? csum_length : (int)read_int(f);
  1211. if (sum->s2length < 0 || sum->s2length > MAX_DIGEST_LEN) {
  1212. rprintf(FERROR, "Invalid checksum length %d [%s]\n",
  1213. sum->s2length, who_am_i());
  1214. exit_cleanup(RERR_PROTOCOL);
  1215. }
  1216. sum->remainder = read_int(f);
  1217. if (sum->remainder < 0 || sum->remainder > sum->blength) {
  1218. rprintf(FERROR, "Invalid remainder length %ld [%s]\n",
  1219. (long)sum->remainder, who_am_i());
  1220. exit_cleanup(RERR_PROTOCOL);
  1221. }
  1222. }
  1223. /* Send the values from a sum_struct over the socket. Set sum to
  1224. * NULL if there are no checksums to send. This is called by both
  1225. * the generator and the sender. */
  1226. void write_sum_head(int f, struct sum_struct *sum)
  1227. {
  1228. static struct sum_struct null_sum;
  1229. if (sum == NULL)
  1230. sum = &null_sum;
  1231. write_int(f, sum->count);
  1232. write_int(f, sum->blength);
  1233. if (protocol_version >= 27)
  1234. write_int(f, sum->s2length);
  1235. write_int(f, sum->remainder);
  1236. }
  1237. /**
  1238. * Sleep after writing to limit I/O bandwidth usage.
  1239. *
  1240. * @todo Rather than sleeping after each write, it might be better to
  1241. * use some kind of averaging. The current algorithm seems to always
  1242. * use a bit less bandwidth than specified, because it doesn't make up
  1243. * for slow periods. But arguably this is a feature. In addition, we
  1244. * ought to take the time used to write the data into account.
  1245. *
  1246. * During some phases of big transfers (file FOO is uptodate) this is
  1247. * called with a small bytes_written every time. As the kernel has to
  1248. * round small waits up to guarantee that we actually wait at least the
  1249. * requested number of microseconds, this can become grossly inaccurate.
  1250. * We therefore keep track of the bytes we've written over time and only
  1251. * sleep when the accumulated delay is at least 1 tenth of a second.
  1252. **/
  1253. static void sleep_for_bwlimit(int bytes_written)
  1254. {
  1255. static struct timeval prior_tv;
  1256. static long total_written = 0;
  1257. struct timeval tv, start_tv;
  1258. long elapsed_usec, sleep_usec;
  1259. #define ONE_SEC 1000000L /* # of microseconds in a second */
  1260. if (!bwlimit_writemax)
  1261. return;
  1262. total_written += bytes_written;
  1263. gettimeofday(&start_tv, NULL);
  1264. if (prior_tv.tv_sec) {
  1265. elapsed_usec = (start_tv.tv_sec - prior_tv.tv_sec) * ONE_SEC
  1266. + (start_tv.tv_usec - prior_tv.tv_usec);
  1267. total_written -= elapsed_usec * bwlimit / (ONE_SEC/1024);
  1268. if (total_written < 0)
  1269. total_written = 0;
  1270. }
  1271. sleep_usec = total_written * (ONE_SEC/1024) / bwlimit;
  1272. if (sleep_usec < ONE_SEC / 10) {
  1273. prior_tv = start_tv;
  1274. return;
  1275. }
  1276. tv.tv_sec = sleep_usec / ONE_SEC;
  1277. tv.tv_usec = sleep_usec % ONE_SEC;
  1278. select(0, NULL, NULL, NULL, &tv);
  1279. gettimeofday(&prior_tv, NULL);
  1280. elapsed_usec = (prior_tv.tv_sec - start_tv.tv_sec) * ONE_SEC
  1281. + (prior_tv.tv_usec - start_tv.tv_usec);
  1282. total_written = (sleep_usec - elapsed_usec) * bwlimit / (ONE_SEC/1024);
  1283. }
  1284. /* Write len bytes to the file descriptor fd, looping as necessary to get
  1285. * the job done and also (in certain circumstances) reading any data on
  1286. * msg_fd_in to avoid deadlock.
  1287. *
  1288. * This function underlies the multiplexing system. The body of the
  1289. * application never calls this function directly. */
  1290. static void writefd_unbuffered(int fd, const char *buf, size_t len)
  1291. {
  1292. size_t n, total = 0;
  1293. fd_set w_fds, r_fds, e_fds;
  1294. int maxfd, count, cnt, using_r_fds;
  1295. int defer_inc = 0;
  1296. struct timeval tv;
  1297. if (no_flush++)
  1298. defer_forwarding_messages++, defer_inc++;
  1299. while (total < len) {
  1300. FD_ZERO(&w_fds);
  1301. FD_SET(fd, &w_fds);
  1302. FD_ZERO(&e_fds);
  1303. FD_SET(fd, &e_fds);
  1304. maxfd = fd;
  1305. if (msg_fd_in >= 0) {
  1306. FD_ZERO(&r_fds);
  1307. FD_SET(msg_fd_in, &r_fds);
  1308. if (msg_fd_in > maxfd)
  1309. maxfd = msg_fd_in;
  1310. using_r_fds = 1;
  1311. } else
  1312. using_r_fds = 0;
  1313. tv.tv_sec = select_timeout;
  1314. tv.tv_usec = 0;
  1315. errno = 0;
  1316. count = select(maxfd + 1, using_r_fds ? &r_fds : NULL,
  1317. &w_fds, &e_fds, &tv);
  1318. if (count <= 0) {
  1319. if (count < 0 && errno == EBADF)
  1320. exit_cleanup(RERR_SOCKETIO);
  1321. check_timeout();
  1322. continue;
  1323. }
  1324. /*if (FD_ISSET(fd, &e_fds))
  1325. rprintf(FINFO, "select exception on fd %d\n", fd); */
  1326. if (using_r_fds && FD_ISSET(msg_fd_in, &r_fds))
  1327. read_msg_fd();
  1328. if (!FD_ISSET(fd, &w_fds))
  1329. continue;
  1330. n = len - total;
  1331. if (bwlimit_writemax && n > bwlimit_writemax)
  1332. n = bwlimit_writemax;
  1333. cnt = write(fd, buf + total, n);
  1334. if (cnt <= 0) {
  1335. if (cnt < 0) {
  1336. if (errno == EINTR)
  1337. continue;
  1338. if (errno == EWOULDBLOCK || errno == EAGAIN) {
  1339. msleep(1);
  1340. continue;
  1341. }
  1342. }
  1343. /* Don't try to write errors back across the stream. */
  1344. if (fd == sock_f_out)
  1345. io_end_multiplex_out();
  1346. /* Don't try to write errors down a failing msg pipe. */
  1347. if (am_server && fd == msg_fd_out)
  1348. exit_cleanup(RERR_STREAMIO);
  1349. rsyserr(FERROR, errno,
  1350. "writefd_unbuffered failed to write %ld bytes [%s]",
  1351. (long)len, who_am_i());
  1352. /* If the other side is sending us error messages, try
  1353. * to grab any messages they sent before they died. */
  1354. while (!am_server && fd == sock_f_out && io_multiplexing_in) {
  1355. char buf[1024];
  1356. set_io_timeout(30);
  1357. ignore_timeout = 0;
  1358. readfd_unbuffered(sock_f_in, buf, sizeof buf);
  1359. }
  1360. exit_cleanup(RERR_STREAMIO);
  1361. }
  1362. total += cnt;
  1363. defer_forwarding_messages++, defer_inc++;
  1364. if (fd == sock_f_out) {
  1365. if (io_timeout || am_generator)
  1366. last_io_out = time(NULL);
  1367. sleep_for_bwlimit(cnt);
  1368. }
  1369. }
  1370. no_flush--;
  1371. if (keep_defer_forwarding)
  1372. defer_inc--;
  1373. if (!(defer_forwarding_messages -= defer_inc) && !no_flush)
  1374. msg_flush();
  1375. }
  1376. int io_flush(int flush_it_all)
  1377. {
  1378. int flushed_something = 0;
  1379. if (no_flush)
  1380. return 0;
  1381. if (iobuf_out_cnt) {
  1382. if (io_multiplexing_out)
  1383. mplex_write(sock_f_out, MSG_DATA, iobuf_out, iobuf_out_cnt, 0);
  1384. else
  1385. writefd_unbuffered(iobuf_f_out, iobuf_out, iobuf_out_cnt);
  1386. iobuf_out_cnt = 0;
  1387. flushed_something = 1;
  1388. }
  1389. if (flush_it_all && !defer_forwarding_messages && msg_queue.head) {
  1390. msg_flush();
  1391. flushed_something = 1;
  1392. }
  1393. return flushed_something;
  1394. }
  1395. static void writefd(int fd, const char *buf, size_t len)
  1396. {
  1397. if (fd == sock_f_out)
  1398. stats.total_written += len;
  1399. if (fd == write_batch_monitor_out) {
  1400. if ((size_t)write(batch_fd, buf, len) != len)
  1401. exit_cleanup(RERR_FILEIO);
  1402. }
  1403. if (!iobuf_out || fd != iobuf_f_out) {
  1404. writefd_unbuffered(fd, buf, len);
  1405. return;
  1406. }
  1407. while (len) {
  1408. int n = MIN((int)len, IO_BUFFER_SIZE - iobuf_out_cnt);
  1409. if (n > 0) {
  1410. memcpy(iobuf_out+iobuf_out_cnt, buf, n);
  1411. buf += n;
  1412. len -= n;
  1413. iobuf_out_cnt += n;
  1414. }
  1415. if (iobuf_out_cnt == IO_BUFFER_SIZE)
  1416. io_flush(NORMAL_FLUSH);
  1417. }
  1418. }
  1419. void write_shortint(int f, unsigned short x)
  1420. {
  1421. char b[2];
  1422. b[0] = (char)x;
  1423. b[1] = (char)(x >> 8);
  1424. writefd(f, b, 2);
  1425. }
  1426. void write_int(int f, int32 x)
  1427. {
  1428. char b[4];
  1429. SIVAL(b, 0, x);
  1430. writefd(f, b, 4);
  1431. }
  1432. void write_varint(int f, int32 x)
  1433. {
  1434. char b[5];
  1435. uchar bit;
  1436. int cnt = 4;
  1437. SIVAL(b, 1, x);
  1438. while (cnt > 1 && b[cnt] == 0)
  1439. cnt--;
  1440. bit = ((uchar)1<<(7-cnt+1));
  1441. if (CVAL(b, cnt) >= bit) {
  1442. cnt++;
  1443. *b = ~(bit-1);
  1444. } else if (cnt > 1)
  1445. *b = b[cnt] | ~(bit*2-1);
  1446. else
  1447. *b = b[cnt];
  1448. writefd(f, b, cnt);
  1449. }
  1450. void write_varlong(int f, int64 x, uchar min_bytes)
  1451. {
  1452. char b[9];
  1453. uchar bit;
  1454. int cnt = 8;
  1455. SIVAL(b, 1, x);
  1456. #if SIZEOF_INT64 >= 8
  1457. SIVAL(b, 5, x >> 32);
  1458. #else
  1459. if (x <= 0x7FFFFFFF && x >= 0)
  1460. memset(b + 5, 0, 4);
  1461. else {
  1462. rprintf(FERROR, "Integer overflow: attempted 64-bit offset\n");
  1463. exit_cleanup(RERR_UNSUPPORTED);
  1464. }
  1465. #endif
  1466. while (cnt > min_bytes && b[cnt] == 0)
  1467. cnt--;
  1468. bit = ((uchar)1<<(7-cnt+min_bytes));
  1469. if (CVAL(b, cnt) >= bit) {
  1470. cnt++;
  1471. *b = ~(bit-1);
  1472. } else if (cnt > min_bytes)
  1473. *b = b[cnt] | ~(bit*2-1);
  1474. else
  1475. *b = b[cnt];
  1476. writefd(f, b, cnt);
  1477. }
  1478. /*
  1479. * Note: int64 may actually be a 32-bit type if ./configure couldn't find any
  1480. * 64-bit types on this platform.
  1481. */
  1482. void write_longint(int f, int64 x)
  1483. {
  1484. char b[12], * const s = b+4;
  1485. SIVAL(s, 0, x);
  1486. if (x <= 0x7FFFFFFF && x >= 0) {
  1487. writefd(f, s, 4);
  1488. return;
  1489. }
  1490. #if SIZEOF_INT64 < 8
  1491. rprintf(FERROR, "Integer overflow: attempted 64-bit offset\n");
  1492. exit_cleanup(RERR_UNSUPPORTED);
  1493. #else
  1494. memset(b, 0xFF, 4);
  1495. SIVAL(s, 4, x >> 32);
  1496. writefd(f, b, 12);
  1497. #endif
  1498. }
  1499. void write_buf(int f, const char *buf, size_t len)
  1500. {
  1501. writefd(f,buf,len);
  1502. }
  1503. /** Write a string to the connection */
  1504. void write_sbuf(int f, const char *buf)
  1505. {
  1506. writefd(f, buf, strlen(buf));
  1507. }
  1508. void write_byte(int f, uchar c)
  1509. {
  1510. writefd(f, (char *)&c, 1);
  1511. }
  1512. void write_vstring(int f, const char *str, int len)
  1513. {
  1514. uchar lenbuf[3], *lb = lenbuf;
  1515. if (len > 0x7F) {
  1516. if (len > 0x7FFF) {
  1517. rprintf(FERROR,
  1518. "attempting to send over-long vstring (%d > %d)\n",
  1519. len, 0x7FFF);
  1520. exit_cleanup(RERR_PROTOCOL);
  1521. }
  1522. *lb++ = len / 0x100 + 0x80;
  1523. }
  1524. *lb = len;
  1525. writefd(f, (char*)lenbuf, lb - lenbuf + 1);
  1526. if (len)
  1527. writefd(f, str, len);
  1528. }
  1529. /* Send a file-list index using a byte-reduction method. */
  1530. void write_ndx(int f, int32 ndx)
  1531. {
  1532. static int32 prev_positive = -1, prev_negative = 1;
  1533. int32 diff, cnt = 0;
  1534. char b[6];
  1535. if (protocol_version < 30 || read_batch) {
  1536. write_int(f, ndx);
  1537. return;
  1538. }
  1539. /* Send NDX_DONE as a single-byte 0 with no side effects. Send
  1540. * negative nums as a positive after sending a leading 0xFF. */
  1541. if (ndx >= 0) {
  1542. diff = ndx - prev_positive;
  1543. prev_positive = ndx;
  1544. } else if (ndx == NDX_DONE) {
  1545. *b = 0;
  1546. writefd(f, b, 1);
  1547. return;
  1548. } else {
  1549. b[cnt++] = (char)0xFF;
  1550. ndx = -ndx;
  1551. diff = ndx - prev_negative;
  1552. prev_negative = ndx;
  1553. }
  1554. /* A diff of 1 - 253 is sent as a one-byte diff; a diff of 254 - 32767
  1555. * or 0 is sent as a 0xFE + a two-byte diff; otherwise we send 0xFE
  1556. * & all 4 bytes of the (non-negative) num with the high-bit set. */
  1557. if (diff < 0xFE && diff > 0)
  1558. b[cnt++] = (char)diff;
  1559. else if (diff < 0 || diff > 0x7FFF) {
  1560. b[cnt++] = (char)0xFE;
  1561. b[cnt++] = (char)((ndx >> 24) | 0x80);
  1562. b[cnt++] = (char)ndx;
  1563. b[cnt++] = (char)(ndx >> 8);
  1564. b[cnt++] = (char)(ndx >> 16);
  1565. } else {
  1566. b[cnt++] = (char)0xFE;
  1567. b[cnt++] = (char)(diff >> 8);
  1568. b[cnt++] = (char)diff;
  1569. }
  1570. writefd(f, b, cnt);
  1571. }
  1572. /* Receive a file-list index using a byte-reduction method. */
  1573. int32 read_ndx(int f)
  1574. {
  1575. static int32 prev_positive = -1, prev_negative = 1;
  1576. int32 *prev_ptr, num;
  1577. char b[4];
  1578. if (protocol_version < 30)
  1579. return read_int(f);
  1580. readfd(f, b, 1);
  1581. if (CVAL(b, 0) == 0xFF) {
  1582. readfd(f, b, 1);
  1583. prev_ptr = &prev_negative;
  1584. } else if (CVAL(b, 0) == 0)
  1585. return NDX_DONE;
  1586. else
  1587. prev_ptr = &prev_positive;
  1588. if (CVAL(b, 0) == 0xFE) {
  1589. readfd(f, b, 2);
  1590. if (CVAL(b, 0) & 0x80) {
  1591. b[3] = CVAL(b, 0) & ~0x80;
  1592. b[0] = b[1];
  1593. readfd(f, b+1, 2);
  1594. num = IVAL(b, 0);
  1595. } else
  1596. num = (UVAL(b,0)<<8) + UVAL(b,1) + *prev_ptr;
  1597. } else
  1598. num = UVAL(b, 0) + *prev_ptr;
  1599. *prev_ptr = num;
  1600. if (prev_ptr == &prev_negative)
  1601. num = -num;
  1602. return num;
  1603. }
  1604. /* Read a line of up to bufsiz-1 characters into buf. Strips
  1605. * the (required) trailing newline and all carriage returns.
  1606. * Returns 1 for success; 0 for I/O error or truncation. */
  1607. int read_line_old(int f, char *buf, size_t bufsiz)
  1608. {
  1609. bufsiz--; /* leave room for the null */
  1610. while (bufsiz > 0) {
  1611. buf[0] = 0;
  1612. read_buf(f, buf, 1);
  1613. if (buf[0] == 0)
  1614. return 0;
  1615. if (buf[0] == '\n')
  1616. break;
  1617. if (buf[0] != '\r') {
  1618. buf++;
  1619. bufsiz--;
  1620. }
  1621. }
  1622. *buf = '\0';
  1623. return bufsiz > 0;
  1624. }
  1625. void io_printf(int fd, const char *format, ...)
  1626. {
  1627. va_list ap;
  1628. char buf[BIGPATHBUFLEN];
  1629. int len;
  1630. va_start(ap, format);
  1631. len = vsnprintf(buf, sizeof buf, format, ap);
  1632. va_end(ap);
  1633. if (len < 0)
  1634. exit_cleanup(RERR_STREAMIO);
  1635. if (len > (int)sizeof buf) {
  1636. rprintf(FERROR, "io_printf() was too long for the buffer.\n");
  1637. exit_cleanup(RERR_STREAMIO);
  1638. }
  1639. write_sbuf(fd, buf);
  1640. }
  1641. /** Setup for multiplexing a MSG_* stream with the data stream. */
  1642. void io_start_multiplex_out(void)
  1643. {
  1644. io_flush(NORMAL_FLUSH);
  1645. io_start_buffering_out(sock_f_out);
  1646. io_multiplexing_out = 1;
  1647. }
  1648. /** Setup for multiplexing a MSG_* stream with the data stream. */
  1649. void io_start_multiplex_in(void)
  1650. {
  1651. io_flush(NORMAL_FLUSH);
  1652. io_start_buffering_in(sock_f_in);
  1653. io_multiplexing_in = 1;
  1654. }
  1655. /** Write an message to the multiplexed data stream. */
  1656. int io_multiplex_write(enum msgcode code, const char *buf, size_t len, int convert)
  1657. {
  1658. if (!io_multiplexing_out)
  1659. return 0;
  1660. io_flush(NORMAL_FLUSH);
  1661. stats.total_written += (len+4);
  1662. mplex_write(sock_f_out, code, buf, len, convert);
  1663. return 1;
  1664. }
  1665. void io_end_multiplex_in(void)
  1666. {
  1667. io_multiplexing_in = 0;
  1668. io_end_buffering_in();
  1669. }
  1670. /** Stop output multiplexing. */
  1671. void io_end_multiplex_out(void)
  1672. {
  1673. io_multiplexing_out = 0;
  1674. io_end_buffering_out();
  1675. }
  1676. void start_write_batch(int fd)
  1677. {
  1678. /* Some communication has already taken place, but we don't
  1679. * enable batch writing until here so that we can write a
  1680. * canonical record of the communication even though the
  1681. * actual communication so far depends on whether a daemon
  1682. * is involved. */
  1683. write_int(batch_fd, protocol_version);
  1684. if (protocol_version >= 30)
  1685. write_byte(batch_fd, inc_recurse);
  1686. write_int(batch_fd, checksum_seed);
  1687. if (am_sender)
  1688. write_batch_monitor_out = fd;
  1689. else
  1690. write_batch_monitor_in = fd;
  1691. }
  1692. void stop_write_batch(void)
  1693. {
  1694. write_batch_monitor_out = -1;
  1695. write_batch_monitor_in = -1;
  1696. }