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,
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,
.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,
.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,
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,
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);
* @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);
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 {
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;
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 <hash key> <raw digest> */
+ 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)
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);
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,
.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",
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)
{
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);
};
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];
};
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;
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;
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
});
}
+ 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 */
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;
+}
{.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},
};
{
/* 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;
}
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);
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++) {
}
}
+ 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");
}
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:
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;
}
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;
RSPAMD_CONTROL_MULTIPATTERN_LOADED,
RSPAMD_CONTROL_REGEXP_MAP_LOADED,
RSPAMD_CONTROL_MEMORY_STAT,
+ RSPAMD_CONTROL_FUZZY_HASH,
RSPAMD_CONTROL_MAX
};
struct {
unsigned int unused;
} fuzzy_sync;
+ struct {
+ unsigned char digest[64];
+ } fuzzy_hash;
struct {
enum {
rspamd_child_offline,
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;
"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 <hex> - 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 *
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";