Skip to content

Commit 181334b

Browse files
committed
feat: full-text search service with tsvector indexes, FTS queries, user search, and analytics
1 parent 987f19f commit 181334b

6 files changed

Lines changed: 148 additions & 27 deletions

File tree

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,67 @@
1+
-- Full-text search: add tsvector columns and GIN indexes for fast FTS
2+
3+
-- 1. Events: searchable on event_type, seller, buyer, trade_id from JSONB
4+
ALTER TABLE events
5+
ADD COLUMN IF NOT EXISTS search_vec tsvector;
6+
7+
UPDATE events
8+
SET search_vec = to_tsvector('english',
9+
coalesce(event_type, '') || ' ' ||
10+
coalesce(data->>'seller', '') || ' ' ||
11+
coalesce(data->>'buyer', '') || ' ' ||
12+
coalesce(data->>'trade_id', '') || ' ' ||
13+
coalesce(data->>'arbitrator', '')
14+
);
15+
16+
CREATE INDEX IF NOT EXISTS idx_events_search_vec ON events USING GIN (search_vec);
17+
18+
-- Keep search_vec in sync on insert/update
19+
CREATE OR REPLACE FUNCTION events_search_vec_update() RETURNS trigger AS $$
20+
BEGIN
21+
NEW.search_vec := to_tsvector('english',
22+
coalesce(NEW.event_type, '') || ' ' ||
23+
coalesce(NEW.data->>'seller', '') || ' ' ||
24+
coalesce(NEW.data->>'buyer', '') || ' ' ||
25+
coalesce(NEW.data->>'trade_id', '') || ' ' ||
26+
coalesce(NEW.data->>'arbitrator', '')
27+
);
28+
RETURN NEW;
29+
END;
30+
$$ LANGUAGE plpgsql;
31+
32+
DROP TRIGGER IF EXISTS trg_events_search_vec ON events;
33+
CREATE TRIGGER trg_events_search_vec
34+
BEFORE INSERT OR UPDATE ON events
35+
FOR EACH ROW EXECUTE FUNCTION events_search_vec_update();
36+
37+
-- 2. User profiles: searchable on address, username_hash, verification
38+
ALTER TABLE user_profiles
39+
ADD COLUMN IF NOT EXISTS search_vec tsvector;
40+
41+
UPDATE user_profiles
42+
SET search_vec = to_tsvector('english',
43+
coalesce(address, '') || ' ' ||
44+
coalesce(username_hash, '') || ' ' ||
45+
coalesce(verification, '')
46+
);
47+
48+
CREATE INDEX IF NOT EXISTS idx_user_profiles_search_vec ON user_profiles USING GIN (search_vec);
49+
50+
CREATE OR REPLACE FUNCTION user_profiles_search_vec_update() RETURNS trigger AS $$
51+
BEGIN
52+
NEW.search_vec := to_tsvector('english',
53+
coalesce(NEW.address, '') || ' ' ||
54+
coalesce(NEW.username_hash, '') || ' ' ||
55+
coalesce(NEW.verification, '')
56+
);
57+
RETURN NEW;
58+
END;
59+
$$ LANGUAGE plpgsql;
60+
61+
DROP TRIGGER IF EXISTS trg_user_profiles_search_vec ON user_profiles;
62+
CREATE TRIGGER trg_user_profiles_search_vec
63+
BEFORE INSERT OR UPDATE ON user_profiles
64+
FOR EACH ROW EXECUTE FUNCTION user_profiles_search_vec_update();
65+
66+
-- 3. GIN index on events.data for fast JSONB key lookups
67+
CREATE INDEX IF NOT EXISTS idx_events_data_gin ON events USING GIN (data);

indexer/src/database.rs

Lines changed: 50 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -469,8 +469,8 @@ impl Database {
469469
let limit = query.limit.unwrap_or(25).clamp(1, 100);
470470
let offset = query.offset.unwrap_or(0).max(0);
471471
let q = query.q.clone().unwrap_or_default();
472-
let q_pattern = format!("%{}%", q);
473472

473+
// Use full-text search when a query term is provided, fall back to no filter
474474
let rows = sqlx::query_as::<_, TradeSearchResult>(
475475
r#"
476476
WITH latest_trade_events AS (
@@ -484,10 +484,11 @@ impl Database {
484484
trade_base AS (
485485
SELECT
486486
(e.data->>'trade_id')::BIGINT AS trade_id,
487-
e.data->>'seller' AS seller,
488-
e.data->>'buyer' AS buyer,
489-
(e.data->>'amount')::BIGINT AS amount,
490-
e.timestamp AS created_at
487+
e.data->>'seller' AS seller,
488+
e.data->>'buyer' AS buyer,
489+
(e.data->>'amount')::BIGINT AS amount,
490+
e.timestamp AS created_at,
491+
e.search_vec
491492
FROM events e
492493
WHERE e.event_type = 'trade_created'
493494
)
@@ -501,18 +502,19 @@ impl Database {
501502
FROM trade_base tb
502503
JOIN latest_trade_events lte ON lte.trade_id = tb.trade_id
503504
WHERE
504-
($1 = '' OR tb.trade_id::TEXT ILIKE $2 OR tb.seller ILIKE $2 OR tb.buyer ILIKE $2)
505-
AND ($3::TEXT IS NULL OR lte.event_type = $3)
506-
AND ($4::TEXT IS NULL OR tb.seller = $4)
507-
AND ($5::TEXT IS NULL OR tb.buyer = $5)
508-
AND ($6::BIGINT IS NULL OR tb.amount >= $6)
509-
AND ($7::BIGINT IS NULL OR tb.amount <= $7)
510-
ORDER BY tb.created_at DESC
511-
LIMIT $8 OFFSET $9
505+
($1 = '' OR tb.search_vec @@ plainto_tsquery('english', $1))
506+
AND ($2::TEXT IS NULL OR lte.event_type = $2)
507+
AND ($3::TEXT IS NULL OR tb.seller = $3)
508+
AND ($4::TEXT IS NULL OR tb.buyer = $4)
509+
AND ($5::BIGINT IS NULL OR tb.amount >= $5)
510+
AND ($6::BIGINT IS NULL OR tb.amount <= $6)
511+
ORDER BY
512+
CASE WHEN $1 = '' THEN 0 ELSE ts_rank(tb.search_vec, plainto_tsquery('english', $1)) END DESC,
513+
tb.created_at DESC
514+
LIMIT $7 OFFSET $8
512515
"#,
513516
)
514517
.bind(q.as_str())
515-
.bind(q_pattern.as_str())
516518
.bind(query.status.as_deref())
517519
.bind(query.seller.as_deref())
518520
.bind(query.buyer.as_deref())
@@ -532,20 +534,19 @@ impl Database {
532534
) -> Result<Vec<DiscoveryResult>, AppError> {
533535
let limit = query.limit.unwrap_or(25).clamp(1, 100);
534536
let q = query.q.clone().unwrap_or_default();
535-
let q_pattern = format!("%{}%", q);
536537

537538
let rows = sqlx::query_as::<_, DiscoveryResult>(
538539
r#"
539540
WITH entities AS (
540-
SELECT data->>'seller' AS address, 'user' AS role, timestamp
541+
SELECT data->>'seller' AS address, 'user' AS role, timestamp, search_vec
541542
FROM events
542543
WHERE event_type = 'trade_created' AND data->>'seller' IS NOT NULL
543544
UNION ALL
544-
SELECT data->>'buyer' AS address, 'user' AS role, timestamp
545+
SELECT data->>'buyer' AS address, 'user' AS role, timestamp, search_vec
545546
FROM events
546547
WHERE event_type = 'trade_created' AND data->>'buyer' IS NOT NULL
547548
UNION ALL
548-
SELECT data->>'arbitrator' AS address, 'arbitrator' AS role, timestamp
549+
SELECT data->>'arbitrator' AS address, 'arbitrator' AS role, timestamp, search_vec
549550
FROM events
550551
WHERE event_type = 'arb_reg' AND data->>'arbitrator' IS NOT NULL
551552
)
@@ -556,15 +557,14 @@ impl Database {
556557
MAX(timestamp) AS last_seen
557558
FROM entities
558559
WHERE
559-
($1 = '' OR address ILIKE $2)
560-
AND ($3::TEXT IS NULL OR role = $3)
560+
($1 = '' OR search_vec @@ plainto_tsquery('english', $1))
561+
AND ($2::TEXT IS NULL OR role = $2)
561562
GROUP BY address, role
562563
ORDER BY seen_count DESC, last_seen DESC
563-
LIMIT $4
564+
LIMIT $3
564565
"#,
565566
)
566567
.bind(q.as_str())
567-
.bind(q_pattern.as_str())
568568
.bind(query.role.as_deref())
569569
.bind(limit)
570570
.fetch_all(&self.pool)
@@ -573,6 +573,34 @@ impl Database {
573573
Ok(rows)
574574
}
575575

576+
/// Full-text search over registered user profiles.
577+
pub async fn search_users(
578+
&self,
579+
q: &str,
580+
limit: i64,
581+
offset: i64,
582+
) -> Result<Vec<crate::models::UserProfile>, AppError> {
583+
let rows = sqlx::query_as::<_, crate::models::UserProfile>(
584+
r#"
585+
SELECT *
586+
FROM user_profiles
587+
WHERE $1 = '' OR search_vec @@ plainto_tsquery('english', $1)
588+
ORDER BY
589+
CASE WHEN $1 = '' THEN 0
590+
ELSE ts_rank(search_vec, plainto_tsquery('english', $1))
591+
END DESC,
592+
registered_at DESC
593+
LIMIT $2 OFFSET $3
594+
"#,
595+
)
596+
.bind(q)
597+
.bind(limit.clamp(1, 100))
598+
.bind(offset.max(0))
599+
.fetch_all(&self.pool)
600+
.await?;
601+
Ok(rows)
602+
}
603+
576604
pub async fn get_search_suggestions(
577605
&self,
578606
prefix: &str,

indexer/src/handlers.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ use crate::models::{
2222
AuditQuery, DiscoveryQuery, Event, EventQuery, EventStats, GlobalSearchQuery,
2323
GlobalSearchResponse, HistoryQuery, IndexerStatus, NewAuditLog, PagedResponse,
2424
PaginatedResponse, ReplayRequest, RetentionRequest, RetentionResponse, StatsResponse,
25-
SuggestionQuery, TradeSearchQuery, WebSocketMessage,
25+
SuggestionQuery, TradeSearchQuery, UserProfile, WebSocketMessage,
2626
};
2727
use crate::websocket::WebSocketManager;
2828

@@ -278,13 +278,15 @@ pub async fn global_search(
278278
let users = state.database.discover_entities(&user_query).await?;
279279
let arbitrators = state.database.discover_entities(&arb_query).await?;
280280
let suggestions = state.database.get_search_suggestions(&params.q, 10).await?;
281+
let profiles = state.database.search_users(&params.q, limit, 0).await?;
281282

282283
state.database.record_search(&params.q, "global").await?;
283284

284285
Ok(Json(GlobalSearchResponse {
285286
trades,
286287
users,
287288
arbitrators,
289+
profiles,
288290
suggestions,
289291
}))
290292
}

indexer/src/main.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -322,6 +322,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
322322
.route("/webhooks/stats", get(get_webhook_stats))
323323
// Users
324324
.route("/users", post(user_handlers::register_user))
325+
.route("/users/search", get(user_handlers::search_users))
325326
.route(
326327
"/users/:address",
327328
get(user_handlers::get_user).patch(user_handlers::update_user),

indexer/src/models.rs

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -283,6 +283,7 @@ pub struct GlobalSearchResponse {
283283
pub trades: Vec<TradeSearchResult>,
284284
pub users: Vec<DiscoveryResult>,
285285
pub arbitrators: Vec<DiscoveryResult>,
286+
pub profiles: Vec<UserProfile>,
286287
pub suggestions: Vec<SearchSuggestion>,
287288
}
288289

@@ -575,3 +576,10 @@ pub struct SetPreferenceRequest {
575576
pub struct SetVerificationRequest {
576577
pub status: String,
577578
}
579+
580+
#[derive(Debug, Clone, Serialize, Deserialize)]
581+
pub struct UserSearchQuery {
582+
pub q: Option<String>,
583+
pub limit: Option<i64>,
584+
pub offset: Option<i64>,
585+
}

indexer/src/user_handlers.rs

Lines changed: 19 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,7 @@
11
use axum::{
2-
extract::{Path, State},
2+
extract::{Path, Query, State},
33
http::StatusCode,
44
response::Json,
5-
routing::{get, patch, post, put},
6-
Router,
75
};
86
use chrono::Utc;
97
use serde_json::json;
@@ -12,7 +10,7 @@ use crate::error::AppError;
1210
use crate::handlers::AppState;
1311
use crate::models::{
1412
RegisterUserRequest, SetPreferenceRequest, SetVerificationRequest, UpdateProfileRequest,
15-
UserAnalyticsRow, UserPreference, UserProfile,
13+
UserAnalyticsRow, UserPreference, UserProfile, UserSearchQuery,
1614
};
1715

1816
/// POST /users — register a new user profile
@@ -181,3 +179,20 @@ pub async fn set_verification(
181179

182180
Ok(Json(json!({ "address": address, "verification": req.status })))
183181
}
182+
183+
/// GET /users/search?q=…&limit=25&offset=0 — full-text search over user profiles
184+
pub async fn search_users(
185+
State(state): State<AppState>,
186+
Query(params): Query<UserSearchQuery>,
187+
) -> Result<Json<Vec<UserProfile>>, AppError> {
188+
let q = params.q.unwrap_or_default();
189+
let limit = params.limit.unwrap_or(25);
190+
let offset = params.offset.unwrap_or(0);
191+
192+
if !q.is_empty() {
193+
state.database.record_search(&q, "users").await?;
194+
}
195+
196+
let results = state.database.search_users(&q, limit, offset).await?;
197+
Ok(Json(results))
198+
}

0 commit comments

Comments
 (0)