diff --git a/interface/src/api/client.ts b/interface/src/api/client.ts index af8f48bd0..d9767b510 100644 --- a/interface/src/api/client.ts +++ b/interface/src/api/client.ts @@ -1199,6 +1199,8 @@ export interface MessagingStatusResponse { webhook: PlatformStatus; twitch: PlatformStatus; email: PlatformStatus; + mattermost: PlatformStatus; + signal: PlatformStatus; instances: AdapterInstanceStatus[]; } @@ -1230,6 +1232,9 @@ export interface CreateMessagingInstanceRequest { webhook_auth_token?: string; mattermost_base_url?: string; mattermost_token?: string; + signal_http_url?: string; + signal_account?: string; + signal_dm_allowed_users?: string; }; } diff --git a/interface/src/components/ChannelEditModal.tsx b/interface/src/components/ChannelEditModal.tsx index 66e3c070b..da86e321f 100644 --- a/interface/src/components/ChannelEditModal.tsx +++ b/interface/src/components/ChannelEditModal.tsx @@ -1,6 +1,7 @@ import {useState} from "react"; import {useMutation, useQuery, useQueryClient} from "@tanstack/react-query"; import {api, type PlatformStatus, type BindingInfo} from "@/api/client"; +import {isValidE164, E164_ERROR_TEXT, validateSignalDmAllowedUsers} from "@/lib/format"; import { Button, Input, @@ -18,7 +19,7 @@ import { import {PlatformIcon} from "@/lib/platformIcons"; import {TagInput} from "@/components/TagInput"; -type Platform = "discord" | "slack" | "telegram" | "twitch" | "email" | "webhook" | "mattermost"; +type Platform = "discord" | "slack" | "telegram" | "twitch" | "email" | "webhook" | "mattermost" | "signal"; interface ChannelEditModalProps { platform: Platform; @@ -168,6 +169,40 @@ export function ChannelEditModal({platform, name, status, open, onOpenChange}: C twitch_client_secret: credentialInputs.twitch_client_secret?.trim(), twitch_refresh_token: credentialInputs.twitch_refresh_token?.trim(), }; + } else if (platform === "signal") { + if (!credentialInputs.signal_http_url?.trim()) { + setMessage({text: "HTTP URL is required", type: "error"}); + return; + } + if (!credentialInputs.signal_account?.trim()) { + setMessage({text: "Account phone number is required", type: "error"}); + return; + } + // Basic E.164 validation + const account = credentialInputs.signal_account.trim(); + if (!isValidE164(account)) { + setMessage({text: E164_ERROR_TEXT, type: "error"}); + return; + } + let dmUsers: string | undefined; + const rawDmUsers = credentialInputs.signal_dm_allowed_users; + if (rawDmUsers !== undefined) { + if (!rawDmUsers.trim()) { + dmUsers = ""; + } else { + const result = validateSignalDmAllowedUsers(rawDmUsers); + if (!result.valid) { + setMessage({text: result.error, type: "error"}); + return; + } + dmUsers = result.entries.length > 0 ? result.entries.join(",") : ""; + } + } + request.platform_credentials = { + signal_http_url: credentialInputs.signal_http_url.trim(), + signal_account: account, + signal_dm_allowed_users: dmUsers, + }; } saveCreds.mutate(request); } @@ -358,13 +393,59 @@ export function ChannelEditModal({platform, name, status, open, onOpenChange}: C )} + {platform === "signal" && ( + <> +
+
+ + setCredentialInputs({...credentialInputs, signal_http_url: e.target.value})} + placeholder="http://127.0.0.1:8686" + onKeyDown={(e) => { if (e.key === "Enter") handleSaveCredentials(); }} + /> +
+
+ + setCredentialInputs({...credentialInputs, signal_account: e.target.value})} + placeholder="+1234567890" + onKeyDown={(e) => { if (e.key === "Enter") handleSaveCredentials(); }} + /> +

+ Your Signal phone number in E.164 format (+ followed by 6-15 digits, first digit 1-9) +

+
+
+ + setCredentialInputs({...credentialInputs, signal_dm_allowed_users: e.target.value})} + placeholder="+1234567890, +1987654321" + onKeyDown={(e) => { if (e.key === "Enter") handleSaveCredentials(); }} + /> +

+ Allowed DM senders: E.164 phone numbers (+1234567890) or uuid:xxx identifiers. Comma-separated. If empty, DMs are blocked. +

+
+
+

+ Need help?{" "} + + Read the Signal setup docs → + +

+ + )} + {platform === "webhook" && (

Webhook receiver requires no additional credentials.

)} - {platform !== "webhook" && Object.values(credentialInputs).some((v) => v?.trim()) && ( + {platform !== "webhook" && (Object.values(credentialInputs).some((v) => v?.trim()) || credentialInputs.signal_dm_allowed_users === "") && ( diff --git a/interface/src/components/ChannelSettingCard.tsx b/interface/src/components/ChannelSettingCard.tsx index f61c4b6bb..7f6947fbc 100644 --- a/interface/src/components/ChannelSettingCard.tsx +++ b/interface/src/components/ChannelSettingCard.tsx @@ -23,11 +23,12 @@ import { Toggle, } from "@/ui"; import {PlatformIcon} from "@/lib/platformIcons"; +import {isValidE164, E164_ERROR_TEXT, validateSignalDmAllowedUsers} from "@/lib/format"; import {TagInput} from "@/components/TagInput"; import {FontAwesomeIcon} from "@fortawesome/react-fontawesome"; import {faChevronDown, faPlus} from "@fortawesome/free-solid-svg-icons"; -type Platform = "discord" | "slack" | "telegram" | "twitch" | "email" | "webhook" | "mattermost"; +type Platform = "discord" | "slack" | "telegram" | "twitch" | "email" | "webhook" | "mattermost" | "signal"; const PLATFORM_LABELS: Record = { discord: "Discord", @@ -37,6 +38,7 @@ const PLATFORM_LABELS: Record = { email: "Email", webhook: "Webhook", mattermost: "Mattermost", + signal: "Signal", }; const DOC_LINKS: Partial> = { @@ -45,6 +47,7 @@ const DOC_LINKS: Partial> = { telegram: "https://docs.spacebot.sh/telegram-setup", twitch: "https://docs.spacebot.sh/twitch-setup", mattermost: "https://docs.spacebot.sh/mattermost-setup", + signal: "https://docs.spacebot.sh/signal-setup", }; // --- Platform Catalog (Left Column) --- @@ -62,6 +65,7 @@ export function PlatformCatalog({onAddInstance}: PlatformCatalogProps) { "email", "webhook", "mattermost", + "signal", ]; const COMING_SOON = [ @@ -650,6 +654,37 @@ export function AddInstanceCard({platform, isDefault, onCancel, onCreated}: AddI } credentials.mattermost_base_url = credentialInputs.mattermost_base_url.trim(); credentials.mattermost_token = credentialInputs.mattermost_token.trim(); + } else if (platform === "signal") { + if (!credentialInputs.signal_http_url?.trim()) { + setMessage({text: "HTTP URL is required", type: "error"}); + return; + } + if (!credentialInputs.signal_account?.trim()) { + setMessage({text: "Account phone number is required", type: "error"}); + return; + } + // Basic E.164 validation (frontend) - match backend rules + const account = credentialInputs.signal_account.trim(); + if (!isValidE164(account)) { + setMessage({ + text: E164_ERROR_TEXT, + type: "error" + }); + return; + } + credentials.signal_http_url = credentialInputs.signal_http_url.trim(); + credentials.signal_account = account; + // Normalize: always omit when blank (empty or undefined) for consistent empty-state behavior + if (credentialInputs.signal_dm_allowed_users?.trim()) { + const result = validateSignalDmAllowedUsers(credentialInputs.signal_dm_allowed_users); + if (!result.valid) { + setMessage({text: result.error, type: "error"}); + return; + } + if (result.entries.length > 0) { + credentials.signal_dm_allowed_users = result.entries.join(','); + } + } } if (!isDefault && !instanceName.trim()) { @@ -964,6 +999,50 @@ export function AddInstanceCard({platform, isDefault, onCancel, onCreated}: AddI )} + {platform === "signal" && ( + <> +
+ + setCredentialInputs({...credentialInputs, signal_http_url: e.target.value})} + placeholder="http://127.0.0.1:8686" + onKeyDown={(e) => { if (e.key === "Enter") handleSave(); }} + /> +

+ URL of your signal-cli daemon (e.g., http://127.0.0.1:8686) +

+
+
+ + setCredentialInputs({...credentialInputs, signal_account: e.target.value})} + placeholder="+1234567890" + onKeyDown={(e) => { if (e.key === "Enter") handleSave(); }} + /> +

+ Your Signal phone number in E.164 format (+ followed by 6-15 digits, first digit 1-9) +

+
+
+ + setCredentialInputs({...credentialInputs, signal_dm_allowed_users: e.target.value})} + placeholder="+1234567890, +1987654321" + onKeyDown={(e) => { if (e.key === "Enter") handleSave(); }} + /> +

+ Allowed DM senders: E.164 phone numbers (+1234567890) or uuid:xxx identifiers. Comma-separated. If empty, DMs are blocked. +

+
+ + )} + {docLink && (

Need help?{" "} diff --git a/interface/src/lib/format.ts b/interface/src/lib/format.ts index a0dd54075..81745fa3b 100644 --- a/interface/src/lib/format.ts +++ b/interface/src/lib/format.ts @@ -57,3 +57,57 @@ export function platformColor(platform: string): string { default: return "bg-gray-500/20 text-gray-400"; } } + +// E.164 Phone Number Validation +// Validates international phone numbers in format: + followed by country code and 6-15 digits +export const E164_REGEX = /^\+[1-9]\d{5,14}$/; + +export const E164_ERROR_TEXT = + "Phone number must be in E.164 format: + followed by 6-15 digits after '+', with the first digit 1-9 (e.g., +1234567890)"; + +export function isValidE164(phoneNumber: string): boolean { + return E164_REGEX.test(phoneNumber.trim()); +} + +export function validateE164(phoneNumber: string): { valid: boolean; error?: string } { + const trimmed = phoneNumber.trim(); + if (!trimmed) { + return { valid: false, error: "Phone number is required" }; + } + if (!E164_REGEX.test(trimmed)) { + return { valid: false, error: E164_ERROR_TEXT }; + } + return { valid: true }; +} + +/** + * Validate Signal DM allowed-users entries. + * Each entry must be E.164 phone or uuid:xxx. + */ +export function validateSignalDmAllowedUsers( + raw: string +): { valid: true; entries: string[] } | { valid: false; error: string } { + const entries = raw.split(',').map(s => s.trim()).filter(s => s.length > 0); + const invalid: string[] = []; + const valid: string[] = []; + + for (const entry of entries) { + if ( + isValidE164(entry) || + (entry.startsWith('uuid:') && entry.length > 5) + ) { + valid.push(entry); + } else { + invalid.push(entry); + } + } + + if (invalid.length > 0) { + return { + valid: false, + error: `Invalid entries: ${invalid.join(', ')}. Must be E.164 phone numbers (+1234567890) or uuid:xxx`, + }; + } + + return { valid: true, entries: valid }; +} diff --git a/interface/src/lib/platformIcons.tsx b/interface/src/lib/platformIcons.tsx index 1193986b4..48c57f7a8 100644 --- a/interface/src/lib/platformIcons.tsx +++ b/interface/src/lib/platformIcons.tsx @@ -18,6 +18,7 @@ export function PlatformIcon({ platform, className = "text-ink-faint", size = "1 email: faEnvelope, mattermost: faServer, whatsapp: faWhatsapp, + signal: faComment, matrix: faComments, imessage: faComment, irc: faComments, diff --git a/interface/src/routes/Settings.tsx b/interface/src/routes/Settings.tsx index 46a230399..b7e834190 100644 --- a/interface/src/routes/Settings.tsx +++ b/interface/src/routes/Settings.tsx @@ -884,7 +884,7 @@ function ThemePreview({ themeId }: { themeId: ThemeId }) { ); } -type Platform = "discord" | "slack" | "telegram" | "twitch" | "email" | "webhook" | "mattermost"; +type Platform = "discord" | "slack" | "telegram" | "twitch" | "email" | "webhook" | "mattermost" | "signal"; function ChannelsSection() { const [expandedKey, setExpandedKey] = useState(null); diff --git a/prompts/en/adapters/cron.md.j2 b/prompts/en/adapters/cron.md.j2 index 7c5828506..dc31f629b 100644 --- a/prompts/en/adapters/cron.md.j2 +++ b/prompts/en/adapters/cron.md.j2 @@ -7,3 +7,5 @@ This is an automated scheduled task, not a live conversation. There is no human - Be concise and data-driven in the final output. Lead with findings, not preamble. Include specifics — numbers, dates, names, links. - Your entire text output will be delivered as-is to the configured channel. Write it like a finished report. - If a worker fails or data is unavailable, say so clearly and include what you were able to gather. +- Use the `reply` tool for your primary output — it will be delivered to the configured destination automatically. +- Only use `send_message_to_another_channel` if the task explicitly requires sending to additional channels beyond the primary delivery target. diff --git a/prompts/en/fragments/conversation_context.md.j2 b/prompts/en/fragments/conversation_context.md.j2 index a3b0e0344..082446e12 100644 --- a/prompts/en/fragments/conversation_context.md.j2 +++ b/prompts/en/fragments/conversation_context.md.j2 @@ -3,6 +3,8 @@ Platform: {{ platform }} Server: {{ server_name }} {%- endif %} {%- if channel_name %} -Channel: #{{ channel_name }} +Channel: {{ channel_name }} ({{ platform }}{% if conversation_id %}, id: `{{ conversation_id }}`{% endif %}) +{%- elif conversation_id %} +Channel ID: `{{ conversation_id }}` {%- endif %} Multiple users may be present. Each message is prefixed with [username]. diff --git a/src/agent/channel.rs b/src/agent/channel.rs index e1c5a4e5e..b153c2457 100644 --- a/src/agent/channel.rs +++ b/src/agent/channel.rs @@ -1257,6 +1257,7 @@ impl Channel { &first.source, server_name, channel_name, + self.conversation_id.as_deref(), )?); } @@ -1736,6 +1737,7 @@ impl Channel { &message.source, server_name, channel_name, + self.conversation_id.as_deref(), )?); } diff --git a/src/api/agents.rs b/src/api/agents.rs index 9ae1060da..f1d1cbfdb 100644 --- a/src/api/agents.rs +++ b/src/api/agents.rs @@ -528,6 +528,18 @@ pub async fn create_agent_internal( // Acquire the config write mutex to prevent concurrent read-modify-write races. let _config_guard = state.config_write_mutex.lock().await; + // Fail early if messaging manager is unavailable — before any config write, + // directory creation, or database init that would leave a half-created agent. + let messaging_manager = { + let guard = state.messaging_manager.read().await; + guard + .as_ref() + .cloned() + .ok_or_else(|| { + "Messaging manager not initialized. Please ensure messaging adapters are configured before creating agents.".to_string() + })? + }; + let content = if config_path.exists() { tokio::fs::read_to_string(&config_path) .await @@ -550,6 +562,19 @@ pub async fn create_agent_internal( .as_array_of_tables_mut() .ok_or_else(|| "agents is not an array of tables in config.toml".to_string())?; + // Revalidate uniqueness under the lock — another request may have written the + // same agent_id to config.toml between our first check and mutex acquisition. + // Check against the parsed TOML document (agents_array) rather than the stale + // in-memory cache to ensure we catch concurrent writes. + if agents_array.iter().any(|t| { + t.get("id") + .and_then(|v| v.as_str()) + .map(|id| id == agent_id) + .unwrap_or(false) + }) { + return Err(format!("Agent '{agent_id}' already exists")); + } + let mut new_table = toml_edit::Table::new(); new_table["id"] = toml_edit::value(&agent_id); if let Some(display_name) = &request.display_name @@ -786,10 +811,7 @@ pub async fn create_agent_internal( event_tx: event_tx.clone(), memory_event_tx: memory_event_tx.clone(), sqlite_pool: db.sqlite.clone(), - messaging_manager: { - let guard = state.messaging_manager.read().await; - guard.as_ref().cloned() - }, + messaging_manager: Some(messaging_manager.clone()), sandbox: sandbox.clone(), links: Arc::new(arc_swap::ArcSwap::from_pointee( (**state.agent_links.load()).clone(), @@ -841,19 +863,14 @@ pub async fn create_agent_internal( deps: deps.clone(), screenshot_dir: agent_config.screenshot_dir(), logs_dir: agent_config.logs_dir(), - messaging_manager: { - let guard = state.messaging_manager.read().await; - guard - .as_ref() - .cloned() - .unwrap_or_else(|| std::sync::Arc::new(crate::messaging::MessagingManager::new())) - }, + messaging_manager: messaging_manager.clone(), store: cron_store.clone(), }; let scheduler = std::sync::Arc::new(crate::cron::Scheduler::new(cron_context)); runtime_config.set_cron(cron_store.clone(), scheduler.clone()); - let cron_tool = crate::tools::CronTool::new(cron_store.clone(), scheduler.clone()); + let cron_tool = + crate::tools::CronTool::new(cron_store.clone(), scheduler.clone(), messaging_manager); let browser_config = (**runtime_config.browser_config.load()).clone(); let brave_search_key = (**runtime_config.brave_search_key.load()).clone(); diff --git a/src/api/channels.rs b/src/api/channels.rs index 0469a546c..cccc7a545 100644 --- a/src/api/channels.rs +++ b/src/api/channels.rs @@ -456,6 +456,7 @@ pub(super) async fn inspect_prompt( &info.platform, server_name, info.display_name.as_deref(), + Some(&info.id), ) .ok() } diff --git a/src/api/messaging.rs b/src/api/messaging.rs index a83154b9e..e0e5b26f7 100644 --- a/src/api/messaging.rs +++ b/src/api/messaging.rs @@ -6,6 +6,9 @@ use axum::http::StatusCode; use serde::{Deserialize, Serialize}; use std::sync::Arc; +// Re-export E.164 validation from messaging target module +pub use crate::messaging::target::is_valid_e164; + #[derive(Serialize, Clone)] pub(super) struct PlatformStatus { configured: bool, @@ -31,6 +34,8 @@ pub(super) struct MessagingStatusResponse { email: PlatformStatus, webhook: PlatformStatus, twitch: PlatformStatus, + mattermost: PlatformStatus, + signal: PlatformStatus, instances: Vec, } @@ -100,6 +105,13 @@ pub(super) struct InstanceCredentials { mattermost_base_url: Option, #[serde(default)] mattermost_token: Option, + // Signal credentials + #[serde(default)] + signal_http_url: Option, + #[serde(default)] + signal_account: Option, + #[serde(default)] + signal_dm_allowed_users: Option, } #[derive(Deserialize)] @@ -185,414 +197,698 @@ fn push_instance_status( }); } +/// Merge incoming Signal credentials with existing TOML values for patch-style updates. +/// Fields omitted from the request are filled from the existing platform table, +/// so callers can update individual fields without resubmitting every credential. +/// +/// **Note on `signal_dm_allowed_users`:** This field intentionally does NOT fall back +/// to existing TOML values. It uses tri-state semantics: +/// - `None` → caller handles as "preserve existing" (do nothing) +/// - `Some("")` → clear the allow-list +/// - `Some(entries)` → set the allow-list +/// +/// The caller (`create_messaging_instance`) must interpret the `None` case correctly. +fn merge_signal_credentials_with_existing( + credentials: &InstanceCredentials, + existing: &toml_edit::Table, +) -> InstanceCredentials { + InstanceCredentials { + signal_http_url: credentials.signal_http_url.clone().or_else(|| { + existing + .get("http_url") + .and_then(|v| v.as_str()) + .filter(|s| !s.is_empty()) + .map(|s| s.to_string()) + }), + signal_account: credentials.signal_account.clone().or_else(|| { + existing + .get("account") + .and_then(|v| v.as_str()) + .filter(|s| !s.is_empty()) + .map(|s| s.to_string()) + }), + signal_dm_allowed_users: credentials.signal_dm_allowed_users.clone(), + ..Default::default() + } +} + +/// Parse and validate Signal credentials from the request. +/// Returns (http_url, account, dm_allowed_users) on success, or an error response on failure. +fn parse_signal_credentials( + credentials: &InstanceCredentials, +) -> Result<(String, String, Option>), MessagingInstanceActionResponse> { + // Validate required fields + let http_url = match credentials + .signal_http_url + .as_ref() + .map(|s| s.trim()) + .filter(|s| !s.is_empty()) + { + Some(url_str) => match reqwest::Url::parse(url_str) { + Ok(url) => { + // Validate URL scheme is http or https + let scheme = url.scheme(); + if scheme != "http" && scheme != "https" { + return Err(MessagingInstanceActionResponse { + success: false, + message: format!( + "signal: URL scheme must be 'http' or 'https', got '{}'", + scheme + ), + }); + } + // url.to_string() normalizes the URL (lowercase host, percent-encoding). + // The stored value may differ from user input (e.g. HTTP://LOCALHOST → http://localhost). + // SignalAdapter::new() additionally strips trailing slashes at runtime. + url.to_string() + } + Err(e) => { + tracing::warn!(%e, "signal: invalid http_url format"); + return Err(MessagingInstanceActionResponse { + success: false, + message: format!("signal: invalid http_url format: {}", e), + }); + } + }, + None => { + tracing::warn!("signal: http_url is required"); + return Err(MessagingInstanceActionResponse { + success: false, + message: "signal: http_url is required (e.g., http://127.0.0.1:8686)".to_string(), + }); + } + }; + + let account = match credentials + .signal_account + .as_ref() + .map(|s| s.trim()) + .filter(|s| !s.is_empty()) + { + Some(account) => account, + None => { + tracing::warn!("signal: account is required"); + return Err(MessagingInstanceActionResponse { + success: false, + message: "signal: account is required".to_string(), + }); + } + }; + + // Validate E.164 format + if !is_valid_e164(account) { + return Err(MessagingInstanceActionResponse { + success: false, + message: format!( + "Invalid Signal account format: '{}'. Must be E.164 format (+1234567890, 6-15 digits after '+', first digit cannot be 0)", + account + ), + }); + } + + // Parse and validate dm_allowed_users if provided + let dm_users = if let Some(dm_str) = credentials.signal_dm_allowed_users.as_ref() { + let entries: Vec = dm_str + .split(',') + .map(|s| s.trim()) + .filter(|s| !s.is_empty()) + .map(|s| s.to_string()) + .collect(); + + // Validate each entry is a valid Signal target format + let invalid_entries: Vec<&String> = entries + .iter() + .filter(|entry| { + // Valid formats: uuid:xxx (non-empty), +E.164 + let is_uuid = entry.starts_with("uuid:") && entry.len() > 5; + !(is_uuid || is_valid_e164(entry)) + }) + .collect(); + + if !invalid_entries.is_empty() { + let invalid_list: String = invalid_entries + .iter() + .map(|s| s.as_str()) + .collect::>() + .join(", "); + return Err(MessagingInstanceActionResponse { + success: false, + message: format!( + "Invalid DM allow-list entries: {}. Must be 'uuid:xxx' or E.164 phone number (+1234567890)", + invalid_list + ), + }); + } + + Some(entries) + } else { + None + }; + + Ok((http_url, account.to_string(), dm_users)) +} + /// Get which messaging platforms are configured and enabled. pub(super) async fn messaging_status( State(state): State>, ) -> Result, StatusCode> { let config_path = state.config_path.read().await.clone(); - let (discord, slack, telegram, email, webhook, twitch, instances) = if config_path.exists() { - let content = tokio::fs::read_to_string(&config_path) - .await - .map_err(|error| { - tracing::warn!(%error, "failed to read config.toml for messaging status"); + let (discord, slack, telegram, email, webhook, twitch, mattermost, signal, instances) = + if config_path.exists() { + let content = tokio::fs::read_to_string(&config_path) + .await + .map_err(|error| { + tracing::warn!(%error, "failed to read config.toml for messaging status"); + StatusCode::INTERNAL_SERVER_ERROR + })?; + let doc: toml_edit::DocumentMut = content.parse().map_err(|error| { + tracing::warn!(%error, "failed to parse config.toml for messaging status"); StatusCode::INTERNAL_SERVER_ERROR })?; - let doc: toml_edit::DocumentMut = content.parse().map_err(|error| { - tracing::warn!(%error, "failed to parse config.toml for messaging status"); - StatusCode::INTERNAL_SERVER_ERROR - })?; - - let mut instances: Vec = Vec::new(); - let bindings = doc - .get("bindings") - .and_then(|value| value.as_array_of_tables()); - - let discord_status = doc - .get("messaging") - .and_then(|m| m.get("discord")) - .map(|d| { - let has_token = d - .get("token") - .and_then(|v| v.as_str()) - .is_some_and(|s| !s.is_empty()); - let enabled = d.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); - - if has_token { - push_instance_status(&mut instances, bindings, "discord", None, true, enabled); - } - if let Some(named_instances) = d - .get("instances") - .and_then(|value| value.as_array_of_tables()) - { - for instance in named_instances { - let instance_name = normalize_adapter_selector( - instance.get("name").and_then(|value| value.as_str()), + let mut instances: Vec = Vec::new(); + let bindings = doc + .get("bindings") + .and_then(|value| value.as_array_of_tables()); + + let discord_status = doc + .get("messaging") + .and_then(|m| m.get("discord")) + .map(|d| { + let has_token = d + .get("token") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let enabled = d.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); + + if has_token { + push_instance_status( + &mut instances, + bindings, + "discord", + None, + true, + enabled, ); - let instance_enabled = instance - .get("enabled") - .and_then(|value| value.as_bool()) - .unwrap_or(true) - && enabled; - let instance_configured = instance - .get("token") - .and_then(|value| value.as_str()) - .is_some_and(|token| !token.is_empty()); - - if let Some(instance_name) = instance_name - && instance_configured - { - push_instance_status( - &mut instances, - bindings, - "discord", - Some(instance_name), - true, - instance_enabled, + } + + if let Some(named_instances) = d + .get("instances") + .and_then(|value| value.as_array_of_tables()) + { + for instance in named_instances { + let instance_name = normalize_adapter_selector( + instance.get("name").and_then(|value| value.as_str()), ); + let instance_enabled = instance + .get("enabled") + .and_then(|value| value.as_bool()) + .unwrap_or(true) + && enabled; + let instance_configured = instance + .get("token") + .and_then(|value| value.as_str()) + .is_some_and(|token| !token.is_empty()); + + if let Some(instance_name) = instance_name + && instance_configured + { + push_instance_status( + &mut instances, + bindings, + "discord", + Some(instance_name), + true, + instance_enabled, + ); + } } } - } - - PlatformStatus { - configured: has_token, - enabled: has_token && enabled, - } - }) - .unwrap_or(PlatformStatus { - configured: false, - enabled: false, - }); - let slack_status = doc - .get("messaging") - .and_then(|m| m.get("slack")) - .map(|s| { - let has_bot_token = s - .get("bot_token") - .and_then(|v| v.as_str()) - .is_some_and(|t| !t.is_empty()); - let has_app_token = s - .get("app_token") - .and_then(|v| v.as_str()) - .is_some_and(|t| !t.is_empty()); - let enabled = s.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); - - if has_bot_token && has_app_token { - push_instance_status(&mut instances, bindings, "slack", None, true, enabled); - } - - if let Some(named_instances) = s - .get("instances") - .and_then(|value| value.as_array_of_tables()) - { - for instance in named_instances { - let instance_name = normalize_adapter_selector( - instance.get("name").and_then(|value| value.as_str()), + PlatformStatus { + configured: has_token, + enabled: has_token && enabled, + } + }) + .unwrap_or(PlatformStatus { + configured: false, + enabled: false, + }); + + let slack_status = doc + .get("messaging") + .and_then(|m| m.get("slack")) + .map(|s| { + let has_bot_token = s + .get("bot_token") + .and_then(|v| v.as_str()) + .is_some_and(|t| !t.is_empty()); + let has_app_token = s + .get("app_token") + .and_then(|v| v.as_str()) + .is_some_and(|t| !t.is_empty()); + let enabled = s.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); + + if has_bot_token && has_app_token { + push_instance_status( + &mut instances, + bindings, + "slack", + None, + true, + enabled, ); - let has_instance_bot = instance - .get("bot_token") - .and_then(|value| value.as_str()) - .is_some_and(|value| !value.is_empty()); - let has_instance_app = instance - .get("app_token") - .and_then(|value| value.as_str()) - .is_some_and(|value| !value.is_empty()); - let instance_enabled = instance - .get("enabled") - .and_then(|value| value.as_bool()) - .unwrap_or(true) - && enabled; + } - if let Some(instance_name) = instance_name - && has_instance_bot - && has_instance_app - { - push_instance_status( - &mut instances, - bindings, - "slack", - Some(instance_name), - true, - instance_enabled, + if let Some(named_instances) = s + .get("instances") + .and_then(|value| value.as_array_of_tables()) + { + for instance in named_instances { + let instance_name = normalize_adapter_selector( + instance.get("name").and_then(|value| value.as_str()), ); + let has_instance_bot = instance + .get("bot_token") + .and_then(|value| value.as_str()) + .is_some_and(|value| !value.is_empty()); + let has_instance_app = instance + .get("app_token") + .and_then(|value| value.as_str()) + .is_some_and(|value| !value.is_empty()); + let instance_enabled = instance + .get("enabled") + .and_then(|value| value.as_bool()) + .unwrap_or(true) + && enabled; + + if let Some(instance_name) = instance_name + && has_instance_bot + && has_instance_app + { + push_instance_status( + &mut instances, + bindings, + "slack", + Some(instance_name), + true, + instance_enabled, + ); + } } } - } - - PlatformStatus { - configured: has_bot_token && has_app_token, - enabled: has_bot_token && has_app_token && enabled, - } - }) - .unwrap_or(PlatformStatus { - configured: false, - enabled: false, - }); - - let webhook_status = doc - .get("messaging") - .and_then(|m| m.get("webhook")) - .map(|w| { - let enabled = w.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); - - push_instance_status(&mut instances, bindings, "webhook", None, true, enabled); - - PlatformStatus { - configured: true, - enabled, - } - }) - .unwrap_or(PlatformStatus { - configured: false, - enabled: false, - }); - - let email_status = doc - .get("messaging") - .and_then(|m| m.get("email")) - .map(|email| { - let has_imap_host = email - .get("imap_host") - .and_then(|v| v.as_str()) - .is_some_and(|s| !s.is_empty()); - let has_imap_username = email - .get("imap_username") - .and_then(|v| v.as_str()) - .is_some_and(|s| !s.is_empty()); - let has_imap_password = email - .get("imap_password") - .and_then(|v| v.as_str()) - .is_some_and(|s| !s.is_empty()); - let has_smtp_host = email - .get("smtp_host") - .and_then(|v| v.as_str()) - .is_some_and(|s| !s.is_empty()); - - let configured = - has_imap_host && has_imap_username && has_imap_password && has_smtp_host; - let enabled = email - .get("enabled") - .and_then(|v| v.as_bool()) - .unwrap_or(false); - - if configured { - push_instance_status(&mut instances, bindings, "email", None, true, enabled); - } - - PlatformStatus { - configured, - enabled: configured && enabled, - } - }) - .unwrap_or(PlatformStatus { - configured: false, - enabled: false, - }); - - let telegram_status = doc - .get("messaging") - .and_then(|m| m.get("telegram")) - .map(|t| { - let has_token = t - .get("token") - .and_then(|v| v.as_str()) - .is_some_and(|s| !s.is_empty()); - let enabled = t.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); + PlatformStatus { + configured: has_bot_token && has_app_token, + enabled: has_bot_token && has_app_token && enabled, + } + }) + .unwrap_or(PlatformStatus { + configured: false, + enabled: false, + }); + + let webhook_status = doc + .get("messaging") + .and_then(|m| m.get("webhook")) + .map(|w| { + let enabled = w.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); + + push_instance_status(&mut instances, bindings, "webhook", None, true, enabled); + + PlatformStatus { + configured: true, + enabled, + } + }) + .unwrap_or(PlatformStatus { + configured: false, + enabled: false, + }); + + let email_status = doc + .get("messaging") + .and_then(|m| m.get("email")) + .map(|email| { + let has_imap_host = email + .get("imap_host") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let has_imap_username = email + .get("imap_username") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let has_imap_password = email + .get("imap_password") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let has_smtp_host = email + .get("smtp_host") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + + let configured = + has_imap_host && has_imap_username && has_imap_password && has_smtp_host; + + let enabled = email + .get("enabled") + .and_then(|v| v.as_bool()) + .unwrap_or(false); - if has_token { - push_instance_status(&mut instances, bindings, "telegram", None, true, enabled); - } + if configured { + push_instance_status( + &mut instances, + bindings, + "email", + None, + true, + enabled, + ); + } - if let Some(named_instances) = t - .get("instances") - .and_then(|value| value.as_array_of_tables()) - { - for instance in named_instances { - let instance_name = normalize_adapter_selector( - instance.get("name").and_then(|value| value.as_str()), + PlatformStatus { + configured, + enabled: configured && enabled, + } + }) + .unwrap_or(PlatformStatus { + configured: false, + enabled: false, + }); + + let telegram_status = doc + .get("messaging") + .and_then(|m| m.get("telegram")) + .map(|t| { + let has_token = t + .get("token") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let enabled = t.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); + + if has_token { + push_instance_status( + &mut instances, + bindings, + "telegram", + None, + true, + enabled, ); - let instance_enabled = instance - .get("enabled") - .and_then(|value| value.as_bool()) - .unwrap_or(true) - && enabled; - let instance_configured = instance - .get("token") - .and_then(|value| value.as_str()) - .is_some_and(|value| !value.is_empty()); - - if let Some(instance_name) = instance_name - && instance_configured - { - push_instance_status( - &mut instances, - bindings, - "telegram", - Some(instance_name), - true, - instance_enabled, + } + + if let Some(named_instances) = t + .get("instances") + .and_then(|value| value.as_array_of_tables()) + { + for instance in named_instances { + let instance_name = normalize_adapter_selector( + instance.get("name").and_then(|value| value.as_str()), ); + let instance_enabled = instance + .get("enabled") + .and_then(|value| value.as_bool()) + .unwrap_or(true) + && enabled; + let instance_configured = instance + .get("token") + .and_then(|value| value.as_str()) + .is_some_and(|value| !value.is_empty()); + + if let Some(instance_name) = instance_name + && instance_configured + { + push_instance_status( + &mut instances, + bindings, + "telegram", + Some(instance_name), + true, + instance_enabled, + ); + } } } - } - - PlatformStatus { - configured: has_token, - enabled: has_token && enabled, - } - }) - .unwrap_or(PlatformStatus { - configured: false, - enabled: false, - }); - - let twitch_status = doc - .get("messaging") - .and_then(|m| m.get("twitch")) - .map(|t| { - let has_username = t - .get("username") - .and_then(|v| v.as_str()) - .is_some_and(|s| !s.is_empty()); - let has_token = t - .get("oauth_token") - .and_then(|v| v.as_str()) - .is_some_and(|s| !s.is_empty()); - let enabled = t.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); - - if has_username && has_token { - push_instance_status(&mut instances, bindings, "twitch", None, true, enabled); - } - if let Some(named_instances) = t - .get("instances") - .and_then(|value| value.as_array_of_tables()) - { - for instance in named_instances { - let instance_name = normalize_adapter_selector( - instance.get("name").and_then(|value| value.as_str()), + PlatformStatus { + configured: has_token, + enabled: has_token && enabled, + } + }) + .unwrap_or(PlatformStatus { + configured: false, + enabled: false, + }); + + let twitch_status = doc + .get("messaging") + .and_then(|m| m.get("twitch")) + .map(|t| { + let has_username = t + .get("username") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let has_token = t + .get("oauth_token") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let enabled = t.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); + + if has_username && has_token { + push_instance_status( + &mut instances, + bindings, + "twitch", + None, + true, + enabled, ); - let instance_enabled = instance - .get("enabled") - .and_then(|value| value.as_bool()) - .unwrap_or(true) - && enabled; - let has_instance_username = instance - .get("username") - .and_then(|value| value.as_str()) - .is_some_and(|value| !value.is_empty()); - let has_instance_token = instance - .get("oauth_token") - .and_then(|value| value.as_str()) - .is_some_and(|value| !value.is_empty()); - - if let Some(instance_name) = instance_name - && has_instance_username - && has_instance_token - { - push_instance_status( - &mut instances, - bindings, - "twitch", - Some(instance_name), - true, - instance_enabled, + } + + if let Some(named_instances) = t + .get("instances") + .and_then(|value| value.as_array_of_tables()) + { + for instance in named_instances { + let instance_name = normalize_adapter_selector( + instance.get("name").and_then(|value| value.as_str()), ); + let instance_enabled = instance + .get("enabled") + .and_then(|value| value.as_bool()) + .unwrap_or(true) + && enabled; + let has_instance_username = instance + .get("username") + .and_then(|value| value.as_str()) + .is_some_and(|value| !value.is_empty()); + let has_instance_token = instance + .get("oauth_token") + .and_then(|value| value.as_str()) + .is_some_and(|value| !value.is_empty()); + + if let Some(instance_name) = instance_name + && has_instance_username + && has_instance_token + { + push_instance_status( + &mut instances, + bindings, + "twitch", + Some(instance_name), + true, + instance_enabled, + ); + } } } - } - - PlatformStatus { - configured: has_username && has_token, - enabled: has_username && has_token && enabled, - } - }) - .unwrap_or(PlatformStatus { - configured: false, - enabled: false, - }); - - // Populate instances for Mattermost (not in the legacy per-platform status fields) - if let Some(mm) = doc.get("messaging").and_then(|m| m.get("mattermost")) { - let has_url = mm - .get("base_url") - .and_then(|v| v.as_str()) - .is_some_and(|s| !s.is_empty()); - let has_token = mm - .get("token") - .and_then(|v| v.as_str()) - .is_some_and(|s| !s.is_empty()); - let enabled = mm.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); - - if has_url && has_token { - push_instance_status(&mut instances, bindings, "mattermost", None, true, enabled); - } - if let Some(named_instances) = mm - .get("instances") - .and_then(|value| value.as_array_of_tables()) - { - for instance in named_instances { - let instance_name = normalize_adapter_selector( - instance.get("name").and_then(|value| value.as_str()), - ); - let instance_enabled = instance - .get("enabled") - .and_then(|value| value.as_bool()) - .unwrap_or(true) - && enabled; - let instance_configured = instance + PlatformStatus { + configured: has_username && has_token, + enabled: has_username && has_token && enabled, + } + }) + .unwrap_or(PlatformStatus { + configured: false, + enabled: false, + }); + + // Populate instances for Mattermost (not in the legacy per-platform status fields) + let mattermost_status = doc + .get("messaging") + .and_then(|m| m.get("mattermost")) + .map(|mm| { + let has_url = mm .get("base_url") - .and_then(|value| value.as_str()) - .is_some_and(|value| !value.is_empty()) - && instance - .get("token") - .and_then(|value| value.as_str()) - .is_some_and(|value| !value.is_empty()); - - if let Some(instance_name) = instance_name - && instance_configured - { + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let has_token = mm + .get("token") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let enabled = mm.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); + + if has_url && has_token { push_instance_status( &mut instances, bindings, "mattermost", - Some(instance_name), + None, true, - instance_enabled, + enabled, ); } - } - } - } - ( - discord_status, - slack_status, - telegram_status, - email_status, - webhook_status, - twitch_status, - instances, - ) - } else { - let default = PlatformStatus { - configured: false, - enabled: false, + if let Some(named_instances) = mm + .get("instances") + .and_then(|value| value.as_array_of_tables()) + { + for instance in named_instances { + let instance_name = normalize_adapter_selector( + instance.get("name").and_then(|value| value.as_str()), + ); + let instance_enabled = instance + .get("enabled") + .and_then(|value| value.as_bool()) + .unwrap_or(true) + && enabled; + let instance_configured = instance + .get("base_url") + .and_then(|value| value.as_str()) + .is_some_and(|value| !value.is_empty()) + && instance + .get("token") + .and_then(|value| value.as_str()) + .is_some_and(|value| !value.is_empty()); + + if let Some(instance_name) = instance_name + && instance_configured + { + push_instance_status( + &mut instances, + bindings, + "mattermost", + Some(instance_name), + true, + instance_enabled, + ); + } + } + } + + PlatformStatus { + configured: has_url && has_token, + enabled: has_url && has_token && enabled, + } + }) + .unwrap_or(PlatformStatus { + configured: false, + enabled: false, + }); + + // Signal status and instances + let signal_status = doc + .get("messaging") + .and_then(|m| m.get("signal")) + .map(|s| { + let has_http_url = s + .get("http_url") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let has_account = s + .get("account") + .and_then(|v| v.as_str()) + .is_some_and(|s| !s.is_empty()); + let enabled = s.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false); + + if has_http_url && has_account { + push_instance_status( + &mut instances, + bindings, + "signal", + None, + true, + enabled, + ); + } + + if let Some(named_instances) = s + .get("instances") + .and_then(|value| value.as_array_of_tables()) + { + for instance in named_instances { + let instance_name = normalize_adapter_selector( + instance.get("name").and_then(|value| value.as_str()), + ); + // Named instances are enabled based only on their own flag (independent of root enabled) + let instance_enabled = instance + .get("enabled") + .and_then(|value| value.as_bool()) + .unwrap_or(true); + let instance_has_http_url = instance + .get("http_url") + .and_then(|value| value.as_str()) + .is_some_and(|value| !value.is_empty()); + let instance_has_account = instance + .get("account") + .and_then(|value| value.as_str()) + .is_some_and(|value| !value.is_empty()); + + if let Some(instance_name) = instance_name + && instance_has_http_url + && instance_has_account + { + push_instance_status( + &mut instances, + bindings, + "signal", + Some(instance_name), + true, + instance_enabled, + ); + } + } + } + + PlatformStatus { + configured: has_http_url && has_account, + enabled: has_http_url && has_account && enabled, + } + }) + .unwrap_or(PlatformStatus { + configured: false, + enabled: false, + }); + + ( + discord_status, + slack_status, + telegram_status, + email_status, + webhook_status, + twitch_status, + mattermost_status, + signal_status, + instances, + ) + } else { + let default = PlatformStatus { + configured: false, + enabled: false, + }; + ( + default.clone(), + default.clone(), + default.clone(), + default.clone(), + default.clone(), + default.clone(), + default.clone(), + default.clone(), + Vec::new(), + ) }; - ( - default.clone(), - default.clone(), - default.clone(), - default.clone(), - default.clone(), - default, - Vec::new(), - ) - }; Ok(Json(MessagingStatusResponse { discord, @@ -601,6 +897,8 @@ pub(super) async fn messaging_status( email, webhook, twitch, + mattermost, + signal, instances, })) } @@ -1145,6 +1443,93 @@ pub(super) async fn toggle_platform( } } } + "signal" => { + if let Some(signal_config) = &new_config.messaging.signal { + match request.adapter.as_ref() { + None => { + // Toggle default adapter only + if !signal_config.http_url.is_empty() + && !signal_config.account.is_empty() + { + let permissions = { + let perms_guard = state.signal_permissions.read().await; + match perms_guard.as_ref() { + Some(existing) => { + // Update existing ArcSwap pointee + let perms = + crate::config::SignalPermissions::from_config( + signal_config, + ); + existing.store(std::sync::Arc::new(perms)); + existing.clone() + } + None => { + drop(perms_guard); + let perms = + crate::config::SignalPermissions::from_config( + signal_config, + ); + let arc_swap = std::sync::Arc::new( + arc_swap::ArcSwap::from_pointee(perms), + ); + state + .set_signal_permissions(arc_swap.clone()) + .await; + arc_swap + } + } + }; + let instance_dir = state.instance_dir.load(); + let tmp_dir = instance_dir.join("tmp"); + let adapter = crate::messaging::signal::SignalAdapter::new( + "signal", + &signal_config.http_url, + &signal_config.account, + signal_config.ignore_stories, + permissions, + tmp_dir, + ); + if let Err(error) = manager.register_and_start(adapter).await { + tracing::error!(%error, "failed to start signal adapter on toggle"); + } + } + } + Some(adapter_name) => { + // Toggle specific named instance only + let adapter_key = adapter_name.trim(); + if let Some(instance) = + signal_config.instances.iter().find(|instance| { + instance.name == adapter_key && instance.enabled + }) + { + let runtime_key = crate::config::binding_runtime_adapter_key( + "signal", + Some(instance.name.as_str()), + ); + let permissions = + std::sync::Arc::new(arc_swap::ArcSwap::from_pointee( + crate::config::SignalPermissions::from_instance_config( + instance, + ), + )); + let instance_dir = state.instance_dir.load(); + let tmp_dir = instance_dir.join("tmp"); + let adapter = crate::messaging::signal::SignalAdapter::new( + runtime_key, + &instance.http_url, + &instance.account, + instance.ignore_stories, + permissions, + tmp_dir, + ); + if let Err(error) = manager.register_and_start(adapter).await { + tracing::error!(%error, adapter = %instance.name, "failed to start named signal adapter on toggle"); + } + } + } + } + } + } _ => {} } } @@ -1156,6 +1541,22 @@ pub(super) async fn toggle_platform( if let Err(error) = manager.remove_adapter(&runtime_key).await { tracing::warn!(%error, adapter = %runtime_key, "failed to shut down named adapter on toggle"); } + } else if platform == "signal" { + // Root toggle disable: only remove the default adapter, not named instances + let runtime_key = crate::config::binding_runtime_adapter_key(platform, None); + if let Err(error) = manager.remove_adapter(&runtime_key).await { + tracing::warn!(%error, adapter = %runtime_key, "failed to shut down default adapter on toggle"); + } + // Clear cached Signal permissions so the allow-list is no longer in memory + let perms_guard = state.signal_permissions.read().await; + if perms_guard.is_some() { + drop(perms_guard); + state + .set_signal_permissions(std::sync::Arc::new(arc_swap::ArcSwap::from_pointee( + crate::config::SignalPermissions::default(), + ))) + .await; + } } else if let Err(error) = manager.remove_platform_adapters(platform).await { tracing::warn!(%error, platform = %platform, "failed to shut down adapters on toggle"); } @@ -1187,7 +1588,7 @@ pub(super) async fn create_messaging_instance( if !matches!( platform.as_str(), - "discord" | "slack" | "telegram" | "twitch" | "email" | "webhook" | "mattermost" + "discord" | "slack" | "telegram" | "twitch" | "email" | "webhook" | "mattermost" | "signal" ) { return Ok(Json(MessagingInstanceActionResponse { success: false, @@ -1210,6 +1611,25 @@ pub(super) async fn create_messaging_instance( message: "instance name cannot contain ':' or spaces".to_string(), })); } + // Block reserved instance name "dm" for Slack/Discord (but allow for Signal) + if trimmed.eq_ignore_ascii_case("dm") && matches!(platform.as_str(), "slack" | "discord") { + return Ok(Json(MessagingInstanceActionResponse { + success: false, + message: "instance name 'dm' is reserved and cannot be used for Slack or Discord" + .to_string(), + })); + } + + // Validate instance name format (rejects all-digits, Slack-style IDs, etc.) + if !crate::messaging::target::is_valid_instance_name(trimmed) { + return Ok(Json(MessagingInstanceActionResponse { + success: false, + message: format!( + "instance name '{}' is invalid. Names must: not be all digits, not match Slack workspace IDs (Txxxxx/Cxxxxx), be 1-20 characters, and contain only alphanumeric characters, underscores, or hyphens", + trimmed + ), + })); + } } let config_path = state.config_path.read().await.clone(); @@ -1348,9 +1768,45 @@ pub(super) async fn create_messaging_instance( platform_table["token"] = toml_edit::value(token.as_str()); } } + "signal" => { + // Merge incoming credentials with existing TOML values for patch-style updates. + // Fields omitted from the request are filled from the current table, + // so callers can update individual fields without resubmitting every credential. + let effective = + merge_signal_credentials_with_existing(credentials, platform_table); + let (http_url, account, dm_users) = match parse_signal_credentials(&effective) { + Ok(result) => result, + Err(response) => return Ok(Json(response)), + }; + + // Store URL as-is; SignalAdapter::new() normalizes trailing slashes. + platform_table["http_url"] = toml_edit::value(http_url.as_str()); + platform_table["account"] = toml_edit::value(account); + + // Store dm_allowed_users: preserve existing when omitted (None), + // set when provided with values, remove only when explicitly cleared (Some(empty)) + match dm_users { + Some(dm_users) if !dm_users.is_empty() => { + let mut dm_array = toml_edit::Array::new(); + for user in dm_users { + dm_array.push(user); + } + platform_table["dm_allowed_users"] = toml_edit::value(dm_array); + } + Some(_) => { + // Explicitly cleared (empty vec) - remove the key + platform_table.remove("dm_allowed_users"); + } + None => { + // Omitted - preserve existing value by doing nothing + } + } + } _ => {} } - platform_table["enabled"] = toml_edit::value(enabled); + if let Some(value) = request.enabled { + platform_table["enabled"] = toml_edit::value(value); + } } Some(ref name) => { // Named instance — add to [[messaging..instances]] @@ -1477,6 +1933,39 @@ pub(super) async fn create_messaging_instance( instance_table["token"] = toml_edit::value(token.as_str()); } } + "signal" => { + // New instance — no existing values to merge, validate directly. + let (http_url, account, dm_users) = match parse_signal_credentials(credentials) + { + Ok(result) => result, + Err(response) => return Ok(Json(response)), + }; + + // Store URL as-is; SignalAdapter::new() normalizes trailing slashes. + instance_table["http_url"] = toml_edit::value(http_url.as_str()); + instance_table["account"] = toml_edit::value(account); + + // Store dm_allowed_users: preserve existing when omitted (None), + // set when provided with values, remove only when explicitly cleared (Some(empty)) + match dm_users { + Some(dm_users) if !dm_users.is_empty() => { + let mut dm_array = toml_edit::Array::new(); + for user in dm_users { + dm_array.push(user); + } + instance_table["dm_allowed_users"] = toml_edit::value(dm_array); + } + Some(_) => { + // Explicitly cleared (empty vec) - remove the key + instance_table.remove("dm_allowed_users"); + } + None => { + // Omitted - don't set the field. When loaded, the absence + // of dm_allowed_users results in an empty Vec which blocks + // all DMs (the safe default). + } + } + } _ => {} } @@ -1542,7 +2031,7 @@ pub(super) async fn delete_messaging_instance( if !matches!( platform.as_str(), - "discord" | "slack" | "telegram" | "twitch" | "email" | "webhook" | "mattermost" + "discord" | "slack" | "telegram" | "twitch" | "email" | "webhook" | "mattermost" | "signal" ) { return Ok(Json(MessagingInstanceActionResponse { success: false, @@ -1646,6 +2135,14 @@ pub(super) async fn delete_messaging_instance( table.remove("dm_allowed_users"); table.remove("max_attachment_bytes"); } + "signal" => { + table.remove("http_url"); + table.remove("account"); + table.remove("dm_allowed_users"); + table.remove("group_ids"); + table.remove("group_allowed_users"); + table.remove("ignore_stories"); + } _ => {} } } @@ -1704,6 +2201,20 @@ pub(super) async fn delete_messaging_instance( { tracing::warn!(%error, adapter = %runtime_key, "failed to shut down adapter during instance delete"); } + drop(manager_guard); + + // Clear cached Signal permissions when deleting the default Signal adapter + if platform == "signal" && adapter_name.is_none() { + let perms_guard = state.signal_permissions.read().await; + if perms_guard.is_some() { + drop(perms_guard); + state + .set_signal_permissions(std::sync::Arc::new(arc_swap::ArcSwap::from_pointee( + crate::config::SignalPermissions::default(), + ))) + .await; + } + } // Clean up twitch token file if applicable if platform == "twitch" { @@ -1739,3 +2250,49 @@ pub(super) async fn delete_messaging_instance( message: format!("{runtime_key} instance deleted"), })) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_is_valid_e164_valid_numbers() { + // Valid E.164 numbers (6-15 digits after +, 7-16 total) + assert!(is_valid_e164("+1234567890")); + assert!(is_valid_e164("+123456")); // Minimum: 6 digits after + + assert!(is_valid_e164("+14155552671")); + assert!(is_valid_e164("+12345678901234")); // 14 digits after + + assert!(is_valid_e164("+123456789012345")); // Maximum: 15 digits after + + } + + #[test] + fn test_is_valid_e164_invalid_numbers() { + // Missing + + assert!(!is_valid_e164("1234567890")); + + // Too short (less than 6 digits after +) + assert!(!is_valid_e164("+12345")); + assert!(!is_valid_e164("+123")); + + // Empty or just + + assert!(!is_valid_e164("+")); + assert!(!is_valid_e164("")); + + // Non-digit characters + assert!(!is_valid_e164("+1234567890a")); + assert!(!is_valid_e164("+123-456-7890")); + assert!(!is_valid_e164("+123 456 7890")); + + // Spaces + assert!(!is_valid_e164(" +1234567890")); + assert!(!is_valid_e164("+1234567890 ")); + + // First digit is 0 (E.164 requires 1-9) + assert!(!is_valid_e164("+0123456")); + assert!(!is_valid_e164("+01234567890")); + + // Too long (more than 15 digits after +) + assert!(!is_valid_e164("+1234567890123456")); + assert!(!is_valid_e164("+12345678901234567")); + } +} diff --git a/src/api/state.rs b/src/api/state.rs index 9500617e1..55416c854 100644 --- a/src/api/state.rs +++ b/src/api/state.rs @@ -3,7 +3,9 @@ use crate::agent::channel::ChannelState; use crate::agent::cortex_chat::CortexChatSession; use crate::agent::status::StatusBlock; -use crate::config::{Binding, DefaultsConfig, DiscordPermissions, RuntimeConfig, SlackPermissions}; +use crate::config::{ + Binding, DefaultsConfig, DiscordPermissions, RuntimeConfig, SignalPermissions, SlackPermissions, +}; use crate::conversation::worker_transcript::{ActionContent, TranscriptStep}; use crate::cron::{CronStore, Scheduler}; use crate::llm::LlmManager; @@ -96,6 +98,8 @@ pub struct ApiState { pub discord_permissions: RwLock>>>, /// Shared reference to the Slack permissions ArcSwap (same instance used by the adapter and file watcher). pub slack_permissions: RwLock>>>, + /// Shared reference to the Signal permissions ArcSwap (same instance used by the adapter and file watcher). + pub signal_permissions: RwLock>>>, /// Shared reference to the bindings ArcSwap (same instance used by the main loop and file watcher). pub bindings: RwLock>>>>, /// Shared messaging manager for runtime adapter addition. @@ -316,6 +320,7 @@ impl ApiState { secrets_store: ArcSwap::from_pointee(None), discord_permissions: RwLock::new(None), slack_permissions: RwLock::new(None), + signal_permissions: RwLock::new(None), bindings: RwLock::new(None), messaging_manager: RwLock::new(None), provider_setup_tx, @@ -786,6 +791,11 @@ impl ApiState { *self.slack_permissions.write().await = Some(permissions); } + /// Share the Signal permissions ArcSwap with the API so reads get hot-reloaded values. + pub async fn set_signal_permissions(&self, permissions: Arc>) { + *self.signal_permissions.write().await = Some(permissions); + } + /// Share the bindings ArcSwap with the API so reads get hot-reloaded values. pub async fn set_bindings(&self, bindings: Arc>>) { *self.bindings.write().await = Some(bindings); diff --git a/src/config/load.rs b/src/config/load.rs index 05a99aedc..fb74e3b3e 100644 --- a/src/config/load.rs +++ b/src/config/load.rs @@ -2217,9 +2217,7 @@ impl Config { .ok() .or_else(|| s.account.as_deref().and_then(resolve_env_value)); - if (http_url.is_none() || account.is_none()) - && !instances.iter().any(|inst| inst.enabled) - { + if (http_url.is_none() || account.is_none()) && instances.is_empty() { return None; } diff --git a/src/config/watcher.rs b/src/config/watcher.rs index ab64d235c..f197897d5 100644 --- a/src/config/watcher.rs +++ b/src/config/watcher.rs @@ -553,61 +553,72 @@ pub fn spawn_file_watcher( } } - // Signal: start default + named instances that are enabled and not already running. - if let Some(signal_config) = &config.messaging.signal - && signal_config.enabled { - if !signal_config.http_url.is_empty() - && !signal_config.account.is_empty() - && !manager.has_adapter("signal").await - { - let permissions = match signal_permissions { - Some(ref existing) => existing.clone(), - None => { - let permissions = SignalPermissions::from_config(signal_config); - Arc::new(arc_swap::ArcSwap::from_pointee(permissions)) - } - }; - let tmp_dir = instance_dir.join("tmp"); - let adapter = crate::messaging::signal::SignalAdapter::new( - "signal", - &signal_config.http_url, - &signal_config.account, - signal_config.ignore_stories, - permissions, - tmp_dir, - ); - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, "failed to hot-start signal adapter from config change"); + // Signal: start default adapter (requires root enabled) and named instances (independent). + // Unlike Discord/Telegram where named instances inherit the root enabled gate, + // Signal named instances start independently when they have valid credentials + // and their own enabled flag is set. This allows running multiple Signal accounts + // without needing a "default" account enabled. + if let Some(signal_config) = &config.messaging.signal { + // Start default adapter only if root is enabled AND has credentials + if signal_config.enabled + && !signal_config.http_url.is_empty() + && !signal_config.account.is_empty() + && !manager.has_adapter("signal").await + { + let permissions = match signal_permissions { + Some(ref existing) => existing.clone(), + None => { + let permissions = SignalPermissions::from_config(signal_config); + Arc::new(arc_swap::ArcSwap::from_pointee(permissions)) } + }; + let tmp_dir = instance_dir.join("tmp"); + let adapter = crate::messaging::signal::SignalAdapter::new( + "signal", + &signal_config.http_url, + &signal_config.account, + signal_config.ignore_stories, + permissions, + tmp_dir, + ); + if let Err(error) = manager.register_and_start(adapter).await { + tracing::error!(%error, "failed to hot-start signal adapter from config change"); } + } - for instance in signal_config.instances.iter().filter(|instance| instance.enabled) { - let runtime_key = binding_runtime_adapter_key( - "signal", - Some(instance.name.as_str()), - ); - if manager.has_adapter(runtime_key.as_str()).await { - // TODO: named instance permissions not hot-updated (see discord block comment) - continue; - } + // Start named instances regardless of root enabled flag (as long as config exists) + for instance in signal_config.instances.iter().filter(|instance| instance.enabled) { + let runtime_key = binding_runtime_adapter_key( + "signal", + Some(instance.name.as_str()), + ); + if manager.has_adapter(runtime_key.as_str()).await { + // TODO: named instance permissions not hot-updated (see discord block comment) + continue; + } + // Skip instances with missing credentials (same gating as cold-start) + if instance.http_url.is_empty() || instance.account.is_empty() { + tracing::warn!(adapter = %instance.name, "skipping enabled signal instance with missing credentials"); + continue; + } - let permissions = Arc::new(arc_swap::ArcSwap::from_pointee( - SignalPermissions::from_instance_config(instance), - )); - let tmp_dir = instance_dir.join("tmp"); - let adapter = crate::messaging::signal::SignalAdapter::new( - runtime_key, - &instance.http_url, - &instance.account, - instance.ignore_stories, - permissions, - tmp_dir, - ); - if let Err(error) = manager.register_and_start(adapter).await { - tracing::error!(%error, adapter = %instance.name, "failed to hot-start named signal adapter from config change"); - } + let permissions = Arc::new(arc_swap::ArcSwap::from_pointee( + SignalPermissions::from_instance_config(instance), + )); + let tmp_dir = instance_dir.join("tmp"); + let adapter = crate::messaging::signal::SignalAdapter::new( + runtime_key, + &instance.http_url, + &instance.account, + instance.ignore_stories, + permissions, + tmp_dir, + ); + if let Err(error) = manager.register_and_start(adapter).await { + tracing::error!(%error, adapter = %instance.name, "failed to hot-start named signal adapter from config change"); } } + } // Mattermost: start default + named instances that are enabled and not already running. if let Some(mattermost_config) = &config.messaging.mattermost diff --git a/src/cron/scheduler.rs b/src/cron/scheduler.rs index 1c4f8b4df..b1a55ab07 100644 --- a/src/cron/scheduler.rs +++ b/src/cron/scheduler.rs @@ -893,10 +893,19 @@ async fn run_cron_job(job: &CronJob, context: &CronContext) -> Result<()> { }); // Send the cron job prompt as a synthetic message + // Derive source from the delivery target's adapter so adapter_selector() can extract + // the platform prefix (e.g., "signal" from "signal:gvoice1"). This ensures the channel + // correctly identifies as being "from" that messaging platform and tools resolve properly. + let source_adapter = job + .delivery_target + .adapter + .split(':') + .next() + .unwrap_or("cron"); let message = InboundMessage { id: uuid::Uuid::new_v4().to_string(), - source: "cron".into(), - adapter: None, + source: source_adapter.into(), + adapter: Some(job.delivery_target.adapter.clone()), conversation_id: format!("cron:{}", job.id), sender_id: "system".into(), agent_id: Some(context.deps.agent_id.clone()), diff --git a/src/main.rs b/src/main.rs index 0502e7508..ca89ccd72 100644 --- a/src/main.rs +++ b/src/main.rs @@ -3303,12 +3303,22 @@ async fn initialize_agents( let perms = spacebot::config::SignalPermissions::from_config(signal_config); Arc::new(ArcSwap::from_pointee(perms)) }); + if let Some(perms) = &*signal_permissions { + api_state.set_signal_permissions(perms.clone()).await; + } - if let Some(signal_config) = &config.messaging.signal - && signal_config.enabled - { - let tmp_dir = config.instance_dir.join("tmp"); - if !signal_config.http_url.is_empty() && !signal_config.account.is_empty() { + // Signal: start default adapter (requires root enabled) and named instances (independent). + // Unlike Discord/Telegram where named instances inherit the root enabled gate, + // Signal named instances start independently when they have valid credentials + // and their own enabled flag is set. This allows running multiple Signal accounts + // without needing a "default" account enabled. + let tmp_dir = config.instance_dir.join("tmp"); + if let Some(signal_config) = &config.messaging.signal { + // Start default adapter only if root is enabled AND has credentials + if signal_config.enabled + && !signal_config.http_url.is_empty() + && !signal_config.account.is_empty() + { let adapter = spacebot::messaging::signal::SignalAdapter::new( "signal", &signal_config.http_url, @@ -3322,6 +3332,7 @@ async fn initialize_agents( new_messaging_manager.register(adapter).await; } + // Start named instances regardless of root enabled flag (as long as config exists) for instance in signal_config .instances .iter() @@ -3449,7 +3460,11 @@ async fn initialize_agents( } // Store cron tool on deps so each channel can register it on its own tool server - let cron_tool = spacebot::tools::CronTool::new(store.clone(), scheduler.clone()); + let cron_tool = spacebot::tools::CronTool::new( + store.clone(), + scheduler.clone(), + messaging_manager.clone(), + ); agent.deps.cron_tool = Some(cron_tool); cron_stores_map.insert(agent_id.to_string(), store); diff --git a/src/messaging/signal.rs b/src/messaging/signal.rs index b9f16a0a6..ff4fe72b5 100644 --- a/src/messaging/signal.rs +++ b/src/messaging/signal.rs @@ -247,7 +247,7 @@ impl SignalAdapter { Self { runtime_key: runtime_key.into(), - http_url: http_url.into(), + http_url: http_url.into().trim_end_matches('/').to_string(), account: account.into(), ignore_stories, permissions, diff --git a/src/messaging/target.rs b/src/messaging/target.rs index c851750ff..e3ddcca24 100644 --- a/src/messaging/target.rs +++ b/src/messaging/target.rs @@ -16,7 +16,26 @@ impl std::fmt::Display for BroadcastTarget { } /// Parse and normalize a delivery target in `adapter:target` format. +/// +/// Signal targets may contain named instance prefixes (e.g., +/// `signal:gvoice1:uuid:xxx`). The generic `split_once(':')` approach +/// cannot distinguish the instance segment from the target, so we +/// delegate to `parse_signal_target_parts` which already handles this. pub fn parse_delivery_target(raw: &str) -> Option { + // Signal needs special handling because named instances add an extra + // colon-separated segment that the generic parser can't distinguish. + if raw.starts_with("signal:") { + let parts: Vec<&str> = raw.split(':').collect(); + return parse_signal_target_parts(parts.get(1..).unwrap_or(&[])); + } + + // Handle other platforms with named instances (telegram, discord, slack) + // Format: platform:: or platform: + if raw.starts_with("telegram:") || raw.starts_with("discord:") || raw.starts_with("slack:") { + let parts: Vec<&str> = raw.split(':').collect(); + return parse_named_instance_target(&parts); + } + let (adapter, raw_target) = raw.split_once(':')?; if adapter.is_empty() || raw_target.is_empty() { return None; @@ -345,6 +364,29 @@ fn normalize_email_target(raw_target: &str) -> Option { None } +/// Validate E.164 phone number format. +/// +/// Requirements: +/// - Must start with '+' +/// - First digit after '+' must be 1-9 (not 0) +/// - Minimum 7 digits total (6 after '+') +/// - Maximum 16 digits total (15 after '+', E.164 standard) +/// - All characters after '+' must be ASCII digits +pub fn is_valid_e164(phone: &str) -> bool { + if let Some(digits) = phone.strip_prefix('+') { + if digits.len() < 6 || digits.len() > 15 { + return false; + } + // First digit must be 1-9 (not 0) + if digits.chars().next().map(|c| c == '0').unwrap_or(true) { + return false; + } + digits.chars().all(|c| c.is_ascii_digit()) + } else { + false + } +} + fn normalize_signal_target(raw_target: &str) -> Option { let target = strip_repeated_prefix(raw_target, "signal"); @@ -364,31 +406,34 @@ fn normalize_signal_target(raw_target: &str) -> Option { return None; } - // Handle e164:+123 or bare +123 format + // Handle e164:+123 format if let Some(phone) = target.strip_prefix("e164:") { - let phone = phone.trim_start_matches('+'); - if !phone.is_empty() && phone.len() >= 7 && phone.chars().all(|c| c.is_ascii_digit()) { - return Some(format!("+{phone}")); + let normalized = format!("+{}", phone.trim_start_matches('+')); + if is_valid_e164(&normalized) { + return Some(normalized); } return None; } // Bare +123 format - if let Some(phone) = target.strip_prefix('+') { - if !phone.is_empty() && phone.len() >= 7 && phone.chars().all(|c| c.is_ascii_digit()) { + if target.starts_with('+') { + if is_valid_e164(target) { return Some(target.to_string()); } return None; } - // Check if it's a valid UUID (contains dashes and alphanumeric) - if target.contains('-') && target.len() > 8 && target.chars().any(|c| c.is_ascii_digit()) { + // Check if it's a bare UUID using strict validation + if uuid::Uuid::parse_str(target).is_ok() { return Some(format!("uuid:{target}")); } - // Check if it's a bare phone number (7+ digits required for E.164) - if target.chars().all(|c| c.is_ascii_digit()) && target.len() >= 7 { - return Some(format!("+{target}")); + // Check if it's a bare phone number (E.164 format) + if target.chars().all(|c| c.is_ascii_digit()) { + let with_plus = format!("+{target}"); + if is_valid_e164(&with_plus) { + return Some(with_plus); + } } None @@ -452,64 +497,191 @@ fn extract_signal_adapter_from_channel_id(channel_id: &str) -> String { pub fn parse_signal_target_parts(parts: &[&str]) -> Option { match parts { // Default adapter: signal:uuid:xxx, signal:group:xxx, signal:e164:+xxx, signal:+xxx - ["uuid", uuid] => Some(BroadcastTarget { + ["uuid", uuid] if !uuid.is_empty() => Some(BroadcastTarget { adapter: "signal".to_string(), target: format!("uuid:{uuid}"), }), - ["group", group_id] => Some(BroadcastTarget { + ["group", group_id] if !group_id.is_empty() => Some(BroadcastTarget { adapter: "signal".to_string(), target: format!("group:{group_id}"), }), // Use normalize_signal_target for phone/e164 to ensure consistent parsing - ["e164", phone] => { - normalize_signal_target(&format!("e164:{phone}")).map(|target| BroadcastTarget { + ["e164", phone] if !phone.is_empty() => normalize_signal_target(&format!("e164:{phone}")) + .map(|target| BroadcastTarget { adapter: "signal".to_string(), target, - }) - } - [phone] if phone.starts_with('+') => { - normalize_signal_target(phone).map(|target| BroadcastTarget { + }), + [phone] if phone.starts_with('+') && !phone.is_empty() => normalize_signal_target(phone) + .map(|target| BroadcastTarget { + adapter: "signal".to_string(), + target, + }), + // Single-part targets: delegate to normalize_signal_target for bare UUIDs/phones + [single] if !single.is_empty() => { + normalize_signal_target(single).map(|target| BroadcastTarget { adapter: "signal".to_string(), target, }) } - // Single-part targets: delegate to normalize_signal_target for bare UUIDs/phones - [single] => normalize_signal_target(single).map(|target| BroadcastTarget { - adapter: "signal".to_string(), - target, - }), // Named adapter: signal:instance:uuid:xxx, signal:instance:group:xxx - [instance, "uuid", uuid] => Some(BroadcastTarget { - adapter: format!("signal:{instance}"), - target: format!("uuid:{uuid}"), - }), - [instance, "group", group_id] => Some(BroadcastTarget { - adapter: format!("signal:{instance}"), - target: format!("group:{group_id}"), - }), + [instance, "uuid", uuid] + if !instance.is_empty() && !uuid.is_empty() && is_valid_instance_name(instance) => + { + Some(BroadcastTarget { + adapter: format!("signal:{instance}"), + target: format!("uuid:{uuid}"), + }) + } + [instance, "group", group_id] + if !instance.is_empty() && !group_id.is_empty() && is_valid_instance_name(instance) => + { + Some(BroadcastTarget { + adapter: format!("signal:{instance}"), + target: format!("group:{group_id}"), + }) + } // Named adapter: signal:instance:e164:+xxx - use normalize_signal_target - [instance, "e164", phone] => { + [instance, "e164", phone] + if !instance.is_empty() && !phone.is_empty() && is_valid_instance_name(instance) => + { normalize_signal_target(&format!("e164:{phone}")).map(|target| BroadcastTarget { adapter: format!("signal:{instance}"), target, }) } // Named adapter: signal:instance:+xxx - use normalize_signal_target - [instance, phone] if phone.starts_with('+') => { + [instance, phone] + if !instance.is_empty() + && phone.starts_with('+') + && !phone.is_empty() + && is_valid_instance_name(instance) => + { normalize_signal_target(phone).map(|target| BroadcastTarget { adapter: format!("signal:{instance}"), target, }) } // Named adapter with single-part target: delegate to normalize_signal_target - [instance, single] => normalize_signal_target(single).map(|target| BroadcastTarget { - adapter: format!("signal:{instance}"), - target, - }), + // Reject all-digit identifiers to avoid misinterpreting numeric IDs as phone numbers + [instance, single] + if !instance.is_empty() + && !single.is_empty() + && !single.chars().all(|c| c.is_ascii_digit()) + && is_valid_instance_name(instance) => + { + normalize_signal_target(single).map(|target| BroadcastTarget { + adapter: format!("signal:{instance}"), + target, + }) + } + _ => None, + } +} + +/// Parse targets for platforms with named instance support (telegram, discord, slack). +/// +/// Handles formats: +/// - Default adapter: ["telegram", target], ["discord", target], ["slack", target] +/// - Legacy format: ["discord", guild_id, channel_id], ["slack", workspace_id, channel_id] +/// - Named adapter: ["telegram", instance, target], ["discord", instance, target], ["slack", instance, target] +/// +/// Returns None for invalid formats. +fn parse_named_instance_target(parts: &[&str]) -> Option { + if parts.len() < 2 { + return None; + } + + let platform = parts[0]; + let is_telegram = platform == "telegram"; + let is_discord = platform == "discord"; + let is_slack = platform == "slack"; + + if !is_telegram && !is_discord && !is_slack { + return None; + } + + // Reject any cases with empty parts + if parts.iter().any(|part| part.is_empty()) { + return None; + } + + match parts { + // "dm" reserved for Slack/Discord (case-insensitive) + [platform @ ("discord" | "slack"), instance, user_id] + if instance.eq_ignore_ascii_case("dm") && !user_id.is_empty() => + { + let normalized = normalize_target(platform, &format!("dm:{user_id}"))?; + Some(BroadcastTarget { + adapter: (*platform).to_string(), + target: normalized, + }) + } + // Named adapter: platform:instance:target + // Heuristic: instance names are simple alphanumeric identifiers (not all digits like guild IDs) + // Reject "dm" as instance name for Discord/Slack (reserved for DM format, case-insensitive) + [platform, instance, target, ..] + if parts.len() >= 3 + && !instance.is_empty() + && !target.is_empty() + && is_valid_instance_name(instance) + && !((*platform == "discord" || *platform == "slack") + && instance.eq_ignore_ascii_case("dm")) + && !instance.eq_ignore_ascii_case(platform) => + { + // Reconstruct target from remaining parts (in case target contains colons) + // Normalize to collapse guild-prefixed IDs and reject malformed inputs + let full_target = parts[2..].join(":"); + let normalized = normalize_target(platform, &full_target)?; + Some(BroadcastTarget { + adapter: format!("{}:{}", platform, instance), + target: normalized, + }) + } + // Legacy multi-part format: platform:workspace_id:target_id (use last part as target) + // Restrict to known legacy adapters (slack, discord) to prevent malformed named-targets + // like "telegram:bad.instance:12345" from being interpreted as "telegram:12345" + [platform @ ("slack" | "discord"), .., target] if parts.len() > 2 && !target.is_empty() => { + let normalized = normalize_target(platform, target)?; + Some(BroadcastTarget { + adapter: (*platform).to_string(), + target: normalized, + }) + } + // Default adapter: platform:target + [platform, target] if !target.is_empty() => { + let normalized = normalize_target(platform, target)?; + Some(BroadcastTarget { + adapter: (*platform).to_string(), + target: normalized, + }) + } _ => None, } } +/// Check if a string looks like a valid instance name (not an ID). +/// Instance names are short identifiers, not numeric IDs or workspace identifiers. +pub fn is_valid_instance_name(name: &str) -> bool { + // Must not be all digits (Discord guild ID, etc.) + if name.chars().all(|c| c.is_ascii_digit()) { + return false; + } + // Must not look like a Slack workspace ID (Txxxxx, Cxxxxx, etc.) + if name.len() > 6 + && name.starts_with(|c: char| c.is_ascii_uppercase()) + && name[1..].chars().all(|c| c.is_ascii_digit()) + { + return false; + } + // Must be reasonably short (instance names are short, not long IDs) + if name.len() > 20 { + return false; + } + // Must contain only alphanumeric characters, underscores, and hyphens + name.chars() + .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-') +} + #[cfg(test)] mod tests { use super::{parse_delivery_target, resolve_broadcast_target}; @@ -790,4 +962,140 @@ mod tests { assert!(super::parse_signal_target_parts(&["uuid"]).is_none()); // missing UUID value assert!(super::parse_signal_target_parts(&["gvoice1", "unknown"]).is_none()); } + + // Tests for parse_named_instance_target + #[test] + fn parse_named_instance_target_telegram() { + let parsed = super::parse_named_instance_target(&["telegram", "mybot", "12345"]); + assert_eq!( + parsed, + Some(super::BroadcastTarget { + adapter: "telegram:mybot".to_string(), + target: "12345".to_string(), + }) + ); + } + + #[test] + fn parse_named_instance_target_discord_named() { + // Named instance: discord:myinstance:channel_id + let parsed = super::parse_named_instance_target(&["discord", "myinstance", "987654321"]); + assert_eq!( + parsed, + Some(super::BroadcastTarget { + adapter: "discord:myinstance".to_string(), + target: "987654321".to_string(), + }) + ); + } + + #[test] + fn parse_named_instance_target_discord_legacy() { + // Legacy format: discord:guild_id:channel_id (all digits = not a named instance) + let parsed = super::parse_named_instance_target(&["discord", "123456789", "987654321"]); + assert_eq!( + parsed, + Some(super::BroadcastTarget { + adapter: "discord".to_string(), + target: "987654321".to_string(), + }) + ); + } + + #[test] + fn parse_named_instance_target_slack_named() { + // Named instance: slack:work:channel_id + let parsed = super::parse_named_instance_target(&["slack", "work", "C012345"]); + assert_eq!( + parsed, + Some(super::BroadcastTarget { + adapter: "slack:work".to_string(), + target: "C012345".to_string(), + }) + ); + } + + #[test] + fn parse_named_instance_target_slack_legacy() { + // Legacy format: slack:workspace_id:channel_id (workspace_id pattern = not named) + let parsed = super::parse_named_instance_target(&["slack", "T012345", "C012345"]); + assert_eq!( + parsed, + Some(super::BroadcastTarget { + adapter: "slack".to_string(), + target: "C012345".to_string(), + }) + ); + } + + #[test] + fn parse_named_instance_target_default_adapter() { + // Default adapter (2 parts): telegram:target + let parsed = super::parse_named_instance_target(&["telegram", "12345"]); + assert_eq!( + parsed, + Some(super::BroadcastTarget { + adapter: "telegram".to_string(), + target: "12345".to_string(), + }) + ); + } + + #[test] + fn parse_named_instance_target_invalid() { + // Too few parts + assert!(super::parse_named_instance_target(&["telegram"]).is_none()); + // Empty instance name + assert!(super::parse_named_instance_target(&["telegram", "", "12345"]).is_none()); + // Empty target + assert!(super::parse_named_instance_target(&["telegram", "work", ""]).is_none()); + // Unsupported platform + assert!(super::parse_named_instance_target(&["unknown", "work", "12345"]).is_none()); + } + + // Tests for is_valid_instance_name + #[test] + fn is_valid_instance_name_valid() { + assert!(super::is_valid_instance_name("work")); + assert!(super::is_valid_instance_name("my_bot")); + assert!(super::is_valid_instance_name("my-bot")); + assert!(super::is_valid_instance_name("instance123")); + assert!(super::is_valid_instance_name("a")); + } + + #[test] + fn is_valid_instance_name_invalid_all_digits() { + // All digits = numeric ID, not an instance name + assert!(!super::is_valid_instance_name("12345")); + assert!(!super::is_valid_instance_name("123456789")); + } + + #[test] + fn is_valid_instance_name_invalid_slack_workspace() { + // Slack workspace ID pattern + assert!(!super::is_valid_instance_name("T012345")); + assert!(!super::is_valid_instance_name("C012345")); + } + + #[test] + fn is_valid_instance_name_invalid_length() { + // Too long (>20 chars) + assert!(!super::is_valid_instance_name( + "this_is_a_very_long_instance_name" + )); + } + + #[test] + fn is_valid_instance_name_invalid_empty() { + // Empty string + assert!(!super::is_valid_instance_name("")); + } + + #[test] + fn is_valid_instance_name_edge_cases() { + // Exactly 20 characters (boundary) + assert!(super::is_valid_instance_name("exactly_twenty_chars")); + // 21 characters (over boundary) + assert!(!super::is_valid_instance_name("exactly_twenty_chars_")); + } } diff --git a/src/prompts/engine.rs b/src/prompts/engine.rs index 5cac008c5..e6ce5aac5 100644 --- a/src/prompts/engine.rs +++ b/src/prompts/engine.rs @@ -240,6 +240,7 @@ impl PromptEngine { platform: &str, server_name: Option<&str>, channel_name: Option<&str>, + conversation_id: Option<&str>, ) -> Result { self.render( "fragments/conversation_context", @@ -247,6 +248,7 @@ impl PromptEngine { platform => platform, server_name => server_name, channel_name => channel_name, + conversation_id => conversation_id, }, ) } diff --git a/src/tools.rs b/src/tools.rs index 59e10e0b6..9e19c4538 100644 --- a/src/tools.rs +++ b/src/tools.rs @@ -434,9 +434,12 @@ pub async fn add_channel_tools( .await?; handle.add_tool(ReactTool::new(response_tx.clone())).await?; if let Some(cron_tool) = cron_tool { - let cron_tool = cron_tool.with_default_delivery_target( - default_delivery_target_for_conversation(&conversation_id, slack_thread_ts), - ); + let cron_tool = cron_tool + .with_default_delivery_target(default_delivery_target_for_conversation( + &conversation_id, + slack_thread_ts, + )) + .with_current_adapter(current_adapter.clone()); handle.add_tool(cron_tool).await?; } if let Some(mut agent_msg) = send_agent_message_tool { diff --git a/src/tools/cron.rs b/src/tools/cron.rs index 6a1ffddfe..628c4366b 100644 --- a/src/tools/cron.rs +++ b/src/tools/cron.rs @@ -2,6 +2,7 @@ use crate::cron::scheduler::{CronConfig, Scheduler}; use crate::cron::store::CronStore; +use crate::messaging::MessagingManager; use rig::completion::ToolDefinition; use rig::tool::Tool; use schemars::JsonSchema; @@ -16,19 +17,38 @@ const MIN_CRON_INTERVAL_SECS: u64 = 60; const MAX_CRON_PROMPT_LENGTH: usize = 10_000; /// Tool for managing cron jobs (scheduled recurring tasks). -#[derive(Debug, Clone)] +#[derive(Clone)] pub struct CronTool { store: Arc, scheduler: Arc, + messaging_manager: Arc, default_delivery_target: Option, + current_adapter: Option, +} + +impl std::fmt::Debug for CronTool { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("CronTool") + .field("store", &self.store) + .field("scheduler", &self.scheduler) + .field("default_delivery_target", &self.default_delivery_target) + .field("current_adapter", &self.current_adapter) + .finish_non_exhaustive() + } } impl CronTool { - pub fn new(store: Arc, scheduler: Arc) -> Self { + pub fn new( + store: Arc, + scheduler: Arc, + messaging_manager: Arc, + ) -> Self { Self { store, scheduler, + messaging_manager, default_delivery_target: None, + current_adapter: None, } } @@ -36,6 +56,11 @@ impl CronTool { self.default_delivery_target = default_delivery_target; self } + + pub fn with_current_adapter(mut self, current_adapter: Option) -> Self { + self.current_adapter = current_adapter; + self + } } #[derive(Debug, thiserror::Error)] @@ -194,12 +219,14 @@ impl CronTool { .filter(|value| !value.is_empty()) .map(ToString::to_string); let interval_secs = args.interval_secs.unwrap_or(3600); - let delivery_target = args + let explicit_delivery_target = args .delivery_target .as_deref() .map(str::trim) .filter(|value| !value.is_empty()) - .map(|value| value.to_string()) + .map(|value| value.to_string()); + let delivery_target = explicit_delivery_target + .clone() .or_else(|| self.default_delivery_target.clone()) .ok_or_else(|| { CronError( @@ -208,6 +235,84 @@ impl CronTool { ) })?; + // Only normalize delivery target when it came from the conversation default, + // not when the user explicitly provided one. An explicit "signal:uuid:xxx" + // targets the default adapter by design; the default fallback from the + // conversation context should be rewritten for the named instance. + // + // For explicit signal: targets, resolve the adapter using the same logic as + // send_message_to_another_channel to handle deployments with only named instances. + let delivery_target = if explicit_delivery_target.is_some() { + let parsed_delivery_target = crate::messaging::target::parse_delivery_target( + &delivery_target, + ) + .ok_or_else(|| CronError(format!("invalid 'delivery_target': '{delivery_target}'")))?; + let should_resolve_signal_default = parsed_delivery_target.adapter == "signal"; + + if should_resolve_signal_default { + let target_part = parsed_delivery_target.target; + let resolved_adapter = + crate::tools::send_message_to_another_channel::resolve_signal_adapter( + &self.messaging_manager, + self.current_adapter.as_deref(), + ) + .await + .map_err(|e| CronError(format!("Failed to resolve Signal adapter: {e}")))?; + if resolved_adapter == "signal" + && !self.messaging_manager.has_adapter("signal").await + { + return Err(CronError( + "No Signal adapter running. Cannot create cron with Signal delivery target." + .into(), + )); + } + format!("{resolved_adapter}:{target_part}") + } else { + if !self + .messaging_manager + .has_adapter(&parsed_delivery_target.adapter) + .await + { + return Err(CronError(if parsed_delivery_target.adapter.contains(':') { + format!("No '{}' adapter running.", parsed_delivery_target.adapter) + } else { + format!( + "No '{}' adapter running. Use '{}::...' to target a specific instance.", + parsed_delivery_target.adapter, parsed_delivery_target.adapter + ) + })); + } + delivery_target + } + } else if let Some(current) = &self.current_adapter + && self.messaging_manager.has_adapter(current).await + { + normalize_delivery_target(&delivery_target, &self.current_adapter) + } else { + // Validate the raw delivery_target before persisting — the adapter may have + // been captured earlier but is no longer available, or was never valid. + let parsed_delivery_target = crate::messaging::target::parse_delivery_target( + &delivery_target, + ) + .ok_or_else(|| CronError(format!("invalid 'delivery_target': '{delivery_target}'")))?; + + if !self + .messaging_manager + .has_adapter(&parsed_delivery_target.adapter) + .await + { + return Err(CronError(if parsed_delivery_target.adapter.contains(':') { + format!("No '{}' adapter running.", parsed_delivery_target.adapter) + } else { + format!( + "No '{}' adapter running. Use '{}::...' to target a specific instance.", + parsed_delivery_target.adapter, parsed_delivery_target.adapter + ) + })); + } + delivery_target + }; + // Validate cron job ID: alphanumeric, hyphens, underscores only if id.is_empty() || id.len() > 50 @@ -411,3 +516,175 @@ fn format_interval(secs: u64) -> String { format!("every {secs} seconds") } } + +/// Normalize delivery target for named instances. +/// +/// If the LLM provided a bare platform adapter (e.g., "signal", "slack") but we're in +/// a named instance conversation (e.g., "signal:gvoice1", "slack:work"), rewrite to +/// include the instance name. This ensures the cron job can find the correct adapter +/// at runtime. +fn normalize_delivery_target(delivery_target: &str, current_adapter: &Option) -> String { + if let Some(parsed) = crate::messaging::target::parse_delivery_target(delivery_target) + && let Some(current_adapter) = current_adapter.as_ref() + { + let expected_prefix = format!("{}:", parsed.adapter); + if current_adapter.starts_with(&expected_prefix) { + // current_adapter is a named instance of the parsed platform + return format!("{current_adapter}:{}", parsed.target); + } + } + delivery_target.to_string() +} + +#[cfg(test)] +mod tests { + use super::normalize_delivery_target; + + #[test] + fn test_normalize_signal_named_instance() { + // Bare signal + named instance context → rewrite + let result = normalize_delivery_target( + "signal:uuid:550e8400-e29b-41d4-a716-446655440000", + &Some("signal:gvoice1".to_string()), + ); + assert_eq!( + result, + "signal:gvoice1:uuid:550e8400-e29b-41d4-a716-446655440000" + ); + } + + #[test] + fn test_normalize_signal_group_named_instance() { + // Signal group + named instance context → rewrite + let result = + normalize_delivery_target("signal:group:grp123", &Some("signal:work".to_string())); + assert_eq!(result, "signal:work:group:grp123"); + } + + #[test] + fn test_normalize_signal_default_instance_no_rewrite() { + // Bare signal + default signal context → no change + let result = normalize_delivery_target( + "signal:uuid:550e8400-e29b-41d4-a716-446655440000", + &Some("signal".to_string()), + ); + assert_eq!(result, "signal:uuid:550e8400-e29b-41d4-a716-446655440000"); + } + + #[test] + fn test_normalize_signal_no_context_no_rewrite() { + // Bare signal + no context → no change + let result = + normalize_delivery_target("signal:uuid:550e8400-e29b-41d4-a716-446655440000", &None); + assert_eq!(result, "signal:uuid:550e8400-e29b-41d4-a716-446655440000"); + } + + #[test] + fn test_normalize_slack_named_instance() { + // Bare slack + named instance context → rewrite + let result = normalize_delivery_target("slack:C123456", &Some("slack:work".to_string())); + assert_eq!(result, "slack:work:C123456"); + } + + #[test] + fn test_normalize_discord_named_instance() { + // Bare discord + named instance context → rewrite + let result = + normalize_delivery_target("discord:987654321", &Some("discord:personal".to_string())); + assert_eq!(result, "discord:personal:987654321"); + } + + #[test] + fn test_normalize_telegram_named_instance() { + // Bare telegram + named instance context → rewrite + let result = + normalize_delivery_target("telegram:-1001234", &Some("telegram:bot1".to_string())); + assert_eq!(result, "telegram:bot1:-1001234"); + } + + #[test] + fn test_normalize_different_platform_no_rewrite() { + // Signal target but in Discord context → no change + let result = normalize_delivery_target( + "signal:uuid:550e8400-e29b-41d4-a716-446655440000", + &Some("discord:general".to_string()), + ); + assert_eq!(result, "signal:uuid:550e8400-e29b-41d4-a716-446655440000"); + } + + #[test] + fn test_normalize_cron_context_no_rewrite() { + // Any target in cron context → no change (cron can't receive delivery) + let result = normalize_delivery_target( + "signal:uuid:550e8400-e29b-41d4-a716-446655440000", + &Some("cron".to_string()), + ); + assert_eq!(result, "signal:uuid:550e8400-e29b-41d4-a716-446655440000"); + } + + #[test] + fn test_normalize_already_qualified_no_rewrite() { + // Already qualified with instance name → no change + let result = normalize_delivery_target( + "signal:gvoice1:uuid:550e8400-e29b-41d4-a716-446655440000", + &Some("signal:gvoice1".to_string()), + ); + assert_eq!( + result, + "signal:gvoice1:uuid:550e8400-e29b-41d4-a716-446655440000" + ); + } + + #[test] + fn test_normalize_email_no_rewrite() { + // Email doesn't use named instances → no change + let result = + normalize_delivery_target("email:alice@example.com", &Some("email".to_string())); + assert_eq!(result, "email:alice@example.com"); + } + + // Regression tests: already-qualified non-Signal targets must not be double-prefixed + #[test] + fn test_normalize_slack_already_qualified_no_rewrite() { + let result = + normalize_delivery_target("slack:team1:userid:123", &Some("slack:team1".to_string())); + assert_eq!(result, "slack:team1:userid:123"); + } + + #[test] + fn test_normalize_discord_already_qualified_no_rewrite() { + let result = normalize_delivery_target( + "discord:guild1:userid:456", + &Some("discord:guild1".to_string()), + ); + assert_eq!(result, "discord:guild1:userid:456"); + } + + #[test] + fn test_normalize_telegram_already_qualified_no_rewrite() { + let result = normalize_delivery_target( + "telegram:bot1:userid:789", + &Some("telegram:bot1".to_string()), + ); + assert_eq!(result, "telegram:bot1:userid:789"); + } + + #[test] + fn test_parse_signal_already_qualified_not_default() { + // Already-qualified Signal targets (signal::...) should not be + // treated as default "signal" adapter. parse_delivery_target extracts + // the adapter correctly. + let parsed = crate::messaging::target::parse_delivery_target("signal:work:+15551234567"); + assert!(parsed.is_some()); + let target = parsed.unwrap(); + assert_eq!(target.adapter, "signal:work"); + assert_eq!(target.target, "+15551234567"); + + // This is different from bare signal: prefix which uses adapter "signal" + let parsed_default = crate::messaging::target::parse_delivery_target("signal:+15551234567"); + assert!(parsed_default.is_some()); + let target_default = parsed_default.unwrap(); + assert_eq!(target_default.adapter, "signal"); + assert_eq!(target_default.target, "+15551234567"); + } +} diff --git a/src/tools/send_message_to_another_channel.rs b/src/tools/send_message_to_another_channel.rs index b2232f637..669312139 100644 --- a/src/tools/send_message_to_another_channel.rs +++ b/src/tools/send_message_to_another_channel.rs @@ -88,12 +88,8 @@ impl Tool for SendMessageTool { async fn definition(&self, _prompt: String) -> ToolDefinition { let email_adapter_available = self.messaging_manager.has_adapter("email").await; - // Check if current adapter is Signal (e.g., "signal:gvoice1" starts with "signal") - let signal_adapter_available = self - .current_adapter - .as_ref() - .map(|adapter| adapter.starts_with("signal")) - .unwrap_or(false); + // Check if any Signal adapter is registered (works in any context, including cron) + let signal_adapter_available = self.messaging_manager.has_platform_adapters("signal").await; let mut description = crate::prompts::text::get("tools/send_message_to_another_channel").to_string(); @@ -109,12 +105,62 @@ impl Tool for SendMessageTool { } if signal_adapter_available { - description.push_str( - " Signal messaging is enabled: you can target `signal:uuid:{uuid}`, `signal:group:{group_id}`, or `signal:+{phone}`.", - ); - target_description.push_str( - " With Signal enabled, explicit targets are also allowed: `signal:uuid:{uuid}`, `signal:group:{group_id}`, `signal:+{phone}`", - ); + // Get actual Signal adapter names to determine if default or named-only + let adapter_names = self.messaging_manager.adapter_names().await; + let signal_adapters: Vec = adapter_names + .into_iter() + .filter(|name| name == "signal" || name.starts_with("signal:")) + .collect(); + let has_default_signal = signal_adapters.iter().any(|name| name == "signal"); + let named_adapters: Vec = signal_adapters + .into_iter() + .filter(|name| name.starts_with("signal:")) + .collect(); + + if has_default_signal { + // Default adapter exists - show generic syntax + description.push_str( + " Signal messaging is enabled: you can target `signal:uuid:{uuid}`, `signal:group:{group_id}`, or `signal:+{phone}`.", + ); + target_description.push_str( + " With Signal enabled, explicit targets are also allowed: `signal:uuid:{uuid}`, `signal:group:{group_id}`, `signal:+{phone}`", + ); + + // Also mention named instances if they exist + if !named_adapters.is_empty() { + let named_examples: Vec = named_adapters + .iter() + .take(2) + .map(|adapter| format!("`{}:+{{phone}}`", adapter)) + .collect(); + description.push_str(&format!( + " Named instances are also available: {}.", + named_examples.join(", ") + )); + target_description + .push_str(&format!(" Named instances: {}", named_examples.join(", "))); + } + } else { + // Only named adapters - show specific instance names + let instance_examples: Vec = named_adapters + .iter() + .take(3) + .map(|adapter| format!("`{}:+{{phone}}`", adapter)) + .collect(); + + if !instance_examples.is_empty() { + description.push_str(&format!( + " Signal messaging is enabled with named instances: target using `{}:{{instance_name}}:{{target}}` format (e.g., {}).", + "signal", + instance_examples.join(", ") + )); + target_description.push_str(&format!( + " With Signal enabled, use named instance format: `{}:{{instance_name}}:{{target}}` (e.g., {})", + "signal", + instance_examples.join(", ") + )); + } + } } ToolDefinition { @@ -147,16 +193,15 @@ impl Tool for SendMessageTool { // Check for explicit signal: prefix first - always honored regardless of current adapter. // This allows users to explicitly target Signal even when in Discord/Telegram/etc. if let Some(mut target) = parse_explicit_signal_prefix(&args.target) { - // If explicit prefix returned default "signal" adapter but we're in a named - // Signal adapter conversation (e.g., signal:gvoice1), use the current adapter - // to ensure the message goes through the correct account. - if target.adapter == "signal" - && let Some(current_adapter) = self - .current_adapter - .as_ref() - .filter(|adapter| adapter.starts_with("signal:")) - { - target.adapter = current_adapter.clone(); + // If explicit prefix returned default "signal" adapter, try to resolve + // to a specific named instance for correct routing. + if target.adapter == "signal" { + target.adapter = resolve_signal_adapter( + &self.messaging_manager, + self.current_adapter.as_deref(), + ) + .await + .map_err(SendMessageError)?; } self.messaging_manager @@ -187,38 +232,57 @@ impl Tool for SendMessageTool { if let Some(current_adapter) = self .current_adapter .as_ref() - .filter(|adapter| adapter.starts_with("signal")) + .filter(|adapter| *adapter == "signal" || adapter.starts_with("signal:")) { - match parse_implicit_signal_shorthand(&args.target, current_adapter) { - Ok(Some(target)) => { - self.messaging_manager - .broadcast( - &target.adapter, - &target.target, - crate::OutboundResponse::Text(args.message), - ) - .await - .map_err(|error| { - SendMessageError(format!("failed to send message: {error}")) - })?; - - tracing::info!( - adapter = %target.adapter, - broadcast_target = %"[REDACTED]", - "message sent via implicit Signal shorthand" - ); - - return Ok(SendMessageOutput { - success: true, - target: target.target, - platform: target.adapter, - }); - } - Err(validation_error) => { - return Err(SendMessageError(validation_error)); - } - Ok(None) => { - // Not a Signal shorthand — fall through to channel-name lookup. + // Verify the cached adapter is still registered before using it. + // The channel's current_adapter is stale if the named adapter was removed. + let live_signal_adapters: Vec = self + .messaging_manager + .adapter_names() + .await + .into_iter() + .filter(|name| name == "signal" || name.starts_with("signal:")) + .collect(); + + if !live_signal_adapters.contains(current_adapter) { + // Adapter was removed — fall through to channel-name lookup instead + // of routing to a dead adapter. + tracing::warn!( + adapter = %current_adapter, + "current_adapter references a removed Signal adapter; skipping implicit shorthand" + ); + } else { + match parse_implicit_signal_shorthand(&args.target, current_adapter) { + Ok(Some(target)) => { + self.messaging_manager + .broadcast( + &target.adapter, + &target.target, + crate::OutboundResponse::Text(args.message), + ) + .await + .map_err(|error| { + SendMessageError(format!("failed to send message: {error}")) + })?; + + tracing::info!( + adapter = %target.adapter, + broadcast_target = %"[REDACTED]", + "message sent via implicit Signal shorthand" + ); + + return Ok(SendMessageOutput { + success: true, + target: target.target, + platform: target.adapter, + }); + } + Err(validation_error) => { + return Err(SendMessageError(validation_error)); + } + Ok(None) => { + // Not a Signal shorthand — fall through to channel-name lookup. + } } } } @@ -337,6 +401,54 @@ fn parse_explicit_signal_prefix(raw: &str) -> Option, +) -> Result { + let all_signal_adapters: Vec = messaging_manager + .adapter_names() + .await + .into_iter() + .filter(|name| name == "signal" || name.starts_with("signal:")) + .collect(); + + let has_default_signal = all_signal_adapters.iter().any(|name| name == "signal"); + let named_adapters: Vec<&str> = all_signal_adapters + .iter() + .filter(|name| name.starts_with("signal:")) + .map(|s| s.as_str()) + .collect(); + + if let Some(adapter) = current_adapter.filter(|adapter| { + adapter.starts_with("signal:") && all_signal_adapters.iter().any(|a| a == adapter) + }) { + if !has_default_signal { + return Ok(adapter.to_string()); + } + // has_default_signal is true — fall through to return "signal" + } else if !has_default_signal && named_adapters.len() == 1 { + // No default, but exactly one named adapter - use it + return Ok(named_adapters[0].to_string()); + } else if !has_default_signal && named_adapters.len() > 1 { + // Multiple named adapters and no default - ambiguity error + return Err(format!( + "Multiple Signal adapters are configured ({}). Please specify which instance to use by targeting 'signal::' instead of 'signal:'.", + named_adapters.join(", ") + )); + } + // has_default_signal is true OR no adapters at all — return default "signal" + Ok("signal".to_string()) +} + /// Parse implicit Signal shorthands - only in Signal conversations. /// Handles bare UUIDs, group:xxx, and +phone without explicit signal: prefix. /// Returns `Ok(Some(...))` on valid shorthand, `Ok(None)` when the input is @@ -369,28 +481,44 @@ fn parse_implicit_signal_shorthand( )); } - // Phone number format: starts with + followed by 7+ digits - if trimmed.starts_with('+') - && trimmed[1..].len() >= 7 - && trimmed[1..].chars().all(|c| c.is_ascii_digit()) - { - return Ok(Some(BroadcastTarget { - adapter: current_adapter.to_string(), - target: trimmed.to_string(), - })); - } - - // Starts with + but doesn't meet phone number requirements. - if let Some(digits) = trimmed.strip_prefix('+') { - if digits.chars().all(|c| c.is_ascii_digit()) { - return Err(format!( - "'{trimmed}' looks like a phone number but is too short. Phone numbers need at least 7 digits after the + prefix." - )); + // Phone number format: use strict E.164 validation + if trimmed.starts_with('+') { + if crate::messaging::target::is_valid_e164(trimmed) { + return Ok(Some(BroadcastTarget { + adapter: current_adapter.to_string(), + target: trimmed.to_string(), + })); } - if !digits.is_empty() { - return Err(format!( - "'{trimmed}' looks like a phone number but contains non-digit characters. Use format: +1234567890" - )); + // Invalid phone number - provide specific error + if let Some(digits) = trimmed.strip_prefix('+') { + if digits.is_empty() { + return Err( + "Phone number cannot be empty after + prefix. Use format: +1234567890" + .to_string(), + ); + } + if !digits.chars().all(|c| c.is_ascii_digit()) { + return Err(format!( + "'{trimmed}' contains non-digit characters. Use format: +1234567890" + )); + } + if digits.len() < 6 { + return Err(format!( + "'{trimmed}' is too short. Phone numbers need 6-15 digits after the + prefix (7-16 total)." + )); + } + if digits.len() > 15 { + return Err(format!( + "'{trimmed}' is too long. Phone numbers need 6-15 digits after the + prefix (7-16 total)." + )); + } + if digits.starts_with('0') { + return Err(format!( + "'{trimmed}' has invalid country code. Country codes cannot start with 0." + )); + } + // Catch-all for any other '+' prefixed input that failed validation + return Err(format!("'{trimmed}' is not a valid E.164 phone number")); } } @@ -620,4 +748,28 @@ mod tests { .expect_err("should be validation error"); assert!(error.contains("requires an ID"), "{error}"); } + + // Tests for resolve_signal_adapter + // Note: These tests use a mock messaging manager to test the resolution logic + + #[tokio::test] + async fn resolve_signal_adapter_returns_signal_when_no_adapters() { + let manager = crate::messaging::MessagingManager::new(); + // No adapters registered - should return "signal" as default + // This lets broadcast() fail with appropriate error if no adapters exist + let result = super::resolve_signal_adapter(&manager, None).await; + assert_eq!(result.unwrap(), "signal"); + } + + #[tokio::test] + async fn resolve_signal_adapter_ignores_current_if_not_registered() { + let manager = crate::messaging::MessagingManager::new(); + // When current_adapter is provided but not registered in manager, + // it falls through to "signal" since we can't verify the adapter exists + let result = super::resolve_signal_adapter(&manager, Some("signal:work")).await; + assert_eq!(result.unwrap(), "signal"); + } + + // TODO: Add test for resolve_signal_adapter ambiguity case once MessagingManager + // supports registering mock adapters or a test helper is available. } diff --git a/tests/context_dump.rs b/tests/context_dump.rs index e84c2f41d..b487bb25e 100644 --- a/tests/context_dump.rs +++ b/tests/context_dump.rs @@ -190,7 +190,7 @@ fn build_channel_system_prompt(rc: &spacebot::config::RuntimeConfig) -> String { .expect("failed to render worker capabilities"); let conversation_context = prompt_engine - .render_conversation_context("discord", Some("Test Server"), Some("#general")) + .render_conversation_context("discord", Some("Test Server"), Some("#general"), None) .ok(); let empty_to_none = |s: String| if s.is_empty() { None } else { Some(s) };