1use crate::api::kasl_server::{AGENT_TOKEN_PROMPT, AGENT_TOKEN_SECRET, KaslServer, UploadError, normalize_url};
18use crate::db::server_outbox::ServerOutbox;
19use crate::db::workdays::Workdays;
20use crate::libs::config::{Config, KaslServerConfig};
21use crate::libs::day_delivery::{Delivered, deliver, record_single};
22use crate::libs::day_upload::build_day_upload;
23use crate::libs::messages::Message;
24use crate::libs::secret::Secret;
25use crate::{msg_error_anyhow, msg_info, msg_print, msg_success, msg_warning};
26use anyhow::{Context, Result};
27use chrono::{Duration, Local, NaiveDate};
28use clap::{Args, Subcommand};
29use dialoguer::{Input, Password, theme::ColorfulTheme};
30use reqwest::StatusCode;
31
32#[derive(Debug, Args)]
34pub struct ServerArgs {
35 #[command(subcommand)]
36 command: ServerCommand,
37}
38
39#[derive(Debug, Subcommand)]
41enum ServerCommand {
42 #[command(about = "Connect this machine to a kasl-server")]
44 Connect(ConnectArgs),
45
46 #[command(about = "Show the current connection to a kasl-server")]
48 Status,
49
50 #[command(about = "Send a day's work to the connected kasl-server")]
52 Push(PushArgs),
53
54 #[command(about = "Send every day still waiting to reach the server")]
56 Flush,
57
58 #[command(about = "Show the days still waiting to reach the server")]
60 Queue,
61
62 #[command(about = "Queue every recorded day in a date range and send them")]
64 Backfill(BackfillArgs),
65
66 #[command(about = "Forget the connection and the stored agent token")]
68 Disconnect,
69}
70
71#[derive(Debug, Args)]
73pub struct BackfillArgs {
74 #[arg(long, value_name = "YYYY-MM-DD")]
76 from: NaiveDate,
77
78 #[arg(long, value_name = "YYYY-MM-DD")]
80 to: Option<NaiveDate>,
81}
82
83#[derive(Debug, Args)]
85pub struct PushArgs {
86 #[arg(long, short, help = "Send the last day instead of today")]
88 last: bool,
89
90 #[arg(long, value_name = "YYYY-MM-DD", conflicts_with = "last")]
92 date: Option<NaiveDate>,
93}
94
95#[derive(Debug, Args)]
97pub struct ConnectArgs {
98 #[arg(long, value_name = "URL")]
100 url: Option<String>,
101
102 #[arg(long, value_name = "PATH")]
104 ca_certificate: Option<String>,
105}
106
107pub async fn cmd(args: ServerArgs) -> Result<()> {
109 match args.command {
110 ServerCommand::Connect(args) => connect(args).await,
111 ServerCommand::Status => status().await,
112 ServerCommand::Push(args) => push(args).await,
113 ServerCommand::Flush => flush().await,
114 ServerCommand::Queue => queue(),
115 ServerCommand::Backfill(args) => backfill(args).await,
116 ServerCommand::Disconnect => disconnect(),
117 }
118}
119
120async fn connect(args: ConnectArgs) -> Result<()> {
126 crate::libs::prompt::ensure_interactive("`kasl server connect` needs a terminal to ask for the agent token")?;
133
134 let mut config = Config::read().unwrap_or_default();
135
136 let url = match args.url {
137 Some(url) => normalize_url(&url),
138 None => {
139 let entered: String = Input::with_theme(&ColorfulTheme::default())
140 .with_prompt(Message::PromptKaslServerUrl.to_string())
141 .with_initial_text(config.kasl_server.as_ref().map(|s| s.url.clone()).unwrap_or_default())
142 .interact_text()?;
143 normalize_url(&entered)
144 }
145 };
146
147 if !url.starts_with("http://") && !url.starts_with("https://") {
150 return Err(msg_error_anyhow!(Message::KaslServerUrlNeedsScheme(url)));
151 }
152
153 let candidate = KaslServerConfig {
154 url: url.clone(),
155 ca_certificate: args
158 .ca_certificate
159 .or_else(|| config.kasl_server.as_ref().and_then(|s| s.ca_certificate.clone())),
160 };
161
162 let client = KaslServer::new(&candidate)?;
163
164 let health = client.health().await?;
167 msg_info!(Message::KaslServerReached {
168 url: url.clone(),
169 version: health.version.clone(),
170 });
171 if health.database != "ok" {
172 msg_warning!(Message::KaslServerDatabaseUnhealthy(health.database.clone()));
175 }
176
177 let secret = Secret::new(AGENT_TOKEN_SECRET, AGENT_TOKEN_PROMPT);
178 let token: String = Password::with_theme(&ColorfulTheme::default())
179 .with_prompt(Message::PromptKaslServerToken.to_string())
180 .interact()?;
181 let token = token.trim().to_string();
182 if token.is_empty() {
183 return Err(msg_error_anyhow!(Message::KaslServerTokenEmpty));
184 }
185
186 let identity = client.identify(&token).await?;
188
189 secret
191 .store(&token)
192 .context("the token was accepted but could not be stored in the OS keyring")?;
193 config.kasl_server = Some(candidate);
194 config.save()?;
195
196 msg_success!(Message::KaslServerConnected {
197 user_name: identity.user_name,
198 agent_name: identity.agent_name,
199 });
200 Ok(())
201}
202
203async fn status() -> Result<()> {
209 let config = Config::read().unwrap_or_default();
210 let Some(server_config) = config.kasl_server else {
211 msg_print!(Message::KaslServerNotConnected);
212 return Ok(());
213 };
214
215 msg_info!(Message::KaslServerConfigured(server_config.url.clone()));
216
217 let secret = Secret::new(AGENT_TOKEN_SECRET, AGENT_TOKEN_PROMPT);
218 let Some(token) = secret.try_get_cached() else {
219 msg_warning!(Message::KaslServerTokenMissing);
222 return Ok(());
223 };
224
225 let client = KaslServer::new(&server_config)?;
226
227 match client.health().await {
228 Ok(health) => msg_info!(Message::KaslServerReached {
229 url: server_config.url.clone(),
230 version: health.version,
231 }),
232 Err(error) => {
233 msg_warning!(Message::KaslServerUnreachable(error.to_string()));
234 return Ok(());
235 }
236 }
237
238 match client.identify(&token).await {
239 Ok(identity) => msg_success!(Message::KaslServerConnected {
240 user_name: identity.user_name,
241 agent_name: identity.agent_name,
242 }),
243 Err(error) => msg_warning!(Message::KaslServerTokenRejected(error.to_string())),
244 }
245
246 Ok(())
247}
248
249async fn push(args: PushArgs) -> Result<()> {
266 let date = match args.date {
267 Some(date) => date,
268 None if args.last => (Local::now() - Duration::days(1)).date_naive(),
269 None => Local::now().date_naive(),
270 };
271
272 let config = Config::read().unwrap_or_default();
273 let Some(server_config) = config.kasl_server else {
274 return Err(msg_error_anyhow!(Message::KaslServerNotConnected));
275 };
276
277 let Some(day) = build_day_upload(date)? else {
278 msg_print!(Message::KaslServerNoDayToPush(date.to_string()));
279 return Ok(());
280 };
281
282 let secret = Secret::new(AGENT_TOKEN_SECRET, AGENT_TOKEN_PROMPT);
283 let Some(token) = secret.try_get_cached() else {
284 return Err(msg_error_anyhow!(Message::KaslServerTokenMissing));
285 };
286
287 let client = KaslServer::new(&server_config)?;
288
289 match client.upload_day(&token, &day).await {
290 Ok(accepted) => {
291 msg_success!(Message::KaslServerDayPushed {
292 date: accepted.date.to_string(),
293 pauses: accepted.pauses,
294 tasks: accepted.tasks,
295 });
296 if accepted.deleted_tasks > 0 {
300 msg_info!(Message::KaslServerTasksDeleted(accepted.deleted_tasks));
301 }
302
303 ServerOutbox::new()?.remove(date)?;
307
308 flush_with(&client, &token).await?;
313 Ok(())
314 }
315 Err(error) => {
321 let outcome = record_single(&mut ServerOutbox::new()?, date, &error)?;
325 if matches!(outcome, Delivered::Deferred { .. }) {
326 msg_info!(Message::KaslServerDayQueued(date.to_string()));
327 }
328
329 match error {
330 error @ UploadError::Rejected {
331 status: StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN,
332 ..
333 } => Err(msg_error_anyhow!(Message::KaslServerPushTokenRejected(error.to_string()))),
334 error @ UploadError::Rejected { .. } => Err(msg_error_anyhow!(Message::KaslServerPushRejected(error.to_string()))),
335 error => Err(msg_error_anyhow!(Message::KaslServerPushRetryable(error.to_string()))),
336 }
337 }
338 }
339}
340
341async fn flush() -> Result<()> {
343 if ServerOutbox::new()?.count()? == 0 {
349 msg_print!(Message::KaslServerQueueEmpty);
350 return Ok(());
351 }
352
353 let (client, token) = connected_client()?;
354 flush_with(&client, &token).await
355}
356
357async fn flush_with(client: &KaslServer, token: &str) -> Result<()> {
362 let mut outbox = ServerOutbox::new()?;
363 let dates: Vec<NaiveDate> = outbox.pending()?.into_iter().map(|owed| owed.date).collect();
364 if dates.is_empty() {
365 return Ok(());
366 }
367
368 msg_info!(Message::KaslServerQueueSending(dates.len()));
369
370 let outcomes = deliver(client, token, &mut outbox, &dates).await?;
371
372 let (mut accepted, mut refused, mut deferred) = (0, 0, 0);
373 for outcome in &outcomes {
374 match outcome {
375 Delivered::Accepted {
376 date,
377 pauses,
378 tasks,
379 deleted_tasks,
380 } => {
381 accepted += 1;
382 msg_success!(Message::KaslServerDayPushed {
383 date: date.to_string(),
384 pauses: *pauses,
385 tasks: *tasks,
386 });
387 if *deleted_tasks > 0 {
388 msg_info!(Message::KaslServerTasksDeleted(*deleted_tasks));
389 }
390 }
391 Delivered::Refused { date, reason } => {
396 refused += 1;
397 msg_warning!(Message::KaslServerDayRefused {
398 date: date.to_string(),
399 reason: reason.clone(),
400 });
401 }
402 Delivered::Deferred { date, reason } => {
403 deferred += 1;
404 msg_warning!(Message::KaslServerDayDeferred {
405 date: date.to_string(),
406 reason: reason.clone(),
407 });
408 }
409 }
410 }
411
412 msg_print!(Message::KaslServerFlushSummary { accepted, refused, deferred });
413 Ok(())
414}
415
416fn queue() -> Result<()> {
421 let outbox = ServerOutbox::new()?;
422 let owed = outbox.pending()?;
423
424 if owed.is_empty() {
425 msg_print!(Message::KaslServerQueueEmpty);
426 return Ok(());
427 }
428
429 msg_info!(Message::KaslServerQueueOwed(owed.len() as i64));
430 for day in &owed {
431 msg_print!(Message::KaslServerQueueEntry {
432 date: day.date.to_string(),
433 attempts: day.attempts,
434 last_error: day.last_error.clone(),
435 });
436 }
437
438 Ok(())
439}
440
441async fn backfill(args: BackfillArgs) -> Result<()> {
448 let to = args.to.unwrap_or_else(|| Local::now().date_naive());
449 if args.from > to {
450 return Err(msg_error_anyhow!(Message::KaslServerBackfillOrderReversed));
451 }
452
453 let (client, token) = connected_client()?;
454
455 let mut workdays = Workdays::new()?;
456 let mut dates = Vec::new();
457 let mut date = args.from;
458 while date <= to {
459 if workdays.fetch(date)?.is_some() {
460 dates.push(date);
461 }
462 date += Duration::days(1);
463 }
464
465 if dates.is_empty() {
466 msg_print!(Message::KaslServerBackfillNoDays {
467 from: args.from.to_string(),
468 to: to.to_string(),
469 });
470 return Ok(());
471 }
472
473 msg_info!(Message::KaslServerBackfillRange {
474 from: args.from.to_string(),
475 to: to.to_string(),
476 days: dates.len(),
477 });
478
479 let mut outbox = ServerOutbox::new()?;
482 for date in &dates {
483 outbox.enqueue(*date, "queued by backfill")?;
484 }
485
486 flush_with(&client, &token).await
487}
488
489fn connected_client() -> Result<(KaslServer, String)> {
495 let config = Config::read().unwrap_or_default();
496 let Some(server_config) = config.kasl_server else {
497 return Err(msg_error_anyhow!(Message::KaslServerNotConnected));
498 };
499
500 let Some(token) = Secret::new(AGENT_TOKEN_SECRET, AGENT_TOKEN_PROMPT).try_get_cached() else {
501 return Err(msg_error_anyhow!(Message::KaslServerTokenMissing));
502 };
503
504 Ok((KaslServer::new(&server_config)?, token))
505}
506
507fn disconnect() -> Result<()> {
513 let mut config = Config::read().unwrap_or_default();
514
515 if let Err(error) = Secret::new(AGENT_TOKEN_SECRET, AGENT_TOKEN_PROMPT).delete() {
522 msg_warning!(Message::KaslServerTokenNotRemoved(error.to_string()));
523 }
524
525 if config.kasl_server.take().is_some() {
526 config.save()?;
527 msg_success!(Message::KaslServerDisconnected);
528 } else {
529 msg_print!(Message::KaslServerNotConnected);
531 }
532
533 Ok(())
534}