use crate::api::kasl_server::{AGENT_TOKEN_PROMPT, AGENT_TOKEN_SECRET, KaslServer, UploadError, normalize_url};
use crate::db::server_outbox::ServerOutbox;
use crate::db::workdays::Workdays;
use crate::libs::config::{Config, KaslServerConfig};
use crate::libs::day_delivery::{Delivered, deliver, record_single};
use crate::libs::day_upload::build_day_upload;
use crate::libs::messages::Message;
use crate::libs::secret::Secret;
use crate::{msg_error_anyhow, msg_info, msg_print, msg_success, msg_warning};
use anyhow::{Context, Result};
use chrono::{Duration, Local, NaiveDate};
use clap::{Args, Subcommand};
use dialoguer::{Input, Password, theme::ColorfulTheme};
use reqwest::StatusCode;
#[derive(Debug, Args)]
pub struct ServerArgs {
#[command(subcommand)]
command: ServerCommand,
}
#[derive(Debug, Subcommand)]
enum ServerCommand {
#[command(about = "Connect this machine to a kasl-server")]
Connect(ConnectArgs),
#[command(about = "Show the current connection to a kasl-server")]
Status,
#[command(about = "Send a day's work to the connected kasl-server")]
Push(PushArgs),
#[command(about = "Send every day still waiting to reach the server")]
Flush,
#[command(about = "Show the days still waiting to reach the server")]
Queue,
#[command(about = "Queue every recorded day in a date range and send them")]
Backfill(BackfillArgs),
#[command(about = "Forget the connection and the stored agent token")]
Disconnect,
}
#[derive(Debug, Args)]
pub struct BackfillArgs {
#[arg(long, value_name = "YYYY-MM-DD")]
from: NaiveDate,
#[arg(long, value_name = "YYYY-MM-DD")]
to: Option<NaiveDate>,
}
#[derive(Debug, Args)]
pub struct PushArgs {
#[arg(long, short, help = "Send the last day instead of today")]
last: bool,
#[arg(long, value_name = "YYYY-MM-DD", conflicts_with = "last")]
date: Option<NaiveDate>,
}
#[derive(Debug, Args)]
pub struct ConnectArgs {
#[arg(long, value_name = "URL")]
url: Option<String>,
#[arg(long, value_name = "PATH")]
ca_certificate: Option<String>,
}
pub async fn cmd(args: ServerArgs) -> Result<()> {
match args.command {
ServerCommand::Connect(args) => connect(args).await,
ServerCommand::Status => status().await,
ServerCommand::Push(args) => push(args).await,
ServerCommand::Flush => flush().await,
ServerCommand::Queue => queue(),
ServerCommand::Backfill(args) => backfill(args).await,
ServerCommand::Disconnect => disconnect(),
}
}
async fn connect(args: ConnectArgs) -> Result<()> {
crate::libs::prompt::ensure_interactive("`kasl server connect` needs a terminal to ask for the agent token")?;
let mut config = Config::read().unwrap_or_default();
let url = match args.url {
Some(url) => normalize_url(&url),
None => {
let entered: String = Input::with_theme(&ColorfulTheme::default())
.with_prompt(Message::PromptKaslServerUrl.to_string())
.with_initial_text(config.kasl_server.as_ref().map(|s| s.url.clone()).unwrap_or_default())
.interact_text()?;
normalize_url(&entered)
}
};
if !url.starts_with("http://") && !url.starts_with("https://") {
return Err(msg_error_anyhow!(Message::KaslServerUrlNeedsScheme(url)));
}
let candidate = KaslServerConfig {
url: url.clone(),
ca_certificate: args
.ca_certificate
.or_else(|| config.kasl_server.as_ref().and_then(|s| s.ca_certificate.clone())),
};
let client = KaslServer::new(&candidate)?;
let health = client.health().await?;
msg_info!(Message::KaslServerReached {
url: url.clone(),
version: health.version.clone(),
});
if health.database != "ok" {
msg_warning!(Message::KaslServerDatabaseUnhealthy(health.database.clone()));
}
let secret = Secret::new(AGENT_TOKEN_SECRET, AGENT_TOKEN_PROMPT);
let token: String = Password::with_theme(&ColorfulTheme::default())
.with_prompt(Message::PromptKaslServerToken.to_string())
.interact()?;
let token = token.trim().to_string();
if token.is_empty() {
return Err(msg_error_anyhow!(Message::KaslServerTokenEmpty));
}
let identity = client.identify(&token).await?;
secret
.store(&token)
.context("the token was accepted but could not be stored in the OS keyring")?;
config.kasl_server = Some(candidate);
config.save()?;
msg_success!(Message::KaslServerConnected {
user_name: identity.user_name,
agent_name: identity.agent_name,
});
Ok(())
}
async fn status() -> Result<()> {
let config = Config::read().unwrap_or_default();
let Some(server_config) = config.kasl_server else {
msg_print!(Message::KaslServerNotConnected);
return Ok(());
};
msg_info!(Message::KaslServerConfigured(server_config.url.clone()));
let secret = Secret::new(AGENT_TOKEN_SECRET, AGENT_TOKEN_PROMPT);
let Some(token) = secret.try_get_cached() else {
msg_warning!(Message::KaslServerTokenMissing);
return Ok(());
};
let client = KaslServer::new(&server_config)?;
match client.health().await {
Ok(health) => msg_info!(Message::KaslServerReached {
url: server_config.url.clone(),
version: health.version,
}),
Err(error) => {
msg_warning!(Message::KaslServerUnreachable(error.to_string()));
return Ok(());
}
}
match client.identify(&token).await {
Ok(identity) => msg_success!(Message::KaslServerConnected {
user_name: identity.user_name,
agent_name: identity.agent_name,
}),
Err(error) => msg_warning!(Message::KaslServerTokenRejected(error.to_string())),
}
Ok(())
}
async fn push(args: PushArgs) -> Result<()> {
let date = match args.date {
Some(date) => date,
None if args.last => (Local::now() - Duration::days(1)).date_naive(),
None => Local::now().date_naive(),
};
let config = Config::read().unwrap_or_default();
let Some(server_config) = config.kasl_server else {
return Err(msg_error_anyhow!(Message::KaslServerNotConnected));
};
let Some(day) = build_day_upload(date)? else {
msg_print!(Message::KaslServerNoDayToPush(date.to_string()));
return Ok(());
};
let secret = Secret::new(AGENT_TOKEN_SECRET, AGENT_TOKEN_PROMPT);
let Some(token) = secret.try_get_cached() else {
return Err(msg_error_anyhow!(Message::KaslServerTokenMissing));
};
let client = KaslServer::new(&server_config)?;
match client.upload_day(&token, &day).await {
Ok(accepted) => {
msg_success!(Message::KaslServerDayPushed {
date: accepted.date.to_string(),
pauses: accepted.pauses,
tasks: accepted.tasks,
});
if accepted.deleted_tasks > 0 {
msg_info!(Message::KaslServerTasksDeleted(accepted.deleted_tasks));
}
ServerOutbox::new()?.remove(date)?;
flush_with(&client, &token).await?;
Ok(())
}
Err(error) => {
let outcome = record_single(&mut ServerOutbox::new()?, date, &error)?;
if matches!(outcome, Delivered::Deferred { .. }) {
msg_info!(Message::KaslServerDayQueued(date.to_string()));
}
match error {
error @ UploadError::Rejected {
status: StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN,
..
} => Err(msg_error_anyhow!(Message::KaslServerPushTokenRejected(error.to_string()))),
error @ UploadError::Rejected { .. } => Err(msg_error_anyhow!(Message::KaslServerPushRejected(error.to_string()))),
error => Err(msg_error_anyhow!(Message::KaslServerPushRetryable(error.to_string()))),
}
}
}
}
async fn flush() -> Result<()> {
if ServerOutbox::new()?.count()? == 0 {
msg_print!(Message::KaslServerQueueEmpty);
return Ok(());
}
let (client, token) = connected_client()?;
flush_with(&client, &token).await
}
async fn flush_with(client: &KaslServer, token: &str) -> Result<()> {
let mut outbox = ServerOutbox::new()?;
let dates: Vec<NaiveDate> = outbox.pending()?.into_iter().map(|owed| owed.date).collect();
if dates.is_empty() {
return Ok(());
}
msg_info!(Message::KaslServerQueueSending(dates.len()));
let outcomes = deliver(client, token, &mut outbox, &dates).await?;
let (mut accepted, mut refused, mut deferred) = (0, 0, 0);
for outcome in &outcomes {
match outcome {
Delivered::Accepted {
date,
pauses,
tasks,
deleted_tasks,
} => {
accepted += 1;
msg_success!(Message::KaslServerDayPushed {
date: date.to_string(),
pauses: *pauses,
tasks: *tasks,
});
if *deleted_tasks > 0 {
msg_info!(Message::KaslServerTasksDeleted(*deleted_tasks));
}
}
Delivered::Refused { date, reason } => {
refused += 1;
msg_warning!(Message::KaslServerDayRefused {
date: date.to_string(),
reason: reason.clone(),
});
}
Delivered::Deferred { date, reason } => {
deferred += 1;
msg_warning!(Message::KaslServerDayDeferred {
date: date.to_string(),
reason: reason.clone(),
});
}
}
}
msg_print!(Message::KaslServerFlushSummary { accepted, refused, deferred });
Ok(())
}
fn queue() -> Result<()> {
let outbox = ServerOutbox::new()?;
let owed = outbox.pending()?;
if owed.is_empty() {
msg_print!(Message::KaslServerQueueEmpty);
return Ok(());
}
msg_info!(Message::KaslServerQueueOwed(owed.len() as i64));
for day in &owed {
msg_print!(Message::KaslServerQueueEntry {
date: day.date.to_string(),
attempts: day.attempts,
last_error: day.last_error.clone(),
});
}
Ok(())
}
async fn backfill(args: BackfillArgs) -> Result<()> {
let to = args.to.unwrap_or_else(|| Local::now().date_naive());
if args.from > to {
return Err(msg_error_anyhow!(Message::KaslServerBackfillOrderReversed));
}
let (client, token) = connected_client()?;
let mut workdays = Workdays::new()?;
let mut dates = Vec::new();
let mut date = args.from;
while date <= to {
if workdays.fetch(date)?.is_some() {
dates.push(date);
}
date += Duration::days(1);
}
if dates.is_empty() {
msg_print!(Message::KaslServerBackfillNoDays {
from: args.from.to_string(),
to: to.to_string(),
});
return Ok(());
}
msg_info!(Message::KaslServerBackfillRange {
from: args.from.to_string(),
to: to.to_string(),
days: dates.len(),
});
let mut outbox = ServerOutbox::new()?;
for date in &dates {
outbox.enqueue(*date, "queued by backfill")?;
}
flush_with(&client, &token).await
}
fn connected_client() -> Result<(KaslServer, String)> {
let config = Config::read().unwrap_or_default();
let Some(server_config) = config.kasl_server else {
return Err(msg_error_anyhow!(Message::KaslServerNotConnected));
};
let Some(token) = Secret::new(AGENT_TOKEN_SECRET, AGENT_TOKEN_PROMPT).try_get_cached() else {
return Err(msg_error_anyhow!(Message::KaslServerTokenMissing));
};
Ok((KaslServer::new(&server_config)?, token))
}
fn disconnect() -> Result<()> {
let mut config = Config::read().unwrap_or_default();
if let Err(error) = Secret::new(AGENT_TOKEN_SECRET, AGENT_TOKEN_PROMPT).delete() {
msg_warning!(Message::KaslServerTokenNotRemoved(error.to_string()));
}
if config.kasl_server.take().is_some() {
config.save()?;
msg_success!(Message::KaslServerDisconnected);
} else {
msg_print!(Message::KaslServerNotConnected);
}
Ok(())
}