Skip to content

Commit 7487ee0

Browse files
authored
Merge pull request #376 from firstJOASH/Search-Service
Search service
2 parents 847b6bb + 01218e3 commit 7487ee0

7 files changed

Lines changed: 217 additions & 23 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
@@ -559,8 +559,8 @@ impl Database {
559559
let limit = query.limit.unwrap_or(25).clamp(1, 100);
560560
let offset = query.offset.unwrap_or(0).max(0);
561561
let q = query.q.clone().unwrap_or_default();
562-
let q_pattern = format!("%{}%", q);
563562

563+
// Use full-text search when a query term is provided, fall back to no filter
564564
let rows = sqlx::query_as::<_, TradeSearchResult>(
565565
r#"
566566
WITH latest_trade_events AS (
@@ -574,10 +574,11 @@ impl Database {
574574
trade_base AS (
575575
SELECT
576576
(e.data->>'trade_id')::BIGINT AS trade_id,
577-
e.data->>'seller' AS seller,
578-
e.data->>'buyer' AS buyer,
579-
(e.data->>'amount')::BIGINT AS amount,
580-
e.timestamp AS created_at
577+
e.data->>'seller' AS seller,
578+
e.data->>'buyer' AS buyer,
579+
(e.data->>'amount')::BIGINT AS amount,
580+
e.timestamp AS created_at,
581+
e.search_vec
581582
FROM events e
582583
WHERE e.event_type = 'trade_created'
583584
)
@@ -591,18 +592,19 @@ impl Database {
591592
FROM trade_base tb
592593
JOIN latest_trade_events lte ON lte.trade_id = tb.trade_id
593594
WHERE
594-
($1 = '' OR tb.trade_id::TEXT ILIKE $2 OR tb.seller ILIKE $2 OR tb.buyer ILIKE $2)
595-
AND ($3::TEXT IS NULL OR lte.event_type = $3)
596-
AND ($4::TEXT IS NULL OR tb.seller = $4)
597-
AND ($5::TEXT IS NULL OR tb.buyer = $5)
598-
AND ($6::BIGINT IS NULL OR tb.amount >= $6)
599-
AND ($7::BIGINT IS NULL OR tb.amount <= $7)
600-
ORDER BY tb.created_at DESC
601-
LIMIT $8 OFFSET $9
595+
($1 = '' OR tb.search_vec @@ plainto_tsquery('english', $1))
596+
AND ($2::TEXT IS NULL OR lte.event_type = $2)
597+
AND ($3::TEXT IS NULL OR tb.seller = $3)
598+
AND ($4::TEXT IS NULL OR tb.buyer = $4)
599+
AND ($5::BIGINT IS NULL OR tb.amount >= $5)
600+
AND ($6::BIGINT IS NULL OR tb.amount <= $6)
601+
ORDER BY
602+
CASE WHEN $1 = '' THEN 0 ELSE ts_rank(tb.search_vec, plainto_tsquery('english', $1)) END DESC,
603+
tb.created_at DESC
604+
LIMIT $7 OFFSET $8
602605
"#,
603606
)
604607
.bind(q.as_str())
605-
.bind(q_pattern.as_str())
606608
.bind(query.status.as_deref())
607609
.bind(query.seller.as_deref())
608610
.bind(query.buyer.as_deref())
@@ -622,20 +624,19 @@ impl Database {
622624
) -> Result<Vec<DiscoveryResult>, AppError> {
623625
let limit = query.limit.unwrap_or(25).clamp(1, 100);
624626
let q = query.q.clone().unwrap_or_default();
625-
let q_pattern = format!("%{}%", q);
626627

627628
let rows = sqlx::query_as::<_, DiscoveryResult>(
628629
r#"
629630
WITH entities AS (
630-
SELECT data->>'seller' AS address, 'user' AS role, timestamp
631+
SELECT data->>'seller' AS address, 'user' AS role, timestamp, search_vec
631632
FROM events
632633
WHERE event_type = 'trade_created' AND data->>'seller' IS NOT NULL
633634
UNION ALL
634-
SELECT data->>'buyer' AS address, 'user' AS role, timestamp
635+
SELECT data->>'buyer' AS address, 'user' AS role, timestamp, search_vec
635636
FROM events
636637
WHERE event_type = 'trade_created' AND data->>'buyer' IS NOT NULL
637638
UNION ALL
638-
SELECT data->>'arbitrator' AS address, 'arbitrator' AS role, timestamp
639+
SELECT data->>'arbitrator' AS address, 'arbitrator' AS role, timestamp, search_vec
639640
FROM events
640641
WHERE event_type = 'arb_reg' AND data->>'arbitrator' IS NOT NULL
641642
)
@@ -646,15 +647,14 @@ impl Database {
646647
MAX(timestamp) AS last_seen
647648
FROM entities
648649
WHERE
649-
($1 = '' OR address ILIKE $2)
650-
AND ($3::TEXT IS NULL OR role = $3)
650+
($1 = '' OR search_vec @@ plainto_tsquery('english', $1))
651+
AND ($2::TEXT IS NULL OR role = $2)
651652
GROUP BY address, role
652653
ORDER BY seen_count DESC, last_seen DESC
653-
LIMIT $4
654+
LIMIT $3
654655
"#,
655656
)
656657
.bind(q.as_str())
657-
.bind(q_pattern.as_str())
658658
.bind(query.role.as_deref())
659659
.bind(limit)
660660
.fetch_all(&self.pool)
@@ -663,6 +663,34 @@ impl Database {
663663
Ok(rows)
664664
}
665665

666+
/// Full-text search over registered user profiles.
667+
pub async fn search_users(
668+
&self,
669+
q: &str,
670+
limit: i64,
671+
offset: i64,
672+
) -> Result<Vec<crate::models::UserProfile>, AppError> {
673+
let rows = sqlx::query_as::<_, crate::models::UserProfile>(
674+
r#"
675+
SELECT *
676+
FROM user_profiles
677+
WHERE $1 = '' OR search_vec @@ plainto_tsquery('english', $1)
678+
ORDER BY
679+
CASE WHEN $1 = '' THEN 0
680+
ELSE ts_rank(search_vec, plainto_tsquery('english', $1))
681+
END DESC,
682+
registered_at DESC
683+
LIMIT $2 OFFSET $3
684+
"#,
685+
)
686+
.bind(q)
687+
.bind(limit.clamp(1, 100))
688+
.bind(offset.max(0))
689+
.fetch_all(&self.pool)
690+
.await?;
691+
Ok(rows)
692+
}
693+
666694
pub async fn get_search_suggestions(
667695
&self,
668696
prefix: &str,

indexer/src/event_monitor.rs

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -347,6 +347,9 @@ impl EventMonitor {
347347
// Update user analytics for trade lifecycle events
348348
self.update_user_analytics(event).await;
349349

350+
// Update user analytics for trade lifecycle events
351+
self.update_user_analytics(event).await;
352+
350353
let mut queue = self.job_queue.lock().await;
351354
if let Err(e) = queue.enqueue(event_job).await {
352355
error!("Failed to enqueue event job for event {}: {}", event.id, e);
@@ -420,6 +423,70 @@ impl EventMonitor {
420423
}
421424
}
422425

426+
async fn update_user_analytics(&self, event: &Event) {
427+
let d = &event.data;
428+
let upsert = |address: &str, seller: bool, buyer: bool, amount: i64, completed: bool, disputed: bool, cancelled: bool| {
429+
let pool = self.database.pool().clone();
430+
let address = address.to_string();
431+
tokio::spawn(async move {
432+
let _ = sqlx::query(
433+
r#"
434+
INSERT INTO user_analytics
435+
(address, total_trades, trades_as_seller, trades_as_buyer,
436+
total_volume, completed_trades, disputed_trades, cancelled_trades, updated_at)
437+
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, NOW())
438+
ON CONFLICT (address) DO UPDATE SET
439+
total_trades = user_analytics.total_trades + EXCLUDED.total_trades,
440+
trades_as_seller = user_analytics.trades_as_seller + EXCLUDED.trades_as_seller,
441+
trades_as_buyer = user_analytics.trades_as_buyer + EXCLUDED.trades_as_buyer,
442+
total_volume = user_analytics.total_volume + EXCLUDED.total_volume,
443+
completed_trades = user_analytics.completed_trades + EXCLUDED.completed_trades,
444+
disputed_trades = user_analytics.disputed_trades + EXCLUDED.disputed_trades,
445+
cancelled_trades = user_analytics.cancelled_trades + EXCLUDED.cancelled_trades,
446+
updated_at = NOW()
447+
"#,
448+
)
449+
.bind(&address)
450+
.bind(if seller || buyer { 1i32 } else { 0 })
451+
.bind(if seller { 1i32 } else { 0 })
452+
.bind(if buyer { 1i32 } else { 0 })
453+
.bind(amount)
454+
.bind(if completed { 1i32 } else { 0 })
455+
.bind(if disputed { 1i32 } else { 0 })
456+
.bind(if cancelled { 1i32 } else { 0 })
457+
.execute(&pool)
458+
.await;
459+
});
460+
};
461+
462+
match event.event_type.as_str() {
463+
"trade_created" => {
464+
let seller = d.get("seller").and_then(|v| v.as_str()).unwrap_or_default();
465+
let buyer = d.get("buyer").and_then(|v| v.as_str()).unwrap_or_default();
466+
let amount = d.get("amount").and_then(|v| v.as_i64()).unwrap_or(0);
467+
upsert(seller, true, false, amount, false, false, false);
468+
upsert(buyer, false, true, amount, false, false, false);
469+
}
470+
"trade_confirmed" => {
471+
let seller = d.get("seller").and_then(|v| v.as_str()).unwrap_or_default();
472+
let buyer = d.get("buyer").and_then(|v| v.as_str()).unwrap_or_default();
473+
upsert(seller, false, false, 0, true, false, false);
474+
upsert(buyer, false, false, 0, true, false, false);
475+
}
476+
"trade_disputed" => {
477+
let seller = d.get("seller").and_then(|v| v.as_str()).unwrap_or_default();
478+
let buyer = d.get("buyer").and_then(|v| v.as_str()).unwrap_or_default();
479+
upsert(seller, false, false, 0, false, true, false);
480+
upsert(buyer, false, false, 0, false, true, false);
481+
}
482+
"trade_cancelled" => {
483+
let seller = d.get("seller").and_then(|v| v.as_str()).unwrap_or_default();
484+
upsert(seller, false, false, 0, false, false, true);
485+
}
486+
_ => {}
487+
}
488+
}
489+
423490
async fn get_latest_ledger(&self) -> Result<i64, AppError> {
424491
let url = format!("{}/ledgers?order=desc&limit=1", self.config.horizon_url);
425492
let response: HorizonResponse<Ledger> = self.client.get(&url).send().await?.json().await?;

indexer/src/handlers.rs

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

@@ -320,13 +320,15 @@ pub async fn global_search(
320320
let users = state.database.discover_entities(&user_query).await?;
321321
let arbitrators = state.database.discover_entities(&arb_query).await?;
322322
let suggestions = state.database.get_search_suggestions(&params.q, 10).await?;
323+
let profiles = state.database.search_users(&params.q, limit, 0).await?;
323324

324325
state.database.record_search(&params.q, "global").await?;
325326

326327
Ok(Json(GlobalSearchResponse {
327328
trades,
328329
users,
329330
arbitrators,
331+
profiles,
330332
suggestions,
331333
}))
332334
}

indexer/src/main.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -365,6 +365,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
365365
.route("/webhooks/stats", get(get_webhook_stats))
366366
// Users
367367
.route("/users", post(user_handlers::register_user))
368+
.route("/users/search", get(user_handlers::search_users))
368369
.route(
369370
"/users/:address",
370371
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
@@ -317,6 +317,7 @@ pub struct GlobalSearchResponse {
317317
pub trades: Vec<TradeSearchResult>,
318318
pub users: Vec<DiscoveryResult>,
319319
pub arbitrators: Vec<DiscoveryResult>,
320+
pub profiles: Vec<UserProfile>,
320321
pub suggestions: Vec<SearchSuggestion>,
321322
}
322323

@@ -630,6 +631,13 @@ pub struct SetPreferenceRequest {
630631
#[derive(Debug, Clone, Serialize, Deserialize)]
631632
pub struct SetVerificationRequest {
632633
pub status: String,
634+
}
635+
636+
#[derive(Debug, Clone, Serialize, Deserialize)]
637+
pub struct UserSearchQuery {
638+
pub q: Option<String>,
639+
pub limit: Option<i64>,
640+
pub offset: Option<i64>,
633641
#[derive(Debug, Clone, Serialize, Deserialize)]
634642
pub struct PushRegistrationRequest {
635643
pub device_token: String,

indexer/src/user_handlers.rs

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,7 @@
11
use axum::{
2+
extract::{Path, Query, State},
3+
http::StatusCode,
4+
response::Json,
25
extract::{Path, State},
36
http::StatusCode,
47
response::Json,
@@ -12,6 +15,7 @@ use crate::error::AppError;
1215
use crate::handlers::AppState;
1316
use crate::models::{
1417
RegisterUserRequest, SetPreferenceRequest, SetVerificationRequest, UpdateProfileRequest,
18+
UserAnalyticsRow, UserPreference, UserProfile, UserSearchQuery,
1519
UserAnalyticsRow, UserPreference, UserProfile,
1620
};
1721

@@ -181,3 +185,20 @@ pub async fn set_verification(
181185

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

0 commit comments

Comments
 (0)