Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion nextest.toml → .config/nextest.toml
Original file line number Diff line number Diff line change
Expand Up @@ -6,4 +6,4 @@ filter = "test(cli::opts)"
test-group = "serial"

[test-groups]
serial.max-threads = 1
serial = { max-threads = 1 }
17 changes: 16 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ rnk job list
| `rnk project activate <project-ref>` | 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

Expand All @@ -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 <id>` | Start an interactive session using a launcher |
| `rnk session stop <session-id>` | Stop a running session |
| `rnk session list` | List currently running sessions |
| `rnk session logs <session-id>` | View session logs |

### Launchers

| Command | Description |
|---------|-------------|
| `rnk launcher list` | List currently running launchers |

### Other
| Command | Description |
|---------|-------------|
| `rnk version` | Show client and server version info |
Expand Down
1 change: 1 addition & 0 deletions src/cli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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?,
};
Expand Down
8 changes: 8 additions & 0 deletions src/cli/cmd.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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 },

Expand All @@ -99,6 +102,11 @@ impl From<job::Error> for CmdError {
CmdError::Job { source }
}
}
impl From<session::Error> for CmdError {
fn from(source: session::Error) -> Self {
CmdError::Session { source }
}
}
impl From<launcher::Error> for CmdError {
fn from(source: launcher::Error) -> Self {
CmdError::Launcher { source }
Expand Down
9 changes: 7 additions & 2 deletions src/cli/cmd/launcher/list.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<SessionMode>,
}

#[derive(Debug, Snafu)]
pub enum Error {
Expand All @@ -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);
}
Expand Down
56 changes: 56 additions & 0 deletions src/cli/cmd/session.rs
Original file line number Diff line number Diff line change
@@ -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),
}
44 changes: 44 additions & 0 deletions src/cli/cmd/session/list.rs
Original file line number Diff line number Diff line change
@@ -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)
}
}
91 changes: 91 additions & 0 deletions src/cli/cmd/session/logs.rs
Original file line number Diff line number Diff line change
@@ -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<usize, httpclient::Error> {
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<bool, httpclient::Error> {
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(());
}
}
}
}
}
}
68 changes: 68 additions & 0 deletions src/cli/cmd/session/start.rs
Original file line number Diff line number Diff line change
@@ -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)
}
}
}
Loading
Loading