You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

worker_util.c 56KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223
  1. /*-
  2. * Copyright 2016 Vsevolod Stakhov
  3. *
  4. * Licensed under the Apache License, Version 2.0 (the "License");
  5. * you may not use this file except in compliance with the License.
  6. * You may obtain a copy of the License at
  7. *
  8. * http://www.apache.org/licenses/LICENSE-2.0
  9. *
  10. * Unless required by applicable law or agreed to in writing, software
  11. * distributed under the License is distributed on an "AS IS" BASIS,
  12. * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  13. * See the License for the specific language governing permissions and
  14. * limitations under the License.
  15. */
  16. #include "config.h"
  17. #include "rspamd.h"
  18. #include "lua/lua_common.h"
  19. #include "worker_util.h"
  20. #include "unix-std.h"
  21. #include "utlist.h"
  22. #include "ottery.h"
  23. #include "rspamd_control.h"
  24. #include "libserver/maps/map.h"
  25. #include "libserver/maps/map_private.h"
  26. #include "libserver/http/http_private.h"
  27. #include "libserver/http/http_router.h"
  28. #include "libutil/rrd.h"
  29. /* sys/resource.h */
  30. #ifdef HAVE_SYS_RESOURCE_H
  31. #include <sys/resource.h>
  32. #endif
  33. /* pwd and grp */
  34. #ifdef HAVE_PWD_H
  35. #include <pwd.h>
  36. #endif
  37. #ifdef HAVE_GRP_H
  38. #include <grp.h>
  39. #endif
  40. #ifdef HAVE_LIBUTIL_H
  41. #include <libutil.h>
  42. #endif
  43. #include "zlib.h"
  44. #ifdef WITH_LIBUNWIND
  45. #define UNW_LOCAL_ONLY 1
  46. #include <libunwind.h>
  47. #define UNWIND_BACKTRACE_DEPTH 256
  48. #endif
  49. #ifdef HAVE_UCONTEXT_H
  50. #include <ucontext.h>
  51. #elif defined(HAVE_SYS_UCONTEXT_H)
  52. #include <sys/ucontext.h>
  53. #endif
  54. #ifdef HAVE_SYS_WAIT_H
  55. #include <sys/wait.h>
  56. #endif
  57. #include "contrib/libev/ev.h"
  58. #include "libstat/stat_api.h"
  59. /* Forward declaration */
  60. static void rspamd_worker_heartbeat_start (struct rspamd_worker *,
  61. struct ev_loop *);
  62. static void rspamd_worker_ignore_signal (struct rspamd_worker_signal_handler *);
  63. /**
  64. * Return worker's control structure by its type
  65. * @param type
  66. * @return worker's control structure or NULL
  67. */
  68. worker_t *
  69. rspamd_get_worker_by_type (struct rspamd_config *cfg, GQuark type)
  70. {
  71. worker_t **pwrk;
  72. pwrk = cfg->compiled_workers;
  73. while (pwrk && *pwrk) {
  74. if (rspamd_check_worker (cfg, *pwrk)) {
  75. if (g_quark_from_string ((*pwrk)->name) == type) {
  76. return *pwrk;
  77. }
  78. }
  79. pwrk++;
  80. }
  81. return NULL;
  82. }
  83. static void
  84. rspamd_worker_check_finished (EV_P_ ev_timer *w, int revents)
  85. {
  86. int *pnchecks = (int *)w->data;
  87. if (*pnchecks > SOFT_SHUTDOWN_TIME * 10) {
  88. msg_warn ("terminating worker before finishing of terminate handlers");
  89. ev_break (EV_A_ EVBREAK_ONE);
  90. }
  91. else {
  92. int refcount = ev_active_cnt (EV_A);
  93. if (refcount == 1) {
  94. ev_break (EV_A_ EVBREAK_ONE);
  95. }
  96. else {
  97. ev_timer_again (EV_A_ w);
  98. }
  99. }
  100. }
  101. static gboolean
  102. rspamd_worker_finalize (gpointer user_data)
  103. {
  104. struct rspamd_task *task = user_data;
  105. if (!(task->flags & RSPAMD_TASK_FLAG_PROCESSING)) {
  106. msg_info_task ("finishing actions has been processed, terminating");
  107. /* ev_break (task->event_loop, EVBREAK_ALL); */
  108. task->worker->state = rspamd_worker_wanna_die;
  109. rspamd_session_destroy (task->s);
  110. return TRUE;
  111. }
  112. return FALSE;
  113. }
  114. gboolean
  115. rspamd_worker_call_finish_handlers (struct rspamd_worker *worker)
  116. {
  117. struct rspamd_task *task;
  118. struct rspamd_config *cfg = worker->srv->cfg;
  119. struct rspamd_abstract_worker_ctx *ctx;
  120. struct rspamd_config_cfg_lua_script *sc;
  121. if (cfg->on_term_scripts) {
  122. ctx = (struct rspamd_abstract_worker_ctx *)worker->ctx;
  123. /* Create a fake task object for async events */
  124. task = rspamd_task_new (worker, cfg, NULL, NULL, ctx->event_loop, FALSE);
  125. task->resolver = ctx->resolver;
  126. task->flags |= RSPAMD_TASK_FLAG_PROCESSING;
  127. task->s = rspamd_session_create (task->task_pool,
  128. rspamd_worker_finalize,
  129. NULL,
  130. (event_finalizer_t) rspamd_task_free,
  131. task);
  132. DL_FOREACH (cfg->on_term_scripts, sc) {
  133. lua_call_finish_script (sc, task);
  134. }
  135. task->flags &= ~RSPAMD_TASK_FLAG_PROCESSING;
  136. if (rspamd_session_pending (task->s)) {
  137. return TRUE;
  138. }
  139. }
  140. return FALSE;
  141. }
  142. static void
  143. rspamd_worker_terminate_handlers (struct rspamd_worker *w)
  144. {
  145. if (w->nconns == 0 &&
  146. (!(w->flags & RSPAMD_WORKER_SCANNER) || w->srv->cfg->on_term_scripts == NULL)) {
  147. /*
  148. * We are here either:
  149. * - No active connections are represented
  150. * - No term scripts are registered
  151. * - Worker is not a scanner, so it can die safely
  152. */
  153. w->state = rspamd_worker_wanna_die;
  154. }
  155. else {
  156. if (w->nconns > 0) {
  157. /*
  158. * Wait until all connections are terminated
  159. */
  160. w->state = rspamd_worker_wait_connections;
  161. }
  162. else {
  163. /*
  164. * Start finish scripts
  165. */
  166. if (w->state != rspamd_worker_wait_final_scripts) {
  167. w->state = rspamd_worker_wait_final_scripts;
  168. if ((w->flags & RSPAMD_WORKER_SCANNER) &&
  169. rspamd_worker_call_finish_handlers (w)) {
  170. msg_info ("performing async finishing actions");
  171. w->state = rspamd_worker_wait_final_scripts;
  172. }
  173. else {
  174. /*
  175. * We are done now
  176. */
  177. msg_info ("no async finishing actions, terminating");
  178. w->state = rspamd_worker_wanna_die;
  179. }
  180. }
  181. }
  182. }
  183. }
  184. static void
  185. rspamd_worker_on_delayed_shutdown (EV_P_ ev_timer *w, int revents)
  186. {
  187. struct rspamd_worker *worker = (struct rspamd_worker *)w->data;
  188. worker->state = rspamd_worker_wanna_die;
  189. ev_timer_stop (EV_A_ w);
  190. ev_break (loop, EVBREAK_ALL);
  191. }
  192. static void
  193. rspamd_worker_shutdown_check (EV_P_ ev_timer *w, int revents)
  194. {
  195. struct rspamd_worker *worker = (struct rspamd_worker *)w->data;
  196. if (worker->state != rspamd_worker_wanna_die) {
  197. rspamd_worker_terminate_handlers (worker);
  198. if (worker->state == rspamd_worker_wanna_die) {
  199. /* We are done, kill event loop */
  200. ev_timer_stop (EV_A_ w);
  201. ev_break (EV_A_ EVBREAK_ALL);
  202. }
  203. else {
  204. /* Try again later */
  205. ev_timer_again (EV_A_ w);
  206. }
  207. }
  208. else {
  209. ev_timer_stop (EV_A_ w);
  210. ev_break (EV_A_ EVBREAK_ALL);
  211. }
  212. }
  213. /*
  214. * Config reload is designed by sending sigusr2 to active workers and pending shutdown of them
  215. */
  216. static gboolean
  217. rspamd_worker_usr2_handler (struct rspamd_worker_signal_handler *sigh, void *arg)
  218. {
  219. /* Do not accept new connections, preparing to end worker's process */
  220. if (sigh->worker->state == rspamd_worker_state_running) {
  221. static ev_timer shutdown_ev, shutdown_check_ev;
  222. ev_tstamp shutdown_ts;
  223. if (sigh->worker->flags & RSPAMD_WORKER_NO_TERMINATE_DELAY) {
  224. shutdown_ts = 0.0;
  225. }
  226. else {
  227. shutdown_ts = MAX (SOFT_SHUTDOWN_TIME,
  228. sigh->worker->srv->cfg->task_timeout * 2.0);
  229. }
  230. rspamd_worker_ignore_signal (sigh);
  231. sigh->worker->state = rspamd_worker_state_terminating;
  232. rspamd_default_log_function (G_LOG_LEVEL_INFO,
  233. sigh->worker->srv->server_pool->tag.tagname,
  234. sigh->worker->srv->server_pool->tag.uid,
  235. G_STRFUNC,
  236. "worker's shutdown is pending in %.2f sec",
  237. shutdown_ts);
  238. /* Soft shutdown timer */
  239. shutdown_ev.data = sigh->worker;
  240. ev_timer_init (&shutdown_ev, rspamd_worker_on_delayed_shutdown,
  241. shutdown_ts, 0.0);
  242. ev_timer_start (sigh->event_loop, &shutdown_ev);
  243. if (!(sigh->worker->flags & RSPAMD_WORKER_NO_TERMINATE_DELAY)) {
  244. /* This timer checks if we are ready to die and is called frequently */
  245. shutdown_check_ev.data = sigh->worker;
  246. ev_timer_init (&shutdown_check_ev, rspamd_worker_shutdown_check,
  247. 0.5, 0.5);
  248. ev_timer_start (sigh->event_loop, &shutdown_check_ev);
  249. }
  250. rspamd_worker_stop_accept (sigh->worker);
  251. }
  252. /* No more signals */
  253. return FALSE;
  254. }
  255. /*
  256. * Reopen log is designed by sending sigusr1 to active workers and pending shutdown of them
  257. */
  258. static gboolean
  259. rspamd_worker_usr1_handler (struct rspamd_worker_signal_handler *sigh, void *arg)
  260. {
  261. struct rspamd_main *rspamd_main = sigh->worker->srv;
  262. rspamd_log_reopen (sigh->worker->srv->logger, rspamd_main->cfg, -1, -1);
  263. msg_info_main ("logging reinitialised");
  264. /* Get more signals */
  265. return TRUE;
  266. }
  267. static gboolean
  268. rspamd_worker_term_handler (struct rspamd_worker_signal_handler *sigh, void *arg)
  269. {
  270. if (sigh->worker->state == rspamd_worker_state_running) {
  271. static ev_timer shutdown_ev, shutdown_check_ev;
  272. ev_tstamp shutdown_ts;
  273. if (sigh->worker->flags & RSPAMD_WORKER_NO_TERMINATE_DELAY) {
  274. shutdown_ts = 0.0;
  275. }
  276. else {
  277. shutdown_ts = MAX (SOFT_SHUTDOWN_TIME,
  278. sigh->worker->srv->cfg->task_timeout * 2.0);
  279. }
  280. rspamd_worker_ignore_signal (sigh);
  281. sigh->worker->state = rspamd_worker_state_terminating;
  282. rspamd_default_log_function (G_LOG_LEVEL_INFO,
  283. sigh->worker->srv->server_pool->tag.tagname,
  284. sigh->worker->srv->server_pool->tag.uid,
  285. G_STRFUNC,
  286. "terminating after receiving signal %s",
  287. g_strsignal (sigh->signo));
  288. rspamd_worker_stop_accept (sigh->worker);
  289. rspamd_worker_terminate_handlers (sigh->worker);
  290. /* Check if we are ready to die */
  291. if (sigh->worker->state != rspamd_worker_wanna_die) {
  292. /* This timer is called when we have no choices but to die */
  293. shutdown_ev.data = sigh->worker;
  294. ev_timer_init (&shutdown_ev, rspamd_worker_on_delayed_shutdown,
  295. shutdown_ts, 0.0);
  296. ev_timer_start (sigh->event_loop, &shutdown_ev);
  297. if (!(sigh->worker->flags & RSPAMD_WORKER_NO_TERMINATE_DELAY)) {
  298. /* This timer checks if we are ready to die and is called frequently */
  299. shutdown_check_ev.data = sigh->worker;
  300. ev_timer_init (&shutdown_check_ev, rspamd_worker_shutdown_check,
  301. 0.5, 0.5);
  302. ev_timer_start (sigh->event_loop, &shutdown_check_ev);
  303. }
  304. }
  305. else {
  306. /* Flag to die has been already set */
  307. ev_break (sigh->event_loop, EVBREAK_ALL);
  308. }
  309. }
  310. /* Stop reacting on signals */
  311. return FALSE;
  312. }
  313. static void
  314. rspamd_worker_signal_handle (EV_P_ ev_signal *w, int revents)
  315. {
  316. struct rspamd_worker_signal_handler *sigh =
  317. (struct rspamd_worker_signal_handler *)w->data;
  318. struct rspamd_worker_signal_handler_elt *cb, *cbtmp;
  319. /* Call all signal handlers registered */
  320. DL_FOREACH_SAFE (sigh->cb, cb, cbtmp) {
  321. if (!cb->handler (sigh, cb->handler_data)) {
  322. DL_DELETE (sigh->cb, cb);
  323. g_free (cb);
  324. }
  325. }
  326. }
  327. static void
  328. rspamd_worker_ignore_signal (struct rspamd_worker_signal_handler *sigh)
  329. {
  330. sigset_t set;
  331. ev_signal_stop (sigh->event_loop, &sigh->ev_sig);
  332. sigemptyset (&set);
  333. sigaddset (&set, sigh->signo);
  334. sigprocmask (SIG_BLOCK, &set, NULL);
  335. }
  336. static void
  337. rspamd_worker_default_signal (int signo)
  338. {
  339. struct sigaction sig;
  340. sigemptyset (&sig.sa_mask);
  341. sigaddset (&sig.sa_mask, signo);
  342. sig.sa_handler = SIG_DFL;
  343. sig.sa_flags = 0;
  344. sigaction (signo, &sig, NULL);
  345. }
  346. static void
  347. rspamd_sigh_free (void *p)
  348. {
  349. struct rspamd_worker_signal_handler *sigh = p;
  350. struct rspamd_worker_signal_handler_elt *cb, *tmp;
  351. DL_FOREACH_SAFE (sigh->cb, cb, tmp) {
  352. DL_DELETE (sigh->cb, cb);
  353. g_free (cb);
  354. }
  355. ev_signal_stop (sigh->event_loop, &sigh->ev_sig);
  356. rspamd_worker_default_signal (sigh->signo);
  357. g_free (sigh);
  358. }
  359. void
  360. rspamd_worker_set_signal_handler (int signo, struct rspamd_worker *worker,
  361. struct ev_loop *event_loop,
  362. rspamd_worker_signal_cb_t handler,
  363. void *handler_data)
  364. {
  365. struct rspamd_worker_signal_handler *sigh;
  366. struct rspamd_worker_signal_handler_elt *cb;
  367. sigh = g_hash_table_lookup (worker->signal_events, GINT_TO_POINTER (signo));
  368. if (sigh == NULL) {
  369. sigh = g_malloc0 (sizeof (*sigh));
  370. sigh->signo = signo;
  371. sigh->worker = worker;
  372. sigh->event_loop = event_loop;
  373. sigh->enabled = TRUE;
  374. sigh->ev_sig.data = sigh;
  375. ev_signal_init (&sigh->ev_sig, rspamd_worker_signal_handle, signo);
  376. ev_signal_start (event_loop, &sigh->ev_sig);
  377. g_hash_table_insert (worker->signal_events,
  378. GINT_TO_POINTER (signo),
  379. sigh);
  380. }
  381. cb = g_malloc0 (sizeof (*cb));
  382. cb->handler = handler;
  383. cb->handler_data = handler_data;
  384. DL_APPEND (sigh->cb, cb);
  385. }
  386. void
  387. rspamd_worker_init_signals (struct rspamd_worker *worker,
  388. struct ev_loop *event_loop)
  389. {
  390. /* A set of terminating signals */
  391. rspamd_worker_set_signal_handler (SIGTERM, worker, event_loop,
  392. rspamd_worker_term_handler, NULL);
  393. rspamd_worker_set_signal_handler (SIGINT, worker, event_loop,
  394. rspamd_worker_term_handler, NULL);
  395. rspamd_worker_set_signal_handler (SIGHUP, worker, event_loop,
  396. rspamd_worker_term_handler, NULL);
  397. /* Special purpose signals */
  398. rspamd_worker_set_signal_handler (SIGUSR1, worker, event_loop,
  399. rspamd_worker_usr1_handler, NULL);
  400. rspamd_worker_set_signal_handler (SIGUSR2, worker, event_loop,
  401. rspamd_worker_usr2_handler, NULL);
  402. }
  403. struct ev_loop *
  404. rspamd_prepare_worker (struct rspamd_worker *worker, const char *name,
  405. rspamd_accept_handler hdl)
  406. {
  407. struct ev_loop *event_loop;
  408. GList *cur;
  409. struct rspamd_worker_listen_socket *ls;
  410. struct rspamd_worker_accept_event *accept_ev;
  411. worker->signal_events = g_hash_table_new_full (g_direct_hash, g_direct_equal,
  412. NULL, rspamd_sigh_free);
  413. event_loop = ev_loop_new (rspamd_config_ev_backend_get (worker->srv->cfg));
  414. worker->srv->event_loop = event_loop;
  415. rspamd_worker_init_signals (worker, event_loop);
  416. rspamd_control_worker_add_default_cmd_handlers (worker, event_loop);
  417. rspamd_worker_heartbeat_start (worker, event_loop);
  418. rspamd_redis_pool_config (worker->srv->cfg->redis_pool,
  419. worker->srv->cfg, event_loop);
  420. /* Accept all sockets */
  421. if (hdl) {
  422. cur = worker->cf->listen_socks;
  423. while (cur) {
  424. ls = cur->data;
  425. if (ls->fd != -1) {
  426. accept_ev = g_malloc0 (sizeof (*accept_ev));
  427. accept_ev->event_loop = event_loop;
  428. accept_ev->accept_ev.data = worker;
  429. ev_io_init (&accept_ev->accept_ev, hdl, ls->fd, EV_READ);
  430. ev_io_start (event_loop, &accept_ev->accept_ev);
  431. DL_APPEND (worker->accept_events, accept_ev);
  432. }
  433. cur = g_list_next (cur);
  434. }
  435. }
  436. return event_loop;
  437. }
  438. void
  439. rspamd_worker_stop_accept (struct rspamd_worker *worker)
  440. {
  441. struct rspamd_worker_accept_event *cur, *tmp;
  442. /* Remove all events */
  443. DL_FOREACH_SAFE (worker->accept_events, cur, tmp) {
  444. if (ev_can_stop (&cur->accept_ev)) {
  445. ev_io_stop (cur->event_loop, &cur->accept_ev);
  446. }
  447. if (ev_can_stop (&cur->throttling_ev)) {
  448. ev_timer_stop (cur->event_loop, &cur->throttling_ev);
  449. }
  450. g_free (cur);
  451. }
  452. /* XXX: we need to do it much later */
  453. #if 0
  454. g_hash_table_iter_init (&it, worker->signal_events);
  455. while (g_hash_table_iter_next (&it, &k, &v)) {
  456. sigh = (struct rspamd_worker_signal_handler *)v;
  457. g_hash_table_iter_steal (&it);
  458. if (sigh->enabled) {
  459. event_del (&sigh->ev);
  460. }
  461. g_free (sigh);
  462. }
  463. g_hash_table_unref (worker->signal_events);
  464. #endif
  465. }
  466. static rspamd_fstring_t *
  467. rspamd_controller_maybe_compress (struct rspamd_http_connection_entry *entry,
  468. rspamd_fstring_t *buf, struct rspamd_http_message *msg)
  469. {
  470. if (entry->support_gzip) {
  471. if (rspamd_fstring_gzip (&buf)) {
  472. rspamd_http_message_add_header (msg, "Content-Encoding", "gzip");
  473. }
  474. }
  475. return buf;
  476. }
  477. void
  478. rspamd_controller_send_error (struct rspamd_http_connection_entry *entry,
  479. gint code, const gchar *error_msg, ...)
  480. {
  481. struct rspamd_http_message *msg;
  482. va_list args;
  483. rspamd_fstring_t *reply;
  484. msg = rspamd_http_new_message (HTTP_RESPONSE);
  485. va_start (args, error_msg);
  486. msg->status = rspamd_fstring_new ();
  487. rspamd_vprintf_fstring (&msg->status, error_msg, args);
  488. va_end (args);
  489. msg->date = time (NULL);
  490. msg->code = code;
  491. reply = rspamd_fstring_sized_new (msg->status->len + 16);
  492. rspamd_printf_fstring (&reply, "{\"error\":\"%V\"}", msg->status);
  493. rspamd_http_message_set_body_from_fstring_steal (msg,
  494. rspamd_controller_maybe_compress (entry, reply, msg));
  495. rspamd_http_connection_reset (entry->conn);
  496. rspamd_http_router_insert_headers (entry->rt, msg);
  497. rspamd_http_connection_write_message (entry->conn,
  498. msg,
  499. NULL,
  500. "application/json",
  501. entry,
  502. entry->rt->timeout);
  503. entry->is_reply = TRUE;
  504. }
  505. void
  506. rspamd_controller_send_openmetrics (struct rspamd_http_connection_entry *entry,
  507. rspamd_fstring_t *str)
  508. {
  509. struct rspamd_http_message *msg;
  510. msg = rspamd_http_new_message (HTTP_RESPONSE);
  511. msg->date = time (NULL);
  512. msg->code = 200;
  513. msg->status = rspamd_fstring_new_init ("OK", 2);
  514. rspamd_http_message_set_body_from_fstring_steal (msg,
  515. rspamd_controller_maybe_compress (entry, str, msg));
  516. rspamd_http_connection_reset (entry->conn);
  517. rspamd_http_router_insert_headers (entry->rt, msg);
  518. rspamd_http_connection_write_message (entry->conn,
  519. msg,
  520. NULL,
  521. "application/openmetrics-text; version=1.0.0; charset=utf-8",
  522. entry,
  523. entry->rt->timeout);
  524. entry->is_reply = TRUE;
  525. }
  526. void
  527. rspamd_controller_send_string (struct rspamd_http_connection_entry *entry,
  528. const gchar *str)
  529. {
  530. struct rspamd_http_message *msg;
  531. rspamd_fstring_t *reply;
  532. msg = rspamd_http_new_message (HTTP_RESPONSE);
  533. msg->date = time (NULL);
  534. msg->code = 200;
  535. msg->status = rspamd_fstring_new_init ("OK", 2);
  536. if (str) {
  537. reply = rspamd_fstring_new_init (str, strlen (str));
  538. }
  539. else {
  540. reply = rspamd_fstring_new_init ("null", 4);
  541. }
  542. rspamd_http_message_set_body_from_fstring_steal (msg,
  543. rspamd_controller_maybe_compress (entry, reply, msg));
  544. rspamd_http_connection_reset (entry->conn);
  545. rspamd_http_router_insert_headers (entry->rt, msg);
  546. rspamd_http_connection_write_message (entry->conn,
  547. msg,
  548. NULL,
  549. "application/json",
  550. entry,
  551. entry->rt->timeout);
  552. entry->is_reply = TRUE;
  553. }
  554. void
  555. rspamd_controller_send_ucl (struct rspamd_http_connection_entry *entry,
  556. ucl_object_t *obj)
  557. {
  558. struct rspamd_http_message *msg;
  559. rspamd_fstring_t *reply;
  560. msg = rspamd_http_new_message (HTTP_RESPONSE);
  561. msg->date = time (NULL);
  562. msg->code = 200;
  563. msg->status = rspamd_fstring_new_init ("OK", 2);
  564. reply = rspamd_fstring_sized_new (BUFSIZ);
  565. rspamd_ucl_emit_fstring (obj, UCL_EMIT_JSON_COMPACT, &reply);
  566. rspamd_http_message_set_body_from_fstring_steal (msg,
  567. rspamd_controller_maybe_compress (entry, reply, msg));
  568. rspamd_http_connection_reset (entry->conn);
  569. rspamd_http_router_insert_headers (entry->rt, msg);
  570. rspamd_http_connection_write_message (entry->conn,
  571. msg,
  572. NULL,
  573. "application/json",
  574. entry,
  575. entry->rt->timeout);
  576. entry->is_reply = TRUE;
  577. }
  578. static void
  579. rspamd_worker_drop_priv (struct rspamd_main *rspamd_main)
  580. {
  581. if (rspamd_main->is_privileged) {
  582. if (setgid (rspamd_main->workers_gid) == -1) {
  583. msg_err_main ("cannot setgid to %d (%s), aborting",
  584. (gint) rspamd_main->workers_gid,
  585. strerror (errno));
  586. exit (-errno);
  587. }
  588. if (rspamd_main->cfg->rspamd_user &&
  589. initgroups (rspamd_main->cfg->rspamd_user,
  590. rspamd_main->workers_gid) == -1) {
  591. msg_err_main ("initgroups failed (%s), aborting", strerror (errno));
  592. exit (-errno);
  593. }
  594. if (setuid (rspamd_main->workers_uid) == -1) {
  595. msg_err_main ("cannot setuid to %d (%s), aborting",
  596. (gint) rspamd_main->workers_uid,
  597. strerror (errno));
  598. exit (-errno);
  599. }
  600. }
  601. }
  602. static void
  603. rspamd_worker_set_limits (struct rspamd_main *rspamd_main,
  604. struct rspamd_worker_conf *cf)
  605. {
  606. struct rlimit rlmt;
  607. if (cf->rlimit_nofile != 0) {
  608. rlmt.rlim_cur = (rlim_t) cf->rlimit_nofile;
  609. rlmt.rlim_max = (rlim_t) cf->rlimit_nofile;
  610. if (setrlimit (RLIMIT_NOFILE, &rlmt) == -1) {
  611. msg_warn_main ("cannot set files rlimit: %L, %s",
  612. cf->rlimit_nofile,
  613. strerror (errno));
  614. }
  615. memset (&rlmt, 0, sizeof (rlmt));
  616. if (getrlimit (RLIMIT_NOFILE, &rlmt) == -1) {
  617. msg_warn_main ("cannot get max files rlimit: %HL, %s",
  618. cf->rlimit_maxcore,
  619. strerror (errno));
  620. }
  621. else {
  622. msg_info_main ("set max file descriptors limit: %HL cur and %HL max",
  623. (guint64) rlmt.rlim_cur,
  624. (guint64) rlmt.rlim_max);
  625. }
  626. }
  627. else {
  628. /* Just report */
  629. if (getrlimit (RLIMIT_NOFILE, &rlmt) == -1) {
  630. msg_warn_main ("cannot get max files rlimit: %HL, %s",
  631. cf->rlimit_maxcore,
  632. strerror (errno));
  633. }
  634. else {
  635. msg_info_main ("use system max file descriptors limit: %HL cur and %HL max",
  636. (guint64) rlmt.rlim_cur,
  637. (guint64) rlmt.rlim_max);
  638. }
  639. }
  640. if (rspamd_main->cores_throttling) {
  641. msg_info_main ("disable core files for the new worker as limits are reached");
  642. rlmt.rlim_cur = 0;
  643. rlmt.rlim_max = 0;
  644. if (setrlimit (RLIMIT_CORE, &rlmt) == -1) {
  645. msg_warn_main ("cannot disable core dumps: error when setting limits: %s",
  646. strerror (errno));
  647. }
  648. }
  649. else {
  650. if (cf->rlimit_maxcore != 0) {
  651. rlmt.rlim_cur = (rlim_t) cf->rlimit_maxcore;
  652. rlmt.rlim_max = (rlim_t) cf->rlimit_maxcore;
  653. if (setrlimit (RLIMIT_CORE, &rlmt) == -1) {
  654. msg_warn_main ("cannot set max core size limit: %HL, %s",
  655. cf->rlimit_maxcore,
  656. strerror (errno));
  657. }
  658. /* Ensure that we did it */
  659. memset (&rlmt, 0, sizeof (rlmt));
  660. if (getrlimit (RLIMIT_CORE, &rlmt) == -1) {
  661. msg_warn_main ("cannot get max core size rlimit: %HL, %s",
  662. cf->rlimit_maxcore,
  663. strerror (errno));
  664. }
  665. else {
  666. if (rlmt.rlim_cur != cf->rlimit_maxcore ||
  667. rlmt.rlim_max != cf->rlimit_maxcore) {
  668. msg_warn_main ("setting of core file limits was unsuccessful: "
  669. "%HL was wanted, "
  670. "but we have %HL cur and %HL max",
  671. cf->rlimit_maxcore,
  672. (guint64) rlmt.rlim_cur,
  673. (guint64) rlmt.rlim_max);
  674. }
  675. else {
  676. msg_info_main ("set max core size limit: %HL cur and %HL max",
  677. (guint64) rlmt.rlim_cur,
  678. (guint64) rlmt.rlim_max);
  679. }
  680. }
  681. }
  682. else {
  683. /* Just report */
  684. if (getrlimit (RLIMIT_CORE, &rlmt) == -1) {
  685. msg_warn_main ("cannot get max core size limit: %HL, %s",
  686. cf->rlimit_maxcore,
  687. strerror (errno));
  688. }
  689. else {
  690. msg_info_main ("use system max core size limit: %HL cur and %HL max",
  691. (guint64) rlmt.rlim_cur,
  692. (guint64) rlmt.rlim_max);
  693. }
  694. }
  695. }
  696. }
  697. static void
  698. rspamd_worker_on_term (EV_P_ ev_child *w, int revents)
  699. {
  700. struct rspamd_worker *wrk = (struct rspamd_worker *)w->data;
  701. if (wrk->ppid == getpid ()) {
  702. if (wrk->term_handler) {
  703. wrk->term_handler (EV_A_ w, wrk->srv, wrk);
  704. }
  705. else {
  706. rspamd_check_termination_clause (wrk->srv, wrk, w->rstatus);
  707. }
  708. }
  709. else {
  710. /* Ignore SIGCHLD for not our children... */
  711. }
  712. }
  713. static void
  714. rspamd_worker_heartbeat_cb (EV_P_ ev_timer *w, int revents)
  715. {
  716. struct rspamd_worker *wrk = (struct rspamd_worker *)w->data;
  717. struct rspamd_srv_command cmd;
  718. memset (&cmd, 0, sizeof (cmd));
  719. cmd.type = RSPAMD_SRV_HEARTBEAT;
  720. rspamd_srv_send_command (wrk, EV_A, &cmd, -1, NULL, NULL);
  721. }
  722. static void
  723. rspamd_worker_heartbeat_start (struct rspamd_worker *wrk, struct ev_loop *event_loop)
  724. {
  725. wrk->hb.heartbeat_ev.data = (void *)wrk;
  726. ev_timer_init (&wrk->hb.heartbeat_ev, rspamd_worker_heartbeat_cb,
  727. 0.0, wrk->srv->cfg->heartbeat_interval);
  728. ev_timer_start (event_loop, &wrk->hb.heartbeat_ev);
  729. }
  730. static void
  731. rspamd_main_heartbeat_cb (EV_P_ ev_timer *w, int revents)
  732. {
  733. struct rspamd_worker *wrk = (struct rspamd_worker *)w->data;
  734. gdouble time_from_last = ev_time ();
  735. struct rspamd_main *rspamd_main;
  736. static struct rspamd_control_command cmd;
  737. struct tm tm;
  738. gchar timebuf[64];
  739. gchar usec_buf[16];
  740. gint r;
  741. time_from_last -= wrk->hb.last_event;
  742. rspamd_main = wrk->srv;
  743. if (wrk->hb.last_event > 0 &&
  744. time_from_last > 0 &&
  745. time_from_last >= rspamd_main->cfg->heartbeat_interval * 2) {
  746. rspamd_localtime (wrk->hb.last_event, &tm);
  747. r = strftime (timebuf, sizeof (timebuf), "%F %H:%M:%S", &tm);
  748. rspamd_snprintf (usec_buf, sizeof (usec_buf), "%.5f",
  749. wrk->hb.last_event - (gdouble)(time_t)wrk->hb.last_event);
  750. rspamd_snprintf (timebuf + r, sizeof (timebuf) - r,
  751. "%s", usec_buf + 1);
  752. if (wrk->hb.nbeats > 0) {
  753. /* First time lost event */
  754. cmd.type = RSPAMD_CONTROL_CHILD_CHANGE;
  755. cmd.cmd.child_change.what = rspamd_child_offline;
  756. cmd.cmd.child_change.pid = wrk->pid;
  757. rspamd_control_broadcast_srv_cmd (rspamd_main, &cmd, wrk->pid);
  758. msg_warn_main ("lost heartbeat from worker type %s with pid %P, "
  759. "last beat on: %s (%L beats received previously)",
  760. g_quark_to_string (wrk->type), wrk->pid,
  761. timebuf,
  762. wrk->hb.nbeats);
  763. wrk->hb.nbeats = -1;
  764. /* TODO: send notify about worker problem */
  765. }
  766. else {
  767. wrk->hb.nbeats --;
  768. msg_warn_main ("lost %L heartbeat from worker type %s with pid %P, "
  769. "last beat on: %s",
  770. -(wrk->hb.nbeats),
  771. g_quark_to_string (wrk->type),
  772. wrk->pid,
  773. timebuf);
  774. if (rspamd_main->cfg->heartbeats_loss_max > 0 &&
  775. -(wrk->hb.nbeats) >= rspamd_main->cfg->heartbeats_loss_max) {
  776. if (-(wrk->hb.nbeats) > rspamd_main->cfg->heartbeats_loss_max + 1) {
  777. msg_err_main ("force kill worker type %s with pid %P, "
  778. "last beat on: %s; %L heartbeat lost",
  779. g_quark_to_string (wrk->type),
  780. wrk->pid,
  781. timebuf,
  782. -(wrk->hb.nbeats));
  783. kill (wrk->pid, SIGKILL);
  784. }
  785. else {
  786. msg_err_main ("terminate worker type %s with pid %P, "
  787. "last beat on: %s; %L heartbeat lost",
  788. g_quark_to_string (wrk->type),
  789. wrk->pid,
  790. timebuf,
  791. -(wrk->hb.nbeats));
  792. kill (wrk->pid, SIGTERM);
  793. }
  794. }
  795. }
  796. }
  797. else if (wrk->hb.nbeats < 0) {
  798. rspamd_localtime (wrk->hb.last_event, &tm);
  799. r = strftime (timebuf, sizeof (timebuf), "%F %H:%M:%S", &tm);
  800. rspamd_snprintf (usec_buf, sizeof (usec_buf), "%.5f",
  801. wrk->hb.last_event - (gdouble)(time_t)wrk->hb.last_event);
  802. rspamd_snprintf (timebuf + r, sizeof (timebuf) - r,
  803. "%s", usec_buf + 1);
  804. cmd.type = RSPAMD_CONTROL_CHILD_CHANGE;
  805. cmd.cmd.child_change.what = rspamd_child_online;
  806. cmd.cmd.child_change.pid = wrk->pid;
  807. rspamd_control_broadcast_srv_cmd (rspamd_main, &cmd, wrk->pid);
  808. msg_info_main ("received heartbeat from worker type %s with pid %P, "
  809. "last beat on: %s (%L beats lost previously)",
  810. g_quark_to_string (wrk->type), wrk->pid,
  811. timebuf,
  812. -(wrk->hb.nbeats));
  813. wrk->hb.nbeats = 1;
  814. /* TODO: send notify about worker restoration */
  815. }
  816. }
  817. static void
  818. rspamd_main_heartbeat_start (struct rspamd_worker *wrk, struct ev_loop *event_loop)
  819. {
  820. wrk->hb.heartbeat_ev.data = (void *)wrk;
  821. ev_timer_init (&wrk->hb.heartbeat_ev, rspamd_main_heartbeat_cb,
  822. 0.0, wrk->srv->cfg->heartbeat_interval * 2);
  823. ev_timer_start (event_loop, &wrk->hb.heartbeat_ev);
  824. }
  825. static bool
  826. rspamd_maybe_reuseport_socket (struct rspamd_worker_listen_socket *ls)
  827. {
  828. gint nfd = -1;
  829. if (ls->is_systemd) {
  830. /* No need to reuseport */
  831. return true;
  832. }
  833. if (ls->fd != -1 && rspamd_inet_address_get_af (ls->addr) == AF_UNIX) {
  834. /* Just try listen */
  835. if (listen (ls->fd, -1) == -1) {
  836. return false;
  837. }
  838. return true;
  839. }
  840. #if defined(SO_REUSEPORT) && defined(SO_REUSEADDR) && defined(LINUX)
  841. if (ls->type == RSPAMD_WORKER_SOCKET_UDP) {
  842. nfd = rspamd_inet_address_listen (ls->addr,
  843. (ls->type == RSPAMD_WORKER_SOCKET_UDP ? SOCK_DGRAM : SOCK_STREAM),
  844. RSPAMD_INET_ADDRESS_LISTEN_ASYNC|RSPAMD_INET_ADDRESS_LISTEN_REUSEPORT,
  845. -1);
  846. if (nfd == -1) {
  847. msg_warn ("cannot create reuseport listen socket for %d: %s",
  848. ls->fd, strerror (errno));
  849. nfd = ls->fd;
  850. }
  851. else {
  852. if (ls->fd != -1) {
  853. close (ls->fd);
  854. }
  855. ls->fd = nfd;
  856. nfd = -1;
  857. }
  858. }
  859. else {
  860. /*
  861. * Reuseport is broken with the current architecture, so it is easier not
  862. * to use it at all
  863. */
  864. nfd = ls->fd;
  865. }
  866. #else
  867. nfd = ls->fd;
  868. #endif
  869. #if 0
  870. /* This needed merely if we have reuseport for tcp, but for now it is disabled */
  871. /* This means that we have an fd with no listening enabled */
  872. if (nfd != -1) {
  873. if (ls->type == RSPAMD_WORKER_SOCKET_TCP) {
  874. if (listen (nfd, -1) == -1) {
  875. return false;
  876. }
  877. }
  878. }
  879. #endif
  880. return true;
  881. }
  882. /**
  883. * Handles worker after fork returned zero
  884. * @param wrk
  885. * @param rspamd_main
  886. * @param cf
  887. * @param listen_sockets
  888. */
  889. static void __attribute__((noreturn))
  890. rspamd_handle_child_fork (struct rspamd_worker *wrk,
  891. struct rspamd_main *rspamd_main,
  892. struct rspamd_worker_conf *cf,
  893. GHashTable *listen_sockets)
  894. {
  895. gint rc;
  896. struct rlimit rlim;
  897. /* Update pid for logging */
  898. rspamd_log_on_fork (cf->type, rspamd_main->cfg, rspamd_main->logger);
  899. wrk->pid = getpid ();
  900. /* Init PRNG after fork */
  901. rc = ottery_init (rspamd_main->cfg->libs_ctx->ottery_cfg);
  902. if (rc != OTTERY_ERR_NONE) {
  903. msg_err_main ("cannot initialize PRNG: %d", rc);
  904. abort ();
  905. }
  906. rspamd_random_seed_fast ();
  907. #ifdef HAVE_EVUTIL_RNG_INIT
  908. evutil_secure_rng_init ();
  909. #endif
  910. /*
  911. * Libev stores all signals in a global table, so
  912. * previous handlers must be explicitly detached and forgotten
  913. * before starting a new loop
  914. */
  915. ev_signal_stop (rspamd_main->event_loop, &rspamd_main->int_ev);
  916. ev_signal_stop (rspamd_main->event_loop, &rspamd_main->term_ev);
  917. ev_signal_stop (rspamd_main->event_loop, &rspamd_main->hup_ev);
  918. ev_signal_stop (rspamd_main->event_loop, &rspamd_main->usr1_ev);
  919. /* Remove the inherited event base */
  920. ev_loop_destroy (rspamd_main->event_loop);
  921. rspamd_main->event_loop = NULL;
  922. /* Close unused sockets */
  923. GHashTableIter it;
  924. gpointer k, v;
  925. g_hash_table_iter_init (&it, listen_sockets);
  926. /*
  927. * Close listen sockets of not our process (inherited from other forks)
  928. */
  929. while (g_hash_table_iter_next (&it, &k, &v)) {
  930. GList *elt = (GList *)v;
  931. GList *our = cf->listen_socks;
  932. if (g_list_position (our, elt) == -1) {
  933. GList *cur = elt;
  934. while (cur) {
  935. struct rspamd_worker_listen_socket *ls =
  936. (struct rspamd_worker_listen_socket *)cur->data;
  937. if (ls->fd != -1 && close (ls->fd) == -1) {
  938. msg_err ("cannot close fd %d (addr = %s): %s",
  939. ls->fd,
  940. rspamd_inet_address_to_string_pretty (ls->addr),
  941. strerror (errno));
  942. }
  943. ls->fd = -1;
  944. cur = g_list_next (cur);
  945. }
  946. }
  947. }
  948. /* Reuseport before dropping privs */
  949. GList *cur = cf->listen_socks;
  950. while (cur) {
  951. struct rspamd_worker_listen_socket *ls =
  952. (struct rspamd_worker_listen_socket *)cur->data;
  953. if (!rspamd_maybe_reuseport_socket (ls)) {
  954. msg_err ("cannot listen on socket %s: %s",
  955. rspamd_inet_address_to_string_pretty (ls->addr),
  956. strerror (errno));
  957. }
  958. cur = g_list_next (cur);
  959. }
  960. /* Drop privileges */
  961. rspamd_worker_drop_priv (rspamd_main);
  962. /* Set limits */
  963. rspamd_worker_set_limits (rspamd_main, cf);
  964. /* Re-set stack limit */
  965. getrlimit (RLIMIT_STACK, &rlim);
  966. rlim.rlim_cur = 100 * 1024 * 1024;
  967. rlim.rlim_max = rlim.rlim_cur;
  968. setrlimit (RLIMIT_STACK, &rlim);
  969. if (cf->bind_conf) {
  970. setproctitle ("%s process (%s)", cf->worker->name,
  971. cf->bind_conf->bind_line);
  972. }
  973. else {
  974. setproctitle ("%s process", cf->worker->name);
  975. }
  976. if (rspamd_main->pfh) {
  977. rspamd_pidfile_close (rspamd_main->pfh);
  978. }
  979. if (rspamd_main->cfg->log_silent_workers) {
  980. rspamd_log_set_log_level (rspamd_main->logger, G_LOG_LEVEL_MESSAGE);
  981. }
  982. wrk->start_time = rspamd_get_calendar_ticks ();
  983. if (cf->bind_conf) {
  984. GString *listen_conf_stringified = g_string_new (NULL);
  985. struct rspamd_worker_bind_conf *cur_conf;
  986. LL_FOREACH (cf->bind_conf, cur_conf) {
  987. if (cur_conf->next) {
  988. rspamd_printf_gstring (listen_conf_stringified, "%s, ",
  989. cur_conf->bind_line);
  990. }
  991. else {
  992. rspamd_printf_gstring (listen_conf_stringified, "%s",
  993. cur_conf->bind_line);
  994. }
  995. }
  996. msg_info_main ("starting %s process %P (%d); listen on: %v",
  997. cf->worker->name,
  998. getpid (), wrk->index, listen_conf_stringified);
  999. g_string_free (listen_conf_stringified, TRUE);
  1000. }
  1001. else {
  1002. msg_info_main ("starting %s process %P (%d); no listen",
  1003. cf->worker->name,
  1004. getpid (), wrk->index);
  1005. }
  1006. /* Close parent part of socketpair */
  1007. close (wrk->control_pipe[0]);
  1008. close (wrk->srv_pipe[0]);
  1009. rspamd_socket_nonblocking (wrk->control_pipe[1]);
  1010. rspamd_socket_nonblocking (wrk->srv_pipe[1]);
  1011. rspamd_main->cfg->cur_worker = wrk;
  1012. /* Execute worker (this function should not return normally!) */
  1013. cf->worker->worker_start_func (wrk);
  1014. /* To distinguish from normal termination */
  1015. exit (EXIT_FAILURE);
  1016. }
  1017. static void
  1018. rspamd_handle_main_fork (struct rspamd_worker *wrk,
  1019. struct rspamd_main *rspamd_main,
  1020. struct rspamd_worker_conf *cf,
  1021. struct ev_loop *ev_base)
  1022. {
  1023. /* Close worker part of socketpair */
  1024. close (wrk->control_pipe[1]);
  1025. close (wrk->srv_pipe[1]);
  1026. rspamd_socket_nonblocking (wrk->control_pipe[0]);
  1027. rspamd_socket_nonblocking (wrk->srv_pipe[0]);
  1028. rspamd_srv_start_watching (rspamd_main, wrk, ev_base);
  1029. /* Child event */
  1030. wrk->cld_ev.data = wrk;
  1031. ev_child_init (&wrk->cld_ev, rspamd_worker_on_term, wrk->pid, 0);
  1032. ev_child_start (rspamd_main->event_loop, &wrk->cld_ev);
  1033. /* Heartbeats */
  1034. rspamd_main_heartbeat_start (wrk, rspamd_main->event_loop);
  1035. /* Insert worker into worker's table, pid is index */
  1036. g_hash_table_insert (rspamd_main->workers,
  1037. GSIZE_TO_POINTER (wrk->pid), wrk);
  1038. #if defined(SO_REUSEPORT) && defined(SO_REUSEADDR) && defined(LINUX)
  1039. /*
  1040. * Close listen sockets in the main process once a child is handling them,
  1041. * if we have reuseport
  1042. */
  1043. GList *cur = cf->listen_socks;
  1044. while (cur) {
  1045. struct rspamd_worker_listen_socket *ls =
  1046. (struct rspamd_worker_listen_socket *)cur->data;
  1047. if (ls->fd != -1 && ls->type == RSPAMD_WORKER_SOCKET_UDP) {
  1048. close (ls->fd);
  1049. ls->fd = -1;
  1050. }
  1051. cur = g_list_next (cur);
  1052. }
  1053. #endif
  1054. }
  1055. struct rspamd_worker *
  1056. rspamd_fork_worker (struct rspamd_main *rspamd_main,
  1057. struct rspamd_worker_conf *cf,
  1058. guint index,
  1059. struct ev_loop *ev_base,
  1060. rspamd_worker_term_cb term_handler,
  1061. GHashTable *listen_sockets)
  1062. {
  1063. struct rspamd_worker *wrk;
  1064. /* Starting worker process */
  1065. wrk = (struct rspamd_worker *) g_malloc0 (sizeof (struct rspamd_worker));
  1066. if (!rspamd_socketpair (wrk->control_pipe, SOCK_DGRAM)) {
  1067. msg_err ("socketpair failure: %s", strerror (errno));
  1068. rspamd_hard_terminate (rspamd_main);
  1069. }
  1070. if (!rspamd_socketpair (wrk->srv_pipe, SOCK_DGRAM)) {
  1071. msg_err ("socketpair failure: %s", strerror (errno));
  1072. rspamd_hard_terminate (rspamd_main);
  1073. }
  1074. if (cf->bind_conf) {
  1075. msg_info_main ("prepare to fork process %s (%d); listen on: %s",
  1076. cf->worker->name,
  1077. index, cf->bind_conf->name);
  1078. }
  1079. else {
  1080. msg_info_main ("prepare to fork process %s (%d), no bind socket",
  1081. cf->worker->name,
  1082. index);
  1083. }
  1084. wrk->srv = rspamd_main;
  1085. wrk->type = cf->type;
  1086. wrk->cf = cf;
  1087. wrk->flags = cf->worker->flags;
  1088. REF_RETAIN (cf);
  1089. wrk->index = index;
  1090. wrk->ctx = cf->ctx;
  1091. wrk->ppid = getpid ();
  1092. wrk->pid = fork ();
  1093. wrk->cores_throttled = rspamd_main->cores_throttling;
  1094. wrk->term_handler = term_handler;
  1095. wrk->control_events_pending = g_hash_table_new_full (g_direct_hash, g_direct_equal,
  1096. NULL, rspamd_pending_control_free);
  1097. switch (wrk->pid) {
  1098. case 0:
  1099. rspamd_handle_child_fork (wrk, rspamd_main, cf, listen_sockets);
  1100. break;
  1101. case -1:
  1102. msg_err_main ("cannot fork main process: %s", strerror (errno));
  1103. if (rspamd_main->pfh) {
  1104. rspamd_pidfile_remove (rspamd_main->pfh);
  1105. }
  1106. rspamd_hard_terminate (rspamd_main);
  1107. break;
  1108. default:
  1109. rspamd_handle_main_fork (wrk, rspamd_main, cf, ev_base);
  1110. break;
  1111. }
  1112. return wrk;
  1113. }
  1114. void
  1115. rspamd_worker_block_signals (void)
  1116. {
  1117. sigset_t set;
  1118. sigemptyset (&set);
  1119. sigaddset (&set, SIGTERM);
  1120. sigaddset (&set, SIGINT);
  1121. sigaddset (&set, SIGHUP);
  1122. sigaddset (&set, SIGUSR1);
  1123. sigaddset (&set, SIGUSR2);
  1124. sigprocmask (SIG_BLOCK, &set, NULL);
  1125. }
  1126. void
  1127. rspamd_worker_unblock_signals (void)
  1128. {
  1129. sigset_t set;
  1130. sigemptyset (&set);
  1131. sigaddset (&set, SIGTERM);
  1132. sigaddset (&set, SIGINT);
  1133. sigaddset (&set, SIGHUP);
  1134. sigaddset (&set, SIGUSR1);
  1135. sigaddset (&set, SIGUSR2);
  1136. sigprocmask (SIG_UNBLOCK, &set, NULL);
  1137. }
  1138. void
  1139. rspamd_hard_terminate (struct rspamd_main *rspamd_main)
  1140. {
  1141. GHashTableIter it;
  1142. gpointer k, v;
  1143. struct rspamd_worker *w;
  1144. sigset_t set;
  1145. /* Block all signals */
  1146. sigemptyset (&set);
  1147. sigaddset (&set, SIGTERM);
  1148. sigaddset (&set, SIGINT);
  1149. sigaddset (&set, SIGHUP);
  1150. sigaddset (&set, SIGUSR1);
  1151. sigaddset (&set, SIGUSR2);
  1152. sigaddset (&set, SIGCHLD);
  1153. sigprocmask (SIG_BLOCK, &set, NULL);
  1154. /* We need to terminate all workers that might be already spawned */
  1155. rspamd_worker_block_signals ();
  1156. g_hash_table_iter_init (&it, rspamd_main->workers);
  1157. while (g_hash_table_iter_next (&it, &k, &v)) {
  1158. w = v;
  1159. msg_err_main ("kill worker %P as Rspamd is terminating due to "
  1160. "an unrecoverable error", w->pid);
  1161. kill (w->pid, SIGKILL);
  1162. }
  1163. msg_err_main ("shutting down Rspamd due to fatal error");
  1164. rspamd_log_close (rspamd_main->logger);
  1165. exit (EXIT_FAILURE);
  1166. }
  1167. gboolean
  1168. rspamd_worker_is_scanner (struct rspamd_worker *w)
  1169. {
  1170. if (w) {
  1171. return !!(w->flags & RSPAMD_WORKER_SCANNER);
  1172. }
  1173. return FALSE;
  1174. }
  1175. gboolean
  1176. rspamd_worker_is_primary_controller (struct rspamd_worker *w)
  1177. {
  1178. if (w) {
  1179. return !!(w->flags & RSPAMD_WORKER_CONTROLLER) && w->index == 0;
  1180. }
  1181. return FALSE;
  1182. }
  1183. struct rspamd_worker_session_elt {
  1184. void *ptr;
  1185. guint *pref;
  1186. const gchar *tag;
  1187. time_t when;
  1188. };
  1189. struct rspamd_worker_session_cache {
  1190. struct ev_loop *ev_base;
  1191. GHashTable *cache;
  1192. struct rspamd_config *cfg;
  1193. struct ev_timer periodic;
  1194. };
  1195. static gint
  1196. rspamd_session_cache_sort_cmp (gconstpointer pa, gconstpointer pb)
  1197. {
  1198. const struct rspamd_worker_session_elt
  1199. *e1 = *(const struct rspamd_worker_session_elt **)pa,
  1200. *e2 = *(const struct rspamd_worker_session_elt **)pb;
  1201. return e2->when < e1->when;
  1202. }
  1203. static void
  1204. rspamd_sessions_cache_periodic (EV_P_ ev_timer *w, int revents)
  1205. {
  1206. struct rspamd_worker_session_cache *c =
  1207. (struct rspamd_worker_session_cache *)w->data;
  1208. GHashTableIter it;
  1209. gchar timebuf[32];
  1210. gpointer k, v;
  1211. struct rspamd_worker_session_elt *elt;
  1212. struct tm tms;
  1213. GPtrArray *res;
  1214. guint i;
  1215. if (g_hash_table_size (c->cache) > c->cfg->max_sessions_cache) {
  1216. res = g_ptr_array_sized_new (g_hash_table_size (c->cache));
  1217. g_hash_table_iter_init (&it, c->cache);
  1218. while (g_hash_table_iter_next (&it, &k, &v)) {
  1219. g_ptr_array_add (res, v);
  1220. }
  1221. msg_err ("sessions cache is overflowed %d elements where %d is limit",
  1222. (gint)res->len, (gint)c->cfg->max_sessions_cache);
  1223. g_ptr_array_sort (res, rspamd_session_cache_sort_cmp);
  1224. PTR_ARRAY_FOREACH (res, i, elt) {
  1225. rspamd_localtime (elt->when, &tms);
  1226. strftime (timebuf, sizeof (timebuf), "%F %H:%M:%S", &tms);
  1227. msg_warn ("redundant session; ptr: %p, "
  1228. "tag: %s, refcount: %d, time: %s",
  1229. elt->ptr, elt->tag ? elt->tag : "unknown",
  1230. elt->pref ? *elt->pref : 0,
  1231. timebuf);
  1232. }
  1233. }
  1234. ev_timer_again (EV_A_ w);
  1235. }
  1236. void *
  1237. rspamd_worker_session_cache_new (struct rspamd_worker *w,
  1238. struct ev_loop *ev_base)
  1239. {
  1240. struct rspamd_worker_session_cache *c;
  1241. static const gdouble periodic_interval = 60.0;
  1242. c = g_malloc0 (sizeof (*c));
  1243. c->ev_base = ev_base;
  1244. c->cache = g_hash_table_new_full (g_direct_hash, g_direct_equal,
  1245. NULL, g_free);
  1246. c->cfg = w->srv->cfg;
  1247. c->periodic.data = c;
  1248. ev_timer_init (&c->periodic, rspamd_sessions_cache_periodic, periodic_interval,
  1249. periodic_interval);
  1250. ev_timer_start (ev_base, &c->periodic);
  1251. return c;
  1252. }
  1253. void
  1254. rspamd_worker_session_cache_add (void *cache, const gchar *tag,
  1255. guint *pref, void *ptr)
  1256. {
  1257. struct rspamd_worker_session_cache *c = cache;
  1258. struct rspamd_worker_session_elt *elt;
  1259. elt = g_malloc0 (sizeof (*elt));
  1260. elt->pref = pref;
  1261. elt->ptr = ptr;
  1262. elt->tag = tag;
  1263. elt->when = time (NULL);
  1264. g_hash_table_insert (c->cache, elt->ptr, elt);
  1265. }
  1266. void
  1267. rspamd_worker_session_cache_remove (void *cache, void *ptr)
  1268. {
  1269. struct rspamd_worker_session_cache *c = cache;
  1270. g_hash_table_remove (c->cache, ptr);
  1271. }
  1272. static void
  1273. rspamd_worker_monitored_on_change (struct rspamd_monitored_ctx *ctx,
  1274. struct rspamd_monitored *m, gboolean alive,
  1275. void *ud)
  1276. {
  1277. struct rspamd_worker *worker = ud;
  1278. struct rspamd_config *cfg = worker->srv->cfg;
  1279. struct ev_loop *ev_base;
  1280. guchar tag[RSPAMD_MONITORED_TAG_LEN];
  1281. static struct rspamd_srv_command srv_cmd;
  1282. rspamd_monitored_get_tag (m, tag);
  1283. ev_base = rspamd_monitored_ctx_get_ev_base (ctx);
  1284. memset (&srv_cmd, 0, sizeof (srv_cmd));
  1285. srv_cmd.type = RSPAMD_SRV_MONITORED_CHANGE;
  1286. rspamd_strlcpy (srv_cmd.cmd.monitored_change.tag, tag,
  1287. sizeof (srv_cmd.cmd.monitored_change.tag));
  1288. srv_cmd.cmd.monitored_change.alive = alive;
  1289. srv_cmd.cmd.monitored_change.sender = getpid ();
  1290. msg_info_config ("broadcast monitored update for %s: %s",
  1291. srv_cmd.cmd.monitored_change.tag, alive ? "alive" : "dead");
  1292. rspamd_srv_send_command (worker, ev_base, &srv_cmd, -1, NULL, NULL);
  1293. }
  1294. void
  1295. rspamd_worker_init_monitored (struct rspamd_worker *worker,
  1296. struct ev_loop *ev_base,
  1297. struct rspamd_dns_resolver *resolver)
  1298. {
  1299. rspamd_monitored_ctx_config (worker->srv->cfg->monitored_ctx,
  1300. worker->srv->cfg, ev_base, resolver->r,
  1301. rspamd_worker_monitored_on_change, worker);
  1302. }
  1303. #ifdef HAVE_SA_SIGINFO
  1304. #ifdef WITH_LIBUNWIND
  1305. static void
  1306. rspamd_print_crash (ucontext_t *uap)
  1307. {
  1308. unw_cursor_t cursor;
  1309. unw_word_t ip, off;
  1310. guint level;
  1311. gint ret;
  1312. if ((ret = unw_init_local (&cursor, uap)) != 0) {
  1313. msg_err ("unw_init_local: %d", ret);
  1314. return;
  1315. }
  1316. level = 0;
  1317. ret = 0;
  1318. for (;;) {
  1319. char name[128];
  1320. if (level >= UNWIND_BACKTRACE_DEPTH) {
  1321. break;
  1322. }
  1323. unw_get_reg (&cursor, UNW_REG_IP, &ip);
  1324. ret = unw_get_proc_name(&cursor, name, sizeof (name), &off);
  1325. if (ret == 0) {
  1326. msg_err ("%d: %p: %s()+0x%xl",
  1327. level, ip, name, (uintptr_t)off);
  1328. } else {
  1329. msg_err ("%d: %p: <unknown>", level, ip);
  1330. }
  1331. level++;
  1332. ret = unw_step (&cursor);
  1333. if (ret <= 0) {
  1334. break;
  1335. }
  1336. }
  1337. if (ret < 0) {
  1338. msg_err ("unw_step_ptr: %d", ret);
  1339. }
  1340. }
  1341. #endif
  1342. static struct rspamd_main *saved_main = NULL;
  1343. static gboolean
  1344. rspamd_crash_propagate (gpointer key, gpointer value, gpointer unused)
  1345. {
  1346. struct rspamd_worker *w = value;
  1347. /* Kill children softly */
  1348. kill (w->pid, SIGTERM);
  1349. return TRUE;
  1350. }
  1351. static void
  1352. rspamd_crash_sig_handler (int sig, siginfo_t *info, void *ctx)
  1353. {
  1354. struct sigaction sa;
  1355. ucontext_t *uap = ctx;
  1356. pid_t pid;
  1357. pid = getpid ();
  1358. msg_err ("caught fatal signal %d(%s), "
  1359. "pid: %P, trace: ",
  1360. sig, strsignal (sig), pid);
  1361. (void)uap;
  1362. #ifdef WITH_LIBUNWIND
  1363. rspamd_print_crash (uap);
  1364. #endif
  1365. msg_err ("please see Rspamd FAQ to learn how to dump core files and how to "
  1366. "fill a bug report");
  1367. if (saved_main) {
  1368. if (pid == saved_main->pid) {
  1369. /*
  1370. * Main process has crashed, propagate crash further to trigger
  1371. * monitoring alerts and mass panic
  1372. */
  1373. g_hash_table_foreach_remove (saved_main->workers,
  1374. rspamd_crash_propagate, NULL);
  1375. }
  1376. }
  1377. /*
  1378. * Invoke signal with the default handler
  1379. */
  1380. sigemptyset (&sa.sa_mask);
  1381. sa.sa_handler = SIG_DFL;
  1382. sa.sa_flags = 0;
  1383. sigaction (sig, &sa, NULL);
  1384. kill (pid, sig);
  1385. }
  1386. #endif
  1387. RSPAMD_NO_SANITIZE void
  1388. rspamd_set_crash_handler (struct rspamd_main *rspamd_main)
  1389. {
  1390. #ifdef HAVE_SA_SIGINFO
  1391. struct sigaction sa;
  1392. #ifdef HAVE_SIGALTSTACK
  1393. void *stack_mem;
  1394. stack_t ss;
  1395. memset (&ss, 0, sizeof ss);
  1396. /*
  1397. * Allocate special stack, NOT freed at the end so far
  1398. * It also cannot be on stack as this memory is used when
  1399. * stack corruption is detected. Leak sanitizer blames about it but
  1400. * I don't know any good ways to stop this behaviour.
  1401. */
  1402. ss.ss_size = MAX (SIGSTKSZ, 8192 * 4);
  1403. stack_mem = g_malloc0 (ss.ss_size);
  1404. ss.ss_sp = stack_mem;
  1405. sigaltstack (&ss, NULL);
  1406. #endif
  1407. saved_main = rspamd_main;
  1408. sigemptyset (&sa.sa_mask);
  1409. sa.sa_sigaction = &rspamd_crash_sig_handler;
  1410. sa.sa_flags = SA_RESTART | SA_SIGINFO | SA_ONSTACK;
  1411. sigaction (SIGSEGV, &sa, NULL);
  1412. sigaction (SIGBUS, &sa, NULL);
  1413. sigaction (SIGABRT, &sa, NULL);
  1414. sigaction (SIGFPE, &sa, NULL);
  1415. sigaction (SIGSYS, &sa, NULL);
  1416. #endif
  1417. }
  1418. static void
  1419. rspamd_enable_accept_event (EV_P_ ev_timer *w, int revents)
  1420. {
  1421. struct rspamd_worker_accept_event *ac_ev =
  1422. (struct rspamd_worker_accept_event *)w->data;
  1423. ev_timer_stop (EV_A_ w);
  1424. ev_io_start (EV_A_ &ac_ev->accept_ev);
  1425. }
  1426. void
  1427. rspamd_worker_throttle_accept_events (gint sock, void *data)
  1428. {
  1429. struct rspamd_worker_accept_event *head, *cur;
  1430. const gdouble throttling = 0.5;
  1431. head = (struct rspamd_worker_accept_event *)data;
  1432. DL_FOREACH (head, cur) {
  1433. ev_io_stop (cur->event_loop, &cur->accept_ev);
  1434. cur->throttling_ev.data = cur;
  1435. ev_timer_init (&cur->throttling_ev, rspamd_enable_accept_event,
  1436. throttling, 0.0);
  1437. ev_timer_start (cur->event_loop, &cur->throttling_ev);
  1438. }
  1439. }
  1440. gboolean
  1441. rspamd_check_termination_clause (struct rspamd_main *rspamd_main,
  1442. struct rspamd_worker *wrk,
  1443. int res)
  1444. {
  1445. gboolean need_refork = TRUE;
  1446. if (wrk->state != rspamd_worker_state_running || rspamd_main->wanna_die ||
  1447. (wrk->flags & RSPAMD_WORKER_OLD_CONFIG)) {
  1448. /* Do not refork workers that are intended to be terminated */
  1449. need_refork = FALSE;
  1450. }
  1451. if (WIFEXITED (res) && WEXITSTATUS (res) == 0) {
  1452. /* Normal worker termination, do not fork one more */
  1453. if (wrk->flags & RSPAMD_WORKER_OLD_CONFIG) {
  1454. /* Never re-fork old workers */
  1455. msg_info_main ("%s process %P terminated normally",
  1456. g_quark_to_string(wrk->type),
  1457. wrk->pid);
  1458. need_refork = FALSE;
  1459. }
  1460. else {
  1461. if (wrk->hb.nbeats < 0 && rspamd_main->cfg->heartbeats_loss_max > 0 &&
  1462. -(wrk->hb.nbeats) >= rspamd_main->cfg->heartbeats_loss_max) {
  1463. msg_info_main ("%s process %P terminated normally, but lost %L "
  1464. "heartbeats, refork it",
  1465. g_quark_to_string(wrk->type),
  1466. wrk->pid,
  1467. -(wrk->hb.nbeats));
  1468. need_refork = TRUE;
  1469. }
  1470. else {
  1471. msg_info_main ("%s process %P terminated normally",
  1472. g_quark_to_string(wrk->type),
  1473. wrk->pid);
  1474. need_refork = FALSE;
  1475. }
  1476. }
  1477. }
  1478. else {
  1479. if (WIFSIGNALED (res)) {
  1480. #ifdef WCOREDUMP
  1481. if (WCOREDUMP (res)) {
  1482. msg_warn_main (
  1483. "%s process %P terminated abnormally by signal: %s"
  1484. " and created core file; please see Rspamd FAQ "
  1485. "to learn how to extract data from core file and "
  1486. "fill a bug report",
  1487. g_quark_to_string (wrk->type),
  1488. wrk->pid,
  1489. g_strsignal (WTERMSIG (res)));
  1490. }
  1491. else {
  1492. #ifdef HAVE_SYS_RESOURCE_H
  1493. struct rlimit rlmt;
  1494. (void) getrlimit (RLIMIT_CORE, &rlmt);
  1495. msg_warn_main (
  1496. "%s process %P terminated abnormally with exit code %d by "
  1497. "signal: %s"
  1498. " but NOT created core file (throttled=%s); "
  1499. "core file limits: %L current, %L max",
  1500. g_quark_to_string (wrk->type),
  1501. wrk->pid,
  1502. WEXITSTATUS (res),
  1503. g_strsignal (WTERMSIG (res)),
  1504. wrk->cores_throttled ? "yes" : "no",
  1505. (gint64) rlmt.rlim_cur,
  1506. (gint64) rlmt.rlim_max);
  1507. #else
  1508. msg_warn_main (
  1509. "%s process %P terminated abnormally with exit code %d by signal: %s"
  1510. " but NOT created core file (throttled=%s); ",
  1511. g_quark_to_string (wrk->type),
  1512. wrk->pid, WEXITSTATUS (res),
  1513. g_strsignal (WTERMSIG (res)),
  1514. wrk->cores_throttled ? "yes" : "no");
  1515. #endif
  1516. }
  1517. #else
  1518. msg_warn_main (
  1519. "%s process %P terminated abnormally with exit code %d by signal: %s",
  1520. g_quark_to_string (wrk->type),
  1521. wrk->pid, WEXITSTATUS (res),
  1522. g_strsignal (WTERMSIG (res)));
  1523. #endif
  1524. if (WTERMSIG (res) == SIGUSR2) {
  1525. /*
  1526. * It is actually race condition when not started process
  1527. * has been requested to be reloaded.
  1528. *
  1529. * We shouldn't refork on this
  1530. */
  1531. need_refork = FALSE;
  1532. }
  1533. }
  1534. else {
  1535. msg_warn_main ("%s process %P terminated abnormally "
  1536. "(but it was not killed by a signal) "
  1537. "with exit code %d",
  1538. g_quark_to_string (wrk->type),
  1539. wrk->pid,
  1540. WEXITSTATUS (res));
  1541. }
  1542. }
  1543. return need_refork;
  1544. }
  1545. #ifdef WITH_HYPERSCAN
  1546. gboolean
  1547. rspamd_worker_hyperscan_ready (struct rspamd_main *rspamd_main,
  1548. struct rspamd_worker *worker, gint fd,
  1549. gint attached_fd,
  1550. struct rspamd_control_command *cmd,
  1551. gpointer ud) {
  1552. struct rspamd_control_reply rep;
  1553. struct rspamd_re_cache *cache = worker->srv->cfg->re_cache;
  1554. memset (&rep, 0, sizeof (rep));
  1555. rep.type = RSPAMD_CONTROL_HYPERSCAN_LOADED;
  1556. if (rspamd_re_cache_is_hs_loaded (cache) != RSPAMD_HYPERSCAN_LOADED_FULL ||
  1557. cmd->cmd.hs_loaded.forced) {
  1558. msg_info ("loading hyperscan expressions after receiving compilation "
  1559. "notice: %s",
  1560. (rspamd_re_cache_is_hs_loaded (cache) != RSPAMD_HYPERSCAN_LOADED_FULL) ?
  1561. "new db" : "forced update");
  1562. rep.reply.hs_loaded.status = rspamd_re_cache_load_hyperscan (
  1563. worker->srv->cfg->re_cache, cmd->cmd.hs_loaded.cache_dir, false);
  1564. }
  1565. if (write (fd, &rep, sizeof (rep)) != sizeof (rep)) {
  1566. msg_err ("cannot write reply to the control socket: %s",
  1567. strerror (errno));
  1568. }
  1569. return TRUE;
  1570. }
  1571. #endif /* With Hyperscan */
  1572. gboolean
  1573. rspamd_worker_check_context (gpointer ctx, guint64 magic)
  1574. {
  1575. struct rspamd_abstract_worker_ctx *actx = (struct rspamd_abstract_worker_ctx*)ctx;
  1576. return actx->magic == magic;
  1577. }
  1578. static gboolean
  1579. rspamd_worker_log_pipe_handler (struct rspamd_main *rspamd_main,
  1580. struct rspamd_worker *worker, gint fd,
  1581. gint attached_fd,
  1582. struct rspamd_control_command *cmd,
  1583. gpointer ud)
  1584. {
  1585. struct rspamd_config *cfg = ud;
  1586. struct rspamd_worker_log_pipe *lp;
  1587. struct rspamd_control_reply rep;
  1588. memset (&rep, 0, sizeof (rep));
  1589. rep.type = RSPAMD_CONTROL_LOG_PIPE;
  1590. if (attached_fd != -1) {
  1591. lp = g_malloc0 (sizeof (*lp));
  1592. lp->fd = attached_fd;
  1593. lp->type = cmd->cmd.log_pipe.type;
  1594. DL_APPEND (cfg->log_pipes, lp);
  1595. msg_info ("added new log pipe");
  1596. }
  1597. else {
  1598. rep.reply.log_pipe.status = ENOENT;
  1599. msg_err ("cannot attach log pipe: invalid fd");
  1600. }
  1601. if (write (fd, &rep, sizeof (rep)) != sizeof (rep)) {
  1602. msg_err ("cannot write reply to the control socket: %s",
  1603. strerror (errno));
  1604. }
  1605. return TRUE;
  1606. }
  1607. static gboolean
  1608. rspamd_worker_monitored_handler (struct rspamd_main *rspamd_main,
  1609. struct rspamd_worker *worker, gint fd,
  1610. gint attached_fd,
  1611. struct rspamd_control_command *cmd,
  1612. gpointer ud)
  1613. {
  1614. struct rspamd_control_reply rep;
  1615. struct rspamd_monitored *m;
  1616. struct rspamd_monitored_ctx *mctx = worker->srv->cfg->monitored_ctx;
  1617. struct rspamd_config *cfg = ud;
  1618. memset (&rep, 0, sizeof (rep));
  1619. rep.type = RSPAMD_CONTROL_MONITORED_CHANGE;
  1620. if (cmd->cmd.monitored_change.sender != getpid ()) {
  1621. m = rspamd_monitored_by_tag (mctx, cmd->cmd.monitored_change.tag);
  1622. if (m != NULL) {
  1623. rspamd_monitored_set_alive (m, cmd->cmd.monitored_change.alive);
  1624. rep.reply.monitored_change.status = 1;
  1625. msg_info_config ("updated monitored status for %s: %s",
  1626. cmd->cmd.monitored_change.tag,
  1627. cmd->cmd.monitored_change.alive ? "alive" : "dead");
  1628. } else {
  1629. msg_err ("cannot find monitored by tag: %*s", 32,
  1630. cmd->cmd.monitored_change.tag);
  1631. rep.reply.monitored_change.status = 0;
  1632. }
  1633. }
  1634. if (write (fd, &rep, sizeof (rep)) != sizeof (rep)) {
  1635. msg_err ("cannot write reply to the control socket: %s",
  1636. strerror (errno));
  1637. }
  1638. return TRUE;
  1639. }
  1640. void
  1641. rspamd_worker_init_scanner (struct rspamd_worker *worker,
  1642. struct ev_loop *ev_base,
  1643. struct rspamd_dns_resolver *resolver,
  1644. struct rspamd_lang_detector **plang_det)
  1645. {
  1646. rspamd_stat_init (worker->srv->cfg, ev_base);
  1647. #ifdef WITH_HYPERSCAN
  1648. rspamd_control_worker_add_cmd_handler (worker,
  1649. RSPAMD_CONTROL_HYPERSCAN_LOADED,
  1650. rspamd_worker_hyperscan_ready,
  1651. NULL);
  1652. #endif
  1653. rspamd_control_worker_add_cmd_handler (worker,
  1654. RSPAMD_CONTROL_LOG_PIPE,
  1655. rspamd_worker_log_pipe_handler,
  1656. worker->srv->cfg);
  1657. rspamd_control_worker_add_cmd_handler (worker,
  1658. RSPAMD_CONTROL_MONITORED_CHANGE,
  1659. rspamd_worker_monitored_handler,
  1660. worker->srv->cfg);
  1661. *plang_det = worker->srv->cfg->lang_det;
  1662. }
  1663. void
  1664. rspamd_controller_store_saved_stats (struct rspamd_main *rspamd_main,
  1665. struct rspamd_config *cfg)
  1666. {
  1667. struct rspamd_stat *stat;
  1668. ucl_object_t *top, *sub;
  1669. struct ucl_emitter_functions *efuncs;
  1670. gint i, fd;
  1671. FILE *fp;
  1672. gchar fpath[PATH_MAX];
  1673. if (cfg->stats_file == NULL) {
  1674. return;
  1675. }
  1676. rspamd_snprintf (fpath, sizeof (fpath), "%s.XXXXXXXX", cfg->stats_file);
  1677. fd = g_mkstemp_full (fpath, O_WRONLY|O_TRUNC, 00644);
  1678. if (fd == -1) {
  1679. msg_err_config ("cannot open for writing controller stats from %s: %s",
  1680. fpath, strerror (errno));
  1681. return;
  1682. }
  1683. fp = fdopen (fd, "w");
  1684. stat = rspamd_main->stat;
  1685. top = ucl_object_typed_new (UCL_OBJECT);
  1686. ucl_object_insert_key (top, ucl_object_fromint (
  1687. stat->messages_scanned), "scanned", 0, false);
  1688. ucl_object_insert_key (top, ucl_object_fromint (
  1689. stat->messages_learned), "learned", 0, false);
  1690. if (stat->messages_scanned > 0) {
  1691. sub = ucl_object_typed_new (UCL_OBJECT);
  1692. for (i = METRIC_ACTION_REJECT; i <= METRIC_ACTION_NOACTION; i++) {
  1693. ucl_object_insert_key (sub,
  1694. ucl_object_fromint (stat->actions_stat[i]),
  1695. rspamd_action_to_str (i), 0, false);
  1696. }
  1697. ucl_object_insert_key (top, sub, "actions", 0, false);
  1698. }
  1699. ucl_object_insert_key (top,
  1700. ucl_object_fromint (stat->connections_count),
  1701. "connections", 0, false);
  1702. ucl_object_insert_key (top,
  1703. ucl_object_fromint (stat->control_connections_count),
  1704. "control_connections", 0, false);
  1705. efuncs = ucl_object_emit_file_funcs (fp);
  1706. if (!ucl_object_emit_full (top, UCL_EMIT_JSON_COMPACT,
  1707. efuncs, NULL)) {
  1708. msg_err_config ("cannot write stats to %s: %s",
  1709. fpath, strerror (errno));
  1710. unlink (fpath);
  1711. }
  1712. else {
  1713. if (rename (fpath, cfg->stats_file) == -1) {
  1714. msg_err_config ("cannot rename stats from %s to %s: %s",
  1715. fpath, cfg->stats_file, strerror (errno));
  1716. }
  1717. }
  1718. ucl_object_unref (top);
  1719. fclose (fp);
  1720. ucl_object_emit_funcs_free (efuncs);
  1721. }
  1722. static ev_timer rrd_timer;
  1723. void
  1724. rspamd_controller_on_terminate (struct rspamd_worker *worker,
  1725. struct rspamd_rrd_file *rrd)
  1726. {
  1727. struct rspamd_abstract_worker_ctx *ctx;
  1728. ctx = (struct rspamd_abstract_worker_ctx *)worker->ctx;
  1729. rspamd_controller_store_saved_stats (worker->srv, worker->srv->cfg);
  1730. if (rrd) {
  1731. ev_timer_stop (ctx->event_loop, &rrd_timer);
  1732. msg_info ("closing rrd file: %s", rrd->filename);
  1733. rspamd_rrd_close (rrd);
  1734. }
  1735. }
  1736. static void
  1737. rspamd_controller_load_saved_stats (struct rspamd_main *rspamd_main,
  1738. struct rspamd_config *cfg)
  1739. {
  1740. struct ucl_parser *parser;
  1741. ucl_object_t *obj;
  1742. const ucl_object_t *elt, *subelt;
  1743. struct rspamd_stat *stat, stat_copy;
  1744. gint i;
  1745. if (cfg->stats_file == NULL) {
  1746. return;
  1747. }
  1748. if (access (cfg->stats_file, R_OK) == -1) {
  1749. msg_err_config ("cannot load controller stats from %s: %s",
  1750. cfg->stats_file, strerror (errno));
  1751. return;
  1752. }
  1753. parser = ucl_parser_new (0);
  1754. if (!ucl_parser_add_file (parser, cfg->stats_file)) {
  1755. msg_err_config ("cannot parse controller stats from %s: %s",
  1756. cfg->stats_file, ucl_parser_get_error (parser));
  1757. ucl_parser_free (parser);
  1758. return;
  1759. }
  1760. obj = ucl_parser_get_object (parser);
  1761. ucl_parser_free (parser);
  1762. stat = rspamd_main->stat;
  1763. memcpy (&stat_copy, stat, sizeof (stat_copy));
  1764. elt = ucl_object_lookup (obj, "scanned");
  1765. if (elt != NULL && ucl_object_type (elt) == UCL_INT) {
  1766. stat_copy.messages_scanned = ucl_object_toint (elt);
  1767. }
  1768. elt = ucl_object_lookup (obj, "learned");
  1769. if (elt != NULL && ucl_object_type (elt) == UCL_INT) {
  1770. stat_copy.messages_learned = ucl_object_toint (elt);
  1771. }
  1772. elt = ucl_object_lookup (obj, "actions");
  1773. if (elt != NULL) {
  1774. for (i = METRIC_ACTION_REJECT; i <= METRIC_ACTION_NOACTION; i++) {
  1775. subelt = ucl_object_lookup (elt, rspamd_action_to_str (i));
  1776. if (subelt && ucl_object_type (subelt) == UCL_INT) {
  1777. stat_copy.actions_stat[i] = ucl_object_toint (subelt);
  1778. }
  1779. }
  1780. }
  1781. elt = ucl_object_lookup (obj, "connections_count");
  1782. if (elt != NULL && ucl_object_type (elt) == UCL_INT) {
  1783. stat_copy.connections_count = ucl_object_toint (elt);
  1784. }
  1785. elt = ucl_object_lookup (obj, "control_connections_count");
  1786. if (elt != NULL && ucl_object_type (elt) == UCL_INT) {
  1787. stat_copy.control_connections_count = ucl_object_toint (elt);
  1788. }
  1789. ucl_object_unref (obj);
  1790. memcpy (stat, &stat_copy, sizeof (stat_copy));
  1791. }
  1792. struct rspamd_controller_periodics_cbdata {
  1793. struct rspamd_worker *worker;
  1794. struct rspamd_rrd_file *rrd;
  1795. struct rspamd_stat *stat;
  1796. ev_timer save_stats_event;
  1797. };
  1798. static void
  1799. rspamd_controller_rrd_update (EV_P_ ev_timer *w, int revents)
  1800. {
  1801. struct rspamd_controller_periodics_cbdata *cbd =
  1802. (struct rspamd_controller_periodics_cbdata *)w->data;
  1803. struct rspamd_stat *stat;
  1804. GArray ar;
  1805. gdouble points[METRIC_ACTION_MAX];
  1806. GError *err = NULL;
  1807. guint i;
  1808. g_assert (cbd->rrd != NULL);
  1809. stat = cbd->stat;
  1810. for (i = METRIC_ACTION_REJECT; i < METRIC_ACTION_MAX; i ++) {
  1811. points[i] = stat->actions_stat[i];
  1812. }
  1813. ar.data = (gchar *)points;
  1814. ar.len = sizeof (points);
  1815. if (!rspamd_rrd_add_record (cbd->rrd, &ar, rspamd_get_calendar_ticks (),
  1816. &err)) {
  1817. msg_err ("cannot update rrd file: %e", err);
  1818. g_error_free (err);
  1819. }
  1820. /* Plan new event */
  1821. ev_timer_again (EV_A_ w);
  1822. }
  1823. static void
  1824. rspamd_controller_stats_save_periodic (EV_P_ ev_timer *w, int revents)
  1825. {
  1826. struct rspamd_controller_periodics_cbdata *cbd =
  1827. (struct rspamd_controller_periodics_cbdata *)w->data;
  1828. rspamd_controller_store_saved_stats (cbd->worker->srv, cbd->worker->srv->cfg);
  1829. ev_timer_again (EV_A_ w);
  1830. }
  1831. void
  1832. rspamd_worker_init_controller (struct rspamd_worker *worker,
  1833. struct rspamd_rrd_file **prrd)
  1834. {
  1835. struct rspamd_abstract_worker_ctx *ctx;
  1836. static const ev_tstamp rrd_update_time = 1.0;
  1837. ctx = (struct rspamd_abstract_worker_ctx *)worker->ctx;
  1838. rspamd_controller_load_saved_stats (worker->srv, worker->srv->cfg);
  1839. if (worker->index == 0) {
  1840. /* Enable periodics and other stuff */
  1841. static struct rspamd_controller_periodics_cbdata cbd;
  1842. const ev_tstamp save_stats_interval = 60; /* 1 minute */
  1843. memset (&cbd, 0, sizeof (cbd));
  1844. cbd.save_stats_event.data = &cbd;
  1845. cbd.worker = worker;
  1846. cbd.stat = worker->srv->stat;
  1847. ev_timer_init (&cbd.save_stats_event,
  1848. rspamd_controller_stats_save_periodic,
  1849. save_stats_interval, save_stats_interval);
  1850. ev_timer_start (ctx->event_loop, &cbd.save_stats_event);
  1851. rspamd_map_watch (worker->srv->cfg, ctx->event_loop,
  1852. ctx->resolver, worker,
  1853. RSPAMD_MAP_WATCH_PRIMARY_CONTROLLER);
  1854. if (prrd != NULL) {
  1855. if (ctx->cfg->rrd_file && worker->index == 0) {
  1856. GError *rrd_err = NULL;
  1857. *prrd = rspamd_rrd_file_default (ctx->cfg->rrd_file, &rrd_err);
  1858. if (*prrd) {
  1859. cbd.rrd = *prrd;
  1860. rrd_timer.data = &cbd;
  1861. ev_timer_init (&rrd_timer, rspamd_controller_rrd_update,
  1862. rrd_update_time, rrd_update_time);
  1863. ev_timer_start (ctx->event_loop, &rrd_timer);
  1864. }
  1865. else if (rrd_err) {
  1866. msg_err ("cannot load rrd from %s: %e", ctx->cfg->rrd_file,
  1867. rrd_err);
  1868. g_error_free (rrd_err);
  1869. }
  1870. else {
  1871. msg_err ("cannot load rrd from %s: unknown error",
  1872. ctx->cfg->rrd_file);
  1873. }
  1874. }
  1875. else {
  1876. *prrd = NULL;
  1877. }
  1878. }
  1879. if (!ctx->cfg->disable_monitored) {
  1880. rspamd_worker_init_monitored (worker,
  1881. ctx->event_loop, ctx->resolver);
  1882. }
  1883. }
  1884. else {
  1885. rspamd_map_watch (worker->srv->cfg, ctx->event_loop,
  1886. ctx->resolver, worker, RSPAMD_MAP_WATCH_SCANNER);
  1887. }
  1888. }