aboutsummaryrefslogtreecommitdiffstats
path: root/crates/daemon/src
diff options
context:
space:
mode:
authorBenedikt Peetz <benedikt.peetz@b-peetz.de>2026-07-20 19:30:40 +0200
committerBenedikt Peetz <benedikt.peetz@b-peetz.de>2026-07-20 19:30:40 +0200
commit966a80c4199a49898cc7d8641012d520ce6b2efa (patch)
tree51029ff75842090fd1eecbea97b6f7c447e3dea9 /crates/daemon/src
parentchore(server): Remove warnings (diff)
downloadatuin-966a80c4199a49898cc7d8641012d520ce6b2efa.zip
chore: Commit
Diffstat (limited to 'crates/daemon/src')
-rw-r--r--crates/daemon/src/aclient/database/mod.rs19
-rw-r--r--crates/daemon/src/aclient/history/builder.rs154
-rw-r--r--crates/daemon/src/aclient/history/mod.rs329
-rw-r--r--crates/daemon/src/aclient/history/store.rs4
-rw-r--r--crates/daemon/src/aclient/mod.rs11
-rw-r--r--crates/daemon/src/aclient/ordering.rs3
-rw-r--r--crates/daemon/src/aclient/secrets.rs223
-rw-r--r--crates/daemon/src/aclient/utils.rs15
-rw-r--r--crates/daemon/src/api/client/mod.rs272
-rw-r--r--crates/daemon/src/api/control.rs (renamed from crates/daemon/src/api/server/control.rs)16
-rw-r--r--crates/daemon/src/api/generated.rs28
-rw-r--r--crates/daemon/src/api/history.rs (renamed from crates/daemon/src/api/server/history.rs)13
-rw-r--r--crates/daemon/src/api/mod.rs8
-rw-r--r--crates/daemon/src/api/server/mod.rs2
-rw-r--r--crates/daemon/src/events.rs2
-rw-r--r--crates/daemon/src/lib.rs161
-rw-r--r--crates/daemon/src/main.rs167
-rw-r--r--crates/daemon/src/server.rs17
18 files changed, 220 insertions, 1224 deletions
diff --git a/crates/daemon/src/aclient/database/mod.rs b/crates/daemon/src/aclient/database/mod.rs
index 5287807c..cdf71065 100644
--- a/crates/daemon/src/aclient/database/mod.rs
+++ b/crates/daemon/src/aclient/database/mod.rs
@@ -3,7 +3,6 @@ use std::{
str::FromStr,
};
-use crate::aclient::utils::setup_db;
use fs_err::{self as fs};
use itertools::Itertools;
use sql_builder::{SqlBuilder, SqlName, bind::Bind, esc, quote};
@@ -13,19 +12,16 @@ use sqlx::{
};
use time::OffsetDateTime;
use tracing::debug;
+use turtle::history::{History, HistoryId, get_host_user};
use turtle_common::utils;
use uuid::Uuid;
use crate::aclient::{
- history::{HistoryId, HistoryStats},
- utils::get_host_user,
-};
-
-use super::{
- history::History,
+ history::HistoryStats,
ordering,
- settings::{FilterMode, SearchMode, Settings},
+ settings::{FilterMode, SearchMode},
};
+use crate::aclient::{settings::Settings, utils::setup_db};
#[derive(Clone)]
pub(crate) struct Context {
@@ -858,10 +854,15 @@ mod test {
}
async fn new_history_item(db: &mut ClientSqlite, cmd: &str) -> Result<()> {
- let mut captured: History = History::capture()
+ const SESSION: &str = "test";
+ const HOSTNAME: &str = "test.host";
+
+ let mut captured: History = History::daemon()
.timestamp(OffsetDateTime::now_utc())
.command(cmd)
.cwd("/home/ellie")
+ .session(SESSION)
+ .hostname(HOSTNAME)
.build()
.into();
diff --git a/crates/daemon/src/aclient/history/builder.rs b/crates/daemon/src/aclient/history/builder.rs
deleted file mode 100644
index ef52637b..00000000
--- a/crates/daemon/src/aclient/history/builder.rs
+++ /dev/null
@@ -1,154 +0,0 @@
-use typed_builder::TypedBuilder;
-
-use super::History;
-
-/// Builder for a history entry that is imported from shell history.
-///
-/// The only two required fields are `timestamp` and `command`.
-#[derive(Debug, Clone, TypedBuilder)]
-pub(crate) struct HistoryImported {
- timestamp: time::OffsetDateTime,
- #[builder(setter(into))]
- command: String,
- #[builder(default = "unknown".into(), setter(into))]
- cwd: String,
- #[builder(default = -1)]
- exit: i64,
- #[builder(default = -1)]
- duration: i64,
- #[builder(default, setter(strip_option, into))]
- session: Option<String>,
- #[builder(default, setter(strip_option, into))]
- hostname: Option<String>,
- #[builder(default, setter(strip_option, into))]
- author: Option<String>,
- #[builder(default, setter(strip_option, into))]
- intent: Option<String>,
-}
-
-impl From<HistoryImported> for History {
- fn from(imported: HistoryImported) -> Self {
- Self::new(
- imported.timestamp,
- imported.command,
- imported.cwd,
- imported.exit,
- imported.duration,
- imported.session,
- imported.hostname,
- imported.author,
- imported.intent,
- None,
- )
- }
-}
-
-/// Builder for a history entry that is captured via hook.
-///
-/// This builder is used only at the `start` step of the hook,
-/// so it doesn't have any fields which are known only after
-/// the command is finished, such as `exit` or `duration`.
-#[derive(Debug, Clone, TypedBuilder)]
-pub struct HistoryCaptured {
- timestamp: time::OffsetDateTime,
- #[builder(setter(into))]
- command: String,
- #[builder(setter(into))]
- cwd: String,
- #[builder(default, setter(strip_option, into))]
- author: Option<String>,
- #[builder(default, setter(strip_option, into))]
- intent: Option<String>,
-}
-
-impl From<HistoryCaptured> for History {
- fn from(captured: HistoryCaptured) -> Self {
- Self::new(
- captured.timestamp,
- captured.command,
- captured.cwd,
- -1,
- -1,
- None,
- None,
- captured.author,
- captured.intent,
- None,
- )
- }
-}
-
-/// Builder for a history entry that is loaded from the database.
-///
-/// All fields are required, as they are all present in the database.
-#[derive(Debug, Clone, TypedBuilder)]
-pub(crate) struct HistoryFromDb {
- id: String,
- timestamp: time::OffsetDateTime,
- command: String,
- cwd: String,
- exit: i64,
- duration: i64,
- session: String,
- hostname: String,
- author: String,
- intent: Option<String>,
- deleted_at: Option<time::OffsetDateTime>,
-}
-
-impl From<HistoryFromDb> for History {
- fn from(from_db: HistoryFromDb) -> Self {
- Self {
- id: from_db.id.into(),
- timestamp: from_db.timestamp,
- exit: from_db.exit,
- command: from_db.command,
- cwd: from_db.cwd,
- duration: from_db.duration,
- session: from_db.session,
- hostname: from_db.hostname,
- author: from_db.author,
- intent: from_db.intent,
- deleted_at: from_db.deleted_at,
- }
- }
-}
-
-/// Builder for a history entry that is captured via hook and sent to the daemon
-///
-/// This builder is similar to Capture, but we just require more information up front.
-/// For the old setup, we could just rely on `History::new` to read some of the missing
-/// data. This is no longer the case.
-#[derive(Debug, Clone, TypedBuilder)]
-pub(crate) struct HistoryDaemonCapture {
- timestamp: time::OffsetDateTime,
- #[builder(setter(into))]
- command: String,
- #[builder(setter(into))]
- cwd: String,
- #[builder(setter(into))]
- session: String,
- #[builder(setter(into))]
- hostname: String,
- #[builder(default, setter(strip_option, into))]
- author: Option<String>,
- #[builder(default, setter(strip_option, into))]
- intent: Option<String>,
-}
-
-impl From<HistoryDaemonCapture> for History {
- fn from(captured: HistoryDaemonCapture) -> Self {
- Self::new(
- captured.timestamp,
- captured.command,
- captured.cwd,
- -1,
- -1,
- Some(captured.session),
- Some(captured.hostname),
- captured.author,
- captured.intent,
- None,
- )
- }
-}
diff --git a/crates/daemon/src/aclient/history/mod.rs b/crates/daemon/src/aclient/history/mod.rs
index 527d1bb6..ea71dcf4 100644
--- a/crates/daemon/src/aclient/history/mod.rs
+++ b/crates/daemon/src/aclient/history/mod.rs
@@ -5,17 +5,15 @@ use rmp::decode::ValueReadError;
use rmp::{Marker, decode::Bytes};
use std::env;
use std::fmt::Display;
+use turtle::history::History;
use turtle_common::record::DecryptedData;
use turtle_common::utils::uuid_v7;
use eyre::{Result, bail, eyre};
-use crate::aclient::secrets::SECRET_PATTERNS_RE;
-use crate::aclient::utils::get_host_user;
use time::OffsetDateTime;
-mod builder;
pub(crate) mod store;
pub(crate) const HISTORY_VERSION_V0: &str = "v0";
@@ -27,72 +25,6 @@ pub(crate) const HISTORY_TAG: &str = "history";
const HISTORY_AUTHOR_ENV: &str = "ATUIN_HISTORY_AUTHOR";
const HISTORY_INTENT_ENV: &str = "ATUIN_HISTORY_INTENT";
-#[derive(Clone, Debug, Eq, PartialEq, Hash)]
-pub struct HistoryId(pub String);
-
-impl Display for HistoryId {
- fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
- write!(f, "{}", self.0)
- }
-}
-
-impl From<String> for HistoryId {
- fn from(s: String) -> Self {
- Self(s)
- }
-}
-
-/// Client-side history entry.
-///
-/// Client stores data unencrypted, and only encrypts it before sending to the server.
-///
-/// To create a new history entry, use one of the builders:
-/// - [`History::import()`] to import an entry from the shell history file
-/// - [`History::capture()`] to capture an entry via hook
-/// - [`History::from_db()`] to create an instance from the database entry
-//
-// ## Implementation Notes
-//
-// New fields must be added to `History::{serialize,deserialize}` in a backwards
-// compatible way (sensible defaults and careful `nfields` handling).
-#[derive(Debug, Clone, PartialEq, Eq, sqlx::FromRow)]
-pub struct History {
- /// A client-generated ID, used to identify the entry when syncing.
- ///
- /// Stored as `client_id` in the database.
- pub id: HistoryId,
-
- /// When the command was run.
- pub timestamp: OffsetDateTime,
-
- /// How long the command took to run.
- pub duration: i64,
-
- /// The exit code of the command.
- pub exit: i64,
-
- /// The command that was run.
- pub command: String,
-
- /// The current working directory when the command was run.
- pub cwd: String,
-
- /// The session ID, associated with a terminal session.
- pub session: String,
-
- /// The hostname of the machine the command was run on.
- pub hostname: String,
-
- /// Who wrote this command (human user or automation/agent identity).
- pub author: String,
-
- /// Optional rationale for why the command was executed.
- pub intent: Option<String>,
-
- /// Timestamp, which is set when the entry is deleted, allowing a soft delete.
- pub deleted_at: Option<OffsetDateTime>,
-}
-
#[derive(Debug, Clone, PartialEq, Eq, sqlx::FromRow)]
pub(crate) struct HistoryStats {
/// The command that was ran after this one in the session
@@ -113,63 +45,17 @@ pub(crate) struct HistoryStats {
pub(crate) duration_over_time: Vec<(String, i64)>,
}
-impl History {
- pub(crate) fn author_from_hostname(hostname: &str) -> String {
- hostname
- .split_once(':')
- .map_or_else(|| hostname.to_owned(), |(_, user)| user.to_owned())
- }
-
- fn normalize_optional_field(field: Option<String>) -> Option<String> {
- field.and_then(|value| {
- let trimmed = value.trim();
- if trimmed.is_empty() {
- None
- } else {
- Some(trimmed.to_owned())
- }
- })
- }
-
- #[expect(clippy::too_many_arguments)]
- fn new(
- timestamp: OffsetDateTime,
- command: String,
- cwd: String,
- exit: i64,
- duration: i64,
- session: Option<String>,
- hostname: Option<String>,
- author: Option<String>,
- intent: Option<String>,
- deleted_at: Option<OffsetDateTime>,
- ) -> Self {
- let session = session
- .or_else(|| env::var("ATUIN_SESSION").ok())
- .unwrap_or_else(|| uuid_v7().as_simple().to_string());
- let hostname = hostname.unwrap_or_else(get_host_user);
- let author = Self::normalize_optional_field(author)
- .or_else(|| Self::normalize_optional_field(env::var(HISTORY_AUTHOR_ENV).ok()))
- .unwrap_or_else(|| Self::author_from_hostname(hostname.as_str()));
- let intent = Self::normalize_optional_field(intent)
- .or_else(|| Self::normalize_optional_field(env::var(HISTORY_INTENT_ENV).ok()));
-
- Self {
- id: uuid_v7().as_simple().to_string().into(),
- timestamp,
- command,
- cwd,
- exit,
- duration,
- session,
- hostname,
- author,
- intent,
- deleted_at,
- }
- }
+pub(crate) trait HistoryExt: Sized {
+ fn serialize(&self) -> Result<DecryptedData>;
+ fn read_optional_string(bytes: &[u8]) -> Result<(Option<String>, &[u8])>;
+ fn deserialize_v0(bytes: &[u8]) -> Result<Self>;
+ fn deserialize_v1(bytes: &[u8]) -> Result<Self>;
+ fn deserialize(bytes: &[u8], version: &str) -> Result<Self>;
+ fn success(&self) -> bool;
+}
- pub(crate) fn serialize(&self) -> Result<DecryptedData> {
+impl HistoryExt for History {
+ fn serialize(&self) -> Result<DecryptedData> {
// This is pretty much the same as what we used for the old history, with one difference -
// it uses integers for timestamps rather than a string format.
@@ -358,7 +244,7 @@ impl History {
})
}
- pub(crate) fn deserialize(bytes: &[u8], version: &str) -> Result<Self> {
+ fn deserialize(bytes: &[u8], version: &str) -> Result<Self> {
match version {
HISTORY_VERSION_V0 => Self::deserialize_v0(bytes),
HISTORY_VERSION_V1 => Self::deserialize_v1(bytes),
@@ -367,118 +253,10 @@ impl History {
}
}
- /// Builder for a history entry that is captured via hook.
- ///
- /// This builder is used only at the `start` step of the hook,
- /// so it doesn't have any fields which are known only after
- /// the command is finished, such as `exit` or `duration`.
- ///
- /// ## Examples
- /// ```rust
- /// use crate::aclient::history::History;
- ///
- /// let history: History = History::capture()
- /// .timestamp(time::OffsetDateTime::now_utc())
- /// .command("ls -la")
- /// .cwd("/home/user")
- /// .build()
- /// .into();
- /// ```
- ///
- /// Command without any required info cannot be captured, which is forced at compile time:
- ///
- /// ```compile_fail
- /// use crate::aclient::history::History;
- ///
- /// // this will not compile because `cwd` is missing
- /// let history: History = History::capture()
- /// .timestamp(time::OffsetDateTime::now_utc())
- /// .command("ls -la")
- /// .build()
- /// .into();
- /// ```
- pub fn capture() -> builder::HistoryCapturedBuilder {
- builder::HistoryCaptured::builder()
- }
-
- /// Builder for a history entry that is captured via hook, and sent to the daemon.
- ///
- /// This builder is used only at the `start` step of the hook,
- /// so it doesn't have any fields which are known only after
- /// the command is finished, such as `exit` or `duration`.
- ///
- /// It does, however, include information that can usually be inferred.
- ///
- /// This is because the daemon we are sending a request to lacks the context of the command
- ///
- /// ## Examples
- /// ```rust
- /// use crate::aclient::history::History;
- ///
- /// let history: History = History::daemon()
- /// .timestamp(time::OffsetDateTime::now_utc())
- /// .command("ls -la")
- /// .cwd("/home/user")
- /// .session("018deb6e8287781f9973ef40e0fde76b")
- /// .hostname("computer:ellie")
- /// .build()
- /// .into();
- /// ```
- ///
- /// Command without any required info cannot be captured, which is forced at compile time:
- ///
- /// ```compile_fail
- /// use crate::aclient::history::History;
- ///
- /// // this will not compile because `hostname` is missing
- /// let history: History = History::daemon()
- /// .timestamp(time::OffsetDateTime::now_utc())
- /// .command("ls -la")
- /// .cwd("/home/user")
- /// .session("018deb6e8287781f9973ef40e0fde76b")
- /// .build()
- /// .into();
- /// ```
- pub(crate) fn daemon() -> builder::HistoryDaemonCaptureBuilder {
- builder::HistoryDaemonCapture::builder()
- }
-
- /// Builder for a history entry that is imported from the database.
- ///
- /// All fields are required, as they are all present in the database.
- ///
- /// ```compile_fail
- /// use crate::aclient::history::History;
- ///
- /// // this will not compile because `id` field is missing
- /// let history: History = History::from_db()
- /// .timestamp(time::OffsetDateTime::now_utc())
- /// .command("ls -la".to_string())
- /// .cwd("/home/user".to_string())
- /// .exit(0)
- /// .duration(100)
- /// .session("somesession".to_string())
- /// .hostname("localhost".to_string())
- /// .author("user".to_string())
- /// .intent(None)
- /// .deleted_at(None)
- /// .build()
- /// .into();
- /// ```
- pub(crate) fn from_db() -> builder::HistoryFromDbBuilder {
- builder::HistoryFromDb::builder()
- }
-
- pub(crate) fn success(&self) -> bool {
+ #[expect(unused)]
+ fn success(&self) -> bool {
self.exit == 0 || self.duration == -1
}
-
- pub fn should_save(&self, filter: SettingsFilter<'_>) -> bool {
- !(self.command.is_empty()
- || filter.history.is_match(&self.command)
- || filter.cwd.is_match(&self.cwd)
- || (filter.secrets && SECRET_PATTERNS_RE.is_match(&self.command)))
- }
}
#[derive(Debug, Copy, Clone)]
@@ -490,89 +268,12 @@ pub struct SettingsFilter<'a> {
#[cfg(test)]
mod tests {
- use regex::RegexSet;
use time::macros::datetime;
- use crate::aclient::{history::HISTORY_VERSION, settings::Settings};
+ use crate::aclient::history::{HISTORY_VERSION, HistoryExt};
use super::History;
- // Test that we don't save history where necessary
- #[test]
- fn privacy_test() {
- let settings = Settings {
- cwd_filter: RegexSet::new(["^/supasecret"]).unwrap(),
- history_filter: RegexSet::new(["^psql"]).unwrap(),
- ..Settings::default()
- };
-
- let normal_command: History = History::capture()
- .timestamp(time::OffsetDateTime::now_utc())
- .command("echo foo")
- .cwd("/")
- .build()
- .into();
-
- let with_space: History = History::capture()
- .timestamp(time::OffsetDateTime::now_utc())
- .command(" echo bar")
- .cwd("/")
- .build()
- .into();
-
- let empty: History = History::capture()
- .timestamp(time::OffsetDateTime::now_utc())
- .command("")
- .cwd("/")
- .build()
- .into();
-
- let stripe_key: History = History::capture()
- .timestamp(time::OffsetDateTime::now_utc())
- .command("curl foo.com/bar?key=sk_test_1234567890abcdefghijklmnop")
- .cwd("/")
- .build()
- .into();
-
- let secret_dir: History = History::capture()
- .timestamp(time::OffsetDateTime::now_utc())
- .command("echo ohno")
- .cwd("/supasecret")
- .build()
- .into();
-
- let with_psql: History = History::capture()
- .timestamp(time::OffsetDateTime::now_utc())
- .command("psql")
- .cwd("/supasecret")
- .build()
- .into();
-
- assert!(normal_command.should_save(&settings));
- assert!(!with_space.should_save(&settings));
- assert!(!empty.should_save(&settings));
- assert!(!stripe_key.should_save(&settings));
- assert!(!secret_dir.should_save(&settings));
- assert!(!with_psql.should_save(&settings));
- }
-
- #[test]
- fn disable_secrets() {
- let settings = Settings {
- secrets_filter: false,
- ..Settings::new().unwrap()
- };
-
- let stripe_key: History = History::capture()
- .timestamp(time::OffsetDateTime::now_utc())
- .command("curl foo.com/bar?key=sk_test_1234567890abcdefghijklmnop")
- .cwd("/")
- .build()
- .into();
-
- assert!(stripe_key.should_save(&settings));
- }
-
#[test]
fn test_serialize_deserialize() {
let history = History {
diff --git a/crates/daemon/src/aclient/history/store.rs b/crates/daemon/src/aclient/history/store.rs
index 216f85d6..c62f9068 100644
--- a/crates/daemon/src/aclient/history/store.rs
+++ b/crates/daemon/src/aclient/history/store.rs
@@ -2,14 +2,16 @@ use std::collections::HashSet;
use eyre::{Result, bail, eyre};
use rmp::decode::Bytes;
+use turtle::history::{History, HistoryId};
use crate::aclient::{
database::ClientSqlite,
+ history::HistoryExt,
record::{encryption::PASETO_V4, sqlite_store::SqliteStore},
};
use turtle_common::record::{DecryptedData, Host, HostId, Record, RecordId, RecordIdx};
-use super::{HISTORY_TAG, HISTORY_VERSION, HISTORY_VERSION_V0, History, HistoryId};
+use super::{HISTORY_TAG, HISTORY_VERSION, HISTORY_VERSION_V0};
#[derive(Debug, Clone)]
pub(crate) struct HistoryStore {
diff --git a/crates/daemon/src/aclient/mod.rs b/crates/daemon/src/aclient/mod.rs
index 3f3709a9..f2d14e01 100644
--- a/crates/daemon/src/aclient/mod.rs
+++ b/crates/daemon/src/aclient/mod.rs
@@ -1,11 +1,10 @@
-pub mod database;
+pub(crate) mod database;
pub mod history;
-pub mod record;
-pub mod settings;
+pub(crate) mod record;
+pub(crate) mod settings;
+pub(crate) mod api_client;
pub(crate) mod encryption;
pub(crate) mod meta;
-pub(crate) mod utils;
pub(crate) mod ordering;
-pub(crate) mod secrets;
-pub(crate) mod api_client;
+pub(crate) mod utils;
diff --git a/crates/daemon/src/aclient/ordering.rs b/crates/daemon/src/aclient/ordering.rs
index 84001f52..8fa6498e 100644
--- a/crates/daemon/src/aclient/ordering.rs
+++ b/crates/daemon/src/aclient/ordering.rs
@@ -1,6 +1,7 @@
use minspan::minspan;
+use turtle::history::History;
-use super::{history::History, settings::SearchMode};
+use super::settings::SearchMode;
pub(crate) fn reorder_fuzzy(mode: SearchMode, query: &str, res: Vec<History>) -> Vec<History> {
match mode {
diff --git a/crates/daemon/src/aclient/secrets.rs b/crates/daemon/src/aclient/secrets.rs
deleted file mode 100644
index 08d24339..00000000
--- a/crates/daemon/src/aclient/secrets.rs
+++ /dev/null
@@ -1,223 +0,0 @@
-// This file will probably trigger a lot of scanners. Sorry.
-
-use regex::RegexSet;
-use std::sync::LazyLock;
-
-#[cfg(test)]
-pub(crate) enum TestValue<'a> {
- Single(&'a str),
- Multiple(&'a [&'a str]),
-}
-
-#[cfg(test)]
-type SpType<'a> = &'a [(&'a str, &'a str, TestValue<'a>)];
-
-#[cfg(not(test))]
-type SpType<'a> = &'a [(&'a str, &'a str)];
-
-/// A list of `(name, regex, test)`, where `test` should match against `regex`.
-pub(crate) static SECRET_PATTERNS: SpType<'_> = &[
- (
- "AWS Access Key ID",
- "A[KS]IA[0-9A-Z]{16}",
- #[cfg(test)]
- TestValue::Single("AKIAIOSFODNN7EXAMPLE"),
- ),
- (
- "AWS Secret Access Key env var",
- "AWS_SECRET_ACCESS_KEY",
- #[cfg(test)]
- TestValue::Single("AWS_SECRET_ACCESS_KEY=KEYDATA"),
- ),
- (
- "AWS Session Token env var",
- "AWS_SESSION_TOKEN",
- #[cfg(test)]
- TestValue::Single("AWS_SESSION_TOKEN=KEYDATA"),
- ),
- (
- "Microsoft Azure secret access key env var",
- "AZURE_.*_KEY",
- #[cfg(test)]
- TestValue::Single("export AZURE_STORAGE_ACCOUNT_KEY=KEYDATA"),
- ),
- (
- "Google cloud platform key env var",
- "GOOGLE_SERVICE_ACCOUNT_KEY",
- #[cfg(test)]
- TestValue::Single("export GOOGLE_SERVICE_ACCOUNT_KEY=KEYDATA"),
- ),
- (
- "Atuin login",
- r"atuin\s+login",
- #[cfg(test)]
- TestValue::Single(
- "atuin login -u mycoolusername -p mycoolpassword -k \"lots of random words\"",
- ),
- ),
- (
- "GitHub PAT (old)",
- "ghp_[a-zA-Z0-9]{36}",
- #[cfg(test)]
- TestValue::Single("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"), // legit, I expired it
- ),
- (
- "GitHub PAT (new)",
- "gh1_[A-Za-z0-9]{21}_[A-Za-z0-9]{59}|github_pat_[0-9][A-Za-z0-9]{21}_[A-Za-z0-9]{59}",
- #[cfg(test)]
- TestValue::Multiple(&[
- "gh1_1234567890abcdefghijk_1234567890abcdefghijklmnopqrstuvwxyz1234567890abcdefghijklm",
- "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", // also legit, also expired
- ]),
- ),
- (
- "GitHub OAuth Access Token",
- "gho_[A-Za-z0-9]{36}",
- #[cfg(test)]
- TestValue::Single("gho_1234567890abcdefghijklmnopqrstuvwx000"), // not a real token
- ),
- (
- "GitHub OAuth Access Token (user)",
- "ghu_[A-Za-z0-9]{36}",
- #[cfg(test)]
- TestValue::Single("ghu_1234567890abcdefghijklmnopqrstuvwx000"), // not a real token
- ),
- (
- "GitHub App Installation Access Token",
- "ghs_[A-Za-z0-9._-]{36,}",
- #[cfg(test)]
- TestValue::Multiple(&[
- "ghs_1234567890abcdefghijklmnopqrstuvwx000", // not a real token
- "ghs_abc-def.ghi_jklMNOP0123456789qrstuv-wxyzABCD", // new token format, fake data
- ]),
- ),
- (
- "GitHub Refresh Token",
- "ghr_[A-Za-z0-9]{76}",
- #[cfg(test)]
- TestValue::Single(
- "ghr_1234567890abcdefghijklmnopqrstuvwx1234567890abcdefghijklmnopqrstuvwx1234567890abcdefghijklmnopqrstuvwx",
- ), // not a real token
- ),
- (
- "GitHub App Installation Access Token v1",
- "v1\\.[0-9A-Fa-f]{40}",
- #[cfg(test)]
- TestValue::Single("v1.1234567890abcdef1234567890abcdef12345678"), // not a real token
- ),
- (
- "GitLab PAT",
- "glpat-[a-zA-Z0-9_]{20}",
- #[cfg(test)]
- TestValue::Single("glpat-RkE_BG5p_bbjML21WSfy"),
- ),
- (
- "Slack OAuth v2 bot",
- "xoxb-[0-9]{11}-[0-9]{11}-[0-9a-zA-Z]{24}",
- #[cfg(test)]
- TestValue::Single("xoxb-17653672481-19874698323-pdFZKVeTuE8sk7oOcBrzbqgy"),
- ),
- (
- "Slack OAuth v2 user token",
- "xoxp-[0-9]{11}-[0-9]{11}-[0-9a-zA-Z]{24}",
- #[cfg(test)]
- TestValue::Single("xoxp-17653672481-19874698323-pdFZKVeTuE8sk7oOcBrzbqgy"),
- ),
- (
- "Slack webhook",
- "T[a-zA-Z0-9_]{8}/B[a-zA-Z0-9_]{8}/[a-zA-Z0-9_]{24}",
- #[cfg(test)]
- TestValue::Single(
- "https://hooks.slack.com/services/T00000000/B00000000/XXXXXXXXXXXXXXXXXXXXXXXX",
- ),
- ),
- (
- "Stripe test key",
- "sk_test_[0-9a-zA-Z]{24}",
- #[cfg(test)]
- TestValue::Single("sk_test_1234567890abcdefghijklmnop"),
- ),
- (
- "Stripe live key",
- "sk_live_[0-9a-zA-Z]{24}",
- #[cfg(test)]
- TestValue::Single("sk_live_1234567890abcdefghijklmnop"),
- ),
- (
- "Netlify authentication token",
- "nf[pcoub]_[0-9a-zA-Z]{36}",
- #[cfg(test)]
- TestValue::Single("nfp_nBh7BdJxUwyaBBwFzpyD29MMFT6pZ9wq5634"),
- ),
- (
- "npm token",
- "npm_[A-Za-z0-9]{36}",
- #[cfg(test)]
- TestValue::Single("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"),
- ),
- (
- "Pulumi personal access token",
- "pul-[0-9a-f]{40}",
- #[cfg(test)]
- TestValue::Single("pul-683c2770662c51d960d72ec27613be7653c5cb26"),
- ),
-];
-
-/// The `regex` expressions from [`SECRET_PATTERNS`] compiled into a `RegexSet`.
-pub(crate) static SECRET_PATTERNS_RE: LazyLock<RegexSet> = LazyLock::new(|| {
- let exprs = SECRET_PATTERNS.iter().map(|f| f.1);
- RegexSet::new(exprs).expect("Failed to build secrets regex")
-});
-
-#[cfg(test)]
-mod tests {
- use regex::Regex;
-
- use crate::aclient::secrets::{SECRET_PATTERNS, TestValue};
-
- #[test]
- fn test_secrets() {
- for (name, regex, test) in SECRET_PATTERNS {
- let re =
- Regex::new(regex).unwrap_or_else(|_| panic!("Failed to compile regex for {name}"));
-
- match test {
- TestValue::Single(test) => {
- assert!(re.is_match(test), "{name} test failed!");
- }
- TestValue::Multiple(tests) => {
- for test_str in tests.iter() {
- assert!(
- re.is_match(test_str),
- "{name} test with value \"{test_str}\" failed!"
- );
- }
- }
- }
- }
- }
-
- #[test]
- fn test_secrets_embedded() {
- for (name, regex, test) in SECRET_PATTERNS {
- let re =
- Regex::new(regex).unwrap_or_else(|_| panic!("Failed to compile regex for {name}"));
-
- match test {
- TestValue::Single(test) => {
- let embedded = format!("some random text {test} some more random text");
- assert!(re.is_match(&embedded), "{name} embedded test failed!");
- }
- TestValue::Multiple(tests) => {
- for test_str in tests.iter() {
- let embedded = format!("some random text {test_str} some more random text");
- assert!(
- re.is_match(&embedded),
- "{name} embedded test with value \"{test_str}\" failed!"
- );
- }
- }
- }
- }
- }
-}
diff --git a/crates/daemon/src/aclient/utils.rs b/crates/daemon/src/aclient/utils.rs
index cf515183..18e732a0 100644
--- a/crates/daemon/src/aclient/utils.rs
+++ b/crates/daemon/src/aclient/utils.rs
@@ -1,18 +1,3 @@
-pub(crate) fn get_hostname() -> String {
- std::env::var("ATUIN_HOST_NAME")
- .unwrap_or_else(|_| whoami::hostname().unwrap_or_else(|_| "unknown-host".to_string()))
-}
-
-pub(crate) fn get_username() -> String {
- std::env::var("ATUIN_HOST_USER")
- .unwrap_or_else(|_| whoami::username().unwrap_or_else(|_| "unknown-user".to_string()))
-}
-
-/// Returns a pair of the hostname and username, separated by a colon.
-pub(crate) fn get_host_user() -> String {
- format!("{}:{}", get_hostname(), get_username())
-}
-
/// Setup a [`SQLite`] database.
///
/// This takes care of correct locking, so that we avoid a race when setting up the database.
diff --git a/crates/daemon/src/api/client/mod.rs b/crates/daemon/src/api/client/mod.rs
deleted file mode 100644
index c588fb09..00000000
--- a/crates/daemon/src/api/client/mod.rs
+++ /dev/null
@@ -1,272 +0,0 @@
-use eyre::{Context as EyreContext, Result};
-use time::OffsetDateTime;
-use tonic::Code;
-use tonic::transport::{Channel, Endpoint, Uri};
-use tower::service_fn;
-
-use hyper_util::rt::TokioIo;
-
-#[cfg(unix)]
-use tokio::net::UnixStream;
-
-use crate::api::generated;
-use crate::api::generated::control::{ForceSyncReply, ForceSyncRequest, PathsReply, PathsRequest};
-use crate::api::generated::history::{HistoryEntry, HistoryRequest};
-use crate::{
- aclient::history::History,
- api::{
- DAEMON_PROTOCOL_VERSION, DAEMON_VERSION,
- generated::{
- control::{
- StatusReply, StatusRequest, control_client::ControlClient as ControlServiceClient,
- },
- history::{
- EndHistoryReply, EndHistoryRequest, StartHistoryReply, StartHistoryRequest,
- TailHistoryRequest, history_client::HistoryClient as HistoryServiceClient,
- },
- },
- },
-};
-
-pub use crate::api::generated::history::{HistoryEventKind, TailHistoryReply};
-
-fn normalize_optional_field(value: &str) -> Option<String> {
- let trimmed = value.trim();
- if trimmed.is_empty() {
- None
- } else {
- Some(trimmed.to_owned())
- }
-}
-
-pub fn history_entry_to_history(entry: HistoryEntry) -> History {
- let timestamp = OffsetDateTime::from_unix_timestamp_nanos(i128::from(entry.timestamp))
- .expect("Daemon history timestamp should always be valid");
-
- History {
- id: entry.id.into(),
- timestamp,
- duration: entry.duration,
- exit: entry.exit,
- command: entry.command,
- cwd: entry.cwd,
- session: entry.session,
- hostname: entry.hostname,
- author: entry.author,
- intent: normalize_optional_field(&entry.intent),
- deleted_at: None,
- }
-}
-
-#[must_use]
-pub fn daemon_matches_expected(version: &str, protocol: u32) -> bool {
- version == DAEMON_VERSION && protocol == DAEMON_PROTOCOL_VERSION
-}
-
-#[must_use]
-pub fn daemon_mismatch_message(version: &str, protocol: u32) -> String {
- if protocol == DAEMON_PROTOCOL_VERSION {
- format!("daemon is out of date: expected {DAEMON_VERSION}, got {version}")
- } else {
- format!("daemon protocol mismatch: expected {DAEMON_PROTOCOL_VERSION}, got {protocol}")
- }
-}
-
-#[derive(Clone, Copy, Debug, Eq, PartialEq)]
-pub enum DaemonClientErrorKind {
- Connect,
- Unavailable,
- Unimplemented,
- Other,
-}
-
-#[must_use]
-pub fn classify_error(error: &eyre::Report) -> DaemonClientErrorKind {
- for cause in error.chain() {
- if cause.downcast_ref::<tonic::transport::Error>().is_some() {
- return DaemonClientErrorKind::Connect;
- }
-
- if let Some(status) = cause.downcast_ref::<tonic::Status>() {
- return match status.code() {
- Code::Unavailable => DaemonClientErrorKind::Unavailable,
- Code::Unimplemented => DaemonClientErrorKind::Unimplemented,
- _ => DaemonClientErrorKind::Other,
- };
- }
- }
-
- DaemonClientErrorKind::Other
-}
-
-#[derive(Debug)]
-pub enum Probe {
- Ready(ControlClient),
- NeedsRestart(String),
- Unreachable(eyre::Report),
-}
-
-/// Check if a client can reach the daemon.
-pub async fn probe(path: String) -> Probe {
- let mut client = match ControlClient::new(path).await {
- Ok(client) => client,
- Err(err) => return Probe::Unreachable(err),
- };
-
- match client.status().await {
- Ok(status) => {
- if daemon_matches_expected(&status.version, status.protocol) {
- Probe::Ready(client)
- } else {
- Probe::NeedsRestart(daemon_mismatch_message(&status.version, status.protocol))
- }
- }
- Err(err) => Probe::Unreachable(err),
- }
-}
-
-// ============================================================================
-// History Client
-// ============================================================================
-
-#[derive(Debug)]
-pub struct HistoryClient {
- client: HistoryServiceClient<Channel>,
-}
-
-pub struct Range {
- pub start: OffsetDateTime,
- pub end: OffsetDateTime,
-}
-
-// Wrap the grpc client
-impl HistoryClient {
- #[cfg(unix)]
- pub async fn new(path: String) -> Result<Self> {
- use eyre::Context;
-
- let log_path = path.clone();
- let channel = Endpoint::try_from("http://atuin_local_daemon:0")?
- .connect_with_connector(service_fn(move |_: Uri| {
- let path = path.clone();
-
- async move {
- Ok::<_, std::io::Error>(TokioIo::new(UnixStream::connect(path.clone()).await?))
- }
- }))
- .await
- .wrap_err_with(|| {
- format!(
- "failed to connect to local atuin daemon at {}. Is it running?",
- &log_path
- )
- })?;
-
- let client = HistoryServiceClient::new(channel);
-
- Ok(Self { client })
- }
-
- pub async fn start_history(&mut self, h: History) -> Result<StartHistoryReply> {
- let req = StartHistoryRequest {
- command: h.command,
- cwd: h.cwd,
- hostname: h.hostname,
- session: h.session,
- timestamp: h.timestamp.unix_timestamp_nanos() as u64,
- author: h.author,
- intent: h.intent.unwrap_or_default(),
- };
-
- Ok(self.client.start_history(req).await?.into_inner())
- }
-
- pub async fn history(&mut self, session: String, range: Option<Range>) -> Result<Vec<History>> {
- let req = HistoryRequest {
- session,
- range: range.map(|r| generated::history::Range {
- start: r.start.unix_timestamp() as u64,
- end: r.end.unix_timestamp() as u64,
- }),
- };
-
- let reply = self.client.history(req).await?.into_inner();
-
- Ok(reply
- .entries
- .into_iter()
- .map(history_entry_to_history)
- .collect())
- }
-
- pub async fn end_history(
- &mut self,
- id: String,
- duration: u64,
- exit: i64,
- ) -> Result<EndHistoryReply> {
- let req = EndHistoryRequest { id, exit, duration };
-
- Ok(self.client.end_history(req).await?.into_inner())
- }
-
- pub async fn tail_history(&mut self) -> Result<tonic::Streaming<TailHistoryReply>> {
- Ok(self
- .client
- .tail_history(TailHistoryRequest {})
- .await?
- .into_inner())
- }
-}
-
-// ============================================================================
-// Control Client
-// ============================================================================
-
-/// Client for the Control gRPC service.
-#[derive(Debug)]
-pub struct ControlClient {
- client: ControlServiceClient<Channel>,
-}
-
-impl ControlClient {
- /// Connect to the daemon's control service.
- pub async fn new(path: String) -> Result<Self> {
- let log_path = path.clone();
- let channel = Endpoint::try_from("http://atuin_local_daemon:0")?
- .connect_with_connector(service_fn(move |_: Uri| {
- let path = path.clone();
-
- async move {
- Ok::<_, std::io::Error>(TokioIo::new(UnixStream::connect(path.clone()).await?))
- }
- }))
- .await
- .wrap_err_with(|| {
- format!(
- "failed to connect to local atuin daemon at {}. Is it running?",
- &log_path
- )
- })?;
-
- let client = ControlServiceClient::new(channel);
-
- Ok(Self { client })
- }
-
- pub async fn paths(&mut self) -> Result<PathsReply> {
- Ok(self.client.paths(PathsRequest {}).await?.into_inner())
- }
-
- pub async fn force_sync(&mut self) -> Result<ForceSyncReply> {
- Ok(self
- .client
- .force_sync(ForceSyncRequest {})
- .await?
- .into_inner())
- }
-
- pub async fn status(&mut self) -> Result<StatusReply> {
- Ok(self.client.status(StatusRequest {}).await?.into_inner())
- }
-}
diff --git a/crates/daemon/src/api/server/control.rs b/crates/daemon/src/api/control.rs
index a5e26355..a9d9cff3 100644
--- a/crates/daemon/src/api/server/control.rs
+++ b/crates/daemon/src/api/control.rs
@@ -6,15 +6,17 @@ use tokio::time::{self, MissedTickBehavior};
use tonic::{Request, Response, Status};
use tracing::{Level, instrument};
+use turtle::generated::{
+ DAEMON_PROTOCOL_VERSION,
+ control::{
+ ForceSyncReply, ForceSyncRequest, PathsReply, PathsRequest, StatusReply, StatusRequest,
+ control_server::{Control, ControlServer},
+ },
+};
+
use crate::{
+ DAEMON_VERSION,
aclient::{history::store::HistoryStore, record::sync, settings::Settings},
- api::{
- DAEMON_PROTOCOL_VERSION, DAEMON_VERSION,
- generated::control::{
- ForceSyncReply, ForceSyncRequest, PathsReply, PathsRequest, StatusReply, StatusRequest,
- control_server::{Control, ControlServer},
- },
- },
daemon::DaemonHandle,
events::DaemonEvent,
};
diff --git a/crates/daemon/src/api/generated.rs b/crates/daemon/src/api/generated.rs
deleted file mode 100644
index 304edcd9..00000000
--- a/crates/daemon/src/api/generated.rs
+++ /dev/null
@@ -1,28 +0,0 @@
-#![expect(
- unreachable_pub,
- unused_qualifications,
- clippy::doc_markdown,
- clippy::default_trait_access,
- clippy::too_many_lines,
- clippy::trivially_copy_pass_by_ref,
- clippy::allow_attributes,
- clippy::derive_partial_eq_without_eq,
- reason = "All of these lints are triggered by the generated code"
-)]
-
-/// Semantic command capture gRPC service types.
-pub(crate) mod semantic {
- tonic::include_proto!("semantic");
-}
-
-/// History module for the daemon gRPC history service.
-///
-/// This module contains the proto-generated types for the history gRPC service.
-pub(crate) mod history {
- tonic::include_proto!("history");
-}
-
-/// Control module for external control.
-pub(crate) mod control {
- tonic::include_proto!("control");
-}
diff --git a/crates/daemon/src/api/server/history.rs b/crates/daemon/src/api/history.rs
index 0edf3b94..bcd2ee5a 100644
--- a/crates/daemon/src/api/server/history.rs
+++ b/crates/daemon/src/api/history.rs
@@ -10,20 +10,23 @@ use tracing::{Level, instrument};
use crate::{
aclient::{
database::{ClientSqlite, current_context},
- history::{History, HistoryId, store::HistoryStore},
+ history::store::HistoryStore,
settings::Settings,
},
- api::{
+ daemon::DaemonHandle,
+ events::DaemonEvent,
+};
+use turtle::{
+ generated::{
DAEMON_PROTOCOL_VERSION,
- generated::history::{
+ history::{
EndHistoryReply, EndHistoryRequest, HistoryEntry, HistoryEventKind, HistoryReply,
HistoryRequest, StartHistoryReply, StartHistoryRequest, TailHistoryReply,
TailHistoryRequest,
history_server::{History as HistorySvc, HistoryServer},
},
},
- daemon::DaemonHandle,
- events::DaemonEvent,
+ history::{History, HistoryId},
};
/// The gRPC service implementation.
diff --git a/crates/daemon/src/api/mod.rs b/crates/daemon/src/api/mod.rs
index b7f82a7b..8d475fe9 100644
--- a/crates/daemon/src/api/mod.rs
+++ b/crates/daemon/src/api/mod.rs
@@ -1,6 +1,2 @@
-pub mod client;
-pub(crate) mod generated;
-pub(crate) mod server;
-
-pub(crate) const DAEMON_VERSION: &str = env!("CARGO_PKG_VERSION");
-const DAEMON_PROTOCOL_VERSION: u32 = 1;
+pub(crate) mod control;
+pub(crate) mod history;
diff --git a/crates/daemon/src/api/server/mod.rs b/crates/daemon/src/api/server/mod.rs
deleted file mode 100644
index 8d475fe9..00000000
--- a/crates/daemon/src/api/server/mod.rs
+++ /dev/null
@@ -1,2 +0,0 @@
-pub(crate) mod control;
-pub(crate) mod history;
diff --git a/crates/daemon/src/events.rs b/crates/daemon/src/events.rs
index e3f6ebc9..1864d224 100644
--- a/crates/daemon/src/events.rs
+++ b/crates/daemon/src/events.rs
@@ -7,7 +7,7 @@
//! External processes (like CLI commands) can also inject events via the
//! Control gRPC service.
-use crate::aclient::history::{History, HistoryId};
+use turtle::history::{History, HistoryId};
use turtle_common::record::RecordId;
/// Events that flow through the daemon's event bus.
diff --git a/crates/daemon/src/lib.rs b/crates/daemon/src/lib.rs
deleted file mode 100644
index c6877096..00000000
--- a/crates/daemon/src/lib.rs
+++ /dev/null
@@ -1,161 +0,0 @@
-#![expect(unused_crate_dependencies, reason = "Didn't remove them yet")]
-
-use std::{
- fs::{self, File, OpenOptions},
- io::Write,
- path::{Path, PathBuf},
- time::{Duration, Instant},
-};
-
-use eyre::{Context, Result, bail, eyre};
-use fs4::fs_std::FileExt;
-use tokio::time::sleep;
-
-use crate::{
- aclient::{database::ClientSqlite, record::sqlite_store::SqliteStore, settings::Settings},
- api::{
- DAEMON_VERSION,
- server::{control::ControlService, history::HistoryService},
- },
- daemon::Daemon,
-};
-
-pub mod aclient;
-
-pub mod api;
-pub(crate) mod daemon;
-pub(crate) mod events;
-pub(crate) mod server;
-
-/// Boot the daemon.
-///
-/// This creates a daemon,
-/// starts the gRPC server with services, and runs the event loop.
-pub async fn boot(settings: Settings, store: SqliteStore, history_db: ClientSqlite) -> Result<()> {
- let pidfile_path = PathBuf::from(&settings.daemon.pidfile_path);
- let _pidfile_guard = PidfileGuard::acquire(&pidfile_path)?;
-
- let mut daemon = Daemon::builder(settings.clone())
- .store(store)
- .history_db(history_db.clone())
- .build()?;
-
- let handle = {
- let handle = daemon.handle();
-
- // Spawn signal handler to emit ShutdownRequested on Ctrl+C/SIGTERM
- let signal_handle = handle.clone();
- tokio::spawn(async move {
- shutdown_signal().await;
- tracing::info!("received shutdown signal");
- signal_handle.shutdown();
- });
-
- handle
- };
-
- let history_service = HistoryService::new(handle.clone(), history_db).await?;
- let control_service = ControlService::new(handle.clone());
-
- server::run_grpc_server(
- &settings,
- history_service.into_server(),
- control_service.into_server(),
- handle,
- )?;
-
- daemon.run_event_loop().await?;
-
- tracing::info!("daemon shut down complete");
- Ok(())
-}
-
-/// Wait for a shutdown signal (Ctrl+C or SIGTERM).
-#[cfg(unix)]
-async fn shutdown_signal() {
- let mut term = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
- .expect("failed to register sigterm handler");
- let mut int = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt())
- .expect("failed to register sigint handler");
-
- tokio::select! {
- _ = term.recv() => {},
- _ = int.recv() => {},
- }
-}
-
-struct PidfileGuard {
- file: File,
-}
-
-impl PidfileGuard {
- fn acquire(path: &Path) -> Result<Self> {
- let mut file = open_lock_file(path)?;
-
- if !file.try_lock_exclusive()? {
- bail!(
- "daemon already running (pidfile lock busy at {})",
- path.display()
- );
- }
-
- file.set_len(0)
- .wrap_err_with(|| format!("could not truncate daemon pidfile {}", path.display()))?;
- writeln!(file, "{}", std::process::id())
- .and_then(|()| writeln!(file, "{DAEMON_VERSION}"))
- .wrap_err_with(|| format!("could not write daemon pidfile {}", path.display()))?;
-
- Ok(Self { file })
- }
-}
-
-impl Drop for PidfileGuard {
- fn drop(&mut self) {
- drop(self.file.unlock());
- }
-}
-
-fn open_lock_file(path: &Path) -> Result<File> {
- if let Some(parent) = path.parent() {
- fs::create_dir_all(parent)
- .wrap_err_with(|| format!("could not create lock directory {}", parent.display()))?;
- }
-
- OpenOptions::new()
- .read(true)
- .write(true)
- .create(true)
- .truncate(false)
- .open(path)
- .wrap_err_with(|| format!("could not open lock file {}", path.display()))
-}
-
-async fn wait_for_lock(path: &Path, timeout: Duration) -> Result<File> {
- const LOCK_POLL: Duration = Duration::from_millis(20);
-
- let file = open_lock_file(path)?;
- let start = Instant::now();
-
- loop {
- match file.try_lock_exclusive() {
- Ok(true) => return Ok(file),
- Ok(false) => {
- if start.elapsed() >= timeout {
- bail!("timed out waiting for lock at {}", path.display());
- }
-
- sleep(LOCK_POLL).await;
- }
- Err(err) => {
- return Err(eyre!("could not lock {}: {err}", path.display()));
- }
- }
- }
-}
-
-async fn wait_for_pidfile_available(path: &Path, timeout: Duration) -> Result<()> {
- let file = wait_for_lock(path, timeout).await?;
- file.unlock()
- .wrap_err_with(|| format!("failed to unlock {}", path.display()))?;
- Ok(())
-}
diff --git a/crates/daemon/src/main.rs b/crates/daemon/src/main.rs
index 46d9ad4d..59d4c7ff 100644
--- a/crates/daemon/src/main.rs
+++ b/crates/daemon/src/main.rs
@@ -1,15 +1,35 @@
-#![expect(unused_crate_dependencies, reason = "Didn't remove them yet")]
+#![expect(unused_crate_dependencies)]
+
+use std::{
+ fs::{self, File, OpenOptions},
+ io::Write,
+ path::{Path, PathBuf},
+ time::{Duration, Instant},
+};
use clap::Parser;
-use eyre::{Result, WrapErr};
-use std::path::PathBuf;
-use turtle_daemon::aclient::{
- database::ClientSqlite, record::sqlite_store::SqliteStore, settings::Settings,
+use eyre::WrapErr;
+use eyre::{Context, Result, bail, eyre};
+use fs4::fs_std::FileExt;
+use tokio::time::sleep;
+
+use crate::{
+ aclient::{database::ClientSqlite, record::sqlite_store::SqliteStore, settings::Settings},
+ api::{control::ControlService, history::HistoryService},
+ daemon::Daemon,
};
+pub(crate) mod aclient;
+pub(crate) mod api;
+pub(crate) mod daemon;
+pub(crate) mod events;
+pub(crate) mod server;
+
+const DAEMON_VERSION: &str = env!("CARGO_PKG_VERSION");
+
#[derive(Parser, Debug)]
#[command(infer_subcommands = true)]
-pub(crate) enum Cmd {
+enum Cmd {
/// Start the daemon server
Start {
/// Also write daemon logs to the console (useful for debugging)
@@ -28,8 +48,139 @@ async fn main() -> Result<()> {
let sqlite_store = SqliteStore::new(record_store_path, settings.local_timeout).await?;
match Cmd::parse() {
- Cmd::Start { show_logs, .. } => {
- turtle_daemon::boot(settings, sqlite_store, history_db).await
+ Cmd::Start { show_logs, .. } => boot(settings, sqlite_store, history_db).await,
+ }
+}
+
+/// Boot the daemon.
+///
+/// This creates a daemon,
+/// starts the gRPC server with services, and runs the event loop.
+async fn boot(settings: Settings, store: SqliteStore, history_db: ClientSqlite) -> Result<()> {
+ let pidfile_path = PathBuf::from(&settings.daemon.pidfile_path);
+ let _pidfile_guard = PidfileGuard::acquire(&pidfile_path)?;
+
+ let mut daemon = Daemon::builder(settings.clone())
+ .store(store)
+ .history_db(history_db.clone())
+ .build()?;
+
+ let handle = {
+ let handle = daemon.handle();
+
+ // Spawn signal handler to emit ShutdownRequested on Ctrl+C/SIGTERM
+ let signal_handle = handle.clone();
+ tokio::spawn(async move {
+ shutdown_signal().await;
+ tracing::info!("received shutdown signal");
+ signal_handle.shutdown();
+ });
+
+ handle
+ };
+
+ let history_service = HistoryService::new(handle.clone(), history_db).await?;
+ let control_service = ControlService::new(handle.clone());
+
+ server::run_grpc_server(
+ &settings,
+ history_service.into_server(),
+ control_service.into_server(),
+ handle,
+ )?;
+
+ daemon.run_event_loop().await?;
+
+ tracing::info!("daemon shut down complete");
+ Ok(())
+}
+
+/// Wait for a shutdown signal (Ctrl+C or SIGTERM).
+#[cfg(unix)]
+async fn shutdown_signal() {
+ let mut term = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
+ .expect("failed to register sigterm handler");
+ let mut int = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt())
+ .expect("failed to register sigint handler");
+
+ tokio::select! {
+ _ = term.recv() => {},
+ _ = int.recv() => {},
+ }
+}
+
+struct PidfileGuard {
+ file: File,
+}
+
+impl PidfileGuard {
+ fn acquire(path: &Path) -> Result<Self> {
+ let mut file = open_lock_file(path)?;
+
+ if !file.try_lock_exclusive()? {
+ bail!(
+ "daemon already running (pidfile lock busy at {})",
+ path.display()
+ );
}
+
+ file.set_len(0)
+ .wrap_err_with(|| format!("could not truncate daemon pidfile {}", path.display()))?;
+ writeln!(file, "{}", std::process::id())
+ .and_then(|()| writeln!(file, "{DAEMON_VERSION}"))
+ .wrap_err_with(|| format!("could not write daemon pidfile {}", path.display()))?;
+
+ Ok(Self { file })
+ }
+}
+
+impl Drop for PidfileGuard {
+ fn drop(&mut self) {
+ drop(self.file.unlock());
+ }
+}
+
+fn open_lock_file(path: &Path) -> Result<File> {
+ if let Some(parent) = path.parent() {
+ fs::create_dir_all(parent)
+ .wrap_err_with(|| format!("could not create lock directory {}", parent.display()))?;
}
+
+ OpenOptions::new()
+ .read(true)
+ .write(true)
+ .create(true)
+ .truncate(false)
+ .open(path)
+ .wrap_err_with(|| format!("could not open lock file {}", path.display()))
+}
+
+async fn wait_for_lock(path: &Path, timeout: Duration) -> Result<File> {
+ const LOCK_POLL: Duration = Duration::from_millis(20);
+
+ let file = open_lock_file(path)?;
+ let start = Instant::now();
+
+ loop {
+ match file.try_lock_exclusive() {
+ Ok(true) => return Ok(file),
+ Ok(false) => {
+ if start.elapsed() >= timeout {
+ bail!("timed out waiting for lock at {}", path.display());
+ }
+
+ sleep(LOCK_POLL).await;
+ }
+ Err(err) => {
+ return Err(eyre!("could not lock {}: {err}", path.display()));
+ }
+ }
+ }
+}
+
+async fn wait_for_pidfile_available(path: &Path, timeout: Duration) -> Result<()> {
+ let file = wait_for_lock(path, timeout).await?;
+ file.unlock()
+ .wrap_err_with(|| format!("failed to unlock {}", path.display()))?;
+ Ok(())
}
diff --git a/crates/daemon/src/server.rs b/crates/daemon/src/server.rs
index 3400ad62..6747f276 100644
--- a/crates/daemon/src/server.rs
+++ b/crates/daemon/src/server.rs
@@ -1,19 +1,14 @@
-use std::os::unix::net::SocketAddr;
-use std::path::PathBuf;
+use std::{os::unix::net::SocketAddr, path::PathBuf};
use eyre::Result;
use eyre::{OptionExt, WrapErr};
-
-#[cfg(unix)]
-use crate::api::server::{control::ControlService, history::HistoryService};
-use crate::{
- aclient::settings::Settings,
- api::generated::{
- control::control_server::ControlServer, history::history_server::HistoryServer,
- },
- daemon::DaemonHandle,
+use turtle::generated::{
+ control::control_server::ControlServer, history::history_server::HistoryServer,
};
+use crate::api::{control::ControlService, history::HistoryService};
+use crate::{aclient::settings::Settings, daemon::DaemonHandle};
+
/// Run the gRPC server with the given services.
///
/// This starts the gRPC server in the background and returns immediately.