use std::collections::BTreeMap;
use std::sync::Arc;
use std::task::Poll;
use datafusion::prelude::{SessionConfig, SessionContext, col, lit};
use datafusion::sql::TableReference;
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_quota_channel::send_crossbeam;
use re_redap_client::{
ClientCredentialsError, ConnectionRegistryHandle, CredentialSource, Credentials,
};
use re_sorbet::ColumnDescriptorRef;
use re_ui::alert::Alert;
use re_ui::{UiExt as _, icons};
use re_uri::DATASET_HIERARCHY_SEPARATOR;
use re_viewer_context::{
AppContext, AsyncRuntimeHandle, EditRedapServerModalCommand, StoreViewContext, ViewStates,
};
use crate::context::Context;
use crate::entries::{Dataset, Entries, Entry, Table};
use crate::server_modal::{LoginFlow, LoginFlowResult, 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()
{
send_crossbeam(
ctx.command_sender,
Command::RefreshCollection(self.origin.clone()),
)
.ok();
}
});
ui.add_space(12.0);
content(ui);
});
}
fn server_ui(
&self,
viewer_ctx: &StoreViewContext<'_>,
ctx: &Context<'_>,
ui: &mut egui::Ui,
inline_login_flow: &mut Option<(re_uri::Origin, Box<LoginFlow>)>,
view_states: &mut ViewStates,
) {
if let Poll::Ready(Err(err)) = self.entries.state() {
self.title_ui(self.origin.host.to_string(), ctx, ui, |ui| {
error_ui(viewer_ctx, ctx, ui, &self.origin, err, inline_login_flow);
});
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, view_states);
}
fn folder_ui(
&self,
viewer_ctx: &StoreViewContext<'_>,
ui: &mut egui::Ui,
origin: &re_uri::Origin,
path_prefix: &str,
) {
use re_viewer_context::{RedapEntryKind, Route, SystemCommand, SystemCommandSender as _};
let command_sender = viewer_ctx.command_sender().clone();
Frame::new().inner_margin(Margin::same(16)).show(ui, |ui| {
ui.horizontal(|ui| {
if ui
.small_icon_button(&icons::ARROW_UP, "Go to parent folder")
.on_hover_text("Go to parent folder")
.clicked()
{
let parent_route = if let Some((parent, _)) =
path_prefix.rsplit_once(DATASET_HIERARCHY_SEPARATOR)
{
Route::RedapEntry {
origin: origin.clone(),
kind: RedapEntryKind::Folder(parent.to_owned()),
}
} else {
Route::RedapServer(origin.clone())
};
if let Some(parent_item) = parent_route.item() {
command_sender.send_system(SystemCommand::set_selection(parent_item));
}
command_sender.send_system(SystemCommand::SetRoute(parent_route));
}
ui.horizontal_centered(|ui| {
ui.heading(RichText::new(path_prefix.to_owned()).strong());
});
});
ui.add_space(12.0);
match self.entries.state() {
Poll::Pending => {
ui.loading_indicator("Loading entries…");
}
Poll::Ready(Err(err)) => {
Alert::error().show_text(
ui,
format!("Error loading entries for folder {path_prefix:?}"),
Some(err.to_string()),
);
}
Poll::Ready(Ok(entries)) => {
crate::folder_card_ui::folder_cards_ui(
ui,
origin,
entries,
path_prefix,
&command_sender,
);
}
}
});
}
fn dataset_entry_ui(
&self,
viewer_ctx: &StoreViewContext<'_>,
ui: &mut egui::Ui,
dataset: &Dataset,
view_states: &mut ViewStates,
) {
const RECORDING_LINK_COLUMN_NAME: &str = "recording link";
re_dataframe_ui::DataFusionTableWidget::new(
self.tables_session_ctx.clone(),
TableReference::bare(dataset.name().to_string()),
)
.title(dataset.name().to_string())
.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, view_states);
}
fn table_entry_ui(
&self,
viewer_ctx: &StoreViewContext<'_>,
ui: &mut egui::Ui,
table: &Table,
view_states: &mut ViewStates,
) {
re_dataframe_ui::DataFusionTableWidget::new(
self.tables_session_ctx.clone(),
TableReference::bare(table.name().to_string()),
)
.title(table.name().to_string())
.url(re_uri::EntryUri::new(table.origin.clone(), table.id()).to_string())
.remote_table(re_uri::EntryUri::new(table.origin.clone(), table.id()))
.show(viewer_ctx, &self.runtime, ui, view_states);
}
}
fn error_ui(
viewer_ctx: &StoreViewContext<'_>,
ctx: &Context<'_>,
ui: &mut egui::Ui,
origin: &re_uri::Origin,
err: &re_redap_client::ApiError,
inline_login_flow: &mut Option<(re_uri::Origin, Box<LoginFlow>)>,
) {
if let Some(conn_err) = err.as_client_credentials_error() {
let message = match conn_err {
ClientCredentialsError::RefreshError { .. }
| ClientCredentialsError::UnauthenticatedMissingToken { .. } => {
"There was an error refreshing your credentials"
}
ClientCredentialsError::SessionExpired => "Your session has expired",
ClientCredentialsError::UnauthenticatedBadToken { credentials, .. } => {
match credentials.source {
CredentialSource::PerOrigin => "The credentials for this origin are invalid",
CredentialSource::Fallback => "The fallback credentials are invalid",
CredentialSource::EnvVar => {
"The credentials provided via environment variable REDAP_TOKEN are invalid"
}
}
}
ClientCredentialsError::HostMismatch(_) => "The token is not allowed for this server",
ClientCredentialsError::NotAuthorized => {
"This server requires authentication to access its data."
}
};
let show_login = match conn_err {
ClientCredentialsError::RefreshError(_)
| ClientCredentialsError::SessionExpired
| ClientCredentialsError::UnauthenticatedMissingToken(_)
| ClientCredentialsError::UnauthenticatedBadToken { .. }
| ClientCredentialsError::NotAuthorized => true,
ClientCredentialsError::HostMismatch(_) => false,
};
let has_active_login_flow = inline_login_flow.as_ref().is_some_and(|(o, _)| o == origin);
if show_login {
Alert::info().show(ui, |ui| {
ui.vertical(|ui| {
ui.strong(message);
if let Some(auth) = viewer_ctx.app_ctx.auth_context {
let identity = if let Some(org) = &auth.org_name {
format!("Logged in as {} ({})", auth.email, org)
} else {
format!("Logged in as {}", auth.email)
};
ui.weak(identity);
}
ui.add_space(8.0);
if has_active_login_flow {
ui.horizontal_centered(|ui| {
let cancel_button_height =
ui.text_style_height(&egui::TextStyle::Button) + 2.0 * 4.0;
ui.set_min_height(cancel_button_height);
ui.loading_indicator("Waiting for login");
ui.label("Waiting for login…");
ui.add_space(8.0);
if ui
.add(
re_ui::ReButton::new(("Cancel", &icons::CLOSE))
.small()
.primary(),
)
.clicked()
{
*inline_login_flow = None;
}
});
} else {
ui.horizontal(|ui| {
if let Some(auth) = viewer_ctx.app_ctx.auth_context {
if ui
.add(
re_ui::ReButton::new(format!("Continue as {}", auth.email))
.primary()
.small(),
)
.clicked()
{
send_crossbeam(
ctx.command_sender,
Command::UseStoredCredentials(origin.clone()),
)
.ok();
}
} else if viewer_ctx.app_ctx.login_enabled {
if ui
.add(re_ui::ReButton::new("Log in").primary().small())
.clicked()
{
match LoginFlow::open_and_start(
ui.ctx(),
viewer_ctx.app_ctx.login_signed_in_url,
) {
Ok(flow) => {
*inline_login_flow =
Some((origin.clone(), Box::new(flow)));
}
Err(err) => {
re_log::error!("Failed to start login: {err}");
}
}
}
}
if ui
.add(re_ui::ReButton::new("Edit connection").small())
.clicked()
{
send_crossbeam(
ctx.command_sender,
Command::OpenEditServerModal(EditRedapServerModalCommand {
origin: origin.clone(),
open_on_success: None,
title: None,
}),
)
.ok();
}
});
}
});
});
} else {
warning_with_edit_button(ctx, ui, origin, message, viewer_ctx.app_ctx.auth_context);
}
} else if matches!(
&err.kind,
re_redap_client::ApiErrorKind::InvalidServer | re_redap_client::ApiErrorKind::Connection
) {
warning_with_edit_button(ctx, ui, origin, &err.to_string(), None);
} else {
ui.error_label(err.to_string());
}
}
fn warning_with_edit_button(
ctx: &Context<'_>,
ui: &mut egui::Ui,
origin: &re_uri::Origin,
message: &str,
auth_context: Option<&re_viewer_context::AuthContext>,
) {
Alert::warning().show(ui, |ui| {
ui.vertical(|ui| {
ui.strong(message);
if let Some(auth) = auth_context {
let identity = if let Some(org) = &auth.org_name {
format!("Logged in as {} ({})", auth.email, org)
} else {
format!("Logged in as {}", auth.email)
};
ui.weak(identity);
}
ui.add_space(8.0);
if ui
.add(re_ui::ReButton::new("Edit connection").small())
.clicked()
{
send_crossbeam(
ctx.command_sender,
Command::OpenEditServerModal(EditRedapServerModalCommand {
origin: origin.clone(),
open_on_success: None,
title: None,
}),
)
.ok();
}
});
});
}
pub struct RedapServers {
servers: BTreeMap<re_uri::Origin, Server>,
pending_servers: Vec<re_uri::Origin>,
command_sender: crossbeam::channel::Sender<Command>,
command_receiver: crossbeam::channel::Receiver<Command>,
server_modal_ui: ServerModal,
inline_login_flow: Option<(re_uri::Origin, Box<LoginFlow>)>,
}
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) = create_channel(256);
Self {
servers: Default::default(),
pending_servers: Default::default(),
command_sender,
command_receiver,
server_modal_ui: Default::default(),
inline_login_flow: None,
}
}
}
fn create_channel<T>(
size: usize,
) -> (
crossbeam::channel::Sender<T>,
crossbeam::channel::Receiver<T>,
) {
cfg_if::cfg_if! {
if #[cfg(target_arch = "wasm32")] {
_ = size;
crossbeam::channel::unbounded() } else {
crossbeam::channel::bounded(size)
}
}
}
pub enum Command {
OpenAddServerModal,
OpenEditServerModal(EditRedapServerModalCommand),
AddServer {
origin: re_uri::Origin,
credentials: Option<re_redap_client::Credentials>,
on_add: Option<Box<dyn FnOnce() + Send>>,
},
RefreshCollection(re_uri::Origin),
UseStoredCredentials(re_uri::Origin),
}
impl std::fmt::Debug for Command {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::OpenAddServerModal => write!(f, "OpenAddServerModal"),
Self::OpenEditServerModal(cmd) => {
f.debug_tuple("OpenEditServerModal").field(cmd).finish()
}
Self::AddServer {
origin,
credentials,
on_add,
} => f
.debug_struct("AddServer")
.field("origin", origin)
.field("credentials", credentials)
.field("on_add", &on_add.as_ref().map(|_| "…"))
.finish(),
Self::RefreshCollection(origin) => {
f.debug_tuple("RefreshCollection").field(origin).finish()
}
Self::UseStoredCredentials(origin) => {
f.debug_tuple("UseStoredCredentials").field(origin).finish()
}
}
}
}
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 remove_server(
&mut self,
origin: &re_uri::Origin,
connection_registry: &re_redap_client::ConnectionRegistryHandle,
) {
self.servers.remove(origin);
connection_registry.remove_credentials(origin);
}
pub fn add_server(&self, origin: re_uri::Origin) {
send_crossbeam(
&self.command_sender,
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) -> Vec<re_uri::Origin> {
self.inline_login_flow = None;
self.server_modal_ui.logout();
let mut origins = Vec::new();
for server in self.servers.values() {
if matches!(
server.connection_registry.credentials(&server.origin),
Some(Credentials::Stored)
) {
origins.push(server.origin.clone());
server
.connection_registry
.remove_credentials(&server.origin);
send_crossbeam(
&self.command_sender,
Command::RefreshCollection(server.origin.clone()),
)
.ok();
}
}
origins
}
pub fn on_frame_start(
&mut self,
connection_registry: &ConnectionRegistryHandle,
runtime: &AsyncRuntimeHandle,
egui_ctx: &egui::Context,
login_enabled: bool,
) {
self.pending_servers.drain(..).for_each(|origin| {
send_crossbeam(
&self.command_sender,
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,
login_enabled,
);
}
if let Some((origin, flow)) = &mut self.inline_login_flow
&& let Some(result) = flow.poll()
{
let origin = origin.clone();
match result {
LoginFlowResult::Success => {
send_crossbeam(&self.command_sender, Command::UseStoredCredentials(origin))
.ok();
}
LoginFlowResult::Failure(err) => {
re_log::warn!("Login failed: {err}");
}
}
self.inline_login_flow = None;
}
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,
login_enabled: bool,
) {
match command {
Command::OpenAddServerModal => {
self.server_modal_ui
.open(ServerModalMode::Add, connection_registry, login_enabled);
}
Command::OpenEditServerModal(origin) => {
self.server_modal_ui.open(
ServerModalMode::Edit(origin),
connection_registry,
login_enabled,
);
}
Command::AddServer {
origin,
credentials,
on_add,
} => {
if let Some(credentials) = credentials {
connection_registry.set_credentials(&origin, credentials);
}
if self.servers.contains_key(&origin) {
re_log::debug!(
"Tried to add pre-existing server at {:?}",
origin.to_string()
);
} else {
self.servers.insert(
origin.clone(),
Server::new(
connection_registry.clone(),
runtime.clone(),
egui_ctx,
origin.clone(),
),
);
}
if let Some(on_add) = on_add {
on_add();
}
}
Command::RefreshCollection(origin) => {
self.servers.entry(origin).and_modify(|server| {
server.refresh_entries(runtime, egui_ctx);
});
}
Command::UseStoredCredentials(origin) => {
connection_registry.set_credentials(&origin, re_redap_client::Credentials::Stored);
send_crossbeam(&self.command_sender, Command::RefreshCollection(origin)).ok();
}
}
}
pub fn server_central_panel_ui(
&mut self,
viewer_ctx: &StoreViewContext<'_>,
ui: &mut egui::Ui,
origin: &re_uri::Origin,
view_states: &mut ViewStates,
) {
if let Some(server) = self.servers.get(origin) {
let ctx = Context {
command_sender: &self.command_sender,
};
server.server_ui(
viewer_ctx,
&ctx,
ui,
&mut self.inline_login_flow,
view_states,
);
} else {
viewer_ctx.revert_to_default_route();
}
}
pub fn folder_central_panel_ui(
&self,
viewer_ctx: &StoreViewContext<'_>,
ui: &mut egui::Ui,
origin: &re_uri::Origin,
path_prefix: &str,
) {
if let Some(server) = self.servers.get(origin) {
server.folder_ui(viewer_ctx, ui, origin, path_prefix);
} else {
viewer_ctx.revert_to_default_route();
}
}
pub fn open_add_server_modal(&self) {
send_crossbeam(&self.command_sender, Command::OpenAddServerModal).ok();
}
pub fn open_edit_server_modal(&self, command: EditRedapServerModalCommand) {
send_crossbeam(&self.command_sender, Command::OpenEditServerModal(command)).ok();
}
pub fn entry_ui(
&self,
ctx: &StoreViewContext<'_>,
ui: &mut egui::Ui,
active_entry: EntryId,
view_states: &mut ViewStates,
) {
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(ctx, ui, dataset, view_states);
return;
}
Ok(crate::entries::EntryInner::Table(table)) => {
server.table_entry_ui(ctx, ui, table, view_states);
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, app_ctx: &AppContext<'_>, ui: &egui::Ui) {
let ctx = Context {
command_sender: &self.command_sender,
};
self.server_modal_ui.ui(app_ctx, &ctx, ui);
}
pub fn send_command(&self, command: Command) {
let result = send_crossbeam(&self.command_sender, command);
if let Err(err) = result {
re_log::warn_once!("Failed to send command: {err}");
}
}
}