123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923 |
- /*-
- * Copyright 2016 Vsevolod Stakhov
- *
- * Licensed under the Apache License, Version 2.0 (the "License");
- * you may not use this file except in compliance with the License.
- * You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
- #include "config.h"
- #include "rspamd.h"
- #include "stat_internal.h"
- #include "upstream.h"
- #include "lua/lua_common.h"
- #include "libserver/mempool_vars_internal.h"
-
- #ifdef WITH_HIREDIS
- #include "hiredis.h"
- #include "adapters/libevent.h"
- #include "ref.h"
-
- #define msg_debug_stat_redis(...) rspamd_conditional_debug_fast (NULL, NULL, \
- rspamd_stat_redis_log_id, "stat_redis", task->task_pool->tag.uid, \
- G_STRFUNC, \
- __VA_ARGS__)
-
- INIT_LOG_MODULE(stat_redis)
-
- #define REDIS_CTX(p) (struct redis_stat_ctx *)(p)
- #define REDIS_RUNTIME(p) (struct redis_stat_runtime *)(p)
- #define REDIS_BACKEND_TYPE "redis"
- #define REDIS_DEFAULT_PORT 6379
- #define REDIS_DEFAULT_OBJECT "%s%l"
- #define REDIS_DEFAULT_USERS_OBJECT "%s%l%r"
- #define REDIS_DEFAULT_TIMEOUT 0.5
- #define REDIS_STAT_TIMEOUT 30
-
- struct redis_stat_ctx {
- lua_State *L;
- struct rspamd_statfile_config *stcf;
- gint conf_ref;
- struct rspamd_stat_async_elt *stat_elt;
- const gchar *redis_object;
- const gchar *password;
- const gchar *dbname;
- gdouble timeout;
- gboolean enable_users;
- gboolean store_tokens;
- gboolean new_schema;
- gboolean enable_signatures;
- guint expiry;
- gint cbref_user;
- };
-
- enum rspamd_redis_connection_state {
- RSPAMD_REDIS_DISCONNECTED = 0,
- RSPAMD_REDIS_CONNECTED,
- RSPAMD_REDIS_REQUEST_SENT,
- RSPAMD_REDIS_TIMEDOUT,
- RSPAMD_REDIS_TERMINATED
- };
-
- struct redis_stat_runtime {
- struct redis_stat_ctx *ctx;
- struct rspamd_task *task;
- struct upstream *selected;
- struct event timeout_event;
- GArray *results;
- struct rspamd_statfile_config *stcf;
- gchar *redis_object_expanded;
- redisAsyncContext *redis;
- guint64 learned;
- gint id;
- gboolean has_event;
- GError *err;
- };
-
- /* Used to get statistics from redis */
- struct rspamd_redis_stat_cbdata;
-
- struct rspamd_redis_stat_elt {
- struct redis_stat_ctx *ctx;
- struct rspamd_stat_async_elt *async;
- struct event_base *ev_base;
- ucl_object_t *stat;
- struct rspamd_redis_stat_cbdata *cbdata;
- };
-
- struct rspamd_redis_stat_cbdata {
- struct rspamd_redis_stat_elt *elt;
- redisAsyncContext *redis;
- ucl_object_t *cur;
- GPtrArray *cur_keys;
- struct upstream *selected;
- guint inflight;
- gboolean wanna_die;
- };
-
- #define GET_TASK_ELT(task, elt) (task == NULL ? NULL : (task)->elt)
-
- static const gchar *M = "redis statistics";
-
- static GQuark
- rspamd_redis_stat_quark (void)
- {
- return g_quark_from_static_string (M);
- }
-
- static inline struct upstream_list *
- rspamd_redis_get_servers (struct redis_stat_ctx *ctx,
- const gchar *what)
- {
- lua_State *L = ctx->L;
- struct upstream_list *res;
-
- lua_rawgeti (L, LUA_REGISTRYINDEX, ctx->conf_ref);
- lua_pushstring (L, what);
- lua_gettable (L, -2);
- res = *((struct upstream_list**)lua_touserdata (L, -1));
- lua_settop (L, 0);
-
- return res;
- }
-
- /*
- * Non-static for lua unit testing
- */
- gsize
- rspamd_redis_expand_object (const gchar *pattern,
- struct redis_stat_ctx *ctx,
- struct rspamd_task *task,
- gchar **target)
- {
- gsize tlen = 0;
- const gchar *p = pattern, *elt;
- gchar *d, *end;
- enum {
- just_char,
- percent_char,
- mod_char
- } state = just_char;
- struct rspamd_statfile_config *stcf;
- lua_State *L = NULL;
- struct rspamd_task **ptask;
- GString *tb;
- const gchar *rcpt = NULL;
- gint err_idx;
-
- g_assert (ctx != NULL);
- stcf = ctx->stcf;
-
- L = task->cfg->lua_state;
-
- if (ctx->enable_users) {
- if (ctx->cbref_user == -1) {
- rcpt = rspamd_task_get_principal_recipient (task);
- }
- else if (L) {
- /* Execute lua function to get userdata */
- lua_pushcfunction (L, &rspamd_lua_traceback);
- err_idx = lua_gettop (L);
-
- lua_rawgeti (L, LUA_REGISTRYINDEX, ctx->cbref_user);
- ptask = lua_newuserdata (L, sizeof (struct rspamd_task *));
- *ptask = task;
- rspamd_lua_setclass (L, "rspamd{task}", -1);
-
- if (lua_pcall (L, 1, 1, err_idx) != 0) {
- tb = lua_touserdata (L, -1);
- msg_err_task ("call to user extraction script failed: %v", tb);
- g_string_free (tb, TRUE);
- }
- else {
- rcpt = rspamd_mempool_strdup (task->task_pool, lua_tostring (L, -1));
- }
-
- /* Result + error function */
- lua_pop (L, 2);
- }
-
- if (rcpt) {
- rspamd_mempool_set_variable (task->task_pool, "stat_user",
- (gpointer)rcpt, NULL);
- }
- }
-
- /* Length calculation */
- while (*p) {
- switch (state) {
- case just_char:
- if (*p == '%') {
- state = percent_char;
- }
- else {
- tlen ++;
- }
- p ++;
- break;
- case percent_char:
- switch (*p) {
- case '%':
- tlen ++;
- state = just_char;
- break;
- case 'u':
- elt = GET_TASK_ELT (task, user);
- if (elt) {
- tlen += strlen (elt);
- }
- break;
- case 'r':
-
- if (rcpt == NULL) {
- elt = rspamd_task_get_principal_recipient (task);
- }
- else {
- elt = rcpt;
- }
-
- if (elt) {
- tlen += strlen (elt);
- }
- break;
- case 'l':
- if (stcf->label) {
- tlen += strlen (stcf->label);
- }
- /* Label miss is OK */
- break;
- case 's':
- if (ctx->new_schema) {
- tlen += sizeof ("RS") - 1;
- }
- else {
- if (stcf->symbol) {
- tlen += strlen (stcf->symbol);
- }
- }
- break;
- default:
- state = just_char;
- tlen ++;
- break;
- }
-
- if (state == percent_char) {
- state = mod_char;
- }
- p ++;
- break;
-
- case mod_char:
- switch (*p) {
- case 'd':
- p ++;
- state = just_char;
- break;
- default:
- state = just_char;
- break;
- }
- break;
- }
- }
-
-
- if (target == NULL || task == NULL) {
- return -1;
- }
-
- *target = rspamd_mempool_alloc (task->task_pool, tlen + 1);
- d = *target;
- end = d + tlen + 1;
- d[tlen] = '\0';
- p = pattern;
- state = just_char;
-
- /* Expand string */
- while (*p && d < end) {
- switch (state) {
- case just_char:
- if (*p == '%') {
- state = percent_char;
- }
- else {
- *d++ = *p;
- }
- p ++;
- break;
- case percent_char:
- switch (*p) {
- case '%':
- *d++ = *p;
- state = just_char;
- break;
- case 'u':
- elt = GET_TASK_ELT (task, user);
- if (elt) {
- d += rspamd_strlcpy (d, elt, end - d);
- }
- break;
- case 'r':
- if (rcpt == NULL) {
- elt = rspamd_task_get_principal_recipient (task);
- }
- else {
- elt = rcpt;
- }
-
- if (elt) {
- d += rspamd_strlcpy (d, elt, end - d);
- }
- break;
- case 'l':
- if (stcf->label) {
- d += rspamd_strlcpy (d, stcf->label, end - d);
- }
- break;
- case 's':
- if (ctx->new_schema) {
- d += rspamd_strlcpy (d, "RS", end - d);
- }
- else {
- if (stcf->symbol) {
- d += rspamd_strlcpy (d, stcf->symbol, end - d);
- }
- }
- break;
- default:
- state = just_char;
- *d++ = *p;
- break;
- }
-
- if (state == percent_char) {
- state = mod_char;
- }
- p ++;
- break;
-
- case mod_char:
- switch (*p) {
- case 'd':
- /* TODO: not supported yet */
- p ++;
- state = just_char;
- break;
- default:
- state = just_char;
- break;
- }
- break;
- }
- }
-
- return tlen;
- }
-
- static void
- rspamd_redis_maybe_auth (struct redis_stat_ctx *ctx, redisAsyncContext *redis)
- {
- if (ctx->password) {
- redisAsyncCommand (redis, NULL, NULL, "AUTH %s", ctx->password);
- }
- if (ctx->dbname) {
- redisAsyncCommand (redis, NULL, NULL, "SELECT %s", ctx->dbname);
- }
- }
-
- static rspamd_fstring_t *
- rspamd_redis_tokens_to_query (struct rspamd_task *task,
- struct redis_stat_runtime *rt,
- GPtrArray *tokens,
- const gchar *command,
- const gchar *prefix,
- gboolean learn,
- gint idx,
- gboolean intvals)
- {
- rspamd_fstring_t *out;
- rspamd_token_t *tok;
- gchar n0[512], n1[64];
- guint i, l0, l1, cmd_len, prefix_len;
- gint ret;
-
- g_assert (tokens != NULL);
-
- cmd_len = strlen (command);
- prefix_len = strlen (prefix);
- out = rspamd_fstring_sized_new (1024);
-
- if (learn) {
- rspamd_printf_fstring (&out, "*1\r\n$5\r\nMULTI\r\n");
-
- ret = redisAsyncFormattedCommand (rt->redis, NULL, NULL,
- out->str, out->len);
-
- if (ret != REDIS_OK) {
- msg_err_task ("call to redis failed: %s", rt->redis->errstr);
- rspamd_fstring_free (out);
-
- return NULL;
- }
-
- out->len = 0;
- }
- else {
- if (rt->ctx->new_schema) {
- /* Multi + HGET */
- rspamd_printf_fstring (&out, "*1\r\n$5\r\nMULTI\r\n");
-
- ret = redisAsyncFormattedCommand (rt->redis, NULL, NULL,
- out->str, out->len);
-
- if (ret != REDIS_OK) {
- msg_err_task ("call to redis failed: %s", rt->redis->errstr);
- rspamd_fstring_free (out);
-
- return NULL;
- }
-
- out->len = 0;
- }
- else {
- rspamd_printf_fstring (&out, ""
- "*%d\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n",
- (tokens->len + 2),
- cmd_len, command,
- prefix_len, prefix);
- }
- }
-
- for (i = 0; i < tokens->len; i ++) {
- tok = g_ptr_array_index (tokens, i);
-
- if (learn) {
- if (intvals) {
- l1 = rspamd_snprintf (n1, sizeof (n1), "%L",
- (gint64) tok->values[idx]);
- } else {
- l1 = rspamd_snprintf (n1, sizeof (n1), "%f",
- tok->values[idx]);
- }
-
- if (rt->ctx->new_schema) {
- /*
- * HINCRBY <prefix_token> <0|1> <value>
- */
- l0 = rspamd_snprintf (n0, sizeof (n0), "%*s_%uL",
- prefix_len, prefix,
- tok->data);
-
- rspamd_printf_fstring (&out, ""
- "*4\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n",
- cmd_len, command,
- l0, n0,
- 1, rt->stcf->is_spam ? "S" : "H",
- l1, n1);
- }
- else {
- l0 = rspamd_snprintf (n0, sizeof (n0), "%uL", tok->data);
-
- /*
- * HINCRBY <prefix> <token> <value>
- */
- rspamd_printf_fstring (&out, ""
- "*4\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n",
- cmd_len, command,
- prefix_len, prefix,
- l0, n0,
- l1, n1);
- }
-
- ret = redisAsyncFormattedCommand (rt->redis, NULL, NULL,
- out->str, out->len);
-
- if (ret != REDIS_OK) {
- msg_err_task ("call to redis failed: %s", rt->redis->errstr);
- rspamd_fstring_free (out);
-
- return NULL;
- }
-
- if (rt->ctx->store_tokens) {
-
- if (!rt->ctx->new_schema) {
- /*
- * We store tokens in form
- * HSET prefix_tokens <token_id> "token_string"
- * ZINCRBY prefix_z 1.0 <token_id>
- */
- if (tok->t1 && tok->t2) {
- redisAsyncCommand (rt->redis, NULL, NULL,
- "HSET %b_tokens %b %b:%b",
- prefix, (size_t) prefix_len,
- n0, (size_t) l0,
- tok->t1->stemmed.begin, tok->t1->stemmed.len,
- tok->t2->stemmed.begin, tok->t2->stemmed.len);
- } else if (tok->t1) {
- redisAsyncCommand (rt->redis, NULL, NULL,
- "HSET %b_tokens %b %b",
- prefix, (size_t) prefix_len,
- n0, (size_t) l0,
- tok->t1->stemmed.begin,
- tok->t1->stemmed.len);
- }
- }
- else {
- /*
- * We store tokens in form
- * HSET <token_id> "tokens" "token_string"
- * ZINCRBY prefix_z 1.0 <token_id>
- */
- if (tok->t1 && tok->t2) {
- redisAsyncCommand (rt->redis, NULL, NULL,
- "HSET %b %s %b:%b",
- n0, (size_t) l0,
- "tokens",
- tok->t1->stemmed.begin, tok->t1->stemmed.len,
- tok->t2->stemmed.begin, tok->t2->stemmed.len);
- } else if (tok->t1) {
- redisAsyncCommand (rt->redis, NULL, NULL,
- "HSET %b %s %b",
- n0, (size_t) l0,
- "tokens",
- tok->t1->stemmed.begin, tok->t1->stemmed.len);
- }
- }
-
- redisAsyncCommand (rt->redis, NULL, NULL,
- "ZINCRBY %b_z %b %b",
- prefix, (size_t)prefix_len,
- n1, (size_t)l1,
- n0, (size_t)l0);
- }
-
- if (rt->ctx->new_schema && rt->ctx->expiry > 0) {
- out->len = 0;
- l1 = rspamd_snprintf (n1, sizeof (n1), "%d",
- rt->ctx->expiry);
-
- rspamd_printf_fstring (&out, ""
- "*3\r\n"
- "$6\r\n"
- "EXPIRE\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n",
- l0, n0,
- l1, n1);
- redisAsyncFormattedCommand (rt->redis, NULL, NULL,
- out->str, out->len);
- }
-
- out->len = 0;
- }
- else {
- if (rt->ctx->new_schema) {
- l0 = rspamd_snprintf (n0, sizeof (n0), "%*s_%uL",
- prefix_len, prefix,
- tok->data);
-
- rspamd_printf_fstring (&out, ""
- "*3\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n",
- cmd_len, command,
- l0, n0,
- 1, rt->stcf->is_spam ? "S" : "H");
-
- ret = redisAsyncFormattedCommand (rt->redis, NULL, NULL,
- out->str, out->len);
-
- if (ret != REDIS_OK) {
- msg_err_task ("call to redis failed: %s", rt->redis->errstr);
- rspamd_fstring_free (out);
-
- return NULL;
- }
-
- out->len = 0;
- }
- else {
- l0 = rspamd_snprintf (n0, sizeof (n0), "%uL", tok->data);
- rspamd_printf_fstring (&out, ""
- "$%d\r\n"
- "%s\r\n", l0, n0);
- }
- }
- }
-
- if (!learn && rt->ctx->new_schema) {
- rspamd_printf_fstring (&out, "*1\r\n$4\r\nEXEC\r\n");
- }
-
- return out;
- }
-
- static void
- rspamd_redis_store_stat_signature (struct rspamd_task *task,
- struct redis_stat_runtime *rt,
- GPtrArray *tokens,
- const gchar *prefix)
- {
- gchar *sig, keybuf[512], nbuf[64];
- rspamd_token_t *tok;
- guint i, blen, klen;
- rspamd_fstring_t *out;
-
- out = rspamd_fstring_sized_new (1024);
- sig = rspamd_mempool_get_variable (task->task_pool,
- RSPAMD_MEMPOOL_STAT_SIGNATURE);
-
- if (sig == NULL) {
- msg_err_task ("cannot get bayes signature");
- return;
- }
-
- klen = rspamd_snprintf (keybuf, sizeof (keybuf), "%s_%s_%s",
- prefix, sig, rt->stcf->is_spam ? "S" : "H");
-
- out->len = 0;
-
- /* Cleanup key */
- rspamd_printf_fstring (&out, ""
- "*2\r\n"
- "$3\r\n"
- "DEL\r\n"
- "$%d\r\n"
- "%s\r\n",
- klen, keybuf);
- redisAsyncFormattedCommand (rt->redis, NULL, NULL,
- out->str, out->len);
- out->len = 0;
-
- rspamd_printf_fstring (&out, ""
- "*%d\r\n"
- "$5\r\n"
- "LPUSH\r\n"
- "$%d\r\n"
- "%s\r\n",
- tokens->len + 2,
- klen, keybuf);
-
- PTR_ARRAY_FOREACH (tokens, i, tok) {
- blen = rspamd_snprintf (nbuf, sizeof (nbuf), "%uL", tok->data);
- rspamd_printf_fstring (&out, ""
- "$%d\r\n"
- "%s\r\n", blen, nbuf);
- }
-
- redisAsyncFormattedCommand (rt->redis, NULL, NULL,
- out->str, out->len);
- out->len = 0;
-
- if (rt->ctx->expiry > 0) {
- out->len = 0;
- blen = rspamd_snprintf (nbuf, sizeof (nbuf), "%d",
- rt->ctx->expiry);
-
- rspamd_printf_fstring (&out, ""
- "*3\r\n"
- "$6\r\n"
- "EXPIRE\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n",
- klen, keybuf,
- blen, nbuf);
- redisAsyncFormattedCommand (rt->redis, NULL, NULL,
- out->str, out->len);
- }
-
- rspamd_fstring_free (out);
- }
-
- static void
- rspamd_redis_async_cbdata_cleanup (struct rspamd_redis_stat_cbdata *cbdata)
- {
- guint i;
- gchar *k;
-
- if (cbdata && !cbdata->wanna_die) {
- /* Avoid double frees */
- cbdata->wanna_die = TRUE;
- redisAsyncFree (cbdata->redis);
-
- for (i = 0; i < cbdata->cur_keys->len; i ++) {
- k = g_ptr_array_index (cbdata->cur_keys, i);
- g_free (k);
- }
-
- g_ptr_array_free (cbdata->cur_keys, TRUE);
-
- if (cbdata->elt) {
- cbdata->elt->cbdata = NULL;
- /* Re-enable parent event */
- cbdata->elt->async->enabled = TRUE;
-
- /* Replace ucl object */
- if (cbdata->cur) {
- if (cbdata->elt->stat) {
- ucl_object_unref (cbdata->elt->stat);
- }
-
- cbdata->elt->stat = cbdata->cur;
- cbdata->cur = NULL;
- }
- }
-
- if (cbdata->cur) {
- ucl_object_unref (cbdata->cur);
- }
-
- g_free (cbdata);
- }
- }
-
- /* Called when we get number of learns for a specific key */
- static void
- rspamd_redis_stat_learns (redisAsyncContext *c, gpointer r, gpointer priv)
- {
- struct rspamd_redis_stat_cbdata *cbdata = priv;
- redisReply *reply = r;
- ucl_object_t *obj;
- gulong num = 0;
-
- if (cbdata->wanna_die) {
- return;
- }
-
- cbdata->inflight --;
-
- if (c->err == 0 && r != NULL) {
- if (G_LIKELY (reply->type == REDIS_REPLY_INTEGER)) {
- num = reply->integer;
- }
- else if (reply->type == REDIS_REPLY_STRING) {
- rspamd_strtoul (reply->str, reply->len, &num);
- }
-
- obj = (ucl_object_t *) ucl_object_lookup (cbdata->cur, "revision");
- if (obj) {
- obj->value.iv += num;
- }
- }
-
- if (cbdata->inflight == 0) {
- rspamd_redis_async_cbdata_cleanup (cbdata);
- }
- }
-
- /* Called when we get number of elements for a specific key */
- static void
- rspamd_redis_stat_key (redisAsyncContext *c, gpointer r, gpointer priv)
- {
- struct rspamd_redis_stat_cbdata *cbdata = priv;
- redisReply *reply = r;
- ucl_object_t *obj;
- glong num = 0;
-
- if (cbdata->wanna_die) {
- return;
- }
-
- cbdata->inflight --;
-
- if (c->err == 0 && r != NULL) {
- if (G_LIKELY (reply->type == REDIS_REPLY_INTEGER)) {
- num = reply->integer;
- }
- else if (reply->type == REDIS_REPLY_STRING) {
- rspamd_strtol (reply->str, reply->len, &num);
- }
-
- if (num < 0) {
- msg_err ("bad learns count: %L", (gint64)num);
- num = 0;
- }
-
- obj = (ucl_object_t *)ucl_object_lookup (cbdata->cur, "used");
- if (obj) {
- obj->value.iv += num;
- }
-
- obj = (ucl_object_t *)ucl_object_lookup (cbdata->cur, "total");
- if (obj) {
- obj->value.iv += num;
- }
-
- obj = (ucl_object_t *)ucl_object_lookup (cbdata->cur, "size");
- if (obj) {
- /* Size of key + size of int64_t */
- obj->value.iv += num * (sizeof (G_STRINGIFY (G_MAXINT64)) +
- sizeof (guint64) + sizeof (gpointer));
- }
- }
-
- if (cbdata->inflight == 0) {
- rspamd_redis_async_cbdata_cleanup (cbdata);
- }
- }
-
- /* Called when we have connected to the redis server and got keys to check */
- static void
- rspamd_redis_stat_keys (redisAsyncContext *c, gpointer r, gpointer priv)
- {
- struct rspamd_redis_stat_cbdata *cbdata = priv;
- redisReply *reply = r, *elt;
- gchar **pk, *k;
- guint i, processed = 0;
-
-
- if (cbdata->wanna_die) {
- return;
- }
-
- cbdata->inflight --;
-
- if (c->err == 0 && r != NULL) {
- if (reply->type == REDIS_REPLY_ARRAY) {
- g_ptr_array_set_size (cbdata->cur_keys, reply->elements);
-
- for (i = 0; i < reply->elements; i ++) {
- elt = reply->element[i];
-
- if (elt->type == REDIS_REPLY_STRING) {
- pk = (gchar **)&g_ptr_array_index (cbdata->cur_keys, i);
- *pk = g_malloc (elt->len + 1);
- rspamd_strlcpy (*pk, elt->str, elt->len + 1);
- processed ++;
- }
- }
-
- if (processed) {
- for (i = 0; i < cbdata->cur_keys->len; i ++) {
- k = (gchar *)g_ptr_array_index (cbdata->cur_keys, i);
-
- if (k) {
- const gchar *learned_key = "learns";
-
- if (cbdata->elt->ctx->new_schema) {
- if (cbdata->elt->ctx->stcf->is_spam) {
- learned_key = "learns_spam";
- }
- else {
- learned_key = "learns_ham";
- }
- redisAsyncCommand (cbdata->redis,
- rspamd_redis_stat_learns,
- cbdata,
- "HGET %s %s",
- k, learned_key);
- cbdata->inflight += 1;
- }
- else {
- redisAsyncCommand (cbdata->redis,
- rspamd_redis_stat_key,
- cbdata,
- "HLEN %s",
- k);
- redisAsyncCommand (cbdata->redis,
- rspamd_redis_stat_learns,
- cbdata,
- "HGET %s %s",
- k, learned_key);
- cbdata->inflight += 2;
- }
- }
- }
- }
- }
-
- /* Set up the required keys */
- ucl_object_insert_key (cbdata->cur,
- ucl_object_typed_new (UCL_INT), "revision", 0, false);
- ucl_object_insert_key (cbdata->cur,
- ucl_object_typed_new (UCL_INT), "used", 0, false);
- ucl_object_insert_key (cbdata->cur,
- ucl_object_typed_new (UCL_INT), "total", 0, false);
- ucl_object_insert_key (cbdata->cur,
- ucl_object_typed_new (UCL_INT), "size", 0, false);
- ucl_object_insert_key (cbdata->cur,
- ucl_object_fromstring (cbdata->elt->ctx->stcf->symbol),
- "symbol", 0, false);
- ucl_object_insert_key (cbdata->cur, ucl_object_fromstring ("redis"),
- "type", 0, false);
- ucl_object_insert_key (cbdata->cur, ucl_object_fromint (0),
- "languages", 0, false);
- ucl_object_insert_key (cbdata->cur, ucl_object_fromint (processed),
- "users", 0, false);
-
- rspamd_upstream_ok (cbdata->selected);
-
- if (cbdata->inflight == 0) {
- rspamd_redis_async_cbdata_cleanup (cbdata);
- }
- }
- else {
- if (c->errstr) {
- msg_err ("cannot get keys to gather stat: %s", c->errstr);
- }
- else {
- msg_err ("cannot get keys to gather stat: unknown error");
- }
-
- rspamd_upstream_fail (cbdata->selected, FALSE);
- rspamd_redis_async_cbdata_cleanup (cbdata);
- }
- }
-
- static void
- rspamd_redis_async_stat_cb (struct rspamd_stat_async_elt *elt, gpointer d)
- {
- struct redis_stat_ctx *ctx;
- struct rspamd_redis_stat_elt *redis_elt = elt->ud;
- struct rspamd_redis_stat_cbdata *cbdata;
- rspamd_inet_addr_t *addr;
- struct upstream_list *ups;
-
- g_assert (redis_elt != NULL);
-
- ctx = redis_elt->ctx;
-
- if (redis_elt->cbdata) {
- /* We have some other process pending */
- rspamd_redis_async_cbdata_cleanup (redis_elt->cbdata);
- }
-
- /* Disable further events unless needed */
- elt->enabled = FALSE;
-
- ups = rspamd_redis_get_servers (ctx, "read_servers");
-
- if (!ups) {
- return;
- }
-
- cbdata = g_malloc0 (sizeof (*cbdata));
-
- cbdata->selected = rspamd_upstream_get (ups,
- RSPAMD_UPSTREAM_ROUND_ROBIN,
- NULL,
- 0);
-
- g_assert (cbdata->selected != NULL);
- addr = rspamd_upstream_addr_next (cbdata->selected);
- g_assert (addr != NULL);
-
- if (rspamd_inet_address_get_af (addr) == AF_UNIX) {
- cbdata->redis = redisAsyncConnectUnix (rspamd_inet_address_to_string (addr));
- }
- else {
- cbdata->redis = redisAsyncConnect (rspamd_inet_address_to_string (addr),
- rspamd_inet_address_get_port (addr));
- }
-
- g_assert (cbdata->redis != NULL);
-
- redisLibeventAttach (cbdata->redis, redis_elt->ev_base);
-
- cbdata->inflight = 1;
- cbdata->cur = ucl_object_typed_new (UCL_OBJECT);
- cbdata->elt = redis_elt;
- cbdata->cur_keys = g_ptr_array_new ();
- redis_elt->cbdata = cbdata;
-
- /* XXX: deal with timeouts maybe */
- /* Get keys in redis that match our symbol */
- rspamd_redis_maybe_auth (ctx, cbdata->redis);
- redisAsyncCommand (cbdata->redis, rspamd_redis_stat_keys, cbdata,
- "SMEMBERS %s_keys",
- ctx->stcf->symbol);
- }
-
- static void
- rspamd_redis_async_stat_fin (struct rspamd_stat_async_elt *elt, gpointer d)
- {
- struct rspamd_redis_stat_elt *redis_elt = elt->ud;
-
- rspamd_redis_async_cbdata_cleanup (redis_elt->cbdata);
- }
-
- /* Called on connection termination */
- static void
- rspamd_redis_fin (gpointer data)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (data);
- redisAsyncContext *redis;
-
- rt->has_event = FALSE;
- /* Stop timeout */
- if (rspamd_event_pending (&rt->timeout_event, EV_TIMEOUT)) {
- event_del (&rt->timeout_event);
- }
-
- if (rt->redis) {
- redis = rt->redis;
- rt->redis = NULL;
- /* This calls for all callbacks pending */
- redisAsyncFree (redis);
- }
- }
-
- static void
- rspamd_redis_fin_learn (gpointer data)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (data);
- redisAsyncContext *redis;
-
- rt->has_event = FALSE;
- /* Stop timeout */
- if (rspamd_event_pending (&rt->timeout_event, EV_TIMEOUT)) {
- event_del (&rt->timeout_event);
- }
-
- if (rt->redis) {
- redis = rt->redis;
- rt->redis = NULL;
- /* This calls for all callbacks pending */
- redisAsyncFree (redis);
- }
- }
-
- static void
- rspamd_redis_timeout (gint fd, short what, gpointer d)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (d);
- struct rspamd_task *task;
- redisAsyncContext *redis;
-
- task = rt->task;
-
- msg_err_task_check ("connection to redis server %s timed out",
- rspamd_upstream_name (rt->selected));
-
- rspamd_upstream_fail (rt->selected, FALSE);
-
- if (rt->redis) {
- redis = rt->redis;
- rt->redis = NULL;
- /* This calls for all callbacks pending */
- redisAsyncFree (redis);
- }
-
- if (!rt->err) {
- g_set_error (&rt->err, rspamd_redis_stat_quark (), ETIMEDOUT,
- "error getting reply from redis server %s: timeout",
- rspamd_upstream_name (rt->selected));
- }
- }
-
- /* Called when we have connected to the redis server and got stats */
- static void
- rspamd_redis_connected (redisAsyncContext *c, gpointer r, gpointer priv)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (priv);
- redisReply *reply = r;
- struct rspamd_task *task;
- glong val = 0;
-
- task = rt->task;
-
- if (c->err == 0) {
- if (r != NULL) {
- if (G_UNLIKELY (reply->type == REDIS_REPLY_INTEGER)) {
- val = reply->integer;
- }
- else if (reply->type == REDIS_REPLY_STRING) {
- rspamd_strtol (reply->str, reply->len, &val);
- }
- else {
- if (reply->type != REDIS_REPLY_NIL) {
- msg_err_task ("bad learned type for %s: %s, nil expected",
- rt->stcf->symbol,
- rspamd_redis_type_to_string (reply->type));
- }
-
- val = 0;
- }
-
- if (val < 0) {
- msg_warn_task ("invalid number of learns for %s: %L",
- rt->stcf->symbol, val);
- val = 0;
- }
-
- rt->learned = val;
- msg_debug_stat_redis ("connected to redis server, tokens learned for %s: %uL",
- rt->redis_object_expanded, rt->learned);
- rspamd_upstream_ok (rt->selected);
- }
- }
- else {
- msg_err_task ("error getting reply from redis server %s: %s",
- rspamd_upstream_name (rt->selected), c->errstr);
- rspamd_upstream_fail (rt->selected, FALSE);
-
- if (!rt->err) {
- g_set_error (&rt->err, rspamd_redis_stat_quark (), c->err,
- "error getting reply from redis server %s: %s",
- rspamd_upstream_name (rt->selected), c->errstr);
- }
- }
-
- }
-
- /* Called when we have received tokens values from redis */
- static void
- rspamd_redis_processed (redisAsyncContext *c, gpointer r, gpointer priv)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (priv);
- redisReply *reply = r, *elt;
- struct rspamd_task *task;
- rspamd_token_t *tok;
- guint i, processed = 0, found = 0;
- gulong val;
- gdouble float_val;
-
- task = rt->task;
-
- if (c->err == 0) {
- if (r != NULL) {
- if (reply->type == REDIS_REPLY_ARRAY) {
-
- if (reply->elements == task->tokens->len) {
- for (i = 0; i < reply->elements; i ++) {
- tok = g_ptr_array_index (task->tokens, i);
- elt = reply->element[i];
-
- if (G_UNLIKELY (elt->type == REDIS_REPLY_INTEGER)) {
- tok->values[rt->id] = elt->integer;
- found ++;
- }
- else if (elt->type == REDIS_REPLY_STRING) {
- if (rt->stcf->clcf->flags &
- RSPAMD_FLAG_CLASSIFIER_INTEGER) {
- rspamd_strtoul (elt->str, elt->len, &val);
- tok->values[rt->id] = val;
- }
- else {
- float_val = strtod (elt->str, NULL);
- tok->values[rt->id] = float_val;
- }
-
- found ++;
- }
- else {
- tok->values[rt->id] = 0;
- }
-
- processed ++;
- }
-
- if (rt->stcf->is_spam) {
- task->flags |= RSPAMD_TASK_FLAG_HAS_SPAM_TOKENS;
- }
- else {
- task->flags |= RSPAMD_TASK_FLAG_HAS_HAM_TOKENS;
- }
- }
- else {
- msg_err_task_check ("got invalid length of reply vector from redis: "
- "%d, expected: %d",
- (gint)reply->elements,
- (gint)task->tokens->len);
- }
- }
- else {
- msg_err_task_check ("got invalid reply from redis: %s, array expected",
- rspamd_redis_type_to_string (reply->type));
- }
-
- msg_debug_stat_redis ("received tokens for %s: %d processed, %d found",
- rt->redis_object_expanded, processed, found);
- rspamd_upstream_ok (rt->selected);
- }
- }
- else {
- msg_err_task ("error getting reply from redis server %s: %s",
- rspamd_upstream_name (rt->selected), c->errstr);
-
- if (rt->redis) {
- rspamd_upstream_fail (rt->selected, FALSE);
- }
-
- if (!rt->err) {
- g_set_error (&rt->err, rspamd_redis_stat_quark (), c->err,
- "cannot get values: error getting reply from redis server %s: %s",
- rspamd_upstream_name (rt->selected), c->errstr);
- }
- }
-
- if (rt->has_event) {
- rspamd_session_remove_event (task->s, rspamd_redis_fin, rt);
- }
- }
-
- /* Called when we have set tokens during learning */
- static void
- rspamd_redis_learned (redisAsyncContext *c, gpointer r, gpointer priv)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (priv);
- struct rspamd_task *task;
-
- task = rt->task;
-
- if (c->err == 0) {
- rspamd_upstream_ok (rt->selected);
- }
- else {
- msg_err_task_check ("error getting reply from redis server %s: %s",
- rspamd_upstream_name (rt->selected), c->errstr);
-
- if (rt->redis) {
- rspamd_upstream_fail (rt->selected, FALSE);
- }
-
- if (!rt->err) {
- g_set_error (&rt->err, rspamd_redis_stat_quark (), c->err,
- "cannot get learned: error getting reply from redis server %s: %s",
- rspamd_upstream_name (rt->selected), c->errstr);
- }
- }
-
- if (rt->has_event) {
- rspamd_session_remove_event (task->s, rspamd_redis_fin_learn, rt);
- }
- }
- static void
- rspamd_redis_parse_classifier_opts (struct redis_stat_ctx *backend,
- const ucl_object_t *obj,
- struct rspamd_config *cfg)
- {
- const gchar *lua_script;
- const ucl_object_t *elt, *users_enabled;
-
- users_enabled = ucl_object_lookup_any (obj, "per_user",
- "users_enabled", NULL);
-
- if (users_enabled != NULL) {
- if (ucl_object_type (users_enabled) == UCL_BOOLEAN) {
- backend->enable_users = ucl_object_toboolean (users_enabled);
- backend->cbref_user = -1;
- }
- else if (ucl_object_type (users_enabled) == UCL_STRING) {
- lua_script = ucl_object_tostring (users_enabled);
-
- if (luaL_dostring (cfg->lua_state, lua_script) != 0) {
- msg_err_config ("cannot execute lua script for users "
- "extraction: %s", lua_tostring (cfg->lua_state, -1));
- }
- else {
- if (lua_type (cfg->lua_state, -1) == LUA_TFUNCTION) {
- backend->enable_users = TRUE;
- backend->cbref_user = luaL_ref (cfg->lua_state,
- LUA_REGISTRYINDEX);
- }
- else {
- msg_err_config ("lua script must return "
- "function(task) and not %s",
- lua_typename (cfg->lua_state, lua_type (
- cfg->lua_state, -1)));
- }
- }
- }
- }
- else {
- backend->enable_users = FALSE;
- backend->cbref_user = -1;
- }
-
- elt = ucl_object_lookup (obj, "prefix");
- if (elt == NULL || ucl_object_type (elt) != UCL_STRING) {
- /* Default non-users statistics */
- if (backend->enable_users || backend->cbref_user != -1) {
- backend->redis_object = REDIS_DEFAULT_USERS_OBJECT;
- }
- else {
- backend->redis_object = REDIS_DEFAULT_OBJECT;
- }
- }
- else {
- /* XXX: sanity check */
- backend->redis_object = ucl_object_tostring (elt);
- }
-
- elt = ucl_object_lookup (obj, "store_tokens");
- if (elt) {
- backend->store_tokens = ucl_object_toboolean (elt);
- }
- else {
- backend->store_tokens = FALSE;
- }
-
- elt = ucl_object_lookup (obj, "new_schema");
- if (elt) {
- backend->new_schema = ucl_object_toboolean (elt);
- }
- else {
- backend->new_schema = FALSE;
-
- msg_warn_config ("you are using old bayes schema for redis statistics, "
- "please consider converting it to a new one "
- "by using 'rspamadm configwizard statistics'");
- }
-
- elt = ucl_object_lookup (obj, "signatures");
- if (elt) {
- backend->enable_signatures = ucl_object_toboolean (elt);
- }
- else {
- backend->enable_signatures = FALSE;
- }
-
- elt = ucl_object_lookup_any (obj, "expiry", "expire", NULL);
- if (elt) {
- backend->expiry = ucl_object_toint (elt);
- }
- else {
- backend->expiry = 0;
- }
- }
-
- gpointer
- rspamd_redis_init (struct rspamd_stat_ctx *ctx,
- struct rspamd_config *cfg, struct rspamd_statfile *st)
- {
- struct redis_stat_ctx *backend;
- struct rspamd_statfile_config *stf = st->stcf;
- struct rspamd_redis_stat_elt *st_elt;
- const ucl_object_t *obj;
- gboolean ret = FALSE;
- gint conf_ref = -1;
- lua_State *L = (lua_State *)cfg->lua_state;
-
- backend = g_malloc0 (sizeof (*backend));
- backend->L = L;
- backend->timeout = REDIS_DEFAULT_TIMEOUT;
-
- /* First search in backend configuration */
- obj = ucl_object_lookup (st->classifier->cfg->opts, "backend");
- if (obj != NULL && ucl_object_type (obj) == UCL_OBJECT) {
- ret = rspamd_lua_try_load_redis (L, obj, cfg, &conf_ref);
- }
-
- /* Now try statfiles config */
- if (!ret && stf->opts) {
- ret = rspamd_lua_try_load_redis (L, stf->opts, cfg, &conf_ref);
- }
-
- /* Now try classifier config */
- if (!ret && st->classifier->cfg->opts) {
- ret = rspamd_lua_try_load_redis (L, st->classifier->cfg->opts, cfg, &conf_ref);
- }
-
- /* Now try global redis settings */
- if (!ret) {
- obj = ucl_object_lookup (cfg->rcl_obj, "redis");
-
- if (obj) {
- const ucl_object_t *specific_obj;
-
- specific_obj = ucl_object_lookup (obj, "statistics");
-
- if (specific_obj) {
- ret = rspamd_lua_try_load_redis (L,
- specific_obj, cfg, &conf_ref);
- }
- else {
- ret = rspamd_lua_try_load_redis (L,
- obj, cfg, &conf_ref);
- }
- }
- }
-
- if (!ret) {
- msg_err_config ("cannot init redis backend for %s", stf->symbol);
- g_free (backend);
- return NULL;
- }
-
- backend->conf_ref = conf_ref;
-
- /* Check some common table values */
- lua_rawgeti (L, LUA_REGISTRYINDEX, conf_ref);
-
- lua_pushstring (L, "timeout");
- lua_gettable (L, -2);
- if (lua_type (L, -1) == LUA_TNUMBER) {
- backend->timeout = lua_tonumber (L, -1);
- }
- lua_pop (L, 1);
-
- lua_pushstring (L, "db");
- lua_gettable (L, -2);
- if (lua_type (L, -1) == LUA_TSTRING) {
- backend->dbname = rspamd_mempool_strdup (cfg->cfg_pool,
- lua_tostring (L, -1));
- }
- lua_pop (L, 1);
-
- lua_pushstring (L, "password");
- lua_gettable (L, -2);
- if (lua_type (L, -1) == LUA_TSTRING) {
- backend->password = rspamd_mempool_strdup (cfg->cfg_pool,
- lua_tostring (L, -1));
- }
- lua_pop (L, 1);
-
- lua_settop (L, 0);
-
- rspamd_redis_parse_classifier_opts (backend, st->classifier->cfg->opts, cfg);
- stf->clcf->flags |= RSPAMD_FLAG_CLASSIFIER_INCREMENTING_BACKEND;
- backend->stcf = stf;
-
- st_elt = g_malloc0 (sizeof (*st_elt));
- st_elt->ev_base = ctx->ev_base;
- st_elt->ctx = backend;
- backend->stat_elt = rspamd_stat_ctx_register_async (
- rspamd_redis_async_stat_cb,
- rspamd_redis_async_stat_fin,
- st_elt,
- REDIS_STAT_TIMEOUT);
- st_elt->async = backend->stat_elt;
-
- return (gpointer)backend;
- }
-
- gpointer
- rspamd_redis_runtime (struct rspamd_task *task,
- struct rspamd_statfile_config *stcf,
- gboolean learn, gpointer c)
- {
- struct redis_stat_ctx *ctx = REDIS_CTX (c);
- struct redis_stat_runtime *rt;
- struct upstream *up;
- struct upstream_list *ups;
- char *object_expanded = NULL;
- rspamd_inet_addr_t *addr;
-
- g_assert (ctx != NULL);
- g_assert (stcf != NULL);
-
- if (learn) {
- ups = rspamd_redis_get_servers (ctx, "write_servers");
-
- if (!ups) {
- msg_err_task ("no write servers defined for %s, cannot learn",
- stcf->symbol);
- return NULL;
- }
- up = rspamd_upstream_get (ups,
- RSPAMD_UPSTREAM_MASTER_SLAVE,
- NULL,
- 0);
- }
- else {
- ups = rspamd_redis_get_servers (ctx, "read_servers");
-
- if (!ups) {
- msg_err_task ("no read servers defined for %s, cannot stat",
- stcf->symbol);
- return NULL;
- }
- up = rspamd_upstream_get (ups,
- RSPAMD_UPSTREAM_ROUND_ROBIN,
- NULL,
- 0);
- }
-
- if (up == NULL) {
- msg_err_task ("no upstreams reachable");
- return NULL;
- }
-
- if (rspamd_redis_expand_object (ctx->redis_object, ctx, task,
- &object_expanded) == 0) {
- msg_err_task ("expansion for learning failed for symbol %s "
- "(maybe learning per user classifier with no user or recipient)",
- stcf->symbol);
- return NULL;
- }
-
- rt = rspamd_mempool_alloc0 (task->task_pool, sizeof (*rt));
- rspamd_mempool_add_destructor (task->task_pool,
- rspamd_gerror_free_maybe, &rt->err);
- rt->selected = up;
- rt->task = task;
- rt->ctx = ctx;
- rt->stcf = stcf;
- rt->redis_object_expanded = object_expanded;
-
- addr = rspamd_upstream_addr_next (up);
- g_assert (addr != NULL);
-
- if (rspamd_inet_address_get_af (addr) == AF_UNIX) {
- rt->redis = redisAsyncConnectUnix (rspamd_inet_address_to_string (addr));
- }
- else {
- rt->redis = redisAsyncConnect (rspamd_inet_address_to_string (addr),
- rspamd_inet_address_get_port (addr));
- }
-
- if (rt->redis == NULL) {
- msg_err_task ("cannot connect redis");
- return NULL;
- }
-
- redisLibeventAttach (rt->redis, task->ev_base);
- rspamd_redis_maybe_auth (ctx, rt->redis);
-
- return rt;
- }
-
- void
- rspamd_redis_close (gpointer p)
- {
- struct redis_stat_ctx *ctx = REDIS_CTX (p);
- lua_State *L = ctx->L;
-
- if (ctx->conf_ref) {
- luaL_unref (L, LUA_REGISTRYINDEX, ctx->conf_ref);
- }
-
- g_free (ctx);
- }
-
- gboolean
- rspamd_redis_process_tokens (struct rspamd_task *task,
- GPtrArray *tokens,
- gint id, gpointer p)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (p);
- rspamd_fstring_t *query;
- struct timeval tv;
- gint ret;
- const gchar *learned_key = "learns";
-
- if (rspamd_session_blocked (task->s)) {
- return FALSE;
- }
-
- if (tokens == NULL || tokens->len == 0 || rt->redis == NULL) {
- return FALSE;
- }
-
- rt->id = id;
-
- if (rt->ctx->new_schema) {
- if (rt->ctx->stcf->is_spam) {
- learned_key = "learns_spam";
- }
- else {
- learned_key = "learns_ham";
- }
- }
-
- if (redisAsyncCommand (rt->redis, rspamd_redis_connected, rt, "HGET %s %s",
- rt->redis_object_expanded, learned_key) == REDIS_OK) {
-
- rspamd_session_add_event (task->s, rspamd_redis_fin, rt, M);
- rt->has_event = TRUE;
-
- if (rspamd_event_pending (&rt->timeout_event, EV_TIMEOUT)) {
- event_del (&rt->timeout_event);
- }
- event_set (&rt->timeout_event, -1, EV_TIMEOUT, rspamd_redis_timeout, rt);
- event_base_set (task->ev_base, &rt->timeout_event);
- double_to_tv (rt->ctx->timeout, &tv);
- event_add (&rt->timeout_event, &tv);
-
- query = rspamd_redis_tokens_to_query (task, rt, tokens,
- rt->ctx->new_schema ? "HGET" : "HMGET",
- rt->redis_object_expanded, FALSE, -1,
- rt->stcf->clcf->flags & RSPAMD_FLAG_CLASSIFIER_INTEGER);
- g_assert (query != NULL);
- rspamd_mempool_add_destructor (task->task_pool,
- (rspamd_mempool_destruct_t)rspamd_fstring_free, query);
-
- ret = redisAsyncFormattedCommand (rt->redis, rspamd_redis_processed, rt,
- query->str, query->len);
-
- if (ret == REDIS_OK) {
- return TRUE;
- }
- else {
- msg_err_task ("call to redis failed: %s", rt->redis->errstr);
- }
- }
-
- return FALSE;
- }
-
- gboolean
- rspamd_redis_finalize_process (struct rspamd_task *task, gpointer runtime,
- gpointer ctx)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (runtime);
- redisAsyncContext *redis;
-
- if (rspamd_event_pending (&rt->timeout_event, EV_TIMEOUT)) {
- event_del (&rt->timeout_event);
- }
-
- if (rt->redis) {
- redis = rt->redis;
- rt->redis = NULL;
- redisAsyncFree (redis);
- }
-
- if (rt->err) {
- return FALSE;
- }
-
- return TRUE;
- }
-
- gboolean
- rspamd_redis_learn_tokens (struct rspamd_task *task, GPtrArray *tokens,
- gint id, gpointer p)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (p);
- struct upstream *up;
- struct upstream_list *ups;
- rspamd_inet_addr_t *addr;
- struct timeval tv;
- rspamd_fstring_t *query;
- const gchar *redis_cmd;
- rspamd_token_t *tok;
- gint ret;
- goffset off;
- const gchar *learned_key = "learns";
-
- if (rspamd_session_blocked (task->s)) {
- return FALSE;
- }
-
- ups = rspamd_redis_get_servers (rt->ctx, "write_servers");
-
- if (!ups) {
- return FALSE;
- }
- up = rspamd_upstream_get (ups,
- RSPAMD_UPSTREAM_MASTER_SLAVE,
- NULL,
- 0);
-
- if (up == NULL) {
- msg_err_task ("no upstreams reachable");
- return FALSE;
- }
-
- rt->selected = up;
-
- if (rt->ctx->new_schema) {
- if (rt->ctx->stcf->is_spam) {
- learned_key = "learns_spam";
- }
- else {
- learned_key = "learns_ham";
- }
- }
-
- addr = rspamd_upstream_addr_next (up);
- g_assert (addr != NULL);
-
- if (rspamd_inet_address_get_af (addr) == AF_UNIX) {
- rt->redis = redisAsyncConnectUnix (rspamd_inet_address_to_string (addr));
- }
- else {
- rt->redis = redisAsyncConnect (rspamd_inet_address_to_string (addr),
- rspamd_inet_address_get_port (addr));
- }
-
- g_assert (rt->redis != NULL);
-
- redisLibeventAttach (rt->redis, task->ev_base);
- rspamd_redis_maybe_auth (rt->ctx, rt->redis);
-
- /*
- * Add the current key to the set of learned keys
- */
- redisAsyncCommand (rt->redis, NULL, NULL, "SADD %s_keys %s",
- rt->stcf->symbol, rt->redis_object_expanded);
-
- if (rt->ctx->new_schema) {
- redisAsyncCommand (rt->redis, NULL, NULL, "HSET %s version 2",
- rt->redis_object_expanded);
- }
-
- if (rt->stcf->clcf->flags & RSPAMD_FLAG_CLASSIFIER_INTEGER) {
- redis_cmd = "HINCRBY";
- }
- else {
- redis_cmd = "HINCRBYFLOAT";
- }
-
- rt->id = id;
- query = rspamd_redis_tokens_to_query (task, rt, tokens,
- redis_cmd, rt->redis_object_expanded, TRUE, id,
- rt->stcf->clcf->flags & RSPAMD_FLAG_CLASSIFIER_INTEGER);
- g_assert (query != NULL);
- query->len = 0;
-
- /*
- * XXX:
- * Dirty hack: we get a token and check if it's value is -1 or 1, so
- * we could understand that we are learning or unlearning
- */
-
- tok = g_ptr_array_index (task->tokens, 0);
-
- if (tok->values[id] > 0) {
- rspamd_printf_fstring (&query, ""
- "*4\r\n"
- "$7\r\n"
- "HINCRBY\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n" /* Learned key */
- "$1\r\n"
- "1\r\n",
- (gint)strlen (rt->redis_object_expanded),
- rt->redis_object_expanded,
- (gint)strlen (learned_key),
- learned_key);
- }
- else {
- rspamd_printf_fstring (&query, ""
- "*4\r\n"
- "$7\r\n"
- "HINCRBY\r\n"
- "$%d\r\n"
- "%s\r\n"
- "$%d\r\n"
- "%s\r\n" /* Learned key */
- "$2\r\n"
- "-1\r\n",
- (gint)strlen (rt->redis_object_expanded),
- rt->redis_object_expanded,
- (gint)strlen (learned_key),
- learned_key);
- }
-
- ret = redisAsyncFormattedCommand (rt->redis, NULL, NULL,
- query->str, query->len);
-
- if (ret != REDIS_OK) {
- msg_err_task ("call to redis failed: %s", rt->redis->errstr);
- rspamd_fstring_free (query);
-
- return FALSE;
- }
-
- off = query->len;
- ret = rspamd_printf_fstring (&query, "*1\r\n$4\r\nEXEC\r\n");
- ret = redisAsyncFormattedCommand (rt->redis, rspamd_redis_learned, rt,
- query->str + off, ret);
- rspamd_mempool_add_destructor (task->task_pool,
- (rspamd_mempool_destruct_t)rspamd_fstring_free, query);
-
- if (ret == REDIS_OK) {
-
- /* Add signature if needed */
- if (rt->ctx->enable_signatures) {
- rspamd_redis_store_stat_signature (task, rt, tokens,
- "RSIG");
- }
-
- rspamd_session_add_event (task->s, rspamd_redis_fin_learn, rt, M);
- rt->has_event = TRUE;
-
- /* Set timeout */
- if (rspamd_event_pending (&rt->timeout_event, EV_TIMEOUT)) {
- event_del (&rt->timeout_event);
- }
- event_set (&rt->timeout_event, -1, EV_TIMEOUT, rspamd_redis_timeout, rt);
- event_base_set (task->ev_base, &rt->timeout_event);
- double_to_tv (rt->ctx->timeout, &tv);
- event_add (&rt->timeout_event, &tv);
-
- return TRUE;
- }
- else {
- msg_err_task ("call to redis failed: %s", rt->redis->errstr);
- }
-
- return FALSE;
- }
-
-
- gboolean
- rspamd_redis_finalize_learn (struct rspamd_task *task, gpointer runtime,
- gpointer ctx, GError **err)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (runtime);
- redisAsyncContext *redis;
-
- if (rspamd_event_pending (&rt->timeout_event, EV_TIMEOUT)) {
- event_del (&rt->timeout_event);
- }
-
- if (rt->redis) {
- redis = rt->redis;
- rt->redis = NULL;
- redisAsyncFree (redis);
- }
-
- if (rt->err) {
- g_propagate_error (err, rt->err);
- rt->err = NULL;
-
- return FALSE;
- }
-
- return TRUE;
- }
-
- gulong
- rspamd_redis_total_learns (struct rspamd_task *task, gpointer runtime,
- gpointer ctx)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (runtime);
-
- return rt->learned;
- }
-
- gulong
- rspamd_redis_inc_learns (struct rspamd_task *task, gpointer runtime,
- gpointer ctx)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (runtime);
-
- /* XXX: may cause races */
- return rt->learned + 1;
- }
-
- gulong
- rspamd_redis_dec_learns (struct rspamd_task *task, gpointer runtime,
- gpointer ctx)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (runtime);
-
- /* XXX: may cause races */
- return rt->learned + 1;
- }
-
- gulong
- rspamd_redis_learns (struct rspamd_task *task, gpointer runtime,
- gpointer ctx)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (runtime);
-
- return rt->learned;
- }
-
- ucl_object_t *
- rspamd_redis_get_stat (gpointer runtime,
- gpointer ctx)
- {
- struct redis_stat_runtime *rt = REDIS_RUNTIME (runtime);
- struct rspamd_redis_stat_elt *st;
- redisAsyncContext *redis;
-
- if (rt->ctx->stat_elt) {
- st = rt->ctx->stat_elt->ud;
-
- if (rt->redis) {
- redis = rt->redis;
- rt->redis = NULL;
- redisAsyncFree (redis);
- }
-
- if (st->stat) {
- return ucl_object_ref (st->stat);
- }
- }
-
- return NULL;
- }
-
- gpointer
- rspamd_redis_load_tokenizer_config (gpointer runtime,
- gsize *len)
- {
- return NULL;
- }
-
- #endif
|