Update analysis hash filters to use WHERE IN (#98)

This commit is contained in:
Naomi Kramer
2026-04-01 12:58:34 -04:00
committed by GitHub
parent 8f9be0291c
commit e101fe2a86
2 changed files with 51 additions and 19 deletions
+19 -19
View File
@@ -217,10 +217,10 @@ func (analyzer *Analyzer) ScoopSNIConns(ctx context.Context, bars *tea.Program)
groupUniqArrayMerge(10)(proxy_ips) AS proxy_ips,
maxMerge(last_seen) AS last_seen,
minMerge(first_seen) as first_seen
FROM usni
RIGHT JOIN unique_sni USING hash
FROM usni
-- Limit query to the last 24 hours of data
WHERE hour >= toStartOfHour(fromUnixTimestamp({min_ts:Int64}))
WHERE hour >= toStartOfHour(fromUnixTimestamp({min_ts:Int64})) AND
hash IN (SELECT hash FROM unique_sni)
GROUP BY hash, src, src_nuid, fqdn, proxy
UNION ALL
@@ -269,15 +269,15 @@ func (analyzer *Analyzer) ScoopSNIConns(ctx context.Context, bars *tea.Program)
historical AS (
SELECT min(first_seen) AS first_seen, fqdn
FROM metadatabase.historical_first_seen
LEFT JOIN sniconns USING fqdn
WHERE fqdn IN (SELECT fqdn FROM sniconns)
GROUP BY fqdn
),
port_proto AS (
SELECT hash, groupUniqArray(20)(port_proto_service) AS port_proto_service FROM (
SELECT DISTINCT hash, concat(po.dst_port, ':', po.proto, ':', po.service) as port_proto_service
FROM port_info po
LEFT JOIN sniconns s ON s.hash = po.hash
WHERE hour >= toStartOfHour(fromUnixTimestamp({min_ts:Int64}))
WHERE hour >= toStartOfHour(fromUnixTimestamp({min_ts:Int64})) AND
hash IN (SELECT hash FROM sniconns)
UNION DISTINCT
SELECT DISTINCT hash, concat(dst_port, ':', proto, ':', service) FROM openhttp
UNION DISTINCT
@@ -412,8 +412,8 @@ func (analyzer *Analyzer) ScoopIPConns(ctx context.Context, bars *tea.Program) e
sniconns AS ( -- usni connections that will be beacons or long connections in this import
SELECT hash, uniqExactMerge(u.unique_ts_count) AS unique_count, countMerge(u.count) AS total_count, sumMerge(total_duration) AS duration
FROM usni u
LEFT SEMI JOIN sniconn_tmp t USING hash
WHERE hour >= toStartOfHour(fromUnixTimestamp({min_ts:Int64}))
WHERE hour >= toStartOfHour(fromUnixTimestamp({min_ts:Int64}))
AND hash IN (SELECT hash FROM sniconn_tmp)
GROUP BY hash
-- is beacon or longconn
HAVING (unique_count >= {unique_connection_threshold:UInt64} AND total_count < 86400) OR
@@ -452,10 +452,10 @@ func (analyzer *Analyzer) ScoopIPConns(ctx context.Context, bars *tea.Program) e
maxMerge(last_seen) as last_seen,
minMerge(first_seen) as first_seen
FROM uconn
-- Limit IP connections to just connections not used by a SNI beacon
RIGHT JOIN filtered_hashes USING hash
-- Limit query to the last 24 hours of data
WHERE hour >= toStartOfHour(fromUnixTimestamp({min_ts:Int64}))
WHERE hour >= toStartOfHour(fromUnixTimestamp({min_ts:Int64}))
-- Limit IP connections to just connections not used by a SNI beacon
AND hash IN (SELECT hash FROM filtered_hashes)
GROUP BY hash, src, src_nuid, dst, dst_nuid, src_local, dst_local
UNION ALL
@@ -475,12 +475,13 @@ func (analyzer *Analyzer) ScoopIPConns(ctx context.Context, bars *tea.Program) e
min(ts) AS first_seen,
max(ts) AS last_seen
FROM openconn
RIGHT JOIN filtered_hashes USING hash -- exclude SNI connections
WHERE hash IN (SELECT hash FROM filtered_hashes) -- exclude SNI connections
GROUP BY hash, src, src_nuid, dst, dst_nuid, src_local, dst_local
),
-- Aggregate data between all union groups
totaled_ipconns AS (
SELECT hash, src, src_nuid, dst, dst_nuid, src_local, dst_local,
multiIf(src_local = true, dst, dst_local = true, src, dst) as target_ip,
sum(missing_host_count) as missing_host_count,
sum(conn_count) as count,
sum(open_count) as open_count,
@@ -506,15 +507,15 @@ func (analyzer *Analyzer) ScoopIPConns(ctx context.Context, bars *tea.Program) e
historical AS (
SELECT min(first_seen) AS first_seen, ip
FROM metadatabase.historical_first_seen h
LEFT JOIN ip_conns i ON h.ip = multiIf(src_local = true, i.dst, dst_local = true, i.src, i.dst)
WHERE ip IN (SELECT target_ip FROM totaled_ipconns)
GROUP BY ip
),
port_proto AS (
SELECT hash, groupUniqArray(20)(port_proto_service) AS port_proto_service FROM (
SELECT DISTINCT hash, if(po.proto = 'icmp', concat(po.proto, ':', po.icmp_type, '/', po.icmp_code), concat(po.dst_port, ':', po.proto, ':', po.service)) as port_proto_service
FROM port_info po
LEFT JOIN ip_conns i ON i.hash = po.hash
WHERE hour >= toStartOfHour(fromUnixTimestamp({min_ts:Int64}))
WHERE hour >= toStartOfHour(fromUnixTimestamp({min_ts:Int64}))
AND hash IN (SELECT hash FROM totaled_ipconns)
UNION DISTINCT
SELECT DISTINCT hash, if(proto = 'icmp', concat(proto, ':', src_port, '/', dst_port), concat(dst_port, ':', proto, ':', service)) as port_proto_service
FROM openconn
@@ -542,10 +543,9 @@ func (analyzer *Analyzer) ScoopIPConns(ctx context.Context, bars *tea.Program) e
po.port_proto_service as port_proto_service
FROM totaled_ipconns i
LEFT JOIN prevalence_counts p ON if(src_local = true, i.dst, i.src) = p.ip
LEFT JOIN metadatabase.threat_intel t ON multiIf(src_local = true, i.dst, dst_local = true, i.src, i.dst) = t.ip
LEFT JOIN metadatabase.threat_intel t ON i.target_ip = t.ip
LEFT JOIN port_proto po ON i.hash = po.hash
LEFT JOIN historical h ON multiIf(src_local = true, i.dst, dst_local = true, i.src, i.dst) = h.ip
LEFT JOIN historical h ON i.target_ip = h.ip
`
rows, err := analyzer.Database.Conn.Query(chCtx, query)
@@ -653,7 +653,7 @@ func (analyzer *Analyzer) ScoopDNS(ctx context.Context, bars *tea.Program) error
SELECT tld, maxMerge(last_seen) as last_seen, minMerge(first_seen) as first_seen from udns
-- limiting the scope to just the domains in this import here
-- has significant performance benefits as opposed to doing it later in the query
RIGHT JOIN unique_tld USING tld
WHERE tld IN (SELECT tld FROM unique_tld)
GROUP BY tld
),
sussy_subdomains AS (
+32
View File
@@ -8,6 +8,7 @@ import (
"strconv"
"testing"
"github.com/activecm/rita/v5/database"
"github.com/activecm/rita/v5/modifier"
"github.com/activecm/rita/v5/util"
@@ -1212,6 +1213,37 @@ func (it *ValidDatasetTestSuite) TestThreatMixtape() {
`).Scan(&count)
require.NoError(t, err)
require.EqualValues(t, 0, count, "all dns entries should have a hash that contains only the fqdn")
// uconns with no connections that weren't part of a SNI connection
verifyHashInMixtape(t, it.db, "0006AF9EDD0596EE5C87E7AAD63DC2FF", false)
verifyHashInMixtape(t, it.db, "000DB22CFD645023D2FE454685A067E5", false)
verifyHashInMixtape(t, it.db, "00281FB4049C4613CA0F4307F7B96932", false)
// uconns with connections that weren't part of a SNI connection
verifyHashInMixtape(t, it.db, "001569CF7AA5791DE4C5DF3498816D34", true)
verifyHashInMixtape(t, it.db, "005BE49BF3FC40ECF0366A79408074B1", true)
}
func verifyHashInMixtape(t *testing.T, db *database.DB, hash string, shouldExist bool) {
t.Helper()
ctx := db.QueryParameters(clickhouse.Parameters{
"hash": hash,
})
var count uint64
err := db.Conn.QueryRow(ctx, `
SELECT count() FROM threat_mixtape
WHERE hash = unhex({hash:String})
`).Scan(&count)
require.NoError(t, err, "querying threat mixtape should not produce an error")
if shouldExist {
require.Positive(t, count, "hash %s should exist in threat mixtape", hash)
} else {
require.EqualValues(t, 0, count, "hash %s should not exist in threat mixtape", hash)
}
}
func (it *ValidDatasetTestSuite) TestBigOlHistogramTable() {