use std::sync::Arc;
use jammi_ai::session::InferenceSession;
use jammi_ai::wire::{columns_from_proto, parse_channel_id};
use jammi_ai::{LocalSession, Session};
use jammi_db::catalog::channel_repo::ChannelSpec;
use tonic::{Request, Response, Status};
use crate::grpc::proto::channel as pb;
use crate::grpc::proto::channel::channel_service_server::ChannelService;
use crate::grpc::wire::{map_engine_error, scoped, session_tenant};
pub struct ChannelServer {
session: Arc<InferenceSession>,
}
impl ChannelServer {
pub fn new(session: Arc<InferenceSession>) -> Self {
Self { session }
}
fn local(&self) -> Session {
Session::Local(LocalSession::new(Arc::clone(&self.session)))
}
}
#[tonic::async_trait]
impl ChannelService for ChannelServer {
async fn register_channel(
&self,
request: Request<pb::RegisterChannelRequest>,
) -> Result<Response<()>, Status> {
let tenant = session_tenant(&request);
let req = request.into_inner();
let id = parse_channel_id(&req.channel_id)?;
let columns = columns_from_proto(req.columns)?;
let spec = ChannelSpec {
id,
priority: req.priority,
columns,
};
let session = self.local();
scoped(&self.session, tenant, || session.register_channel(&spec))
.await
.map_err(map_engine_error)?;
Ok(Response::new(()))
}
async fn add_channel_columns(
&self,
request: Request<pb::AddChannelColumnsRequest>,
) -> Result<Response<()>, Status> {
let tenant = session_tenant(&request);
let req = request.into_inner();
let id = parse_channel_id(&req.channel_id)?;
let columns = columns_from_proto(req.columns)?;
let session = self.local();
scoped(&self.session, tenant, || {
session.add_channel_columns(&id, &columns)
})
.await
.map_err(map_engine_error)?;
Ok(Response::new(()))
}
}