Ви не можете вибрати більше 25 тем Теми мають розпочинатися з літери або цифри, можуть містити дефіси (-) і не повинні перевищувати 35 символів.

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