diff --git a/.beads/interactions.jsonl b/.beads/interactions.jsonl index c30ec60..a0adcc2 100644 --- a/.beads/interactions.jsonl +++ b/.beads/interactions.jsonl @@ -86,3 +86,4 @@ {"id":"int-a150f81b","kind":"field_change","created_at":"2026-06-30T22:08:23.925241993Z","actor":"dawn","issue_id":"klbr-7e0","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Removed unregistered old memory tool modules and updated memory docs"}} {"id":"int-c575a8f2","kind":"field_change","created_at":"2026-06-30T22:20:51.577686369Z","actor":"dawn","issue_id":"klbr-x1i","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Replaced klbr-bench raw argv dispatch with clap-derived commands and structured runner options"}} {"id":"int-96f16ee3","kind":"field_change","created_at":"2026-06-30T22:27:11.180142339Z","actor":"dawn","issue_id":"klbr-rds","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Split klbr-bench main.rs into cli, continuous_loop, and longmem modules while preserving bench behavior"}} +{"id":"int-5c780dfb","kind":"field_change","created_at":"2026-06-30T22:34:06.198966004Z","actor":"dawn","issue_id":"klbr-f8w","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Split Discord formatting, inbox storage, and tool schema builders into modules and deduplicated batch row rendering prep"}} diff --git a/.beads/issues.jsonl b/.beads/issues.jsonl index b04b8a7..235f7b4 100644 --- a/.beads/issues.jsonl +++ b/.beads/issues.jsonl @@ -45,7 +45,7 @@ {"_type":"issue","id":"klbr-noi","title":"nix: fix broken flake.nix referencing non-existent default.nix","description":"flake.nix references ./default.nix to build the default package, but default.nix does not exist in the repository. Refactor flake.nix to output the derivation directly using nix-cargo-integration outputs.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:42:19Z","created_by":"dawn","updated_at":"2026-06-30T22:01:41Z","started_at":"2026-06-30T22:00:06Z","closed_at":"2026-06-30T22:01:41Z","close_reason":"Replaced the missing ./default.nix package with config.nci.outputs.\"klbr-daemon\".packages.release. Verified nix eval .#packages.x86_64-linux.default.name returns klbr-daemon and nix flake show succeeds without the default.nix error.","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-x1i","title":"bench: replace manual CLI argument parsing with clap","description":"klbr-bench main.rs parses command-line arguments using raw std::env::args and index offsets. This is fragile and lacks auto-help or completions. Port CLI arguments parsing to clap.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:42:13Z","created_by":"dawn","updated_at":"2026-06-30T22:20:52Z","started_at":"2026-06-30T22:09:02Z","closed_at":"2026-06-30T22:20:52Z","close_reason":"Replaced klbr-bench raw argv dispatch with clap-derived commands and structured runner options","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-rds","title":"bench: split 6700-line main.rs monolith into logical submodules","description":"klbr-bench main.rs is currently over 6700 lines of code. It contains benchmark suites for passive recall, continuous loop, router, sweep, and command line parsing. Needs to be split into domain submodules (cli.rs, sweep.rs, passive_recall.rs, continuous_loop.rs, etc.).","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:42:03Z","created_by":"dawn","updated_at":"2026-06-30T22:27:11Z","started_at":"2026-06-30T22:21:17Z","closed_at":"2026-06-30T22:27:11Z","close_reason":"Split klbr-bench main.rs into cli, continuous_loop, and longmem modules while preserving bench behavior","dependency_count":0,"dependent_count":0,"comment_count":0} -{"_type":"issue","id":"klbr-f8w","title":"discord: decompose monolithic lib.rs and deduplicate batch formatting","description":"klbr-discord lib.rs is a 2150-line monolith containing gateway loop, state database, formatting styles, and tools. Split into logical submodules, and deduplicate bracket/xml string formatting loops.","status":"open","priority":2,"issue_type":"task","owner":"90008@klbr.net","created_at":"2026-06-30T21:40:14Z","created_by":"dawn","updated_at":"2026-06-30T21:40:14Z","dependency_count":0,"dependent_count":0,"comment_count":0} +{"_type":"issue","id":"klbr-f8w","title":"discord: decompose monolithic lib.rs and deduplicate batch formatting","description":"klbr-discord lib.rs is a 2150-line monolith containing gateway loop, state database, formatting styles, and tools. Split into logical submodules, and deduplicate bracket/xml string formatting loops.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:40:14Z","created_by":"dawn","updated_at":"2026-06-30T22:34:06Z","started_at":"2026-06-30T22:27:44Z","closed_at":"2026-06-30T22:34:06Z","close_reason":"Split Discord formatting, inbox storage, and tool schema builders into modules and deduplicated batch row rendering prep","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-41z","title":"daemon: fix blocking std::fs calls and secure default WS bind address","description":"DumpMemories handler in daemon.rs blocks the tokio runtime thread using sync std::fs::write. Fix to tokio::fs::write. Also, change default WS bind address to 127.0.0.1 for local-only safety.","status":"closed","priority":2,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-30T21:40:08Z","created_by":"dawn","updated_at":"2026-06-30T22:04:39Z","started_at":"2026-06-30T22:03:35Z","closed_at":"2026-06-30T22:04:39Z","close_reason":"Changed daemon default websocket bind from 0.0.0.0:8765 to 127.0.0.1:8765 and replaced DumpMemories std::fs::write with tokio::fs::write(...).await. Verified rg for old patterns, cargo fmt --check, and cargo test -p klbr-daemon.","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-5h3","title":"ipc/daemon: eliminate duplicate DTO structs and field mapping boilerplate","description":"HistoryEntry, ToolCall, CompactionRecord, and ResolutionEventDto are identical duplicates between klbr-ipc and klbr-core. The daemon has extensive boilerplate mapping them field-by-field. Share or re-export these types.","status":"open","priority":2,"issue_type":"task","owner":"90008@klbr.net","created_at":"2026-06-30T21:40:01Z","created_by":"dawn","updated_at":"2026-06-30T21:40:01Z","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-2kc","title":"core: reduce KDL config parser boilerplate in parser.rs","description":"parser.rs contains ~20 copy-pasted optional_*_node helpers that share identical structure. Simplify with macros or generic helpers. Also resolve duplicate deserialize_usize.","status":"open","priority":2,"issue_type":"task","owner":"90008@klbr.net","created_at":"2026-06-30T21:39:55Z","created_by":"dawn","updated_at":"2026-06-30T21:39:55Z","dependency_count":0,"dependent_count":0,"comment_count":0} diff --git a/klbr-discord/src/formatting.rs b/klbr-discord/src/formatting.rs new file mode 100644 index 0000000..ad3d13e --- /dev/null +++ b/klbr-discord/src/formatting.rs @@ -0,0 +1,272 @@ +use super::*; + +pub(super) fn format_history( + channel_id: Id, + conversation_id: Option<&str>, + messages: &[Message], + style: FormatStyle, +) -> String { + match style { + FormatStyle::Brackets => { + let conversation = conversation_id + .map(|id| format!("conv:{id}")) + .unwrap_or_else(|| format!("channel:{channel_id}")); + let mut lines = vec![ + "[discord history]".to_string(), + conversation, + "messages:".to_string(), + ]; + + for message in messages.iter().rev() { + let author = message + .author + .global_name + .as_deref() + .unwrap_or(&message.author.name); + lines.push(format!( + "[msg:{} {}] {}: {}", + message.id, + format_discord_timestamp(message.timestamp), + author, + one_line(&message.content) + )); + } + + lines.join("\n") + } + FormatStyle::Xml => { + let mut lines = vec![]; + let conversation_attr = conversation_id + .map(|id| format!(" conversation_id=\"{id}\"")) + .unwrap_or_else(|| format!(" channel_id=\"{channel_id}\"")); + lines.push(format!("", conversation_attr)); + + for message in messages.iter().rev() { + let author = message + .author + .global_name + .as_deref() + .unwrap_or(&message.author.name); + lines.push(format!( + " {}", + message.id, + format_discord_timestamp(message.timestamp), + author, + one_line(&message.content) + )); + } + + lines.push("".to_string()); + lines.join("\n") + } + } +} + +pub(super) fn format_discord_batch( + messages: &[DiscordIncomingMessage], + pending: &[DiscordInboxItem], + include_instructions: bool, + style: FormatStyle, +) -> String { + let rows = batch_message_rows(messages, pending); + let older_pending = older_pending_items(messages, pending); + + match style { + FormatStyle::Brackets => { + let mut lines = vec!["[discord message batch]".to_string()]; + if include_instructions { + lines.push("messages are context, not direct requests".to_string()); + lines.push("for dm: reply if a response is needed".to_string()); + lines.push( + "for guild: channel context only — replying is optional, use judgment" + .to_string(), + ); + lines.push("use discord_send to respond (marks pending items acted); use discord_mark only when not sending".to_string()); + } + lines.push("messages:".to_string()); + + for row in &rows { + let message = row.message; + lines.push(format!( + "[conv:{} msg:{} {} {}] {}: {}", + message.conversation_id, + message.message_id, + row.context, + format_discord_batch_time(message.timestamp), + message.author_name, + one_line(&message.content) + )); + } + + if !older_pending.is_empty() { + lines.push("pending:".to_string()); + for item in older_pending { + lines.push(format!( + "[msg:{} {} conv:{} {}] {}: {}", + item.item_id, + item.bucket, + item.conversation_id, + format_discord_batch_time(item.timestamp), + item.author_name, + one_line(&item.content) + )); + } + } + + lines.join("\n") + } + FormatStyle::Xml => { + let mut lines = vec!["".to_string()]; + if include_instructions { + lines.push(" ".to_string()); + lines.push(" messages are context, not direct requests".to_string()); + lines.push(" for dm: reply if a response is needed".to_string()); + lines.push( + " for guild: channel context only — replying is optional, use judgment" + .to_string(), + ); + lines.push(" use discord_send to respond (marks pending items acted); use discord_mark only when not sending".to_string()); + lines.push(" ".to_string()); + } + + lines.push(" ".to_string()); + for row in &rows { + let message = row.message; + let item_attr = row + .pending_item + .map(|item| format!(" item_id=\"{}\"", item.item_id)) + .unwrap_or_default(); + + lines.push(format!( + " {}", + message.message_id, + item_attr, + message.conversation_id, + row.context, + format_discord_batch_time(message.timestamp), + message.author_name, + one_line(&message.content) + )); + } + lines.push(" ".to_string()); + + if !older_pending.is_empty() { + lines.push(" ".to_string()); + for item in older_pending { + lines.push(format!( + " {}", + item.message_id, + item.item_id, + item.conversation_id, + item.bucket, + format_discord_batch_time(item.timestamp), + item.author_name, + one_line(&item.content) + )); + } + lines.push(" ".to_string()); + } + + lines.push("".to_string()); + lines.join("\n") + } + } +} + +struct BatchMessageRow<'a> { + message: &'a DiscordIncomingMessage, + pending_item: Option<&'a DiscordInboxItem>, + context: &'a str, +} + +fn batch_message_rows<'a>( + messages: &'a [DiscordIncomingMessage], + pending: &'a [DiscordInboxItem], +) -> Vec> { + let pending_by_message: HashMap<&str, &DiscordInboxItem> = pending + .iter() + .map(|item| (item.message_id.as_str(), item)) + .collect(); + + messages + .iter() + .map(|message| { + let pending_item = pending_by_message.get(message.message_id.as_str()).copied(); + let context = pending_item + .map(|item| item.bucket.as_str()) + .unwrap_or(if message.is_dm { "dm" } else { "guild" }); + BatchMessageRow { + message, + pending_item, + context, + } + }) + .collect() +} + +fn older_pending_items<'a>( + messages: &[DiscordIncomingMessage], + pending: &'a [DiscordInboxItem], +) -> Vec<&'a DiscordInboxItem> { + let batch_msg_ids: HashSet<&str> = messages.iter().map(|m| m.message_id.as_str()).collect(); + pending + .iter() + .filter(|item| !batch_msg_ids.contains(item.message_id.as_str())) + .collect() +} + +pub(super) fn format_discord_timestamp(timestamp: Timestamp) -> String { + format_discord_timestamp_secs(timestamp.as_secs()) +} + +pub(super) fn format_discord_timestamp_secs(secs: i64) -> String { + let Some(datetime) = chrono::DateTime::from_timestamp(secs, 0) else { + return secs.to_string(); + }; + datetime.format("%Y-%m-%d %H:%M:%S").to_string() +} + +pub(super) fn format_discord_batch_time(secs: i64) -> String { + let Some(datetime) = chrono::DateTime::from_timestamp(secs, 0) else { + return secs.to_string(); + }; + datetime.format("%Y-%m-%d %H:%M:%S").to_string() +} + +pub(super) fn one_line(text: &str) -> String { + use std::sync::OnceLock; + static DATA_URL_RE: OnceLock = OnceLock::new(); + let re = DATA_URL_RE.get_or_init(|| { + regex::Regex::new(r"data:([a-zA-Z0-9\-+\.]+/[a-zA-Z0-9\-+\.]+);base64,([a-zA-Z0-9/+=]+)") + .unwrap() + }); + + let mut media_urls = Vec::new(); + let mut cleaned_text = text.to_string(); + + while let Some(mat) = re.find(&cleaned_text) { + media_urls.push(mat.as_str().to_string()); + cleaned_text.replace_range(mat.range(), ""); + } + + let normalized = cleaned_text + .split_whitespace() + .collect::>() + .join(" "); + let mut out = if normalized.chars().count() <= 500 { + normalized + } else { + let mut truncated = normalized.chars().take(500).collect::(); + truncated.push_str("..."); + truncated + }; + + for url in media_urls { + if !out.is_empty() { + out.push(' '); + } + out.push_str(&url); + } + + out +} diff --git a/klbr-discord/src/inbox.rs b/klbr-discord/src/inbox.rs new file mode 100644 index 0000000..f522145 --- /dev/null +++ b/klbr-discord/src/inbox.rs @@ -0,0 +1,237 @@ +use super::*; + +#[derive(Clone, Debug)] +pub(super) struct DiscordInboxStore { + pub(super) conn: Arc>, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub(super) struct DiscordInboxItem { + pub(super) item_id: String, + pub(super) conversation_id: String, + pub(super) channel_id: String, + pub(super) message_id: String, + pub(super) author_name: String, + pub(super) content: String, + pub(super) bucket: String, + pub(super) status: String, + pub(super) timestamp: i64, +} + +impl DiscordInboxStore { + pub(super) fn open(path: impl AsRef) -> Result { + let conn = Connection::open(path)?; + let store = Self { + conn: Arc::new(Mutex::new(conn)), + }; + store.init_schema()?; + Ok(store) + } + + pub(super) fn init_schema(&self) -> Result<()> { + self.conn.lock().unwrap().execute_batch( + "CREATE TABLE IF NOT EXISTS discord_inbox_items ( + item_id TEXT PRIMARY KEY, + conversation_id TEXT NOT NULL, + channel_id TEXT NOT NULL, + message_id TEXT NOT NULL UNIQUE, + author_name TEXT NOT NULL, + content TEXT NOT NULL, + bucket TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'pending', + action TEXT NOT NULL DEFAULT 'none', + note TEXT, + message_ts INTEGER NOT NULL, + first_seen_ts INTEGER NOT NULL DEFAULT (unixepoch()), + last_seen_ts INTEGER NOT NULL DEFAULT (unixepoch()), + decision_ts INTEGER + ); + CREATE INDEX IF NOT EXISTS idx_discord_inbox_status + ON discord_inbox_items(status, last_seen_ts DESC); + CREATE INDEX IF NOT EXISTS idx_discord_inbox_channel + ON discord_inbox_items(channel_id, status, last_seen_ts DESC);", + )?; + Ok(()) + } + + pub(super) fn upsert_pending( + &self, + message: &DiscordIncomingMessage, + ) -> Result> { + let Some(bucket) = pending_bucket(message) else { + return Ok(None); + }; + + let conn = self.conn.lock().unwrap(); + conn.execute( + "INSERT INTO discord_inbox_items ( + item_id, conversation_id, channel_id, message_id, author_name, + content, bucket, status, action, message_ts + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, 'pending', 'none', ?8) + ON CONFLICT(item_id) DO UPDATE SET + conversation_id = excluded.conversation_id, + channel_id = excluded.channel_id, + author_name = excluded.author_name, + content = excluded.content, + bucket = excluded.bucket, + last_seen_ts = unixepoch()", + params![ + message.message_id, + message.conversation_id, + message.channel_id, + message.message_id, + message.author_name, + message.content, + bucket, + message.timestamp, + ], + )?; + Ok(Some(message.message_id.clone())) + } + + pub(super) fn pending_items(&self, limit: usize) -> Result> { + let conn = self.conn.lock().unwrap(); + let mut stmt = conn.prepare( + "SELECT item_id, conversation_id, channel_id, message_id, author_name, + content, bucket, status, message_ts + FROM discord_inbox_items + WHERE status = 'pending' + ORDER BY message_ts DESC, last_seen_ts DESC + LIMIT ?1", + )?; + let rows = stmt.query_map(params![limit as i64], inbox_item_from_row)?; + rows.collect::, _>>() + .map_err(Into::into) + } + + pub(super) fn demote_stale_pending(&self, older_than_secs: i64) -> Result { + let changed = self.conn.lock().unwrap().execute( + "UPDATE discord_inbox_items + SET status = 'seen', + action = CASE WHEN action = 'none' THEN 'noted' ELSE action END, + decision_ts = unixepoch(), + last_seen_ts = unixepoch() + WHERE status = 'pending' + AND message_ts <= unixepoch() - ?1", + params![older_than_secs], + )?; + Ok(changed) + } + + pub(super) fn pending_item_for_message(&self, message_id: &str) -> Result> { + self.conn + .lock() + .unwrap() + .query_row( + "SELECT item_id + FROM discord_inbox_items + WHERE message_id = ?1 AND status = 'pending' + LIMIT 1", + params![message_id], + |row| row.get(0), + ) + .optional() + .map_err(Into::into) + } + + pub(super) fn conversation_summaries(&self) -> Result> { + // returns (conversation_id, channel_id, bucket, latest_ts) per conversation + let conn = self.conn.lock().unwrap(); + let mut stmt = conn.prepare( + "SELECT conversation_id, channel_id, + CASE WHEN SUM(bucket = 'dm') > 0 THEN 'dm' ELSE 'guild' END AS bucket, + MAX(message_ts) AS latest_ts + FROM discord_inbox_items + GROUP BY conversation_id + ORDER BY latest_ts DESC + LIMIT 50", + )?; + let rows = stmt.query_map([], |row| { + Ok(( + row.get::<_, String>(0)?, + row.get::<_, String>(1)?, + row.get::<_, String>(2)?, + row.get::<_, i64>(3)?, + )) + })?; + rows.collect::, _>>() + .map_err(Into::into) + } + + pub(super) fn pending_dm_item_for_channel(&self, channel_id: &str) -> Result> { + self.conn + .lock() + .unwrap() + .query_row( + "SELECT item_id + FROM discord_inbox_items + WHERE channel_id = ?1 + AND bucket = 'dm' + AND status = 'pending' + ORDER BY message_ts DESC + LIMIT 1", + params![channel_id], + |row| row.get(0), + ) + .optional() + .map_err(Into::into) + } + + pub(super) fn mark( + &self, + item_id: &str, + status: &str, + action: &str, + note: Option<&str>, + ) -> Result<()> { + validate_inbox_status(status)?; + validate_inbox_action(action)?; + let item_id = normalize_item_id(item_id); + let changed = self.conn.lock().unwrap().execute( + "UPDATE discord_inbox_items + SET status = ?1, + action = ?2, + note = ?3, + decision_ts = unixepoch(), + last_seen_ts = unixepoch() + WHERE item_id = ?4", + params![status, action, note, item_id], + )?; + if changed == 0 { + anyhow::bail!("unknown Discord inbox item {item_id:?}"); + } + Ok(()) + } +} + +fn inbox_item_from_row(row: &Row<'_>) -> rusqlite::Result { + Ok(DiscordInboxItem { + item_id: row.get(0)?, + conversation_id: row.get(1)?, + channel_id: row.get(2)?, + message_id: row.get(3)?, + author_name: row.get(4)?, + content: row.get(5)?, + bucket: row.get(6)?, + status: row.get(7)?, + timestamp: row.get(8)?, + }) +} + +fn validate_inbox_status(status: &str) -> Result<()> { + match status { + "pending" | "seen" | "acted" | "ignored" => Ok(()), + other => anyhow::bail!( + "invalid Discord inbox status {other:?}; expected pending, seen, acted, or ignored" + ), + } +} + +fn validate_inbox_action(action: &str) -> Result<()> { + match action { + "none" | "replied" | "messaged" | "dismissed" | "noted" => Ok(()), + other => anyhow::bail!( + "invalid Discord inbox action {other:?}; expected none, replied, messaged, dismissed, or noted" + ), + } +} diff --git a/klbr-discord/src/lib.rs b/klbr-discord/src/lib.rs index a3e6032..b87fa06 100644 --- a/klbr-discord/src/lib.rs +++ b/klbr-discord/src/lib.rs @@ -8,12 +8,11 @@ use anyhow::{anyhow, Context, Result}; use klbr_core::{ harness_block::{FormatStyle, DEFAULT_FORMAT_STYLE}, interrupt::{ExternalEvent, Interrupt}, - models::ToolDef, tools::{Subroutines, Tool, ToolFuture}, AgentEvent, }; use rusqlite::{params, Connection, OptionalExtension, Row}; -use serde_json::{json, Value}; +use serde_json::Value; use tokio::{ sync::{broadcast, mpsc}, task::JoinHandle, @@ -29,6 +28,21 @@ use twilight_model::{ util::Timestamp, }; +mod formatting; +mod inbox; +mod tool_defs; + +use formatting::{ + format_discord_batch, format_discord_batch_time, format_discord_timestamp, format_history, +}; +#[cfg(test)] +use formatting::{format_discord_timestamp_secs, one_line}; +use inbox::{DiscordInboxItem, DiscordInboxStore}; +use tool_defs::{ + discord_fetch_history_def, discord_list_conversations_def, discord_mark_def, discord_react_def, + discord_send_def, +}; + #[derive(Clone, Debug)] pub struct DiscordConfig { pub token: String, @@ -82,24 +96,6 @@ struct DiscordIncomingMessage { author_is_bot: bool, } -#[derive(Clone, Debug)] -struct DiscordInboxStore { - conn: Arc>, -} - -#[derive(Clone, Debug, PartialEq, Eq)] -struct DiscordInboxItem { - item_id: String, - conversation_id: String, - channel_id: String, - message_id: String, - author_name: String, - content: String, - bucket: String, - status: String, - timestamp: i64, -} - impl DiscordConfig { pub fn from_env() -> Result> { let enabled = env_bool("KLBR_DISCORD_ENABLED"); @@ -143,183 +139,6 @@ impl DiscordConfig { } } -impl DiscordInboxStore { - fn open(path: impl AsRef) -> Result { - let conn = Connection::open(path)?; - let store = Self { - conn: Arc::new(Mutex::new(conn)), - }; - store.init_schema()?; - Ok(store) - } - - fn init_schema(&self) -> Result<()> { - self.conn.lock().unwrap().execute_batch( - "CREATE TABLE IF NOT EXISTS discord_inbox_items ( - item_id TEXT PRIMARY KEY, - conversation_id TEXT NOT NULL, - channel_id TEXT NOT NULL, - message_id TEXT NOT NULL UNIQUE, - author_name TEXT NOT NULL, - content TEXT NOT NULL, - bucket TEXT NOT NULL, - status TEXT NOT NULL DEFAULT 'pending', - action TEXT NOT NULL DEFAULT 'none', - note TEXT, - message_ts INTEGER NOT NULL, - first_seen_ts INTEGER NOT NULL DEFAULT (unixepoch()), - last_seen_ts INTEGER NOT NULL DEFAULT (unixepoch()), - decision_ts INTEGER - ); - CREATE INDEX IF NOT EXISTS idx_discord_inbox_status - ON discord_inbox_items(status, last_seen_ts DESC); - CREATE INDEX IF NOT EXISTS idx_discord_inbox_channel - ON discord_inbox_items(channel_id, status, last_seen_ts DESC);", - )?; - Ok(()) - } - - fn upsert_pending(&self, message: &DiscordIncomingMessage) -> Result> { - let Some(bucket) = pending_bucket(message) else { - return Ok(None); - }; - - let conn = self.conn.lock().unwrap(); - conn.execute( - "INSERT INTO discord_inbox_items ( - item_id, conversation_id, channel_id, message_id, author_name, - content, bucket, status, action, message_ts - ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, 'pending', 'none', ?8) - ON CONFLICT(item_id) DO UPDATE SET - conversation_id = excluded.conversation_id, - channel_id = excluded.channel_id, - author_name = excluded.author_name, - content = excluded.content, - bucket = excluded.bucket, - last_seen_ts = unixepoch()", - params![ - message.message_id, - message.conversation_id, - message.channel_id, - message.message_id, - message.author_name, - message.content, - bucket, - message.timestamp, - ], - )?; - Ok(Some(message.message_id.clone())) - } - - fn pending_items(&self, limit: usize) -> Result> { - let conn = self.conn.lock().unwrap(); - let mut stmt = conn.prepare( - "SELECT item_id, conversation_id, channel_id, message_id, author_name, - content, bucket, status, message_ts - FROM discord_inbox_items - WHERE status = 'pending' - ORDER BY message_ts DESC, last_seen_ts DESC - LIMIT ?1", - )?; - let rows = stmt.query_map(params![limit as i64], inbox_item_from_row)?; - rows.collect::, _>>() - .map_err(Into::into) - } - - fn demote_stale_pending(&self, older_than_secs: i64) -> Result { - let changed = self.conn.lock().unwrap().execute( - "UPDATE discord_inbox_items - SET status = 'seen', - action = CASE WHEN action = 'none' THEN 'noted' ELSE action END, - decision_ts = unixepoch(), - last_seen_ts = unixepoch() - WHERE status = 'pending' - AND message_ts <= unixepoch() - ?1", - params![older_than_secs], - )?; - Ok(changed) - } - - fn pending_item_for_message(&self, message_id: &str) -> Result> { - self.conn - .lock() - .unwrap() - .query_row( - "SELECT item_id - FROM discord_inbox_items - WHERE message_id = ?1 AND status = 'pending' - LIMIT 1", - params![message_id], - |row| row.get(0), - ) - .optional() - .map_err(Into::into) - } - - fn conversation_summaries(&self) -> Result> { - // returns (conversation_id, channel_id, bucket, latest_ts) per conversation - let conn = self.conn.lock().unwrap(); - let mut stmt = conn.prepare( - "SELECT conversation_id, channel_id, - CASE WHEN SUM(bucket = 'dm') > 0 THEN 'dm' ELSE 'guild' END AS bucket, - MAX(message_ts) AS latest_ts - FROM discord_inbox_items - GROUP BY conversation_id - ORDER BY latest_ts DESC - LIMIT 50", - )?; - let rows = stmt.query_map([], |row| { - Ok(( - row.get::<_, String>(0)?, - row.get::<_, String>(1)?, - row.get::<_, String>(2)?, - row.get::<_, i64>(3)?, - )) - })?; - rows.collect::, _>>() - .map_err(Into::into) - } - - fn pending_dm_item_for_channel(&self, channel_id: &str) -> Result> { - self.conn - .lock() - .unwrap() - .query_row( - "SELECT item_id - FROM discord_inbox_items - WHERE channel_id = ?1 - AND bucket = 'dm' - AND status = 'pending' - ORDER BY message_ts DESC - LIMIT 1", - params![channel_id], - |row| row.get(0), - ) - .optional() - .map_err(Into::into) - } - - fn mark(&self, item_id: &str, status: &str, action: &str, note: Option<&str>) -> Result<()> { - validate_inbox_status(status)?; - validate_inbox_action(action)?; - let item_id = normalize_item_id(item_id); - let changed = self.conn.lock().unwrap().execute( - "UPDATE discord_inbox_items - SET status = ?1, - action = ?2, - note = ?3, - decision_ts = unixepoch(), - last_seen_ts = unixepoch() - WHERE item_id = ?4", - params![status, action, note, item_id], - )?; - if changed == 0 { - anyhow::bail!("unknown Discord inbox item {item_id:?}"); - } - Ok(()) - } -} - impl DiscordRuntime { pub async fn new(config: DiscordConfig, db_path: impl AsRef) -> Result { let _ = rustls::crypto::ring::default_provider().install_default(); @@ -994,128 +813,6 @@ impl DiscordRuntime { } } -fn discord_send_def() -> ToolDef { - ToolDef::function( - "discord_send", - "send a Discord message. use this when a Discord response is useful; many ambient batches need no reply. plain assistant text only goes to the local chat. conversation_id is optional when the current Discord batch has a single conversation. pass source_item_id when responding to a pending item; successful sends mark that item acted, so do not call discord_mark afterward.", - json!({ - "type": "object", - "properties": { - "conversation_id": { - "type": "string", - "description": "Discord conversation id from the batch, like d1 or conv:d1; omit only for a single-conversation batch" - }, - "content": { - "type": "string", - "description": "message content to send, max 2000 characters" - }, - "source_item_id": { - "type": "string", - "description": "optional pending inbox item id shown as msg:; marks it acted after sending" - } - }, - "required": ["content"] - }), - ) -} - -fn discord_fetch_history_def() -> ToolDef { - ToolDef::function( - "discord_fetch_history", - "fetch recent Discord message history for a conversation. use this when the current batch lacks enough context. conversation_id is optional when the current batch has a single conversation.", - json!({ - "type": "object", - "properties": { - "conversation_id": { - "type": "string", - "description": "Discord conversation id from the batch, like d1 or conv:d1; omit only for a single-conversation batch" - }, - "before_message_id": { - "type": "string", - "description": "optional msg: from the batch; fetch messages before it" - }, - "after_message_id": { - "type": "string", - "description": "optional msg: from the batch; fetch messages after it" - }, - "limit": { - "type": "integer", - "description": "number of messages to fetch, 1-100; default 20" - } - }, - "required": [] - }), - ) -} - -fn discord_react_def() -> ToolDef { - ToolDef::function( - "discord_react", - "add a reaction to a Discord message. prefer real unicode emoji like 😍, 👍, or ❤️. only use a custom emoji string like <:name:id> when the exact string came from Discord; never invent custom emoji ids. conversation_id is optional when the current Discord batch has a single conversation.", - json!({ - "type": "object", - "properties": { - "conversation_id": { - "type": "string", - "description": "Discord conversation id from the batch, like d1 or conv:d1; omit only for a single-conversation batch" - }, - "message_id": { - "type": "string", - "description": "msg: of the message to react to, as shown in the batch" - }, - "emoji": { - "type": "string", - "description": "unicode emoji is preferred; custom emoji must be an exact Discord form <:name:id> / , copied from Discord rather than guessed" - } - }, - "required": ["message_id", "emoji"] - }), - ) -} - -fn discord_list_conversations_def() -> ToolDef { - ToolDef::function( - "discord_list_conversations", - "list known Discord conversations with their ids, type (dm/guild), and last activity. use this to find a conversation_id when it is not in the current batch.", - json!({ - "type": "object", - "properties": {}, - "required": [] - }), - ) -} - -fn discord_mark_def() -> ToolDef { - ToolDef::function( - "discord_mark", - "mark a pending Discord inbox item as seen, acted, or ignored so future batches remember the decision. use this when you are not sending a Discord response, especially ignored when you deliberately choose not to respond. successful discord_send already marks the item acted.", - json!({ - "type": "object", - "properties": { - "item_id": { - "type": "string", - "description": "pending Discord inbox message id, either raw snowflake or the displayed msg:" - }, - "status": { - "type": "string", - "enum": ["pending", "seen", "acted", "ignored"], - "description": "new item status" - }, - "action": { - "type": "string", - "enum": ["none", "replied", "messaged", "dismissed", "noted"], - "description": "optional decision/action label" - }, - "note": { - "type": "string", - "description": "optional short decision note" - } - }, - "required": ["item_id", "status"] - }), - ) -} - fn env_bool(name: &str) -> Option { std::env::var(name).ok().map(|value| { matches!( @@ -1248,38 +945,6 @@ fn normalize_item_id(raw: &str) -> &str { .unwrap_or(s) } -fn inbox_item_from_row(row: &Row<'_>) -> rusqlite::Result { - Ok(DiscordInboxItem { - item_id: row.get(0)?, - conversation_id: row.get(1)?, - channel_id: row.get(2)?, - message_id: row.get(3)?, - author_name: row.get(4)?, - content: row.get(5)?, - bucket: row.get(6)?, - status: row.get(7)?, - timestamp: row.get(8)?, - }) -} - -fn validate_inbox_status(status: &str) -> Result<()> { - match status { - "pending" | "seen" | "acted" | "ignored" => Ok(()), - other => anyhow::bail!( - "invalid Discord inbox status {other:?}; expected pending, seen, acted, or ignored" - ), - } -} - -fn validate_inbox_action(action: &str) -> Result<()> { - match action { - "none" | "replied" | "messaged" | "dismissed" | "noted" => Ok(()), - other => anyhow::bail!( - "invalid Discord inbox action {other:?}; expected none, replied, messaged, dismissed, or noted" - ), - } -} - fn resolve_channel_arg_from_state( state: &DiscordState, args: &Value, @@ -1425,258 +1090,6 @@ fn normalize_emoji_shortcode(input: &str) -> Option { Some(name.replace('-', "_").to_ascii_lowercase()) } -fn format_history( - channel_id: Id, - conversation_id: Option<&str>, - messages: &[Message], - style: FormatStyle, -) -> String { - match style { - FormatStyle::Brackets => { - let conversation = conversation_id - .map(|id| format!("conv:{id}")) - .unwrap_or_else(|| format!("channel:{channel_id}")); - let mut lines = vec![ - "[discord history]".to_string(), - conversation, - "messages:".to_string(), - ]; - - for message in messages.iter().rev() { - let author = message - .author - .global_name - .as_deref() - .unwrap_or(&message.author.name); - lines.push(format!( - "[msg:{} {}] {}: {}", - message.id, - format_discord_timestamp(message.timestamp), - author, - one_line(&message.content) - )); - } - - lines.join("\n") - } - FormatStyle::Xml => { - let mut lines = vec![]; - let conversation_attr = conversation_id - .map(|id| format!(" conversation_id=\"{id}\"")) - .unwrap_or_else(|| format!(" channel_id=\"{channel_id}\"")); - lines.push(format!("", conversation_attr)); - - for message in messages.iter().rev() { - let author = message - .author - .global_name - .as_deref() - .unwrap_or(&message.author.name); - lines.push(format!( - " {}", - message.id, - format_discord_timestamp(message.timestamp), - author, - one_line(&message.content) - )); - } - - lines.push("".to_string()); - lines.join("\n") - } - } -} - -fn format_discord_batch( - messages: &[DiscordIncomingMessage], - pending: &[DiscordInboxItem], - include_instructions: bool, - style: FormatStyle, -) -> String { - // pending items mapped by message_id - let batch_pending: HashMap<&str, &DiscordInboxItem> = pending - .iter() - .map(|item| (item.message_id.as_str(), item)) - .collect(); - - match style { - FormatStyle::Brackets => { - let mut lines = vec!["[discord message batch]".to_string()]; - if include_instructions { - lines.push("messages are context, not direct requests".to_string()); - lines.push("for dm: reply if a response is needed".to_string()); - lines.push( - "for guild: channel context only — replying is optional, use judgment" - .to_string(), - ); - lines.push("use discord_send to respond (marks pending items acted); use discord_mark only when not sending".to_string()); - } - lines.push("messages:".to_string()); - - for message in messages { - let tag = batch_pending - .get(message.message_id.as_str()) - .map(|item| item.bucket.as_str()) - .unwrap_or(if message.is_dm { "dm" } else { "guild" }); - lines.push(format!( - "[conv:{} msg:{} {} {}] {}: {}", - message.conversation_id, - message.message_id, - tag, - format_discord_batch_time(message.timestamp), - message.author_name, - one_line(&message.content) - )); - } - - let batch_msg_ids: HashSet<&str> = - messages.iter().map(|m| m.message_id.as_str()).collect(); - let older_pending: Vec<&DiscordInboxItem> = pending - .iter() - .filter(|item| !batch_msg_ids.contains(item.message_id.as_str())) - .collect(); - - if !older_pending.is_empty() { - lines.push("pending:".to_string()); - for item in older_pending { - lines.push(format!( - "[msg:{} {} conv:{} {}] {}: {}", - item.item_id, - item.bucket, - item.conversation_id, - format_discord_batch_time(item.timestamp), - item.author_name, - one_line(&item.content) - )); - } - } - - lines.join("\n") - } - FormatStyle::Xml => { - let mut lines = vec!["".to_string()]; - if include_instructions { - lines.push(" ".to_string()); - lines.push(" messages are context, not direct requests".to_string()); - lines.push(" for dm: reply if a response is needed".to_string()); - lines.push( - " for guild: channel context only — replying is optional, use judgment" - .to_string(), - ); - lines.push(" use discord_send to respond (marks pending items acted); use discord_mark only when not sending".to_string()); - lines.push(" ".to_string()); - } - - lines.push(" ".to_string()); - for message in messages { - let pending_item = batch_pending.get(message.message_id.as_str()); - let tag = pending_item - .map(|item| item.bucket.as_str()) - .unwrap_or(if message.is_dm { "dm" } else { "guild" }); - - let item_attr = pending_item - .map(|item| format!(" item_id=\"{}\"", item.item_id)) - .unwrap_or_default(); - - lines.push(format!( - " {}", - message.message_id, - item_attr, - message.conversation_id, - tag, - format_discord_batch_time(message.timestamp), - message.author_name, - one_line(&message.content) - )); - } - lines.push(" ".to_string()); - - let batch_msg_ids: HashSet<&str> = - messages.iter().map(|m| m.message_id.as_str()).collect(); - let older_pending: Vec<&DiscordInboxItem> = pending - .iter() - .filter(|item| !batch_msg_ids.contains(item.message_id.as_str())) - .collect(); - - if !older_pending.is_empty() { - lines.push(" ".to_string()); - for item in older_pending { - lines.push(format!( - " {}", - item.message_id, - item.item_id, - item.conversation_id, - item.bucket, - format_discord_batch_time(item.timestamp), - item.author_name, - one_line(&item.content) - )); - } - lines.push(" ".to_string()); - } - - lines.push("".to_string()); - lines.join("\n") - } - } -} - -fn format_discord_timestamp(timestamp: Timestamp) -> String { - format_discord_timestamp_secs(timestamp.as_secs()) -} - -fn format_discord_timestamp_secs(secs: i64) -> String { - let Some(datetime) = chrono::DateTime::from_timestamp(secs, 0) else { - return secs.to_string(); - }; - datetime.format("%Y-%m-%d %H:%M:%S").to_string() -} - -fn format_discord_batch_time(secs: i64) -> String { - let Some(datetime) = chrono::DateTime::from_timestamp(secs, 0) else { - return secs.to_string(); - }; - datetime.format("%Y-%m-%d %H:%M:%S").to_string() -} - -fn one_line(text: &str) -> String { - use std::sync::OnceLock; - static DATA_URL_RE: OnceLock = OnceLock::new(); - let re = DATA_URL_RE.get_or_init(|| { - regex::Regex::new(r"data:([a-zA-Z0-9\-+\.]+/[a-zA-Z0-9\-+\.]+);base64,([a-zA-Z0-9/+=]+)") - .unwrap() - }); - - let mut media_urls = Vec::new(); - let mut cleaned_text = text.to_string(); - - while let Some(mat) = re.find(&cleaned_text) { - media_urls.push(mat.as_str().to_string()); - cleaned_text.replace_range(mat.range(), ""); - } - - let normalized = cleaned_text - .split_whitespace() - .collect::>() - .join(" "); - let mut out = if normalized.chars().count() <= 500 { - normalized - } else { - let mut truncated = normalized.chars().take(500).collect::(); - truncated.push_str("..."); - truncated - }; - - for url in media_urls { - if !out.is_empty() { - out.push(' '); - } - out.push_str(&url); - } - - out -} - #[cfg(test)] mod tests { use super::{ diff --git a/klbr-discord/src/tool_defs.rs b/klbr-discord/src/tool_defs.rs new file mode 100644 index 0000000..b5e357f --- /dev/null +++ b/klbr-discord/src/tool_defs.rs @@ -0,0 +1,124 @@ +use klbr_core::models::ToolDef; +use serde_json::json; + +pub(super) fn discord_send_def() -> ToolDef { + ToolDef::function( + "discord_send", + "send a Discord message. use this when a Discord response is useful; many ambient batches need no reply. plain assistant text only goes to the local chat. conversation_id is optional when the current Discord batch has a single conversation. pass source_item_id when responding to a pending item; successful sends mark that item acted, so do not call discord_mark afterward.", + json!({ + "type": "object", + "properties": { + "conversation_id": { + "type": "string", + "description": "Discord conversation id from the batch, like d1 or conv:d1; omit only for a single-conversation batch" + }, + "content": { + "type": "string", + "description": "message content to send, max 2000 characters" + }, + "source_item_id": { + "type": "string", + "description": "optional pending inbox item id shown as msg:; marks it acted after sending" + } + }, + "required": ["content"] + }), + ) +} + +pub(super) fn discord_fetch_history_def() -> ToolDef { + ToolDef::function( + "discord_fetch_history", + "fetch recent Discord message history for a conversation. use this when the current batch lacks enough context. conversation_id is optional when the current batch has a single conversation.", + json!({ + "type": "object", + "properties": { + "conversation_id": { + "type": "string", + "description": "Discord conversation id from the batch, like d1 or conv:d1; omit only for a single-conversation batch" + }, + "before_message_id": { + "type": "string", + "description": "optional msg: from the batch; fetch messages before it" + }, + "after_message_id": { + "type": "string", + "description": "optional msg: from the batch; fetch messages after it" + }, + "limit": { + "type": "integer", + "description": "number of messages to fetch, 1-100; default 20" + } + }, + "required": [] + }), + ) +} + +pub(super) fn discord_react_def() -> ToolDef { + ToolDef::function( + "discord_react", + "add a reaction to a Discord message. prefer real unicode emoji like 😍, 👍, or ❤️. only use a custom emoji string like <:name:id> when the exact string came from Discord; never invent custom emoji ids. conversation_id is optional when the current Discord batch has a single conversation.", + json!({ + "type": "object", + "properties": { + "conversation_id": { + "type": "string", + "description": "Discord conversation id from the batch, like d1 or conv:d1; omit only for a single-conversation batch" + }, + "message_id": { + "type": "string", + "description": "msg: of the message to react to, as shown in the batch" + }, + "emoji": { + "type": "string", + "description": "unicode emoji is preferred; custom emoji must be an exact Discord form <:name:id> / , copied from Discord rather than guessed" + } + }, + "required": ["message_id", "emoji"] + }), + ) +} + +pub(super) fn discord_list_conversations_def() -> ToolDef { + ToolDef::function( + "discord_list_conversations", + "list known Discord conversations with their ids, type (dm/guild), and last activity. use this to find a conversation_id when it is not in the current batch.", + json!({ + "type": "object", + "properties": {}, + "required": [] + }), + ) +} + +pub(super) fn discord_mark_def() -> ToolDef { + ToolDef::function( + "discord_mark", + "mark a pending Discord inbox item as seen, acted, or ignored so future batches remember the decision. use this when you are not sending a Discord response, especially ignored when you deliberately choose not to respond. successful discord_send already marks the item acted.", + json!({ + "type": "object", + "properties": { + "item_id": { + "type": "string", + "description": "pending Discord inbox message id, either raw snowflake or the displayed msg:" + }, + "status": { + "type": "string", + "enum": ["pending", "seen", "acted", "ignored"], + "description": "new item status" + }, + "action": { + "type": "string", + "enum": ["none", "replied", "messaged", "dismissed", "noted"], + "description": "optional decision/action label" + }, + "note": { + "type": "string", + "description": "optional short decision note" + } + }, + "required": ["item_id", "status"] + }), + ) +}