diff --git a/Cargo.lock b/Cargo.lock index 609fcded11deb..d0ffda774df71 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8076,7 +8076,6 @@ dependencies = [ "turborepo-task-id", "turborepo-turbo-json", "turborepo-types", - "turborepo-ui", "wax", "webbrowser", ] @@ -8096,7 +8095,6 @@ dependencies = [ "turborepo-signals", "turborepo-turbo-json", "turborepo-types", - "turborepo-ui", ] [[package]] @@ -8476,7 +8474,6 @@ version = "0.1.0" dependencies = [ "anyhow", "async-graphql", - "axum", "base64 0.22.1", "chrono", "clipboard-win", diff --git a/crates/turborepo-lib/src/cli/snapshots/turborepo_lib__cli__test__turbo_long_help.snap b/crates/turborepo-lib/src/cli/snapshots/turborepo_lib__cli__test__turbo_long_help.snap index 02df5e4ccf4ed..86a44c6216d96 100644 --- a/crates/turborepo-lib/src/cli/snapshots/turborepo_lib__cli__test__turbo_long_help.snap +++ b/crates/turborepo-lib/src/cli/snapshots/turborepo_lib__cli__test__turbo_long_help.snap @@ -55,7 +55,6 @@ Options: - tui: Use the terminal user interface - stream: Use the standard output stream - stream-with-experimental-timestamps: Use the standard output stream with timestamps. Note: This feature is experimental and may change or be removed at any time - - web: Use the web user interface. Note: This feature is undocumented, experimental, and not meant to be used. It may change or be removed at any time --login Override the login endpoint diff --git a/crates/turborepo-lib/src/cli/snapshots/turborepo_lib__cli__test__turbo_short_help.snap b/crates/turborepo-lib/src/cli/snapshots/turborepo_lib__cli__test__turbo_short_help.snap index d341e3f90939f..2e0b682be61c1 100644 --- a/crates/turborepo-lib/src/cli/snapshots/turborepo_lib__cli__test__turbo_short_help.snap +++ b/crates/turborepo-lib/src/cli/snapshots/turborepo_lib__cli__test__turbo_short_help.snap @@ -42,7 +42,7 @@ Options: --heap Specify a file to save a pprof heap profile --ui - Specify whether to use the streaming UI or TUI [possible values: tui, stream, stream-with-experimental-timestamps, web] + Specify whether to use the streaming UI or TUI [possible values: tui, stream, stream-with-experimental-timestamps] --login Override the login endpoint --no-color diff --git a/crates/turborepo-lib/src/commands/run.rs b/crates/turborepo-lib/src/commands/run.rs index b08c8a9f32bf0..f3cd392762551 100644 --- a/crates/turborepo-lib/src/commands/run.rs +++ b/crates/turborepo-lib/src/commands/run.rs @@ -169,7 +169,6 @@ pub async fn run( analytics_handle.close_with_timeout().await; } - // We only stop if it's the TUI, for the web UI we don't need to stop if let Some(UISender::Tui(sender)) = sender { sender.stop().await; } diff --git a/crates/turborepo-lib/src/lib.rs b/crates/turborepo-lib/src/lib.rs index e11c76e297f42..55589fcd2e49f 100644 --- a/crates/turborepo-lib/src/lib.rs +++ b/crates/turborepo-lib/src/lib.rs @@ -62,9 +62,8 @@ pub fn get_version() -> &'static str { /// Main entry point for the turborepo CLI. /// /// `query_server` provides the GraphQL query execution layer. When `None`, -/// the `turbo query` command returns an error and the Web UI mode falls -/// back silently. Pass `Some(...)` with a [`QueryServer`] implementation -/// to enable the full query subsystem. +/// the `turbo query` command returns an error. Pass `Some(...)` with a +/// [`QueryServer`] implementation to enable the full query subsystem. pub fn main( query_server: Option>, ) -> Result { diff --git a/crates/turborepo-lib/src/run/mod.rs b/crates/turborepo-lib/src/run/mod.rs index 5f3f38c6fc03f..30a8e599398c3 100644 --- a/crates/turborepo-lib/src/run/mod.rs +++ b/crates/turborepo-lib/src/run/mod.rs @@ -6,7 +6,6 @@ pub(crate) mod package_discovery; pub(crate) mod scope; pub mod task_access; pub(crate) mod task_filter; -mod ui; pub mod watch; use std::{ @@ -44,7 +43,7 @@ use turborepo_task_hash::{ }; use turborepo_telemetry::events::generic::GenericEventBuilder; use turborepo_types::{EnvMode, UIMode}; -use turborepo_ui::{sender::UISender, tui, tui::TuiSender, wui::sender::WebUISender, ColorConfig}; +use turborepo_ui::{sender::UISender, tui, tui::TuiSender, ColorConfig}; pub use crate::run::error::Error; use crate::{ @@ -106,7 +105,6 @@ pub struct Run { type UIResult = Result>)>, Error>; -type WuiResult = UIResult; type TuiResult = UIResult; #[derive(Debug, Clone, Copy)] @@ -525,22 +523,8 @@ impl Run { .start_terminal_ui() .map(|res| res.map(|(sender, handle)| (UISender::Tui(sender), handle))), UIMode::Stream | UIMode::StreamWithTimestamps => Ok(None), - UIMode::Web => self - .start_web_ui() - .map(|res| res.map(|(sender, handle)| (UISender::Wui(sender), handle))), } } - fn start_web_ui(self: &Arc) -> WuiResult { - let Some(query_server) = self.query_server.clone() else { - tracing::warn!("Web UI requires a query server implementation"); - return Ok(None); - }; - let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); - - let handle = tokio::spawn(ui::start_web_ui_server(rx, self.clone(), query_server)); - - Ok(Some((WebUISender { tx }, handle))) - } #[allow(clippy::type_complexity)] fn start_terminal_ui(&self) -> TuiResult { diff --git a/crates/turborepo-lib/src/run/ui.rs b/crates/turborepo-lib/src/run/ui.rs deleted file mode 100644 index da3316b209660..0000000000000 --- a/crates/turborepo-lib/src/run/ui.rs +++ /dev/null @@ -1,27 +0,0 @@ -use std::sync::Arc; - -use turborepo_query_api::QueryServer; -use turborepo_ui::wui::{event::WebUIEvent, query::SharedState}; - -use crate::run::Run; - -pub async fn start_web_ui_server( - rx: tokio::sync::mpsc::UnboundedReceiver, - run: Arc, - query_server: Arc, -) -> Result<(), turborepo_ui::Error> { - let state = SharedState::default(); - let subscriber = turborepo_ui::wui::subscriber::Subscriber::new(rx); - tokio::spawn(subscriber.watch(state.clone())); - - let run: Arc = run; - query_server - .run_web_ui_server(state, run) - .await - .map_err(|e| { - let wui_err = turborepo_ui::wui::Error::Server(std::io::Error::other(e)); - turborepo_ui::Error::Wui(wui_err) - })?; - - Ok(()) -} diff --git a/crates/turborepo-query-api/Cargo.toml b/crates/turborepo-query-api/Cargo.toml index b33e107faa59f..5f5b50bac3894 100644 --- a/crates/turborepo-query-api/Cargo.toml +++ b/crates/turborepo-query-api/Cargo.toml @@ -16,7 +16,6 @@ turborepo-scope = { workspace = true } turborepo-signals = { workspace = true } turborepo-turbo-json = { workspace = true } turborepo-types = { workspace = true } -turborepo-ui = { workspace = true } [lints] workspace = true diff --git a/crates/turborepo-query-api/src/lib.rs b/crates/turborepo-query-api/src/lib.rs index ec71575542af9..592cc12abad92 100644 --- a/crates/turborepo-query-api/src/lib.rs +++ b/crates/turborepo-query-api/src/lib.rs @@ -91,8 +91,6 @@ pub enum Error { #[error(transparent)] #[diagnostic(transparent)] Path(#[from] turbopath::PathError), - #[error(transparent)] - UI(#[from] turborepo_ui::Error), #[error("Failed to calculate affected packages: {0}")] AffectedPackages(#[from] AffectedPackagesError), #[error(transparent)] @@ -146,15 +144,6 @@ pub trait QueryServer: Send + Sync { run: Arc, signal: turborepo_signals::SignalHandler, ) -> Pin> + Send + '_>>; - - /// Start the Web UI server that serves the TUI-integrated query interface. - /// - /// The shared state is used to stream build events to the UI. - fn run_web_ui_server( - &self, - state: turborepo_ui::wui::query::SharedState, - run: Arc, - ) -> Pin> + Send + '_>>; } // Compile-time assertions that both traits remain object-safe. diff --git a/crates/turborepo-query/Cargo.toml b/crates/turborepo-query/Cargo.toml index 64410bfea9d76..706c5c372949f 100644 --- a/crates/turborepo-query/Cargo.toml +++ b/crates/turborepo-query/Cargo.toml @@ -37,7 +37,6 @@ turborepo-signals = { workspace = true } turborepo-task-id = { workspace = true } turborepo-turbo-json = { path = "../turborepo-turbo-json" } turborepo-types = { workspace = true } -turborepo-ui = { workspace = true } wax = { workspace = true } webbrowser = { workspace = true } diff --git a/crates/turborepo-query/src/lib.rs b/crates/turborepo-query/src/lib.rs index a22be32faa280..f640e99f80119 100644 --- a/crates/turborepo-query/src/lib.rs +++ b/crates/turborepo-query/src/lib.rs @@ -72,11 +72,6 @@ impl From for Error { Error::Api(e.into()) } } -impl From for Error { - fn from(e: turborepo_ui::Error) -> Self { - Error::Api(e.into()) - } -} impl From for Error { fn from(e: AffectedPackagesError) -> Self { Error::Api(e.into()) @@ -812,7 +807,7 @@ pub async fn run_query_server(run: Arc, signal: SignalHandler) -> println!("Shutting down GraphQL server"); return Ok(()); } - result = server::run_server(None, run) => { + result = server::run_server(run) => { result?; } } diff --git a/crates/turborepo-query/src/server.rs b/crates/turborepo-query/src/server.rs index 7a2659a647518..0111d890633a6 100644 --- a/crates/turborepo-query/src/server.rs +++ b/crates/turborepo-query/src/server.rs @@ -1,43 +1,27 @@ use std::sync::Arc; -use async_graphql::{EmptyMutation, EmptySubscription, MergedObject, Schema}; +use async_graphql::{EmptyMutation, EmptySubscription, Schema}; use async_graphql_axum::GraphQL; use axum::{http::Method, routing::get, Router}; use tokio::net::TcpListener; use tower_http::cors::{Any, CorsLayer}; -use turborepo_ui::wui::query::SharedState; use crate::{graphiql, QueryRun, RepositoryQuery}; -#[derive(MergedObject)] -struct Query(turborepo_ui::wui::RunQuery, RepositoryQuery); - -pub async fn run_server( - state: Option, - run: Arc, -) -> Result<(), turborepo_ui::Error> { +pub async fn run_server(run: Arc) -> std::io::Result<()> { let cors = CorsLayer::new() .allow_methods([Method::GET, Method::POST]) .allow_headers(Any) .allow_origin(Any); - let web_ui_query = turborepo_ui::wui::RunQuery::new(state.clone()); let turbo_query = RepositoryQuery::new(run); - let combined_query = Query(web_ui_query, turbo_query); - let schema = Schema::new(combined_query, EmptyMutation, EmptySubscription); + let schema = Schema::new(turbo_query, EmptyMutation, EmptySubscription); let app = Router::new() .route("/", get(graphiql).post_service(GraphQL::new(schema))) .layer(cors); - axum::serve( - TcpListener::bind("127.0.0.1:8000") - .await - .map_err(turborepo_ui::wui::Error::Server)?, - app, - ) - .await - .map_err(turborepo_ui::wui::Error::Server)?; + axum::serve(TcpListener::bind("127.0.0.1:8000").await?, app).await?; Ok(()) } diff --git a/crates/turborepo-types/src/lib.rs b/crates/turborepo-types/src/lib.rs index 92a0dc5607f9d..ff491f9f0ff40 100644 --- a/crates/turborepo-types/src/lib.rs +++ b/crates/turborepo-types/src/lib.rs @@ -185,12 +185,6 @@ pub enum UIMode { #[schemars(rename = "stream-with-experimental-timestamps")] #[value(name = "stream-with-experimental-timestamps")] StreamWithTimestamps, - /// Use the web user interface. - /// Note: This feature is undocumented, experimental, and not meant to be - /// used. It may change or be removed at any time. - #[schemars(skip)] - #[ts(skip)] - Web, } impl fmt::Display for UIMode { @@ -199,7 +193,6 @@ impl fmt::Display for UIMode { UIMode::Tui => write!(f, "tui"), UIMode::Stream => write!(f, "stream"), UIMode::StreamWithTimestamps => write!(f, "stream-with-experimental-timestamps"), - UIMode::Web => write!(f, "web"), } } } @@ -210,9 +203,9 @@ impl UIMode { } /// Returns true if the UI mode has a sender, - /// i.e. web or tui but not stream + /// i.e. tui but not stream pub fn has_sender(&self) -> bool { - matches!(self, Self::Tui | Self::Web) + matches!(self, Self::Tui) } /// Returns true if this UI mode should include timestamps in the prefix diff --git a/crates/turborepo-ui/Cargo.toml b/crates/turborepo-ui/Cargo.toml index c432928619c62..15b41360c6723 100644 --- a/crates/turborepo-ui/Cargo.toml +++ b/crates/turborepo-ui/Cargo.toml @@ -17,7 +17,6 @@ workspace = true [dependencies] async-graphql = { workspace = true } -axum = { workspace = true, features = ["ws"] } base64 = "0.22" chrono = { workspace = true } console = { workspace = true } diff --git a/crates/turborepo-ui/README.md b/crates/turborepo-ui/README.md index fb009f9c07ed9..f828dacea675e 100644 --- a/crates/turborepo-ui/README.md +++ b/crates/turborepo-ui/README.md @@ -12,8 +12,7 @@ turborepo-ui ├── PrefixedUI - Output with task prefixes ├── ColorSelector - Assign colors to concurrent tasks ├── LogWriter - Task log handling - ├── tui/ - Interactive terminal UI (ratatui) - └── wui/ - Web UI server + └── tui/ - Interactive terminal UI (ratatui) ``` Key components: diff --git a/crates/turborepo-ui/src/lib.rs b/crates/turborepo-ui/src/lib.rs index f297ff753b28e..ae1fb98750712 100644 --- a/crates/turborepo-ui/src/lib.rs +++ b/crates/turborepo-ui/src/lib.rs @@ -10,7 +10,6 @@ pub mod sender; mod terminal_sink; pub mod tui; mod tui_sink; -pub mod wui; use std::{borrow::Cow, env, f64::consts::PI, io::IsTerminal, sync::LazyLock, time::Duration}; @@ -50,8 +49,6 @@ pub use crate::{ pub enum Error { #[error(transparent)] Tui(#[from] tui::Error), - #[error(transparent)] - Wui(#[from] wui::Error), #[error("Cannot read logs: {0}")] CannotReadLogs(#[source] std::io::Error), #[error("Cannot write logs: {0}")] diff --git a/crates/turborepo-ui/src/sender.rs b/crates/turborepo-ui/src/sender.rs index d17cbf2de1a21..3fba6412ce92e 100644 --- a/crates/turborepo-ui/src/sender.rs +++ b/crates/turborepo-ui/src/sender.rs @@ -3,55 +3,47 @@ use std::sync::{Arc, Mutex}; use crate::{ tui, tui::event::{CacheResult, OutputLogs, PaneSize, TaskResult}, - wui::sender, }; -/// Enum to abstract over sending events to either the Tui or the Web UI +/// Enum to abstract over sending events to the TUI. #[derive(Debug, Clone)] pub enum UISender { Tui(tui::TuiSender), - Wui(sender::WebUISender), } impl UISender { pub fn start_task(&self, task: String, output_logs: OutputLogs) { match self { UISender::Tui(sender) => sender.start_task(task, output_logs), - UISender::Wui(sender) => sender.start_task(task, output_logs), } } pub fn restart_tasks(&self, tasks: Vec) -> Result<(), crate::Error> { match self { UISender::Tui(sender) => sender.restart_tasks(tasks), - UISender::Wui(sender) => sender.restart_tasks(tasks), } } pub fn end_task(&self, task: String, result: TaskResult) { match self { UISender::Tui(sender) => sender.end_task(task, result), - UISender::Wui(sender) => sender.end_task(task, result), } } pub fn status(&self, task: String, status: String, result: CacheResult) { match self { UISender::Tui(sender) => sender.status(task, status, result), - UISender::Wui(sender) => sender.status(task, status, result), } } fn set_stdin(&self, task: String, stdin: Box) { match self { UISender::Tui(sender) => sender.set_stdin(task, stdin), - UISender::Wui(sender) => sender.set_stdin(task, stdin), } } pub fn output(&self, task: String, output: Vec) -> Result<(), crate::Error> { match self { UISender::Tui(sender) => sender.output(task, output), - UISender::Wui(sender) => sender.output(task, output), } } @@ -59,27 +51,22 @@ impl UISender { pub fn task(&self, task: String) -> TaskSender { match self { UISender::Tui(sender) => sender.task(task), - UISender::Wui(sender) => sender.task(task), } } pub async fn stop(&self) { match self { UISender::Tui(sender) => sender.stop().await, - UISender::Wui(sender) => sender.stop(), } } pub fn update_tasks(&self, tasks: Vec) -> Result<(), crate::Error> { match self { UISender::Tui(sender) => sender.update_tasks(tasks), - UISender::Wui(sender) => sender.update_tasks(tasks), } } pub async fn pane_size(&self) -> Option { match self { UISender::Tui(sender) => sender.pane_size().await, - // Not applicable to the web UI - UISender::Wui(_) => None, } } } diff --git a/crates/turborepo-ui/src/wui/event.rs b/crates/turborepo-ui/src/wui/event.rs deleted file mode 100644 index df9cb12eb2c53..0000000000000 --- a/crates/turborepo-ui/src/wui/event.rs +++ /dev/null @@ -1,34 +0,0 @@ -use serde::Serialize; - -use crate::tui::event::{CacheResult, OutputLogs, TaskResult}; - -/// Specific events that the GraphQL server can send to the client, -/// not all the `Event` types from the TUI. -#[derive(Debug, Clone, Serialize)] -#[serde(tag = "type", content = "payload")] -pub enum WebUIEvent { - StartTask { - task: String, - output_logs: OutputLogs, - }, - TaskOutput { - task: String, - output: Vec, - }, - EndTask { - task: String, - result: TaskResult, - }, - CacheStatus { - task: String, - message: String, - result: CacheResult, - }, - UpdateTasks { - tasks: Vec, - }, - RestartTasks { - tasks: Vec, - }, - Stop, -} diff --git a/crates/turborepo-ui/src/wui/mod.rs b/crates/turborepo-ui/src/wui/mod.rs deleted file mode 100644 index e3c7a5a40ac18..0000000000000 --- a/crates/turborepo-ui/src/wui/mod.rs +++ /dev/null @@ -1,25 +0,0 @@ -//! Web UI for Turborepo. Creates a WebSocket server that can be subscribed to -//! by a web client to display the status of tasks. - -pub mod event; -pub mod query; -pub mod sender; -pub mod subscriber; - -use event::WebUIEvent; -pub use query::RunQuery; -use thiserror::Error; - -#[derive(Debug, Error)] -pub enum Error { - #[error("Failed to start server.")] - Server(#[from] std::io::Error), - #[error("Failed to start websocket server: {0}")] - WebSocket(#[source] axum::Error), - #[error("Failed to serialize message: {0}")] - Serde(#[from] serde_json::Error), - #[error("Failed to send message.")] - Send(#[from] axum::Error), - #[error("Failed to send message through channel.")] - Broadcast(#[from] tokio::sync::mpsc::error::SendError), -} diff --git a/crates/turborepo-ui/src/wui/query.rs b/crates/turborepo-ui/src/wui/query.rs deleted file mode 100644 index 1d862dec83344..0000000000000 --- a/crates/turborepo-ui/src/wui/query.rs +++ /dev/null @@ -1,62 +0,0 @@ -use std::sync::Arc; - -use async_graphql::{Object, SimpleObject}; -use serde::Serialize; -use tokio::sync::Mutex; - -use crate::wui::subscriber::{TaskState, WebUIState}; - -#[derive(Debug, Clone, Serialize, SimpleObject)] -struct RunTask { - name: String, - state: TaskState, -} - -struct CurrentRun<'a> { - state: &'a SharedState, -} - -#[Object] -impl CurrentRun<'_> { - async fn tasks(&self) -> Vec { - self.state - .lock() - .await - .tasks() - .iter() - .map(|(task, state)| RunTask { - name: task.clone(), - state: state.clone(), - }) - .collect() - } -} - -/// We keep the state in a `Arc>>` so both `Subscriber` and -/// `Query` can access it, with `Subscriber` mutating it and `Query` only -/// reading it. -pub type SharedState = Arc>; - -/// The query for actively running tasks. -/// -/// (As opposed to the query for general repository state `RepositoryQuery` -/// in `turborepo_lib::query`) -/// This is `None` when we're not actually running a task (e.g. `turbo query`) -pub struct RunQuery { - state: Option, -} - -impl RunQuery { - pub fn new(state: Option) -> Self { - Self { state } - } -} - -#[Object] -impl RunQuery { - async fn current_run(&self) -> Option> { - Some(CurrentRun { - state: self.state.as_ref()?, - }) - } -} diff --git a/crates/turborepo-ui/src/wui/sender.rs b/crates/turborepo-ui/src/wui/sender.rs deleted file mode 100644 index cc1edce79843f..0000000000000 --- a/crates/turborepo-ui/src/wui/sender.rs +++ /dev/null @@ -1,78 +0,0 @@ -use std::io::Write; - -use tracing::warn; - -use crate::{ - sender::{TaskSender, UISender}, - tui::event::{CacheResult, OutputLogs, TaskResult}, - wui::{Error, event::WebUIEvent}, -}; - -#[derive(Debug, Clone)] -pub struct WebUISender { - pub tx: tokio::sync::mpsc::UnboundedSender, -} - -impl WebUISender { - pub fn new(tx: tokio::sync::mpsc::UnboundedSender) -> Self { - Self { tx } - } - pub fn start_task(&self, task: String, output_logs: OutputLogs) { - self.tx - .send(WebUIEvent::StartTask { task, output_logs }) - .ok(); - } - - pub fn restart_tasks(&self, tasks: Vec) -> Result<(), crate::Error> { - self.tx - .send(WebUIEvent::RestartTasks { tasks }) - .map_err(Error::Broadcast)?; - Ok(()) - } - - pub fn end_task(&self, task: String, result: TaskResult) { - self.tx.send(WebUIEvent::EndTask { task, result }).ok(); - } - - pub fn status(&self, task: String, message: String, result: CacheResult) { - self.tx - .send(WebUIEvent::CacheStatus { - task, - message, - result, - }) - .ok(); - } - - pub fn set_stdin(&self, _: String, _: Box) { - warn!("stdin is not supported (yet) in web ui"); - } - - pub fn task(&self, task: String) -> TaskSender { - TaskSender { - name: task, - handle: UISender::Wui(self.clone()), - logs: Default::default(), - } - } - - pub fn stop(&self) { - self.tx.send(WebUIEvent::Stop).ok(); - } - - pub fn update_tasks(&self, tasks: Vec) -> Result<(), crate::Error> { - self.tx - .send(WebUIEvent::UpdateTasks { tasks }) - .map_err(Error::Broadcast)?; - - Ok(()) - } - - pub fn output(&self, task: String, output: Vec) -> Result<(), crate::Error> { - self.tx - .send(WebUIEvent::TaskOutput { task, output }) - .map_err(Error::Broadcast)?; - - Ok(()) - } -} diff --git a/crates/turborepo-ui/src/wui/subscriber.rs b/crates/turborepo-ui/src/wui/subscriber.rs deleted file mode 100644 index 18c42f2dd8f77..0000000000000 --- a/crates/turborepo-ui/src/wui/subscriber.rs +++ /dev/null @@ -1,282 +0,0 @@ -use std::{collections::BTreeMap, sync::Arc}; - -use async_graphql::{Enum, SimpleObject}; -use serde::Serialize; -use tokio::sync::Mutex; - -use crate::{ - tui::event::{CacheResult, TaskResult}, - wui::{event::WebUIEvent, query::SharedState}, -}; - -/// Subscribes to the Web UI events and updates the state -pub struct Subscriber { - rx: tokio::sync::mpsc::UnboundedReceiver, -} - -impl Subscriber { - pub fn new(rx: tokio::sync::mpsc::UnboundedReceiver) -> Self { - Self { rx } - } - - pub async fn watch( - self, - // We use a tokio::sync::Mutex here because we want this future to be Send. - #[allow(clippy::type_complexity)] state: SharedState, - ) { - let mut rx = self.rx; - while let Some(event) = rx.recv().await { - Self::add_message(&state, event).await; - } - } - - async fn add_message(state: &Arc>, event: WebUIEvent) { - let mut state = state.lock().await; - - match event { - WebUIEvent::StartTask { - task, - output_logs: _, - } => { - state.tasks.insert( - task, - TaskState { - output: Vec::new(), - status: TaskStatus::Running, - cache_result: None, - cache_message: None, - }, - ); - } - WebUIEvent::TaskOutput { task, output } => { - if let Some(task) = state.tasks.get_mut(&task) { - task.output.extend(output); - } - } - WebUIEvent::EndTask { task, result } => { - if let Some(task) = state.tasks.get_mut(&task) { - task.status = TaskStatus::from(result); - } - } - WebUIEvent::CacheStatus { - task, - result, - message, - } => { - if let Some(task) = state.tasks.get_mut(&task) { - if result == CacheResult::Hit { - task.status = TaskStatus::Cached; - } - task.cache_result = Some(result); - task.cache_message = Some(message); - } - } - WebUIEvent::Stop => { - // TODO: stop watching - } - WebUIEvent::UpdateTasks { tasks } => { - state.tasks = tasks - .into_iter() - .map(|task| { - ( - task, - TaskState { - output: Vec::new(), - status: TaskStatus::Pending, - cache_result: None, - cache_message: None, - }, - ) - }) - .collect(); - } - WebUIEvent::RestartTasks { tasks } => { - state.tasks = tasks - .into_iter() - .map(|task| { - ( - task, - TaskState { - output: Vec::new(), - status: TaskStatus::Running, - cache_result: None, - cache_message: None, - }, - ) - }) - .collect(); - } - } - } -} - -#[derive(Debug, Clone, Copy, Serialize, PartialEq, Eq, Enum)] -pub enum TaskStatus { - Pending, - Running, - Cached, - Failed, - Succeeded, -} - -impl From for TaskStatus { - fn from(result: TaskResult) -> Self { - match result { - TaskResult::Success => Self::Succeeded, - TaskResult::CacheHit => Self::Cached, - TaskResult::Failure => Self::Failed, - } - } -} - -#[derive(Debug, Clone, Serialize, SimpleObject)] -pub struct TaskState { - output: Vec, - status: TaskStatus, - cache_result: Option, - /// The message for the cache status, i.e. `cache hit, replaying logs` - cache_message: Option, -} - -#[derive(Debug, Default, Clone, Serialize)] -pub struct WebUIState { - tasks: BTreeMap, -} - -impl WebUIState { - pub fn tasks(&self) -> &BTreeMap { - &self.tasks - } -} - -#[cfg(test)] -mod test { - use async_graphql::{EmptyMutation, EmptySubscription, Schema}; - - use super::*; - use crate::{ - tui::event::OutputLogs, - wui::{query::RunQuery, sender::WebUISender}, - }; - - #[tokio::test] - async fn test_web_ui_state() -> Result<(), crate::Error> { - let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); - let state = Arc::new(Mutex::new(WebUIState::default())); - let subscriber = Subscriber::new(rx); - - let sender = WebUISender::new(tx); - - // Start a successful task - sender.start_task("task".to_string(), OutputLogs::Full); - sender.output("task".to_string(), b"this is my output".to_vec())?; - sender.end_task("task".to_string(), TaskResult::Success); - - // Start a cached task - sender.start_task("task2".to_string(), OutputLogs::Full); - sender.status("task2".to_string(), "status".to_string(), CacheResult::Hit); - - // Start a failing task - sender.start_task("task3".to_string(), OutputLogs::Full); - sender.end_task("task3".to_string(), TaskResult::Failure); - - // Drop the sender so the subscriber can terminate - drop(sender); - - // Run the subscriber blocking - subscriber.watch(state.clone()).await; - - let state_handle = state.lock().await.clone(); - assert_eq!(state_handle.tasks().len(), 3); - assert_eq!( - state_handle.tasks().get("task2").unwrap().status, - TaskStatus::Cached - ); - assert_eq!( - state_handle.tasks().get("task").unwrap().status, - TaskStatus::Succeeded - ); - assert_eq!( - state_handle.tasks().get("task").unwrap().output, - b"this is my output" - ); - assert_eq!( - state_handle.tasks().get("task3").unwrap().status, - TaskStatus::Failed - ); - - // Now let's check with the GraphQL API - let schema = Schema::new(RunQuery::new(Some(state)), EmptyMutation, EmptySubscription); - let result = schema - .execute("query { currentRun { tasks { name state { status } } } }") - .await; - assert!(result.errors.is_empty()); - assert_eq!( - result.data, - async_graphql::Value::from_json(serde_json::json!({ - "currentRun": { - "tasks": [ - { - "name": "task", - "state": { - "status": "SUCCEEDED" - } - }, - { - "name": "task2", - "state": { - "status": "CACHED" - } - }, - { - "name": "task3", - "state": { - "status": "FAILED" - } - } - ] - } - })) - .unwrap() - ); - - Ok(()) - } - - #[tokio::test] - async fn test_restart_tasks() -> Result<(), crate::Error> { - let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); - let state = Arc::new(Mutex::new(WebUIState::default())); - let subscriber = Subscriber::new(rx); - - let sender = WebUISender::new(tx); - - // Start a successful task - sender.start_task("task".to_string(), OutputLogs::Full); - sender.output("task".to_string(), b"this is my output".to_vec())?; - sender.end_task("task".to_string(), TaskResult::Success); - - // Start a cached task - sender.start_task("task2".to_string(), OutputLogs::Full); - sender.status("task2".to_string(), "status".to_string(), CacheResult::Hit); - - // Restart a task - sender.restart_tasks(vec!["task".to_string()])?; - - // Drop the sender so the subscriber can terminate - drop(sender); - - // Run the subscriber blocking - subscriber.watch(state.clone()).await; - - let state_handle = state.lock().await.clone(); - assert_eq!(state_handle.tasks().len(), 1); - assert_eq!( - state_handle.tasks().get("task").unwrap().status, - TaskStatus::Running - ); - assert!(state_handle.tasks().get("task").unwrap().output.is_empty()); - - Ok(()) - } -} diff --git a/crates/turborepo/ARCHITECTURE.md b/crates/turborepo/ARCHITECTURE.md index 514e960f61da3..19ef8d9a6ef49 100644 --- a/crates/turborepo/ARCHITECTURE.md +++ b/crates/turborepo/ARCHITECTURE.md @@ -383,7 +383,7 @@ The summary module is responsible for any time of summary: ### 8. Query Subsystem The query subsystem powers `turbo query` (GraphQL introspection of the -package/task graph) and the Web UI mode (`--ui=web`). +package/task graph). **Crate layout:** diff --git a/crates/turborepo/src/main.rs b/crates/turborepo/src/main.rs index d8f89482c54ae..f6ab81491cfe7 100644 --- a/crates/turborepo/src/main.rs +++ b/crates/turborepo/src/main.rs @@ -58,18 +58,6 @@ impl turborepo_query_api::QueryServer for TurboQueryServer { .map_err(Into::into) }) } - - fn run_web_ui_server( - &self, - state: turborepo_ui::wui::query::SharedState, - run: Arc, - ) -> Pin> + Send + '_>> { - Box::pin(async move { - turborepo_query::run_server(Some(state), run) - .await - .map_err(Into::into) - }) - } } // This function should not expanded. Please add any logic to