use std::collections::BTreeMap;
use std::sync::Arc;
use std::sync::mpsc::{Receiver, Sender};
use std::task::Poll;
use datafusion::prelude::{SessionConfig, SessionContext, col, lit};
use egui::{Frame, Margin, RichText};
use re_dataframe_ui::{ColumnBlueprint, default_display_name_for_column};
use re_log_types::{EntityPathPart, EntryId};
use re_protos::cloud::v1alpha1::{EntryKind, ScanSegmentTableResponse};
use re_redap_client::ConnectionRegistryHandle;
use re_sorbet::ColumnDescriptorRef;
use re_ui::alert::Alert;
use re_ui::{UiExt as _, icons};
use re_viewer_context::{
AsyncRuntimeHandle, EditRedapServerModalCommand, GlobalContext, ViewerContext,
};
use crate::context::Context;
use crate::entries::{Dataset, Entries, Entry, Table};
use crate::server_modal::{ServerModal, ServerModalMode};
pub struct Server {
origin: re_uri::Origin,
entries: Entries,
tables_session_ctx: Arc<SessionContext>,
connection_registry: re_redap_client::ConnectionRegistryHandle,
runtime: AsyncRuntimeHandle,
}
impl Server {
fn new(
connection_registry: re_redap_client::ConnectionRegistryHandle,
runtime: AsyncRuntimeHandle,
egui_ctx: &egui::Context,
origin: re_uri::Origin,
) -> Self {
let tables_session_ctx = Self::session_context();
let entries = Entries::new(
connection_registry.clone(),
&runtime,
egui_ctx,
origin.clone(),
tables_session_ctx.clone(),
);
Self {
origin,
entries,
tables_session_ctx,
connection_registry,
runtime,
}
}
fn session_context() -> Arc<SessionContext> {
let session_ctx = SessionContext::new_with_config(
SessionConfig::new()
.with_coalesce_batches(false),
);
Arc::new(session_ctx)
}
fn refresh_entries(&mut self, runtime: &AsyncRuntimeHandle, egui_ctx: &egui::Context) {
self.tables_session_ctx = Self::session_context();
self.entries = Entries::new(
self.connection_registry.clone(),
runtime,
egui_ctx,
self.origin.clone(),
self.tables_session_ctx.clone(),
);
}
#[inline]
pub fn origin(&self) -> &re_uri::Origin {
&self.origin
}
#[inline]
pub fn entries(&self) -> &Entries {
&self.entries
}
fn on_frame_start(&mut self) {
self.entries.on_frame_start();
}
fn find_entry(&self, entry_id: EntryId) -> Option<&Entry> {
self.entries.find_entry(entry_id)
}
fn title_ui(
&self,
title: String,
ctx: &Context<'_>,
ui: &mut egui::Ui,
content: impl FnOnce(&mut egui::Ui),
) {
Frame::new().inner_margin(Margin::same(16)).show(ui, |ui| {
ui.horizontal(|ui| {
ui.heading(RichText::new(title).strong());
if ui
.small_icon_button(&icons::RESET, "Refresh collection")
.clicked()
{
ctx.command_sender
.send(Command::RefreshCollection(self.origin.clone()))
.ok();
}
});
ui.add_space(12.0);
content(ui);
});
}
fn server_ui(&self, viewer_ctx: &ViewerContext<'_>, ctx: &Context<'_>, ui: &mut egui::Ui) {
if let Poll::Ready(Err(err)) = self.entries.state() {
self.title_ui(self.origin.host.to_string(), ctx, ui, |ui| {
if let Some(conn_err) = err.as_client_credentials_error() {
let message = if conn_err.is_missing_token() {
"This server requires authentication to access its data."
} else {
"The provided credentials are invalid for this server."
};
let edit_message = if conn_err.is_missing_token() {
"Add credentials"
} else {
"Edit credentials"
};
Alert::warning().show(ui, |ui| {
ui.vertical(|ui| {
ui.strong(message);
if ui
.link(RichText::new(edit_message).strong().underline())
.clicked()
{
ctx.command_sender
.send(Command::OpenEditServerModal(
EditRedapServerModalCommand {
origin: self.origin.clone(),
open_on_success: None,
title: None,
},
))
.ok();
}
});
});
} else {
ui.error_label(err.to_string());
}
});
return;
}
const ENTRY_LINK_COLUMN_NAME: &str = "link";
re_dataframe_ui::DataFusionTableWidget::new(self.tables_session_ctx.clone(), "__entries")
.title(self.origin.host.to_string())
.column_blueprint(|desc| {
let mut blueprint = ColumnBlueprint::default();
if let ColumnDescriptorRef::Component(component) = desc
&& component.component == "entry_kind"
{
blueprint = blueprint.variant_ui(re_component_ui::REDAP_ENTRY_KIND_VARIANT);
}
let column_sort_key = match desc.display_name().as_str() {
"name" => 0,
ENTRY_LINK_COLUMN_NAME => 1,
_ => 2,
};
blueprint = blueprint.sort_key(column_sort_key);
if desc.display_name().as_str() == ENTRY_LINK_COLUMN_NAME {
blueprint = blueprint.variant_ui(re_component_ui::REDAP_URI_BUTTON_VARIANT);
}
blueprint
})
.generate_entry_links(ENTRY_LINK_COLUMN_NAME, "id", self.origin.clone())
.prefilter(
col("entry_kind")
.in_list(
vec![lit(EntryKind::Table as i32), lit(EntryKind::Dataset as i32)],
false,
)
.and(col("name").not_eq(lit("__entries"))),
)
.show(viewer_ctx, &self.runtime, ui);
}
fn dataset_entry_ui(
&self,
viewer_ctx: &ViewerContext<'_>,
ui: &mut egui::Ui,
dataset: &Dataset,
) {
const RECORDING_LINK_COLUMN_NAME: &str = "recording link";
re_dataframe_ui::DataFusionTableWidget::new(
self.tables_session_ctx.clone(),
dataset.name(),
)
.title(dataset.name())
.url(re_uri::EntryUri::new(dataset.origin.clone(), dataset.id()).to_string())
.column_blueprint(|desc| {
let mut name = default_display_name_for_column(desc);
name = name
.strip_prefix("rerun_")
.map(|name| name.replace('_', " "))
.unwrap_or(name);
let default_visible = if desc.entity_path().is_some_and(|entity_path| {
entity_path.starts_with(&std::iter::once(EntityPathPart::properties()).collect())
}) {
true
} else {
matches!(
desc.display_name().as_str(),
RECORDING_LINK_COLUMN_NAME | ScanSegmentTableResponse::FIELD_SEGMENT_ID
)
};
let column_sort_key = match desc.display_name().as_str() {
ScanSegmentTableResponse::FIELD_SEGMENT_ID => 0,
RECORDING_LINK_COLUMN_NAME => 1,
_ => 2,
};
let mut blueprint = ColumnBlueprint::default()
.display_name(name)
.default_visibility(default_visible)
.sort_key(column_sort_key);
if desc.display_name().as_str() == RECORDING_LINK_COLUMN_NAME {
blueprint = blueprint.variant_ui(re_component_ui::REDAP_URI_BUTTON_VARIANT);
}
blueprint
})
.generate_segment_links(
RECORDING_LINK_COLUMN_NAME,
ScanSegmentTableResponse::FIELD_SEGMENT_ID,
self.origin.clone(),
dataset.id(),
)
.show(viewer_ctx, &self.runtime, ui);
}
fn table_entry_ui(&self, viewer_ctx: &ViewerContext<'_>, ui: &mut egui::Ui, table: &Table) {
re_dataframe_ui::DataFusionTableWidget::new(self.tables_session_ctx.clone(), table.name())
.title(table.name())
.url(re_uri::EntryUri::new(table.origin.clone(), table.id()).to_string())
.show(viewer_ctx, &self.runtime, ui);
}
}
pub struct RedapServers {
servers: BTreeMap<re_uri::Origin, Server>,
pending_servers: Vec<re_uri::Origin>,
command_sender: Sender<Command>,
command_receiver: Receiver<Command>,
server_modal_ui: ServerModal,
}
impl serde::Serialize for RedapServers {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
self.servers
.keys()
.collect::<Vec<_>>()
.serialize(serializer)
}
}
impl<'de> serde::Deserialize<'de> for RedapServers {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let origins = Vec::<re_uri::Origin>::deserialize(deserializer)?;
let mut servers = Self::default();
for origin in origins {
servers.pending_servers.push(origin);
}
Ok(servers)
}
}
impl Default for RedapServers {
fn default() -> Self {
let (command_sender, command_receiver) = std::sync::mpsc::channel();
Self {
servers: Default::default(),
pending_servers: Default::default(),
command_sender,
command_receiver,
server_modal_ui: Default::default(),
}
}
}
pub enum Command {
OpenAddServerModal,
OpenEditServerModal(EditRedapServerModalCommand),
AddServer {
origin: re_uri::Origin,
credentials: Option<re_redap_client::Credentials>,
on_add: Option<Box<dyn FnOnce() + Send>>,
},
RemoveServer(re_uri::Origin),
RefreshCollection(re_uri::Origin),
}
impl RedapServers {
pub fn is_empty(&self) -> bool {
self.servers.is_empty() && self.pending_servers.is_empty()
}
pub fn has_server(&self, origin: &re_uri::Origin) -> bool {
self.servers.contains_key(origin) || self.pending_servers.contains(origin)
}
pub fn add_server(&self, origin: re_uri::Origin) {
self.command_sender
.send(Command::AddServer {
origin,
credentials: None,
on_add: None,
})
.ok();
}
pub fn iter_servers(&self) -> impl Iterator<Item = &Server> {
self.servers.values()
}
pub fn is_authenticated(&self, origin: &re_uri::Origin) -> bool {
self.servers
.get(origin)
.and_then(|server| server.connection_registry.credentials(origin))
.is_some()
}
pub fn logout(&mut self) {
self.server_modal_ui.logout();
}
pub fn on_frame_start(
&mut self,
connection_registry: &ConnectionRegistryHandle,
runtime: &AsyncRuntimeHandle,
egui_ctx: &egui::Context,
) {
self.pending_servers.drain(..).for_each(|origin| {
self.command_sender
.send(Command::AddServer {
origin,
credentials: None,
on_add: None,
})
.ok();
});
while let Ok(command) = self.command_receiver.try_recv() {
self.handle_command(connection_registry, runtime, egui_ctx, command);
}
for server in self.servers.values_mut() {
server.on_frame_start();
}
}
fn handle_command(
&mut self,
connection_registry: &re_redap_client::ConnectionRegistryHandle,
runtime: &AsyncRuntimeHandle,
egui_ctx: &egui::Context,
command: Command,
) {
match command {
Command::OpenAddServerModal => {
self.server_modal_ui
.open(ServerModalMode::Add, connection_registry);
}
Command::OpenEditServerModal(origin) => {
self.server_modal_ui
.open(ServerModalMode::Edit(origin), connection_registry);
}
Command::AddServer {
origin,
credentials,
on_add,
} => {
if let Some(credentials) = credentials {
connection_registry.set_credentials(&origin, credentials);
}
if !self.servers.contains_key(&origin) {
self.servers.insert(
origin.clone(),
Server::new(
connection_registry.clone(),
runtime.clone(),
egui_ctx,
origin.clone(),
),
);
} else {
re_log::debug!(
"Tried to add pre-existing server at {:?}",
origin.to_string()
);
}
if let Some(on_add) = on_add {
on_add();
}
}
Command::RemoveServer(origin) => {
self.servers.remove(&origin);
connection_registry.remove_credentials(&origin);
}
Command::RefreshCollection(origin) => {
self.servers.entry(origin).and_modify(|server| {
server.refresh_entries(runtime, egui_ctx);
});
}
}
}
pub fn server_central_panel_ui(
&self,
viewer_ctx: &ViewerContext<'_>,
ui: &mut egui::Ui,
origin: &re_uri::Origin,
) {
if let Some(server) = self.servers.get(origin) {
self.with_ctx(|ctx| {
server.server_ui(viewer_ctx, ctx, ui);
});
} else {
viewer_ctx.revert_to_default_display_mode();
}
}
pub fn open_add_server_modal(&self) {
self.command_sender.send(Command::OpenAddServerModal).ok();
}
pub fn open_edit_server_modal(&self, command: EditRedapServerModalCommand) {
self.command_sender
.send(Command::OpenEditServerModal(command))
.ok();
}
pub fn entry_ui(
&self,
viewer_ctx: &ViewerContext<'_>,
ui: &mut egui::Ui,
active_entry: EntryId,
) {
for server in self.servers.values() {
if let Some(entry) = server.find_entry(active_entry) {
match entry.inner() {
Ok(crate::entries::EntryInner::Dataset(dataset)) => {
server.dataset_entry_ui(viewer_ctx, ui, dataset);
return;
}
Ok(crate::entries::EntryInner::Table(table)) => {
server.table_entry_ui(viewer_ctx, ui, table);
return;
}
Err(err) => {
Frame::new().inner_margin(16.0).show(ui, |ui| {
Alert::error().show_text(
ui,
format!("Error loading entry {}", entry.name()),
Some(err.to_string()),
);
});
}
}
}
}
}
pub fn modals_ui(&mut self, global_ctx: &GlobalContext<'_>, ui: &egui::Ui) {
let ctx = Context {
command_sender: &self.command_sender,
};
self.server_modal_ui.ui(global_ctx, &ctx, ui);
}
pub fn send_command(&self, command: Command) {
let result = self.command_sender.send(command);
if let Err(err) = result {
re_log::warn_once!("Failed to send command: {}", err);
}
}
#[inline]
fn with_ctx<R>(&self, func: impl FnOnce(&Context<'_>) -> R) -> R {
let ctx = Context {
command_sender: &self.command_sender,
};
func(&ctx)
}
}