From a08f3245b275c8f21af357bffaab5eb3283f4b5a Mon Sep 17 00:00:00 2001 From: Ralf Grubenmann Date: Tue, 22 Sep 2026 09:00:14 +0200 Subject: [PATCH 1/3] feat: allow starting interactive sessions --- src/cli.rs | 1 + src/cli/cmd.rs | 8 ++ src/cli/cmd/launcher/list.rs | 9 ++- src/cli/cmd/session.rs | 56 ++++++++++++++ src/cli/cmd/session/list.rs | 44 +++++++++++ src/cli/cmd/session/logs.rs | 91 ++++++++++++++++++++++ src/cli/cmd/session/start.rs | 68 +++++++++++++++++ src/cli/cmd/session/stop.rs | 45 +++++++++++ src/cli/complete.rs | 144 +++++++++++++++++++++-------------- src/cli/opts.rs | 2 + src/cli/sink.rs | 1 + src/httpclient/data.rs | 49 ++++++++++-- 12 files changed, 452 insertions(+), 66 deletions(-) create mode 100644 src/cli/cmd/session.rs create mode 100644 src/cli/cmd/session/list.rs create mode 100644 src/cli/cmd/session/logs.rs create mode 100644 src/cli/cmd/session/start.rs create mode 100644 src/cli/cmd/session/stop.rs diff --git a/src/cli.rs b/src/cli.rs index e81a4b1..d033fd1 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -30,6 +30,7 @@ pub async fn execute_cmd(opts: MainOpts) -> Result<(), CmdError> { SubCommand::Dataset(input) => input.exec(ctx).await?, SubCommand::Job(input) => input.exec(ctx).await?, + SubCommand::Session(input) => input.exec(ctx).await?, SubCommand::Launcher(input) => input.exec(ctx).await?, SubCommand::Logout(input) => input.exec(&ctx).await?, }; diff --git a/src/cli/cmd.rs b/src/cli/cmd.rs index d51429c..39413ed 100644 --- a/src/cli/cmd.rs +++ b/src/cli/cmd.rs @@ -4,6 +4,7 @@ pub mod launcher; pub mod login; pub mod logout; pub mod project; +pub mod session; pub mod update; #[cfg(feature = "user-doc")] pub mod userdoc; @@ -87,6 +88,8 @@ pub enum CmdError { #[snafu(display("Job - {}", source))] Job { source: job::Error }, + #[snafu(display("Session - {}", source))] + Session { source: session::Error }, #[snafu(display("Launcher - {}", source))] Launcher { source: launcher::Error }, @@ -99,6 +102,11 @@ impl From for CmdError { CmdError::Job { source } } } +impl From for CmdError { + fn from(source: session::Error) -> Self { + CmdError::Session { source } + } +} impl From for CmdError { fn from(source: launcher::Error) -> Self { CmdError::Launcher { source } diff --git a/src/cli/cmd/launcher/list.rs b/src/cli/cmd/launcher/list.rs index 5292c4d..c6ce232 100644 --- a/src/cli/cmd/launcher/list.rs +++ b/src/cli/cmd/launcher/list.rs @@ -12,7 +12,10 @@ use snafu::{ResultExt, Snafu}; /// /// List currently running launchers. #[derive(Parser, Debug)] -pub struct Input {} +pub struct Input { + #[arg(long, short)] + pub mode: Option, +} #[derive(Debug, Snafu)] pub enum Error { @@ -27,7 +30,9 @@ impl Input { pub async fn exec(&self, ctx: Context) -> Result<(), Error> { let mut result = ctx.client.list_launchers().await.context(HttpClientSnafu)?; - result.retain(|e| e.launcher_type == SessionMode::NonInteractive); + if let Some(mode) = self.mode { + result.retain(|e| e.launcher_type == mode); + } if let Ok(Some(project)) = ctx.resolve_project_context().await { result.retain(|v| v.project_id == project.id); } diff --git a/src/cli/cmd/session.rs b/src/cli/cmd/session.rs new file mode 100644 index 0000000..27622ad --- /dev/null +++ b/src/cli/cmd/session.rs @@ -0,0 +1,56 @@ +pub mod list; +pub mod logs; +pub mod start; +pub mod stop; + +use super::Context; +use clap::Parser; +use snafu::{ResultExt, Snafu}; + +#[derive(Debug, Snafu)] +pub enum Error { + #[snafu(display("Error starting session: {}", source))] + Start { source: start::Error }, + + #[snafu(display("Error stopping session: {}", source))] + Stop { source: stop::Error }, + + #[snafu(display("Error listing sessions: {}", source))] + List { source: list::Error }, + + #[snafu(display("Error getting logs: {}", source))] + Logs { source: logs::Error }, +} + +/// Sub command for managing sessions [alias: s] +#[derive(Parser, Debug)] +pub struct Input { + #[command(subcommand)] + pub subcmd: SessionCommand, +} + +impl Input { + pub async fn exec(&self, ctx: Context) -> Result<(), Error> { + match &self.subcmd { + SessionCommand::Start(input) => input.exec(ctx).await.context(StartSnafu), + SessionCommand::Stop(input) => input.exec(ctx).await.context(StopSnafu), + SessionCommand::List(input) => input.exec(ctx).await.context(ListSnafu), + SessionCommand::Logs(input) => input.exec(ctx).await.context(LogsSnafu), + } + } +} + +#[derive(Parser, Debug)] +pub enum SessionCommand { + #[command()] + Start(start::Input), + + #[command()] + Stop(stop::Input), + + #[command(alias = "ls")] + List(list::Input), + + #[command()] + Logs(logs::Input), +} diff --git a/src/cli/cmd/session/list.rs b/src/cli/cmd/session/list.rs new file mode 100644 index 0000000..824a586 --- /dev/null +++ b/src/cli/cmd/session/list.rs @@ -0,0 +1,44 @@ +use super::Context; +use crate::{ + cli::sink::Error as SinkError, + httpclient::{ + self, + data::{InteractiveSessionList, SessionMode}, + }, +}; + +use clap::Parser; + +use snafu::{ResultExt, Snafu}; + +/// Listing sessions [alias: ls]. +/// +/// List currently running sessions. +#[derive(Parser, Debug)] +pub struct Input {} + +#[derive(Debug, Snafu)] +pub enum Error { + #[snafu(display("Error writing data: {}", source))] + WriteResult { source: SinkError }, + + #[snafu(display("Http error: {}", source))] + HttpClient { source: httpclient::Error }, +} + +impl Input { + pub async fn exec(&self, ctx: Context) -> Result<(), Error> { + let mut result = ctx + .client + .list_sessions(Some(SessionMode::Interactive)) + .await + .context(HttpClientSnafu)?; + + if let Ok(Some(project)) = ctx.resolve_project_context().await { + result.retain(|v| v.project_id == project.id); + } + let result = InteractiveSessionList(result); + + ctx.write_result(&result).await.context(WriteResultSnafu) + } +} diff --git a/src/cli/cmd/session/logs.rs b/src/cli/cmd/session/logs.rs new file mode 100644 index 0000000..7484bc6 --- /dev/null +++ b/src/cli/cmd/session/logs.rs @@ -0,0 +1,91 @@ +use super::Context; +use crate::cli::complete::complete_interactive_session_name; +use crate::{cli::sink::Error as SinkError, httpclient}; +use clap::{Parser, ValueHint}; +use std::time::Duration; +use tokio::signal; +use tokio::time::sleep; + +use clap_complete::ArgValueCompleter; +use snafu::{ResultExt, Snafu}; + +/// List the logs of a session. +#[derive(Parser, Debug)] +pub struct Input { + /// The session name/id to get logs for. + #[arg(value_hint=ValueHint::Other, add = ArgValueCompleter::new(complete_interactive_session_name))] + pub session_id: String, + + /// Periodically retrieves logs, it will stop when the session finished. + #[arg(long, short, default_value_t = false)] + pub follow: bool, + + /// The interval in seconds to wait between calls for logs. + #[arg(long, default_value_t = 2)] + pub follow_interval: u8, +} + +#[derive(Debug, Snafu)] +pub enum Error { + #[snafu(display("Error writing data: {}", source))] + WriteResult { source: SinkError }, + + #[snafu(display("Http error: {}", source))] + HttpClient { source: httpclient::Error }, +} + +impl Input { + pub async fn exec(&self, ctx: Context) -> Result<(), Error> { + if self.follow { + self.follow_logs(ctx).await.context(HttpClientSnafu) + } else { + self.show_logs(&ctx, 0).await.context(HttpClientSnafu)?; + Ok(()) + } + } + + async fn show_logs(&self, ctx: &Context, seen: usize) -> Result { + let result = ctx.client.session_logs(&self.session_id).await?; + if let Some(lines_blob) = result.0.get("amalthea-session") { + let lines: Vec<&str> = lines_blob.lines().collect(); + if lines.len() > seen { + for line in &lines[seen..] { + println!("{}", line); + } + return Ok(lines.len()); + } + } + Ok(seen) + } + + async fn is_session_finished(&self, ctx: &Context) -> Result { + let details = ctx.client.get_session(&self.session_id).await?; + + match &details { + None => Ok(true), + Some(d) => Ok(!d.status.state.is_running()), + } + } + + pub async fn follow_logs(&self, ctx: Context) -> Result<(), httpclient::Error> { + let mut seen: usize = self.show_logs(&ctx, 0).await?; + if self.is_session_finished(&ctx).await? { + return Ok(()); + } + + loop { + tokio::select! { + _ = signal::ctrl_c() => { + eprintln!("Interrupted, exiting."); + break Ok(()); + } + _ = sleep(Duration::from_secs(self.follow_interval as u64)) => { + seen = self.show_logs(&ctx, seen).await?; + if self.is_session_finished(&ctx).await? { + break Ok(()); + } + } + } + } + } +} diff --git a/src/cli/cmd/session/start.rs b/src/cli/cmd/session/start.rs new file mode 100644 index 0000000..ba31fc8 --- /dev/null +++ b/src/cli/cmd/session/start.rs @@ -0,0 +1,68 @@ +use crate::{ + cli::{cmd::session::logs, complete::complete_session_launcher_id}, + data::simple_message::SimpleMessage, + httpclient::{self, data::SessionStartRequest}, +}; + +use super::Context; +use crate::cli::sink::Error as SinkError; + +use clap::{Parser, ValueHint}; +use clap_complete::ArgValueCompleter; +use ulid::Ulid; + +use snafu::{ResultExt, Snafu}; + +/// Start a session. +/// +/// Starts an interactive session using a pre-configured session launcher. +#[derive(Parser, Debug)] +pub struct Input { + /// The launcher to use for launching the session. + #[arg(long, value_hint=ValueHint::Other, add = ArgValueCompleter::new(complete_session_launcher_id))] + pub launcher: Ulid, + + /// Start the session and show the logs until it ends or the user cancels with Ctrl-C. + #[arg(long, default_value_t = false)] + pub wait: bool, +} + +#[derive(Debug, Snafu)] +pub enum Error { + #[snafu(display("Error writing data: {}", source))] + WriteResult { source: SinkError }, + + #[snafu(display("Http error: {}", source))] + HttpClient { source: httpclient::Error }, +} + +impl Input { + pub async fn exec(&self, ctx: Context) -> Result<(), Error> { + let req = SessionStartRequest { + launcher_id: self.launcher.to_string(), + session_type: "interactive".into(), + ..Default::default() + }; + let result = ctx + .client + .start_session(req) + .await + .context(HttpClientSnafu)?; + + if self.wait { + ctx.write_result(&SimpleMessage { + message: format!("Started session {}. Waiting for logs...", result.name,), + }) + .await + .context(WriteResultSnafu)?; + let log_input = logs::Input { + session_id: result.name, + follow: true, + follow_interval: 2, + }; + log_input.follow_logs(ctx).await.context(HttpClientSnafu) + } else { + ctx.write_result(&result).await.context(WriteResultSnafu) + } + } +} diff --git a/src/cli/cmd/session/stop.rs b/src/cli/cmd/session/stop.rs new file mode 100644 index 0000000..b7fa296 --- /dev/null +++ b/src/cli/cmd/session/stop.rs @@ -0,0 +1,45 @@ +use crate::{ + cli::complete::complete_interactive_session_name, data::simple_message::SimpleMessage, + httpclient, +}; + +use super::Context; +use crate::cli::sink::Error as SinkError; + +use clap::{Parser, ValueHint}; + +use clap_complete::ArgValueCompleter; +use snafu::{ResultExt, Snafu}; + +/// Stop a session. +/// +/// Stop a running interactive session. +#[derive(Parser, Debug)] +pub struct Input { + /// The id of the session to stop + #[arg(value_hint=ValueHint::Other, add = ArgValueCompleter::new(complete_interactive_session_name))] + pub session_id: String, +} + +#[derive(Debug, Snafu)] +pub enum Error { + #[snafu(display("Error writing data: {}", source))] + WriteResult { source: SinkError }, + + #[snafu(display("Http error: {}", source))] + HttpClient { source: httpclient::Error }, +} + +impl Input { + pub async fn exec(&self, ctx: Context) -> Result<(), Error> { + ctx.client + .stop_session(&self.session_id) + .await + .context(HttpClientSnafu)?; + ctx.write_result(&SimpleMessage { + message: "Session is being removed.".into(), + }) + .await + .context(WriteResultSnafu) + } +} diff --git a/src/cli/complete.rs b/src/cli/complete.rs index b749773..cc6b46c 100644 --- a/src/cli/complete.rs +++ b/src/cli/complete.rs @@ -128,70 +128,96 @@ async fn resolve_project_id(client: &Client, id: ProjectId) -> Option { /// Complete a job session launcher id pub fn complete_job_launcher_id(current: &ffi::OsStr) -> Vec { make_sync_completer(current, async |client, opts| { - let launchers = match client.list_launchers().await { - Err(msg) => { - eprintln!( - "Completions failed: Error getting list of launchers: {}", - msg - ); - return vec![]; - } - Ok(res) => res, - }; - let mut result: Vec = vec![]; - let project_ctx = opts.get_project_context().ok().flatten(); - let project_id = match project_ctx { - Some(id) => resolve_project_id(&client, id).await, - None => None, - }; - for launcher in launchers - .0 - .iter() - .filter(|e| e.launcher_type == SessionMode::NonInteractive) - .filter(|e| match &project_id { - Some(id) => id == &e.project_id, - None => true, - }) - { - let cc = make_launcher_completion_candidate(&client, launcher).await; - result.push(cc); - } - if result.is_empty() { - eprintln!("No job launchers found."); - } - result + complete_launcher_id(client, opts, SessionMode::NonInteractive).await }) } -/// Complete a job name -pub fn complete_job_name(current: &ffi::OsStr) -> Vec { +/// Complete a interactive session launcher id +pub fn complete_session_launcher_id(current: &ffi::OsStr) -> Vec { make_sync_completer(current, async |client, opts| { - let jobs = match client - .list_sessions(Some(SessionMode::NonInteractive)) - .await - { - Err(msg) => { - eprintln!("Completions failed: Error getting list of jobs: {}", msg); - return vec![]; - } - Ok(res) => res, - }; - let mut result: Vec = vec![]; - let project_ctx = opts.get_project_context().ok().flatten(); - let project_id = match project_ctx { - Some(id) => resolve_project_id(&client, id).await, - None => None, - }; - for job in jobs.0.iter().filter(|e| match &project_id { + complete_launcher_id(client, opts, SessionMode::Interactive).await + }) +} + +async fn complete_launcher_id( + client: Client, + opts: CommonOpts, + mode: SessionMode, +) -> Vec { + let launchers = match client.list_launchers().await { + Err(msg) => { + eprintln!( + "Completions failed: Error getting list of launchers: {}", + msg + ); + return vec![]; + } + Ok(res) => res, + }; + let mut result: Vec = vec![]; + let project_ctx = opts.get_project_context().ok().flatten(); + let project_id = match project_ctx { + Some(id) => resolve_project_id(&client, id).await, + None => None, + }; + for launcher in launchers + .0 + .iter() + .filter(|e| e.launcher_type == mode) + .filter(|e| match &project_id { Some(id) => id == &e.project_id, None => true, - }) { - let cc = make_job_name_completion_candidate(&client, job).await; - result.push(cc); - } - if result.is_empty() { - eprintln!("No job launchers found."); - } - result + }) + { + let cc = make_launcher_completion_candidate(&client, launcher).await; + result.push(cc); + } + if result.is_empty() { + eprintln!("No job launchers found."); + } + result +} + +/// Complete a job name +pub fn complete_job_name(current: &ffi::OsStr) -> Vec { + make_sync_completer(current, async |client, opts| { + complete_session_name(client, opts, SessionMode::NonInteractive).await }) } +/// Complete a interactive session name +pub fn complete_interactive_session_name(current: &ffi::OsStr) -> Vec { + make_sync_completer(current, async |client, opts| { + complete_session_name(client, opts, SessionMode::NonInteractive).await + }) +} + +async fn complete_session_name( + client: Client, + opts: CommonOpts, + mode: SessionMode, +) -> Vec { + let jobs = match client.list_sessions(Some(mode)).await { + Err(msg) => { + eprintln!("Completions failed: Error getting list of jobs: {}", msg); + return vec![]; + } + Ok(res) => res, + }; + let mut result: Vec = vec![]; + let project_ctx = opts.get_project_context().ok().flatten(); + let project_id = match project_ctx { + Some(id) => resolve_project_id(&client, id).await, + None => None, + }; + for job in jobs.0.iter().filter(|e| match &project_id { + Some(id) => id == &e.project_id, + None => true, + }) { + let cc = make_job_name_completion_candidate(&client, job).await; + result.push(cc); + } + if result.is_empty() { + eprintln!("No job launchers found."); + } + result +} diff --git a/src/cli/opts.rs b/src/cli/opts.rs index 3bde8e3..0eec05c 100644 --- a/src/cli/opts.rs +++ b/src/cli/opts.rs @@ -175,6 +175,8 @@ pub enum SubCommand { #[command(alias = "j")] Job(job::Input), + #[command(alias = "s")] + Session(session::Input), #[command(alias = "l")] Launcher(launcher::Input), diff --git a/src/cli/sink.rs b/src/cli/sink.rs index 3966f7a..6162f98 100644 --- a/src/cli/sink.rs +++ b/src/cli/sink.rs @@ -64,6 +64,7 @@ impl Sink for UserCode {} impl Sink for Response {} impl Sink for SessionStartResponse {} impl Sink for SessionList {} +impl Sink for InteractiveSessionList {} impl Sink for LauncherList {} impl Sink for SessionLogs {} impl Sink for VersionInfo {} diff --git a/src/httpclient/data.rs b/src/httpclient/data.rs index ee34f26..f75a5ff 100644 --- a/src/httpclient/data.rs +++ b/src/httpclient/data.rs @@ -2,6 +2,7 @@ //! `De/Serialize` instances. use crate::data::{renku_url::RenkuUrl, submission_id::SubmissionId}; +use clap::ValueEnum; use iso8601_timestamp::Timestamp; use serde::{Deserialize, Serialize}; use std::{collections::HashMap, fmt}; @@ -37,7 +38,7 @@ impl fmt::Display for SessionLogs { } } -#[derive(Debug, Serialize, Deserialize, PartialEq)] +#[derive(Debug, Serialize, Deserialize, PartialEq, Clone, Copy, ValueEnum)] pub enum SessionMode { #[serde(rename = "interactive")] Interactive, @@ -60,7 +61,7 @@ impl SessionMode { } } -#[derive(Debug, Serialize, Deserialize)] +#[derive(Debug, Serialize, Deserialize, Default)] pub struct SessionStartRequest { pub launcher_id: String, pub session_type: String, @@ -105,10 +106,11 @@ where { let mut builder = Builder::default(); for r in data { - let data = vec![&r.name, r.id.as_str(), &r.project_id]; + let launcher_type = r.launcher_type.to_string(); + let data = vec![&r.name, r.id.as_str(), &launcher_type, &r.project_id]; builder.push_record(data); } - builder.insert_record(0, vec!["Launcher", "Id", "Project Id"]); + builder.insert_record(0, vec!["Launcher", "Id", "Type", "Project Id"]); let mut table = builder.build(); let settings = Settings::default().with(Style::sharp()); @@ -174,7 +176,7 @@ where impl fmt::Display for SessionList { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { if self.0.is_empty() { - write!(f, "No jobs found.") + write!(f, "No sessions found.") } else { let table = create_session_table(&self.0); write!(f, "{}", table) @@ -182,6 +184,41 @@ impl fmt::Display for SessionList { } } +// New Type to distinguish formatting for job lists and interactive session lists +#[derive(Debug, Serialize, Deserialize)] +pub struct InteractiveSessionList(pub SessionList); + +fn create_interactive_session_table<'a, I>(data: I) -> Table +where + I: IntoIterator, +{ + let mut builder = Builder::default(); + for r in data { + let started = r.started.format(); + let status = r.status.state.to_str(); + let data = vec![&r.name, &r.project_id, status, &started, &r.url]; + builder.push_record(data); + } + builder.insert_record(0, vec!["Name", "Project Id", "Status", "Started", "Url"]); + + let mut table = builder.build(); + let settings = Settings::default().with(Style::sharp()); + + table.with(settings); + table +} + +impl fmt::Display for InteractiveSessionList { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + if self.0.0.is_empty() { + write!(f, "No sessions found.") + } else { + let table = create_interactive_session_table(&self.0.0); + write!(f, "{}", table) + } + } +} + #[derive(Debug, Serialize, Deserialize)] #[serde(rename_all = "lowercase")] pub enum SessionState { @@ -237,6 +274,8 @@ pub struct SessionStartResponse { pub submission_id: Option, pub status: SessionStatus, pub started: Timestamp, + pub session_type: SessionMode, + pub url: String, } impl fmt::Display for SessionStartResponse { From 82864557fec6869cc7a2d6138dfb58f8a19252ea Mon Sep 17 00:00:00 2001 From: Ralf Grubenmann Date: Tue, 22 Sep 2026 11:11:27 +0200 Subject: [PATCH 2/3] update readme --- README.md | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index c57ac47..a4b0ca6 100644 --- a/README.md +++ b/README.md @@ -55,6 +55,7 @@ rnk job list | `rnk project activate ` | Set the active project for the current user (short form: `rnk p a`) | | `rnk project deactivate` | Unset the active project (short form: `rnk p d`) | | `rnk project current` | Show the currently active project (short form: `rnk p c`) | +| `rnk project list` | List projects (supports `--all` / `-a` to show all projects, `--n-results` / `-n` to limit results) | ### Datasets @@ -71,8 +72,22 @@ rnk job list | `rnk job stop` | Stop a job | | `rnk job logs` | View job logs | -### Other +### Sessions + +| Command | Description | +|---------|-------------| +| `rnk session start --launcher ` | Start an interactive session using a launcher | +| `rnk session stop ` | Stop a running session | +| `rnk session list` | List currently running sessions | +| `rnk session logs ` | View session logs | + +### Launchers +| Command | Description | +|---------|-------------| +| `rnk launcher list` | List currently running launchers | + +### Other | Command | Description | |---------|-------------| | `rnk version` | Show client and server version info | From 7367417b3c21af032e2a8c00985b805a26251983 Mon Sep 17 00:00:00 2001 From: Ralf Grubenmann Date: Wed, 23 Sep 2026 10:32:29 +0200 Subject: [PATCH 3/3] fix nextest running tests serially --- nextest.toml => .config/nextest.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) rename nextest.toml => .config/nextest.toml (91%) diff --git a/nextest.toml b/.config/nextest.toml similarity index 91% rename from nextest.toml rename to .config/nextest.toml index 7626158..560c7cc 100644 --- a/nextest.toml +++ b/.config/nextest.toml @@ -6,4 +6,4 @@ filter = "test(cli::opts)" test-group = "serial" [test-groups] -serial.max-threads = 1 \ No newline at end of file +serial = { max-threads = 1 }