Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
64fe9ba
fix: prevent UTF-8 panics in byte-index string truncation (#1688)
zmanian Mar 29, 2026
70214c4
fix(bedrock): strip tool blocks from messages when toolConfig is abse…
rajulbhatnagar Mar 29, 2026
de384b0
feat: support custom LLM provider configuration via web UI (#1340)
italic-jinxin Mar 29, 2026
4f277c9
Handle empty tool completions in autonomous jobs (#1720)
henrypark133 Mar 29, 2026
368d2f5
feat(gateway): OIDC JWT authentication for reverse-proxy deployments …
synner88 Mar 29, 2026
8acdd08
fix(wasm): inject Content-Length: 0 for bodyless mutating HTTP reques…
werdan Mar 30, 2026
c75dea0
fix(auth): make shared Google tool status scope-aware (#1532)
G7CNF Mar 30, 2026
10d5a53
fix: resolve 11 test failures from multi-tenant bootstrap and sandbox…
ilblackdragon Mar 30, 2026
d0f7862
fix(slack): respond to thread replies without requiring @mention (#1405)
synner88 Mar 30, 2026
d567d94
fix(routines): clone Arc before await in web handler event cache refr…
ilblackdragon Mar 30, 2026
21f613f
test(e2e): align WASM reinstall expectation with uninstall cleanup (#…
henrypark133 Mar 30, 2026
8be61c2
Merge pull request #1763 from nearai/staging-promote/21f613ff-2376117…
henrypark133 Mar 30, 2026
ac79803
Merge pull request #1761 from nearai/staging-promote/d567d94c-2375525…
henrypark133 Mar 30, 2026
30a3e59
Merge pull request #1755 from nearai/staging-promote/d0f7862a-2373313…
henrypark133 Mar 30, 2026
e6c568a
Merge pull request #1753 from nearai/staging-promote/c75dea0e-2373121…
henrypark133 Mar 30, 2026
527b88e
Merge pull request #1747 from nearai/staging-promote/8acdd080-2372603…
henrypark133 Mar 30, 2026
afc8342
Merge pull request #1744 from nearai/staging-promote/368d2f52-2372018…
henrypark133 Mar 30, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
353 changes: 210 additions & 143 deletions Cargo.lock

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,7 @@ hyper-util = { version = "0.1", features = ["server", "tokio", "http1", "http2"]
http-body-util = "0.1"
bytes = "1"
base64 = "0.22.1"
jsonwebtoken = "9"
mime_guess = "2.0.5"
clap_complete = "4.5.0"
lru = "0.16.3"
Expand Down
2 changes: 1 addition & 1 deletion FEATURE_PARITY.md
Original file line number Diff line number Diff line change
Expand Up @@ -112,7 +112,7 @@ This document tracks feature parity between IronClaw (Rust implementation) and O
|---------|----------|----------|-------|
| Streaming draft replies | ✅ | ❌ | Partial replies via draft message updates |
| Configurable stream modes | ✅ | ❌ | Per-channel stream behavior |
| Thread ownership | ✅ | ❌ | Thread-level ownership tracking plus reply participation memory |
| Thread ownership | ✅ | 🚧 | Reply participation memory now persists with TTL-bounded tracking; full thread-level ownership tracking is still missing |
| Download-file action | ✅ | ❌ | On-demand attachment downloads via message actions |

### Mattermost-Specific Features (since Mar 2026)
Expand Down
202 changes: 184 additions & 18 deletions channels-src/slack/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ wit_bindgen::generate!({
});

use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;

// Re-export generated types
use exports::near::agent::channel::{
Expand Down Expand Up @@ -129,9 +130,17 @@ const OWNER_ID_PATH: &str = "state/owner_id";
const DM_POLICY_PATH: &str = "state/dm_policy";
/// Workspace path for persisting allow_from (JSON array) across WASM callbacks.
const ALLOW_FROM_PATH: &str = "state/allow_from";
/// Workspace path for tracking recently active Slack threads.
const ACTIVE_THREADS_PATH: &str = "state/active_threads.json";
/// Recently active threads expire after 24 hours to avoid reviving stale threads forever.
const ACTIVE_THREAD_TTL_MS: u64 = 24 * 60 * 60 * 1000;
/// Cap stored thread markers so the workspace state stays bounded.
const ACTIVE_THREAD_MAX_ENTRIES: usize = 256;
/// Channel name for pairing store (used by pairing host APIs).
const CHANNEL_NAME: &str = "slack";

type ActiveThreads = BTreeMap<String, u64>;

/// Channel configuration from capabilities file.
#[derive(Debug, Deserialize)]
struct SlackConfig {
Expand Down Expand Up @@ -263,9 +272,9 @@ impl Guest for SlackChannel {
"text": response.content,
});

// Add thread_ts for threaded replies
if let Some(thread_ts) = response.thread_id.or(metadata.thread_ts) {
payload["thread_ts"] = serde_json::Value::String(thread_ts);
let thread_ts = response.thread_id.or(metadata.thread_ts);
if let Some(ref thread_ts) = thread_ts {
payload["thread_ts"] = serde_json::Value::String(thread_ts.clone());
}

let payload_bytes = serde_json::to_vec(&payload)
Expand Down Expand Up @@ -308,6 +317,10 @@ impl Guest for SlackChannel {
));
}

if let Some(thread_ts) = thread_ts {
track_active_thread(&metadata.channel, &thread_ts)?;
}

channel_host::log(
channel_host::LogLevel::Debug,
&format!(
Expand Down Expand Up @@ -452,13 +465,14 @@ fn download_and_store_slack_files(attachments: &[InboundAttachment]) {
}
}

/// Handle a Slack event and emit message if applicable.
fn handle_slack_event(event: SlackEvent, team_id: Option<String>, _event_id: Option<String>) {
let attachments = extract_slack_attachments(&event.files);

// Download and store file attachments for host-side processing
fn prepare_inbound_attachments(files: &Option<Vec<SlackFile>>) -> Vec<InboundAttachment> {
let attachments = extract_slack_attachments(files);
download_and_store_slack_files(&attachments);
attachments
}

/// Handle a Slack event and emit message if applicable.
fn handle_slack_event(event: SlackEvent, team_id: Option<String>, _event_id: Option<String>) {
match event.event_type.as_str() {
// Direct mention of the bot (always in a channel, not a DM)
"app_mention" => {
Expand All @@ -472,6 +486,7 @@ fn handle_slack_event(event: SlackEvent, team_id: Option<String>, _event_id: Opt
if !check_sender_permission(&user, &channel, false) {
return;
}
let attachments = prepare_inbound_attachments(&event.files);
emit_message(
user,
text,
Expand All @@ -483,7 +498,7 @@ fn handle_slack_event(event: SlackEvent, team_id: Option<String>, _event_id: Opt
}
}

// Direct message to the bot
// Direct message or thread follow-up to the bot
"message" => {
// Skip messages from bots (including ourselves)
if event.bot_id.is_some() || event.subtype.is_some() {
Expand All @@ -496,11 +511,20 @@ fn handle_slack_event(event: SlackEvent, team_id: Option<String>, _event_id: Opt
event.text,
event.ts.clone(),
) {
// Only process DMs (channel IDs starting with D)
if channel.starts_with('D') {
if !check_sender_permission(&user, &channel, true) {
let is_dm = channel.starts_with('D');

// Check if this is a reply in a thread where we previously participated
let is_active_thread = !is_dm
&& event
.thread_ts
.as_ref()
.is_some_and(|thread_ts| is_active_thread(&channel, thread_ts));

if is_dm || is_active_thread {
if !check_sender_permission(&user, &channel, is_dm) {
return;
}
let attachments = prepare_inbound_attachments(&event.files);
emit_message(
user,
text,
Expand All @@ -522,6 +546,93 @@ fn handle_slack_event(event: SlackEvent, team_id: Option<String>, _event_id: Opt
}
}

fn active_thread_key(channel: &str, thread_ts: &str) -> String {
format!("{channel}/{thread_ts}")
}

fn is_thread_marker_fresh(last_seen_millis: u64, now_millis: u64) -> bool {
now_millis.saturating_sub(last_seen_millis) <= ACTIVE_THREAD_TTL_MS
}

fn prune_active_threads(active_threads: &mut ActiveThreads, now_millis: u64) -> bool {
let mut changed = false;
active_threads.retain(|_, last_seen_millis| {
let keep = is_thread_marker_fresh(*last_seen_millis, now_millis);
if !keep {
changed = true;
}
keep
});

if active_threads.len() > ACTIVE_THREAD_MAX_ENTRIES {
let mut oldest_first: Vec<_> = active_threads
.iter()
.map(|(key, last_seen_millis)| (key.clone(), *last_seen_millis))
.collect();
oldest_first.sort_by_key(|(_, last_seen_millis)| *last_seen_millis);

for (key, _) in oldest_first
.into_iter()
.take(active_threads.len() - ACTIVE_THREAD_MAX_ENTRIES)
{
active_threads.remove(&key);
changed = true;
}
}

changed
}

fn load_active_threads() -> ActiveThreads {
let Some(raw) = channel_host::workspace_read(ACTIVE_THREADS_PATH) else {
return ActiveThreads::new();
};

match serde_json::from_str(&raw) {
Ok(active_threads) => active_threads,
Err(e) => {
channel_host::log(
channel_host::LogLevel::Warn,
&format!("Failed to parse active thread state: {e}"),
);
ActiveThreads::new()
}
}
}

fn persist_active_threads(active_threads: &ActiveThreads) -> Result<(), String> {
let serialized = serde_json::to_string(active_threads)
.map_err(|e| format!("Failed to serialize active thread state: {e}"))?;
channel_host::workspace_write(ACTIVE_THREADS_PATH, &serialized)
.map_err(|e| format!("Failed to persist active thread state: {e}"))
}

fn track_active_thread(channel: &str, thread_ts: &str) -> Result<(), String> {
let now_millis = channel_host::now_millis();
let mut active_threads = load_active_threads();
prune_active_threads(&mut active_threads, now_millis);
active_threads.insert(active_thread_key(channel, thread_ts), now_millis);
prune_active_threads(&mut active_threads, now_millis);
persist_active_threads(&active_threads)
}

fn is_active_thread(channel: &str, thread_ts: &str) -> bool {
let now_millis = channel_host::now_millis();
let mut active_threads = load_active_threads();
let changed = prune_active_threads(&mut active_threads, now_millis);

if changed {
if let Err(e) = persist_active_threads(&active_threads) {
channel_host::log(
channel_host::LogLevel::Warn,
&format!("Failed to prune active thread state: {e}"),
);
}
}

active_threads.contains_key(&active_thread_key(channel, thread_ts))
}

/// Emit a message to the agent.
fn emit_message(
user_id: String,
Expand Down Expand Up @@ -606,8 +717,7 @@ fn check_sender_permission(user_id: &str, channel_id: &str, is_dm: bool) -> bool
}

// 4. Check sender (Slack events only have user ID, not username)
let is_allowed =
allowed.contains(&"*".to_string()) || allowed.contains(&user_id.to_string());
let is_allowed = allowed.contains(&"*".to_string()) || allowed.contains(&user_id.to_string());

if is_allowed {
return true;
Expand All @@ -625,10 +735,7 @@ fn check_sender_permission(user_id: &str, channel_id: &str, is_dm: bool) -> bool
Ok(result) => {
channel_host::log(
channel_host::LogLevel::Info,
&format!(
"Pairing request for user {}: code {}",
user_id, result.code
),
&format!("Pairing request for user {}: code {}", user_id, result.code),
);
if result.created {
let _ = send_pairing_reply(channel_id, &result.code);
Expand Down Expand Up @@ -826,4 +933,63 @@ mod tests {
// Verify the constant is 20 MB
assert_eq!(MAX_DOWNLOAD_SIZE_BYTES, 20 * 1024 * 1024);
}

#[test]
fn test_active_thread_key_scopes_by_channel_and_thread() {
assert_eq!(
active_thread_key("C123", "1742486400.000100"),
"C123/1742486400.000100"
);
}

#[test]
fn test_prune_active_threads_removes_expired_entries() {
let now_millis = ACTIVE_THREAD_TTL_MS + 1_000;
let mut active_threads = ActiveThreads::from([
(
"C1/expired".to_string(),
now_millis - ACTIVE_THREAD_TTL_MS - 1,
),
("C1/fresh".to_string(), now_millis - ACTIVE_THREAD_TTL_MS),
]);

let changed = prune_active_threads(&mut active_threads, now_millis);

assert!(changed);
assert!(!active_threads.contains_key("C1/expired"));
assert!(active_threads.contains_key("C1/fresh"));
}

#[test]
fn test_prune_active_threads_trims_oldest_entries_when_over_limit() {
let now_millis = ACTIVE_THREAD_TTL_MS + 1_000;
let mut active_threads = ActiveThreads::new();

for i in 0..=ACTIVE_THREAD_MAX_ENTRIES {
active_threads.insert(format!("C1/{i}"), now_millis + i as u64);
}

let changed = prune_active_threads(
&mut active_threads,
now_millis + ACTIVE_THREAD_MAX_ENTRIES as u64,
);

assert!(changed);
assert_eq!(active_threads.len(), ACTIVE_THREAD_MAX_ENTRIES);
assert!(!active_threads.contains_key("C1/0"));
assert!(active_threads.contains_key(&format!("C1/{ACTIVE_THREAD_MAX_ENTRIES}")));
}

#[test]
fn test_is_thread_marker_fresh_respects_ttl_boundary() {
let now_millis = ACTIVE_THREAD_TTL_MS + 1_000;
assert!(is_thread_marker_fresh(
now_millis - ACTIVE_THREAD_TTL_MS,
now_millis
));
assert!(!is_thread_marker_fresh(
now_millis - ACTIVE_THREAD_TTL_MS - 1,
now_millis
));
}
}
Loading
Loading