use std::io::{BufRead, BufReader, Write};
use std::os::unix::net::{UnixListener, UnixStream};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use crate::protocol::{Request, Response, TimerInfo, TimerStatus, socket_path};
use crate::storage::{data_dir, load_tasks, save_tasks};
use crate::task::Task;
use crate::timer::{Timer, TimerState};
struct ServiceState {
tasks: Vec<Task>,
timer: Timer,
active_task: Option<usize>,
data_dir: std::path::PathBuf,
}
impl ServiceState {
fn new() -> Self {
let data_dir = data_dir();
let tasks = load_tasks(&data_dir);
Self {
tasks,
timer: Timer::new(),
active_task: None,
data_dir,
}
}
fn save(&self) {
save_tasks(&self.data_dir, &self.tasks);
}
fn timer_info(&self) -> TimerInfo {
TimerInfo {
state: match self.timer.state {
TimerState::Idle => TimerStatus::Idle,
TimerState::Running => TimerStatus::Running,
TimerState::Paused => TimerStatus::Paused,
TimerState::Finished => TimerStatus::Finished,
},
task_index: self.active_task,
remaining_secs: self.timer.remaining_secs(),
total_secs: self.timer.total_secs,
}
}
#[allow(clippy::too_many_lines)]
fn handle(&mut self, req: Request) -> Response {
match req {
Request::ListTasks | Request::Status => Response::ok()
.with_tasks(self.tasks.clone())
.with_timer(self.timer_info()),
Request::AddTask {
name,
duration_minutes,
} => {
self.tasks.push(Task::new(name, duration_minutes));
self.save();
Response::ok()
.with_tasks(self.tasks.clone())
.with_timer(self.timer_info())
}
Request::RemoveTask { index } => {
if index >= self.tasks.len() {
return Response::error(format!("Index {index} out of range"));
}
if self.active_task == Some(index) {
self.timer.stop();
self.active_task = None;
} else if let Some(ai) = self.active_task
&& ai > index
{
self.active_task = Some(ai - 1);
}
let removed = self.tasks.remove(index);
self.save();
Response::ok_msg(format!("Removed: \"{}\"", removed.name))
.with_tasks(self.tasks.clone())
.with_timer(self.timer_info())
}
Request::ToggleDone { index } => {
if index >= self.tasks.len() {
return Response::error(format!("Index {index} out of range"));
}
self.tasks[index].completed = !self.tasks[index].completed;
self.save();
Response::ok()
.with_tasks(self.tasks.clone())
.with_timer(self.timer_info())
}
Request::MarkDone { index } => {
if index >= self.tasks.len() {
return Response::error(format!("Index {index} out of range"));
}
self.tasks[index].completed = true;
self.save();
Response::ok_msg(format!("Marked done: \"{}\"", self.tasks[index].name))
.with_tasks(self.tasks.clone())
.with_timer(self.timer_info())
}
Request::ClearCompleted => {
let before = self.tasks.len();
// Adjust active_task index if needed
if let Some(ai) = self.active_task {
if self.tasks[ai].completed {
self.timer.stop();
self.active_task = None;
} else {
let removed_before =
self.tasks[..ai].iter().filter(|t| t.completed).count();
self.active_task = Some(ai - removed_before);
}
}
self.tasks.retain(|t| !t.completed);
let cleared = before - self.tasks.len();
self.save();
Response::ok_msg(format!("Cleared {cleared} completed task(s)"))
.with_tasks(self.tasks.clone())
.with_timer(self.timer_info())
}
Request::StartTimer { index } => {
if index >= self.tasks.len() {
return Response::error(format!("Index {index} out of range"));
}
let secs = self.tasks[index].duration_secs();
self.timer.start(secs);
self.active_task = Some(index);
Response::ok()
.with_tasks(self.tasks.clone())
.with_timer(self.timer_info())
}
Request::PauseTimer => {
self.timer.pause();
Response::ok().with_timer(self.timer_info())
}
Request::ResumeTimer => {
self.timer.resume();
Response::ok().with_timer(self.timer_info())
}
Request::StopTimer => {
self.timer.stop();
self.active_task = None;
Response::ok().with_timer(self.timer_info())
}
Request::Shutdown => Response::ok_msg("Shutting down"),
}
}
}
/// Run the service in the foreground (blocks forever).
pub fn start() {
let sock_path = socket_path();
// Remove stale socket if no service is listening
if sock_path.exists() {
if UnixStream::connect(&sock_path).is_ok() {
eprintln!("Service is already running");
return;
}
let _ = std::fs::remove_file(&sock_path);
}
let listener = match UnixListener::bind(&sock_path) {
Ok(l) => l,
Err(e) => {
eprintln!("Failed to bind socket: {e}");
std::process::exit(1);
}
};
let state = Arc::new(Mutex::new(ServiceState::new()));
// Timer tick thread — checks every second for expiry and sends notifications
let tick_state = Arc::clone(&state);
std::thread::spawn(move || timer_tick_loop(&tick_state));
// Accept connections
for stream in listener.incoming() {
match stream {
Ok(stream) => {
let st = Arc::clone(&state);
std::thread::spawn(move || handle_connection(&stream, &st));
}
Err(e) => eprintln!("Connection error: {e}"),
}
}
let _ = std::fs::remove_file(&sock_path);
}
fn timer_tick_loop(state: &Arc<Mutex<ServiceState>>) {
loop {
std::thread::sleep(Duration::from_secs(1));
let mut s = state.lock().expect("lock poisoned");
if s.timer.tick() {
// Timer just finished
let task_name = s
.active_task
.and_then(|i| s.tasks.get(i))
.map(|t| t.name.clone());
if let Some(idx) = s.active_task
&& idx < s.tasks.len()
{
s.tasks[idx].completed = true;
s.save();
}
// Drop lock before showing notification (could block briefly)
drop(s);
let body = task_name.map_or_else(
|| "Timer finished.".to_string(),
|name| format!("Finished: {name}. Take a break or start the next task."),
);
let _ = notify_rust::Notification::new()
.summary("\u{23f0} Time's up!")
.body(&body)
.timeout(notify_rust::Timeout::Milliseconds(10_000))
.show();
}
}
}
fn handle_connection(stream: &UnixStream, state: &Arc<Mutex<ServiceState>>) {
let mut reader = BufReader::new(stream);
let mut line = String::new();
if reader.read_line(&mut line).is_err() {
return;
}
let request: Request = match serde_json::from_str(&line) {
Ok(r) => r,
Err(e) => {
let resp = Response::error(format!("Invalid request: {e}"));
let _ = write_response(stream, &resp);
return;
}
};
let is_shutdown = matches!(request, Request::Shutdown);
let response = {
let mut s = state.lock().expect("lock poisoned");
s.handle(request)
};
let _ = write_response(stream, &response);
if is_shutdown {
let _ = std::fs::remove_file(socket_path());
std::process::exit(0);
}
}
fn write_response(mut stream: &UnixStream, resp: &Response) -> std::io::Result<()> {
let mut payload = serde_json::to_string(resp).unwrap_or_default();
payload.push('\n');
stream.write_all(payload.as_bytes())
}
/// CLI entry point for `hyperfocus service <start|stop|status>`.
pub fn run_cli(args: &[String]) {
match args.first().map(String::as_str) {
Some("start") => start(),
Some("stop") => match crate::client::Client::send(&Request::Shutdown) {
Ok(_) => println!("Service stopped"),
Err(e) => eprintln!("Failed to stop service: {e}"),
},
Some("status") => {
if UnixStream::connect(socket_path()).is_ok() {
println!("Service is running");
} else {
println!("Service is not running");
}
}
_ => {
println!("Usage: hyperfocus service <start|stop|status>");
}
}
}