From: Vsevolod Stakhov Date: Fri, 24 Jul 2026 18:18:47 +0000 (+0100) Subject: [Feature] fuzzy: per-hash introspection via rspamadm control fuzzyhash X-Git-Tag: 4.1.3~23^2~1 X-Git-Url: http://git.ipfire.org/gitweb/index.cgi?a=commitdiff_plain;h=f94ecd1de68116a54da9523082a78d64d87f9db0;p=thirdparty%2Frspamd.git [Feature] fuzzy: per-hash introspection via rspamadm control fuzzyhash rspamadm control fuzzyhash <128 hex digest> asks all fuzzy workers for storage-side diagnostics of a specific hash: flag slots and values, creation time, remaining ttl and, using the persisted shingle set, slot ownership (owned/foreign/vacant out of total). Ownership directly answers the key false-positive triage question: can this hash still produce fuzzy matches, or has its shingle anchor decayed? Implemented as a new inspect backend API: the redis backend gathers everything in a single server-side script round trip; sqlite queries the digests and shingles tables synchronously. The worker control handler replies asynchronously from the backend callback (the control pipe is persistent) using the same fd-attachment mechanism as fuzzy_stat. --- diff --git a/src/libserver/fuzzy_backend/fuzzy_backend.c b/src/libserver/fuzzy_backend/fuzzy_backend.c index d5155d9101..1055240497 100644 --- a/src/libserver/fuzzy_backend/fuzzy_backend.c +++ b/src/libserver/fuzzy_backend/fuzzy_backend.c @@ -49,6 +49,10 @@ static void rspamd_fuzzy_backend_version_sqlite(struct rspamd_fuzzy_backend *bk, void *subr_ud); static const char *rspamd_fuzzy_backend_id_sqlite(struct rspamd_fuzzy_backend *bk, void *subr_ud); +static void rspamd_fuzzy_backend_inspect_sqlite_wrapper(struct rspamd_fuzzy_backend *bk, + const unsigned char *digest, + rspamd_fuzzy_inspect_cb cb, void *ud, + void *subr_ud); static void rspamd_fuzzy_backend_expire_sqlite(struct rspamd_fuzzy_backend *bk, void *subr_ud); static void rspamd_fuzzy_backend_close_sqlite(struct rspamd_fuzzy_backend *bk, @@ -69,6 +73,10 @@ struct rspamd_fuzzy_backend_subr { void (*count)(struct rspamd_fuzzy_backend *bk, rspamd_fuzzy_count_cb cb, void *ud, void *subr_ud); + void (*inspect)(struct rspamd_fuzzy_backend *bk, + const unsigned char *digest, + rspamd_fuzzy_inspect_cb cb, void *ud, + void *subr_ud); void (*version)(struct rspamd_fuzzy_backend *bk, const char *src, rspamd_fuzzy_version_cb cb, void *ud, @@ -84,6 +92,7 @@ static const struct rspamd_fuzzy_backend_subr fuzzy_subrs[] = { .check = rspamd_fuzzy_backend_check_sqlite, .update = rspamd_fuzzy_backend_update_sqlite, .count = rspamd_fuzzy_backend_count_sqlite, + .inspect = rspamd_fuzzy_backend_inspect_sqlite_wrapper, .version = rspamd_fuzzy_backend_version_sqlite, .id = rspamd_fuzzy_backend_id_sqlite, .periodic = rspamd_fuzzy_backend_expire_sqlite, @@ -94,6 +103,7 @@ static const struct rspamd_fuzzy_backend_subr fuzzy_subrs[] = { .check = rspamd_fuzzy_backend_check_redis, .update = rspamd_fuzzy_backend_update_redis, .count = rspamd_fuzzy_backend_count_redis, + .inspect = rspamd_fuzzy_backend_inspect_redis, .version = rspamd_fuzzy_backend_version_redis, .id = rspamd_fuzzy_backend_id_redis, .periodic = rspamd_fuzzy_backend_expire_redis, @@ -468,6 +478,33 @@ void rspamd_fuzzy_backend_count(struct rspamd_fuzzy_backend *bk, bk->subr->count(bk, cb, ud, bk->subr_ud); } +void rspamd_fuzzy_backend_inspect(struct rspamd_fuzzy_backend *bk, + const unsigned char *digest, + rspamd_fuzzy_inspect_cb cb, void *ud) +{ + g_assert(bk != NULL); + + if (bk->subr->inspect) { + bk->subr->inspect(bk, digest, cb, ud, bk->subr_ud); + } + else if (cb) { + cb(NULL, ud); + } +} + +static void +rspamd_fuzzy_backend_inspect_sqlite_wrapper(struct rspamd_fuzzy_backend *bk, + const unsigned char *digest, + rspamd_fuzzy_inspect_cb cb, void *ud, + void *subr_ud) +{ + struct rspamd_fuzzy_backend_sqlite *sq = subr_ud; + + if (cb) { + cb(rspamd_fuzzy_backend_sqlite_inspect(sq, digest), ud); + } +} + void rspamd_fuzzy_backend_version(struct rspamd_fuzzy_backend *bk, const char *src, diff --git a/src/libserver/fuzzy_backend/fuzzy_backend.h b/src/libserver/fuzzy_backend/fuzzy_backend.h index ac66c77546..47addac27d 100644 --- a/src/libserver/fuzzy_backend/fuzzy_backend.h +++ b/src/libserver/fuzzy_backend/fuzzy_backend.h @@ -51,6 +51,8 @@ typedef void (*rspamd_fuzzy_update_cb)(gboolean success, typedef void (*rspamd_fuzzy_version_cb)(uint64_t rev, void *ud); typedef void (*rspamd_fuzzy_count_cb)(uint64_t count, void *ud); +/* Called with an ucl object describing the hash (ownership transferred) or NULL */ +typedef void (*rspamd_fuzzy_inspect_cb)(ucl_object_t *res, void *ud); typedef gboolean (*rspamd_fuzzy_periodic_cb)(void *ud); @@ -93,6 +95,14 @@ void rspamd_fuzzy_backend_process_updates(struct rspamd_fuzzy_backend *bk, * @param cb * @param ud */ +/** + * Get diagnostic information about a specific hash: flags, values, creation + * time, ttl and shingle slot ownership (where the backend supports it) + */ +void rspamd_fuzzy_backend_inspect(struct rspamd_fuzzy_backend *bk, + const unsigned char *digest, + rspamd_fuzzy_inspect_cb cb, void *ud); + void rspamd_fuzzy_backend_count(struct rspamd_fuzzy_backend *bk, rspamd_fuzzy_count_cb cb, void *ud); diff --git a/src/libserver/fuzzy_backend/fuzzy_backend_redis.c b/src/libserver/fuzzy_backend/fuzzy_backend_redis.c index 4f7a8540ea..2eed8a0f0a 100644 --- a/src/libserver/fuzzy_backend/fuzzy_backend_redis.c +++ b/src/libserver/fuzzy_backend/fuzzy_backend_redis.c @@ -67,7 +67,8 @@ struct rspamd_fuzzy_backend_redis { enum rspamd_fuzzy_redis_command { RSPAMD_FUZZY_REDIS_COMMAND_COUNT, RSPAMD_FUZZY_REDIS_COMMAND_VERSION, - RSPAMD_FUZZY_REDIS_COMMAND_CHECK + RSPAMD_FUZZY_REDIS_COMMAND_CHECK, + RSPAMD_FUZZY_REDIS_COMMAND_INSPECT }; struct rspamd_fuzzy_redis_session { @@ -86,6 +87,7 @@ struct rspamd_fuzzy_redis_session { rspamd_fuzzy_check_cb cb_check; rspamd_fuzzy_version_cb cb_version; rspamd_fuzzy_count_cb cb_count; + rspamd_fuzzy_inspect_cb cb_inspect; } callback; void *cbdata; @@ -855,6 +857,196 @@ rspamd_fuzzy_redis_count_callback(redisAsyncContext *c, gpointer r, rspamd_fuzzy_redis_session_dtor(session, FALSE); } +/* + * Runs entirely on the redis side to gather per-hash diagnostics in a single + * round trip: flag slots, ttl and shingle slot ownership computed from the + * persisted 'S' field (see fuzzy_update.lua) + */ +static const char *rspamd_fuzzy_redis_inspect_script = + "local key = KEYS[1]\n" + "if redis.call('EXISTS', key) == 0 then return cjson.encode({found=false}) end\n" + "local o = {found=true, ttl=redis.call('TTL', key)}\n" + "local data = redis.call('HGETALL', key)\n" + "local fields = {}\n" + "for i=1,#data,2 do fields[data[i]] = data[i+1] end\n" + "o.flag = tonumber(fields['F'])\n" + "o.value = tonumber(fields['V'])\n" + "o.created = tonumber(fields['C'])\n" + "local extra = {}\n" + "for i=1,7 do\n" + " if fields['F'..i] and fields['V'..i] then\n" + " extra[#extra+1] = {flag=tonumber(fields['F'..i]), value=tonumber(fields['V'..i])}\n" + " end\n" + "end\n" + "if #extra > 0 then o.extra_flags = extra end\n" + "if fields['S'] then\n" + " local owned, total, vacant = 0, 0, 0\n" + " local prefix = string.sub(key, 1, #key - #ARGV[1])\n" + " for suf in string.gmatch(fields['S'], '[^,]+') do\n" + " total = total + 1\n" + " local owner = redis.call('GET', prefix .. '_' .. suf)\n" + " if owner == ARGV[1] then owned = owned + 1\n" + " elseif owner == false then vacant = vacant + 1 end\n" + " end\n" + " o.shingles = {total=total, owned=owned, vacant=vacant, foreign=total-owned-vacant}\n" + "end\n" + "return cjson.encode(o)\n"; + +static void +rspamd_fuzzy_redis_inspect_callback(redisAsyncContext *c, gpointer r, + gpointer priv) +{ + struct rspamd_fuzzy_redis_session *session = priv; + redisReply *reply = r; + ucl_object_t *res = NULL; + + ev_timer_stop(session->event_loop, &session->timeout); + + if (c->err == 0 && reply != NULL) { + rspamd_upstream_ok(session->up); + + if (reply->type == REDIS_REPLY_STRING) { + struct ucl_parser *parser = ucl_parser_new(UCL_PARSER_SAFE_FLAGS); + + if (ucl_parser_add_chunk(parser, reply->str, reply->len)) { + res = ucl_parser_get_object(parser); + } + else { + msg_err_redis_session("cannot parse inspect reply: %s", + ucl_parser_get_error(parser)); + } + + ucl_parser_free(parser); + } + else if (reply->type == REDIS_REPLY_ERROR) { + msg_err_redis_session("fuzzy backend redis error: \"%s\"", + reply->str); + } + } + else { + if (c->errstr) { + msg_err_redis_session("error inspecting hash on %s: %s", + rspamd_inet_address_to_string_pretty(rspamd_upstream_addr_cur(session->up)), + c->errstr); + rspamd_upstream_fail(session->up, FALSE, c->errstr); + } + } + + if (session->callback.cb_inspect) { + session->callback.cb_inspect(res, session->cbdata); + } + else if (res) { + ucl_object_unref(res); + } + + rspamd_fuzzy_redis_session_dtor(session, FALSE); +} + +void rspamd_fuzzy_backend_inspect_redis(struct rspamd_fuzzy_backend *bk, + const unsigned char *digest, + rspamd_fuzzy_inspect_cb cb, void *ud, + void *subr_ud) +{ + struct rspamd_fuzzy_backend_redis *backend = subr_ud; + struct rspamd_fuzzy_redis_session *session; + struct upstream *up; + struct upstream_list *ups; + rspamd_inet_addr_t *addr; + GString *key; + + g_assert(backend != NULL); + + ups = rspamd_redis_get_servers(backend, "read_servers"); + if (!ups) { + if (cb) { + cb(NULL, ud); + } + + return; + } + + session = g_malloc0(sizeof(*session)); + session->backend = backend; + REF_RETAIN(session->backend); + + session->callback.cb_inspect = cb; + session->cbdata = ud; + session->command = RSPAMD_FUZZY_REDIS_COMMAND_INSPECT; + session->event_loop = rspamd_fuzzy_backend_event_base(bk); + + /* EVAL script 1 */ + session->nargs = 5; + session->argv = g_malloc0(sizeof(char *) * session->nargs); + session->argv_lens = g_malloc0(sizeof(gsize) * session->nargs); + session->argv[0] = g_strdup("EVAL"); + session->argv_lens[0] = 4; + session->argv[1] = g_strdup(rspamd_fuzzy_redis_inspect_script); + session->argv_lens[1] = strlen(rspamd_fuzzy_redis_inspect_script); + session->argv[2] = g_strdup("1"); + session->argv_lens[2] = 1; + key = g_string_new(backend->redis_object); + g_string_append_len(key, digest, rspamd_cryptobox_HASHBYTES); + session->argv[3] = key->str; + session->argv_lens[3] = key->len; + g_string_free(key, FALSE); /* Do not free underlying array */ + session->argv[4] = g_malloc(rspamd_cryptobox_HASHBYTES); + memcpy(session->argv[4], digest, rspamd_cryptobox_HASHBYTES); + session->argv_lens[4] = rspamd_cryptobox_HASHBYTES; + + up = rspamd_upstream_get(ups, + RSPAMD_UPSTREAM_ROUND_ROBIN, + NULL, + 0); + + if (up == NULL) { + msg_err_redis_session("cannot select fuzzy redis upstream for inspect: " + "all backends are dead or pending DNS resolution"); + rspamd_fuzzy_redis_session_dtor(session, TRUE); + if (cb) { + cb(NULL, ud); + } + return; + } + + session->up = rspamd_upstream_ref(up); + addr = rspamd_upstream_addr_next(up); + g_assert(addr != NULL); + session->ctx = rspamd_redis_pool_connect(backend->pool, + backend->dbname, + backend->username, backend->password, + rspamd_inet_address_to_string(addr), + rspamd_inet_address_get_port(addr)); + + if (session->ctx == NULL) { + rspamd_upstream_fail(up, TRUE, strerror(errno)); + rspamd_fuzzy_redis_session_dtor(session, TRUE); + + if (cb) { + cb(NULL, ud); + } + } + else { + if (redisAsyncCommandArgv(session->ctx, rspamd_fuzzy_redis_inspect_callback, + session, session->nargs, + (const char **) session->argv, session->argv_lens) != REDIS_OK) { + rspamd_fuzzy_redis_session_dtor(session, TRUE); + + if (cb) { + cb(NULL, ud); + } + } + else { + /* Add timeout */ + session->timeout.data = session; + ev_now_update_if_cheap((struct ev_loop *) session->event_loop); + ev_timer_init(&session->timeout, + rspamd_fuzzy_redis_timeout, + session->backend->timeout, 0.0); + ev_timer_start(session->event_loop, &session->timeout); + } + } +} + void rspamd_fuzzy_backend_count_redis(struct rspamd_fuzzy_backend *bk, rspamd_fuzzy_count_cb cb, void *ud, void *subr_ud) diff --git a/src/libserver/fuzzy_backend/fuzzy_backend_redis.h b/src/libserver/fuzzy_backend/fuzzy_backend_redis.h index 0a536c2fa2..109597edda 100644 --- a/src/libserver/fuzzy_backend/fuzzy_backend_redis.h +++ b/src/libserver/fuzzy_backend/fuzzy_backend_redis.h @@ -41,6 +41,10 @@ void rspamd_fuzzy_backend_update_redis(struct rspamd_fuzzy_backend *bk, rspamd_fuzzy_update_cb cb, void *ud, void *subr_ud); +void rspamd_fuzzy_backend_inspect_redis(struct rspamd_fuzzy_backend *bk, + const unsigned char *digest, + rspamd_fuzzy_inspect_cb cb, void *ud, + void *subr_ud); void rspamd_fuzzy_backend_count_redis(struct rspamd_fuzzy_backend *bk, rspamd_fuzzy_count_cb cb, void *ud, void *subr_ud); diff --git a/src/libserver/fuzzy_backend/fuzzy_backend_sqlite.c b/src/libserver/fuzzy_backend/fuzzy_backend_sqlite.c index f3c1690e06..a953fc81ce 100644 --- a/src/libserver/fuzzy_backend/fuzzy_backend_sqlite.c +++ b/src/libserver/fuzzy_backend/fuzzy_backend_sqlite.c @@ -97,6 +97,8 @@ enum rspamd_fuzzy_statement_idx { RSPAMD_FUZZY_BACKEND_GET_DIGEST_BY_ID, RSPAMD_FUZZY_BACKEND_DELETE, RSPAMD_FUZZY_BACKEND_COUNT, + RSPAMD_FUZZY_BACKEND_INSPECT, + RSPAMD_FUZZY_BACKEND_SHINGLES_COUNT, RSPAMD_FUZZY_BACKEND_EXPIRE, RSPAMD_FUZZY_BACKEND_VACUUM, RSPAMD_FUZZY_BACKEND_DELETE_ORPHANED, @@ -177,6 +179,16 @@ static struct rspamd_fuzzy_stmts { .args = "", .stmt = NULL, .result = SQLITE_ROW}, + {.idx = RSPAMD_FUZZY_BACKEND_INSPECT, + .sql = "SELECT id, value, time, flag FROM digests WHERE digest==?1;", + .args = "D", + .stmt = NULL, + .result = SQLITE_ROW}, + {.idx = RSPAMD_FUZZY_BACKEND_SHINGLES_COUNT, + .sql = "SELECT COUNT(*) FROM shingles WHERE digest_id=?1;", + .args = "I", + .stmt = NULL, + .result = SQLITE_ROW}, {.idx = RSPAMD_FUZZY_BACKEND_EXPIRE, .sql = "DELETE FROM digests WHERE id IN (SELECT id FROM digests WHERE time < ?1 LIMIT ?2);", .args = "II", @@ -999,6 +1011,63 @@ gsize rspamd_fuzzy_backend_sqlite_count(struct rspamd_fuzzy_backend_sqlite *back return 0; } +ucl_object_t * +rspamd_fuzzy_backend_sqlite_inspect(struct rspamd_fuzzy_backend_sqlite *backend, + const unsigned char *digest) +{ + ucl_object_t *res; + int rc; + + if (backend == NULL) { + return NULL; + } + + res = ucl_object_typed_new(UCL_OBJECT); + + rc = rspamd_fuzzy_backend_sqlite_run_stmt(backend, FALSE, + RSPAMD_FUZZY_BACKEND_INSPECT, digest); + + if (rc == SQLITE_OK) { + int64_t id = sqlite3_column_int64( + prepared_stmts[RSPAMD_FUZZY_BACKEND_INSPECT].stmt, 0); + + ucl_object_insert_key(res, ucl_object_frombool(true), "found", 0, false); + ucl_object_insert_key(res, + ucl_object_fromint(sqlite3_column_int64( + prepared_stmts[RSPAMD_FUZZY_BACKEND_INSPECT].stmt, 1)), + "value", 0, false); + ucl_object_insert_key(res, + ucl_object_fromint(sqlite3_column_int64( + prepared_stmts[RSPAMD_FUZZY_BACKEND_INSPECT].stmt, 2)), + "created", 0, false); + ucl_object_insert_key(res, + ucl_object_fromint(sqlite3_column_int( + prepared_stmts[RSPAMD_FUZZY_BACKEND_INSPECT].stmt, 3)), + "flag", 0, false); + rspamd_fuzzy_backend_sqlite_cleanup_stmt(backend, RSPAMD_FUZZY_BACKEND_INSPECT); + + if (rspamd_fuzzy_backend_sqlite_run_stmt(backend, FALSE, + RSPAMD_FUZZY_BACKEND_SHINGLES_COUNT, id) == SQLITE_OK) { + int64_t nshingles = sqlite3_column_int64( + prepared_stmts[RSPAMD_FUZZY_BACKEND_SHINGLES_COUNT].stmt, 0); + ucl_object_t *sgl = ucl_object_typed_new(UCL_OBJECT); + + /* Shingles are bound to the digest id in sqlite, so all are owned */ + ucl_object_insert_key(sgl, ucl_object_fromint(nshingles), "total", 0, false); + ucl_object_insert_key(sgl, ucl_object_fromint(nshingles), "owned", 0, false); + ucl_object_insert_key(res, sgl, "shingles", 0, false); + } + + rspamd_fuzzy_backend_sqlite_cleanup_stmt(backend, RSPAMD_FUZZY_BACKEND_SHINGLES_COUNT); + } + else { + ucl_object_insert_key(res, ucl_object_frombool(false), "found", 0, false); + rspamd_fuzzy_backend_sqlite_cleanup_stmt(backend, RSPAMD_FUZZY_BACKEND_INSPECT); + } + + return res; +} + int rspamd_fuzzy_backend_sqlite_version(struct rspamd_fuzzy_backend_sqlite *backend, const char *source) { diff --git a/src/libserver/fuzzy_backend/fuzzy_backend_sqlite.h b/src/libserver/fuzzy_backend/fuzzy_backend_sqlite.h index 1ace52f350..5286f33858 100644 --- a/src/libserver/fuzzy_backend/fuzzy_backend_sqlite.h +++ b/src/libserver/fuzzy_backend/fuzzy_backend_sqlite.h @@ -94,6 +94,13 @@ void rspamd_fuzzy_backend_sqlite_close(struct rspamd_fuzzy_backend_sqlite *backe gsize rspamd_fuzzy_backend_sqlite_count(struct rspamd_fuzzy_backend_sqlite *backend); +/** + * Get diagnostic info for a specific digest (or NULL if the backend is not available): + * flags, values, creation time and the number of shingles stored for the digest + */ +ucl_object_t *rspamd_fuzzy_backend_sqlite_inspect(struct rspamd_fuzzy_backend_sqlite *backend, + const unsigned char *digest); + int rspamd_fuzzy_backend_sqlite_version(struct rspamd_fuzzy_backend_sqlite *backend, const char *source); gsize rspamd_fuzzy_backend_sqlite_expired(struct rspamd_fuzzy_backend_sqlite *backend); diff --git a/src/libserver/fuzzy_storage_internal.h b/src/libserver/fuzzy_storage_internal.h index 3c491895ec..3353d21838 100644 --- a/src/libserver/fuzzy_storage_internal.h +++ b/src/libserver/fuzzy_storage_internal.h @@ -82,7 +82,7 @@ struct rspamd_leaky_bucket_elt { }; struct rspamd_fuzzy_dynamic_ban { - double expire_ts; /* monotonic clock; 0.0 = never expires */ + double expire_ts; /* monotonic clock; 0.0 = never expires */ int32_t response_code; /* fuzzy reply code; 0 = use default (503) */ char reason[64]; }; @@ -182,6 +182,8 @@ struct rspamd_fuzzy_storage_ctx { struct rspamd_http_context *http_ctx; rspamd_lru_hash_t *errors_ips; rspamd_lru_hash_t *ratelimit_buckets; + /* Aggregate + per-IP stats for unkeyed clients (allowed by allow_update etc) */ + struct fuzzy_key_stat *unkeyed_stat; struct rspamd_fuzzy_backend *backend; GArray *updates_pending; unsigned int updates_failed; @@ -222,7 +224,7 @@ struct fuzzy_session { rspamd_inet_addr_t *addr; struct rspamd_fuzzy_storage_ctx *ctx; - struct rspamd_fuzzy_shingle_cmd cmd; /* Can handle both shingles and non-shingles */ + struct rspamd_fuzzy_shingle_cmd cmd; /* Can handle both shingles and non-shingles */ union { struct rspamd_fuzzy_encrypted_reply v1; struct rspamd_fuzzy_encrypted_reply_v2 v2; @@ -349,4 +351,10 @@ gboolean rspamd_fuzzy_storage_stat(struct rspamd_main *rspamd_main, struct rspamd_control_command *cmd, gpointer ud); +gboolean rspamd_fuzzy_storage_hash_info(struct rspamd_main *rspamd_main, + struct rspamd_worker *worker, int fd, + int attached_fd, + struct rspamd_control_command *cmd, + gpointer ud); + #endif diff --git a/src/libserver/fuzzy_storage_stat.c b/src/libserver/fuzzy_storage_stat.c index ef61001c74..e169844a23 100644 --- a/src/libserver/fuzzy_storage_stat.c +++ b/src/libserver/fuzzy_storage_stat.c @@ -135,6 +135,29 @@ rspamd_fuzzy_stat_to_ucl(struct rspamd_fuzzy_storage_ctx *ctx, gboolean ip_stat) }); } + if (ctx->unkeyed_stat) { + /* Pseudo-key entry for unkeyed clients (e.g. allowed by IP) */ + elt = rspamd_fuzzy_storage_stat_key(ctx->unkeyed_stat); + + if (ctx->unkeyed_stat->last_ips && ip_stat) { + int i = 0; + gpointer k, v; + + ip_elt = ucl_object_typed_new(UCL_OBJECT); + + while ((i = rspamd_lru_hash_foreach(ctx->unkeyed_stat->last_ips, + i, &k, &v)) != -1) { + ucl_object_insert_key(ip_elt, + rspamd_fuzzy_storage_stat_key(v), + rspamd_inet_address_to_string(k), 0, true); + } + + ucl_object_insert_key(elt, ip_elt, "ips", 0, false); + } + + ucl_object_insert_key(keys_obj, elt, "unkeyed", 0, false); + } + ucl_object_insert_key(obj, keys_obj, "keys", 0, false); /* Now generic stats */ @@ -298,3 +321,124 @@ rspamd_fuzzy_storage_stat(struct rspamd_main *rspamd_main, return TRUE; } + +struct rspamd_fuzzy_hash_info_cbdata { + struct rspamd_main *rspamd_main; + struct rspamd_fuzzy_storage_ctx *ctx; + uint64_t id; + int fd; +}; + +static void +rspamd_fuzzy_hash_info_cb(ucl_object_t *res, void *ud) +{ + struct rspamd_fuzzy_hash_info_cbdata *cbd = ud; + struct rspamd_main *rspamd_main = cbd->rspamd_main; + struct rspamd_control_reply rep; + struct ucl_emitter_functions *emit_subr; + unsigned char fdspace[CMSG_SPACE(sizeof(int))]; + struct iovec iov; + struct msghdr msg; + struct cmsghdr *cmsg; + int outfd = -1; + char tmppath[PATH_MAX]; + + memset(&rep, 0, sizeof(rep)); + rep.type = RSPAMD_CONTROL_FUZZY_HASH; + rep.id = cbd->id; + + if (res == NULL) { + rep.reply.fuzzy_hash.status = ENOENT; + } + else { + const char *backend_id = rspamd_fuzzy_backend_id(cbd->ctx->backend); + + if (backend_id) { + memcpy(rep.reply.fuzzy_hash.storage_id, + backend_id, + sizeof(rep.reply.fuzzy_hash.storage_id)); + } + + rspamd_snprintf(tmppath, sizeof(tmppath), "%s%c%s-XXXXXXXXXX", + rspamd_main->cfg->temp_dir, G_DIR_SEPARATOR, "fuzzy-hash"); + + if ((outfd = mkstemp(tmppath)) == -1) { + rep.reply.fuzzy_hash.status = errno; + msg_info_main("cannot make temporary file for fuzzy hash info: %s", + strerror(errno)); + } + else { + rep.reply.fuzzy_hash.status = 0; + emit_subr = ucl_object_emit_fd_funcs(outfd); + ucl_object_emit_full(res, UCL_EMIT_JSON_COMPACT, emit_subr, NULL); + ucl_object_emit_funcs_free(emit_subr); + /* Rewind output file */ + close(outfd); + outfd = open(tmppath, O_RDONLY); + unlink(tmppath); + } + + ucl_object_unref(res); + } + + memset(&msg, 0, sizeof(msg)); + + if (outfd != -1) { + memset(fdspace, 0, sizeof(fdspace)); + msg.msg_control = fdspace; + msg.msg_controllen = sizeof(fdspace); + cmsg = CMSG_FIRSTHDR(&msg); + + if (cmsg) { + cmsg->cmsg_level = SOL_SOCKET; + cmsg->cmsg_type = SCM_RIGHTS; + cmsg->cmsg_len = CMSG_LEN(sizeof(int)); + memcpy(CMSG_DATA(cmsg), &outfd, sizeof(int)); + } + } + + iov.iov_base = &rep; + iov.iov_len = sizeof(rep); + msg.msg_iov = &iov; + msg.msg_iovlen = 1; + + if (sendmsg(cbd->fd, &msg, 0) == -1) { + msg_err_main("cannot send fuzzy hash info: %s", strerror(errno)); + } + + if (outfd != -1) { + close(outfd); + } + + g_free(cbd); +} + +gboolean +rspamd_fuzzy_storage_hash_info(struct rspamd_main *rspamd_main, + struct rspamd_worker *worker, int fd, + int attached_fd, + struct rspamd_control_command *cmd, + gpointer ud) +{ + struct rspamd_fuzzy_storage_ctx *ctx = ud; + struct rspamd_fuzzy_hash_info_cbdata *cbd; + + if (attached_fd != -1) { + close(attached_fd); + } + + cbd = g_malloc0(sizeof(*cbd)); + cbd->rspamd_main = rspamd_main; + cbd->ctx = ctx; + cbd->id = cmd->id; + cbd->fd = fd; + + /* + * The reply is sent from the backend callback: the control pipe is + * persistent, so a deferred reply is safe here + */ + rspamd_fuzzy_backend_inspect(ctx->backend, cmd->cmd.fuzzy_hash.digest, + rspamd_fuzzy_hash_info_cb, cbd); + + return TRUE; +} diff --git a/src/libserver/rspamd_control.c b/src/libserver/rspamd_control.c index cc6e92e47e..55e4874c6e 100644 --- a/src/libserver/rspamd_control.c +++ b/src/libserver/rspamd_control.c @@ -84,6 +84,7 @@ static const struct rspamd_control_cmd_match { {.name = {.begin = "/recompile", .len = sizeof("/recompile") - 1}, .type = RSPAMD_CONTROL_RECOMPILE}, {.name = {.begin = "/fuzzystat", .len = sizeof("/fuzzystat") - 1}, .type = RSPAMD_CONTROL_FUZZY_STAT}, {.name = {.begin = "/fuzzysync", .len = sizeof("/fuzzysync") - 1}, .type = RSPAMD_CONTROL_FUZZY_SYNC}, + {.name = {.begin = "/fuzzyhash", .len = sizeof("/fuzzyhash") - 1}, .type = RSPAMD_CONTROL_FUZZY_HASH}, {.name = {.begin = "/compositesstats", .len = sizeof("/compositesstats") - 1}, .type = RSPAMD_CONTROL_COMPOSITES_STATS}, {.name = {.begin = "/memstat", .len = sizeof("/memstat") - 1}, .type = RSPAMD_CONTROL_MEMORY_STAT}, }; @@ -198,7 +199,8 @@ rspamd_control_write_reply(struct rspamd_control_session *session) { /* Skip incompatible worker for fuzzy_stat */ if ((session->cmd.type == RSPAMD_CONTROL_FUZZY_STAT || - session->cmd.type == RSPAMD_CONTROL_FUZZY_SYNC) && + session->cmd.type == RSPAMD_CONTROL_FUZZY_SYNC || + session->cmd.type == RSPAMD_CONTROL_FUZZY_HASH) && elt->wrk_type != g_quark_from_static_string("fuzzy")) { continue; } @@ -279,6 +281,34 @@ rspamd_control_write_reply(struct rspamd_control_session *session) case RSPAMD_CONTROL_FUZZY_SYNC: ucl_object_insert_key(cur, ucl_object_fromint(elt->reply.reply.fuzzy_sync.status), "status", 0, false); break; + case RSPAMD_CONTROL_FUZZY_HASH: + ucl_object_insert_key(cur, + ucl_object_fromint(elt->reply.reply.fuzzy_hash.status), + "status", 0, false); + + if (elt->attached_fd != -1) { + parser = ucl_parser_new(UCL_PARSER_SAFE_FLAGS); + + if (ucl_parser_add_fd(parser, elt->attached_fd)) { + ucl_object_insert_key(cur, ucl_parser_get_object(parser), + "data", 0, false); + } + else { + ucl_object_insert_key(cur, + ucl_object_fromstring(ucl_parser_get_error(parser)), + "error", 0, false); + } + + ucl_parser_free(parser); + ucl_object_insert_key(cur, + ucl_object_fromlstring( + elt->reply.reply.fuzzy_hash.storage_id, + MEMPOOL_UID_LEN - 1), + "id", + 0, + false); + } + break; case RSPAMD_CONTROL_COMPOSITES_STATS: ucl_object_insert_key(cur, ucl_object_fromint(elt->reply.reply.composites_stats.checked_slow), "checked_slow", 0, false); @@ -741,6 +771,16 @@ rspamd_control_finish_handler(struct rspamd_http_connection *conn, srch.begin = msg->url->str; srch.len = msg->url->len; + /* Some commands carry arguments in the query string */ + const char *qpos = memchr(msg->url->str, '?', msg->url->len); + rspamd_ftok_t query = {.begin = NULL, .len = 0}; + + if (qpos != NULL) { + query.begin = qpos + 1; + query.len = msg->url->len - (qpos - msg->url->str) - 1; + srch.len = qpos - msg->url->str; + } + session->is_reply = TRUE; for (i = 0; i < G_N_ELEMENTS(cmd_matches); i++) { @@ -751,6 +791,25 @@ rspamd_control_finish_handler(struct rspamd_http_connection *conn, } } + if (found && session->cmd.type == RSPAMD_CONTROL_FUZZY_HASH) { + /* Parse digest=<128 hex chars> */ + static const char digest_pfx[] = "digest="; + gsize hexlen = sizeof(session->cmd.cmd.fuzzy_hash.digest) * 2; + + if (query.len == sizeof(digest_pfx) - 1 + hexlen && + memcmp(query.begin, digest_pfx, sizeof(digest_pfx) - 1) == 0 && + rspamd_decode_hex_buf(query.begin + sizeof(digest_pfx) - 1, hexlen, + session->cmd.cmd.fuzzy_hash.digest, + sizeof(session->cmd.cmd.fuzzy_hash.digest)) != -1) { + /* Parsed fine */ + } + else { + rspamd_control_send_error(session, 400, + "fuzzyhash requires ?digest=<%z hex characters>", hexlen); + return 0; + } + } + if (!found) { rspamd_control_send_error(session, 404, "Command not defined"); } @@ -857,6 +916,7 @@ rspamd_control_default_cmd_handler(int fd, case RSPAMD_CONTROL_MONITORED_CHANGE: case RSPAMD_CONTROL_FUZZY_STAT: case RSPAMD_CONTROL_FUZZY_SYNC: + case RSPAMD_CONTROL_FUZZY_HASH: case RSPAMD_CONTROL_LOG_PIPE: case RSPAMD_CONTROL_CHILD_CHANGE: case RSPAMD_CONTROL_FUZZY_BLOCKED: @@ -1803,6 +1863,9 @@ rspamd_control_command_from_string(const char *str) else if (g_ascii_strcasecmp(str, "fuzzy_sync") == 0) { ret = RSPAMD_CONTROL_FUZZY_SYNC; } + else if (g_ascii_strcasecmp(str, "fuzzy_hash") == 0) { + ret = RSPAMD_CONTROL_FUZZY_HASH; + } else if (g_ascii_strcasecmp(str, "monitored_change") == 0) { ret = RSPAMD_CONTROL_MONITORED_CHANGE; } @@ -1849,6 +1912,9 @@ rspamd_control_command_to_string(enum rspamd_control_type cmd) case RSPAMD_CONTROL_FUZZY_SYNC: reply = "fuzzy_sync"; break; + case RSPAMD_CONTROL_FUZZY_HASH: + reply = "fuzzy_hash"; + break; case RSPAMD_CONTROL_MONITORED_CHANGE: reply = "monitored_change"; break; diff --git a/src/libserver/rspamd_control.h b/src/libserver/rspamd_control.h index 39835d1c21..32faa3cfb7 100644 --- a/src/libserver/rspamd_control.h +++ b/src/libserver/rspamd_control.h @@ -42,6 +42,7 @@ enum rspamd_control_type { RSPAMD_CONTROL_MULTIPATTERN_LOADED, RSPAMD_CONTROL_REGEXP_MAP_LOADED, RSPAMD_CONTROL_MEMORY_STAT, + RSPAMD_CONTROL_FUZZY_HASH, RSPAMD_CONTROL_MAX }; @@ -106,6 +107,9 @@ struct rspamd_control_command { struct { unsigned int unused; } fuzzy_sync; + struct { + unsigned char digest[64]; + } fuzzy_hash; struct { enum { rspamd_child_offline, @@ -168,6 +172,10 @@ struct rspamd_control_reply { unsigned int status; char storage_id[MEMPOOL_UID_LEN]; } fuzzy_stat; + struct { + unsigned int status; + char storage_id[MEMPOOL_UID_LEN]; + } fuzzy_hash; struct { unsigned int status; } fuzzy_sync; diff --git a/src/rspamadm/control.c b/src/rspamadm/control.c index e90697702e..71c437c6eb 100644 --- a/src/rspamadm/control.c +++ b/src/rspamadm/control.c @@ -63,15 +63,16 @@ static GOptionEntry entries[] = { "Set IO timeout (1s by default)", NULL}, {NULL, 0, 0, G_OPTION_ARG_NONE, NULL, NULL, NULL}}; -#define RSPAMADM_CONTROL_COMMAND_LIST \ - "Supported commands:\n" \ - " stat - show statistics\n" \ - " reload - reload workers dynamic data\n" \ - " reresolve - resolve upstreams addresses\n" \ - " recompile - recompile hyperscan regexes\n" \ - " fuzzystat - show fuzzy statistics\n" \ - " fuzzysync - immediately sync fuzzy database to storage\n" \ - " compositesstats - show composites processing statistics\n" \ +#define RSPAMADM_CONTROL_COMMAND_LIST \ + "Supported commands:\n" \ + " stat - show statistics\n" \ + " reload - reload workers dynamic data\n" \ + " reresolve - resolve upstreams addresses\n" \ + " recompile - recompile hyperscan regexes\n" \ + " fuzzystat - show fuzzy statistics\n" \ + " fuzzysync - immediately sync fuzzy database to storage\n" \ + " fuzzyhash - show storage info for a fuzzy hash (128 hex chars)\n" \ + " compositesstats - show composites processing statistics\n" \ " memstat - show memory usage statistics across all workers\n" static const char * @@ -232,6 +233,20 @@ rspamadm_control(int argc, char **argv, const struct rspamadm_command *_cmd) g_ascii_strcasecmp(cmd, "fuzzy_sync") == 0) { path = "/fuzzysync"; } + else if (g_ascii_strcasecmp(cmd, "fuzzyhash") == 0 || + g_ascii_strcasecmp(cmd, "fuzzy_hash") == 0) { + static char hash_path_buf[256]; + + if (argc < 3 || strlen(argv[2]) != 128) { + rspamd_fprintf(stderr, + "fuzzyhash requires a 128 characters hex digest argument\n"); + exit(EXIT_FAILURE); + } + + rspamd_snprintf(hash_path_buf, sizeof(hash_path_buf), + "/fuzzyhash?digest=%s", argv[2]); + path = hash_path_buf; + } else if (g_ascii_strcasecmp(cmd, "compositesstats") == 0 || g_ascii_strcasecmp(cmd, "composites_stats") == 0) { path = "/compositesstats";