Skip to content

Commit 5197e79

Browse files
committed
feat: allow custom indexer
1 parent aa61469 commit 5197e79

3 files changed

Lines changed: 56 additions & 45 deletions

File tree

src/main.rs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,8 @@ use tokio_util::sync::CancellationToken;
33
use tracing::{info, level_filters::LevelFilter};
44
use tracing_subscriber::fmt::format::FmtSpan;
55

6-
use pulsebeam_server_foss::proto::signaling_server::SignalingServer;
76
use pulsebeam_server_foss::server::Server;
7+
use pulsebeam_server_foss::{manager::IndexManager, proto::signaling_server::SignalingServer};
88
use std::time::Duration;
99
use tonic::service::LayerExt;
1010
use tower_http::cors::{AllowOrigin, CorsLayer};
@@ -29,15 +29,15 @@ async fn main() -> anyhow::Result<()> {
2929
let token = CancellationToken::new();
3030
let (mut health_reporter, health_service) = tonic_health::server::health_reporter();
3131
health_reporter
32-
.set_serving::<SignalingServer<crate::Server>>()
32+
.set_serving::<SignalingServer<crate::Server<IndexManager>>>()
3333
.await;
3434
// https://github.com/hyperium/tonic/discussions/1784
3535
// switch to v1 after tooling supports it
3636
let reflector = tonic_reflection::server::Builder::configure()
3737
.register_encoded_file_descriptor_set(pulsebeam_server_foss::proto::FILE_DESCRIPTOR_SET)
3838
.register_encoded_file_descriptor_set(tonic_health::pb::FILE_DESCRIPTOR_SET)
3939
.build_v1()?;
40-
let server = Server::spawn(token, CONNECTION_CAPACITY);
40+
let server = Server::spawn_default(token, CONNECTION_CAPACITY);
4141
let grpc_server = tower::ServiceBuilder::new()
4242
.layer(cors)
4343
.layer(tonic_web::GrpcWebLayer::new())

src/manager.rs

Lines changed: 41 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,26 @@ pub enum ConnEvent {
9393
Removed(PeerInfo),
9494
}
9595

96+
pub trait Indexer: Send + Sync + 'static {
97+
fn select(&self, range: impl RangeBounds<PeerInfo>) -> Vec<(PeerInfo, PeerStats)>;
98+
fn select_one(&self, range: impl RangeBounds<PeerInfo>) -> Option<PeerInfo>;
99+
100+
fn select_group(&self, group_id: GroupId) -> Vec<(PeerInfo, PeerStats)> {
101+
let start = PeerInfo {
102+
group_id: group_id.clone(),
103+
peer_id: "".to_string(),
104+
conn_id: u32::MIN,
105+
};
106+
107+
let end = PeerInfo {
108+
group_id,
109+
peer_id: "~".to_string(),
110+
conn_id: u32::MAX,
111+
};
112+
self.select(start..=end)
113+
}
114+
}
115+
96116
#[derive(Clone)]
97117
pub struct IndexManager {
98118
state: Arc<RwLock<IndexManagerState>>,
@@ -108,6 +128,27 @@ impl Default for IndexManager {
108128
}
109129
}
110130

131+
impl Indexer for IndexManager {
132+
fn select(&self, range: impl RangeBounds<PeerInfo>) -> Vec<(PeerInfo, PeerStats)> {
133+
let state = self.state.read();
134+
let mut result = Vec::new();
135+
for (k, v) in state.index.range(range) {
136+
result.push((k.clone(), v.clone()))
137+
}
138+
result
139+
}
140+
141+
fn select_one(&self, range: impl RangeBounds<PeerInfo>) -> Option<PeerInfo> {
142+
let state = self.state.read();
143+
// pick the youngest connection
144+
let found = state
145+
.index
146+
.range(range)
147+
.max_by_key(|(_, p)| p.inserted_at)?;
148+
Some(found.0.clone())
149+
}
150+
}
151+
111152
impl IndexManager {
112153
pub async fn run_until_cancelled_owned(self, mut event_ch: mpsc::UnboundedReceiver<ConnEvent>) {
113154
tracing::info!("spawned index worker");
@@ -128,40 +169,6 @@ impl IndexManager {
128169
}
129170
}
130171
}
131-
132-
pub fn select(&self, range: impl RangeBounds<PeerInfo>) -> Vec<(PeerInfo, PeerStats)> {
133-
let state = self.state.read();
134-
let mut result = Vec::new();
135-
for (k, v) in state.index.range(range) {
136-
result.push((k.clone(), v.clone()))
137-
}
138-
result
139-
}
140-
141-
pub fn select_one(&self, range: impl RangeBounds<PeerInfo>) -> Option<PeerInfo> {
142-
let state = self.state.read();
143-
// pick the youngest connection
144-
let found = state
145-
.index
146-
.range(range)
147-
.max_by_key(|(_, p)| p.inserted_at)?;
148-
Some(found.0.clone())
149-
}
150-
151-
pub fn select_group(&self, group_id: GroupId) -> Vec<(PeerInfo, PeerStats)> {
152-
let start = PeerInfo {
153-
group_id: group_id.clone(),
154-
peer_id: "".to_string(),
155-
conn_id: u32::MIN,
156-
};
157-
158-
let end = PeerInfo {
159-
group_id,
160-
peer_id: "~".to_string(),
161-
conn_id: u32::MAX,
162-
};
163-
self.select(start..=end)
164-
}
165172
}
166173

167174
#[derive(Clone, Debug, Serialize)]

src/server.rs

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -9,21 +9,21 @@ use tokio_util::sync::CancellationToken;
99
use tracing::field::valuable;
1010
use valuable::Enumerable;
1111

12-
use crate::manager::{IndexManager, Manager};
12+
use crate::manager::{IndexManager, Indexer, Manager};
1313
const RESERVED_CONN_ID_DISCOVERY: u32 = 0;
1414
const RECV_STREAM_BUFFER: usize = 8;
1515
const KEEP_ALIVE_INTERVAL: Duration = Duration::from_secs(45);
1616

1717
#[derive(Clone)]
18-
pub struct Server {
18+
pub struct Server<I> {
1919
pub manager: Manager,
20-
pub index: IndexManager,
20+
pub index: I,
2121
}
2222

2323
pub type MessageStream = Pin<Box<dyn Stream<Item = proto::Message> + Send>>;
2424

25-
impl Server {
26-
pub fn spawn(token: CancellationToken, capacity: u64) -> Self {
25+
impl Server<IndexManager> {
26+
pub fn spawn_default(token: CancellationToken, capacity: u64) -> Self {
2727
let event_ch = mpsc::unbounded_channel();
2828
let manager = Manager::new(capacity, event_ch.0);
2929
let index = IndexManager::default();
@@ -35,7 +35,9 @@ impl Server {
3535
}
3636
Self { manager, index }
3737
}
38+
}
3839

40+
impl<I: Indexer> Server<I> {
3941
pub fn insert_recv_stream(&self, src: PeerInfo) -> MessageStream {
4042
let conn = self.manager.allocate(src);
4143
let payload_stream = ReceiverStream::new(conn);
@@ -55,7 +57,7 @@ impl Server {
5557
pub type RecvStream = Pin<Box<dyn Stream<Item = Result<proto::RecvResp, tonic::Status>> + Send>>;
5658

5759
#[tonic::async_trait]
58-
impl Signaling for Server {
60+
impl<I: Indexer> Signaling for Server<I> {
5961
async fn prepare(
6062
&self,
6163
_req: tonic::Request<proto::PrepareReq>,
@@ -185,6 +187,8 @@ impl Signaling for Server {
185187
mod test {
186188
use std::iter::zip;
187189

190+
use crate::manager::IndexManager;
191+
188192
use super::*;
189193
use proto::*;
190194

@@ -220,8 +224,8 @@ mod test {
220224
.await
221225
}
222226

223-
fn setup() -> (Server, PeerInfo, PeerInfo) {
224-
let s = Server::spawn(CancellationToken::new(), 65536);
227+
fn setup() -> (Server<IndexManager>, PeerInfo, PeerInfo) {
228+
let s = Server::spawn_default(CancellationToken::new(), 65536);
225229
let peer1 = PeerInfo {
226230
group_id: String::from("default"),
227231
peer_id: String::from("peer1"),

0 commit comments

Comments
 (0)