mirror of
https://github.com/outbackdingo/optimclaw.git
synced 2026-08-25 14:53:34 +00:00
Fix REPL single-message hang and cap CI test duration (#1643)
* Fix REPL single-message hang and cap CI test duration * Fix Clippy nested-if lint in REPL startup * Fix single-message approval flow * Handle empty single-message REPL exits * Wait for one-shot event routines before exit
This commit is contained in:
@@ -12,6 +12,7 @@ jobs:
|
|||||||
tests:
|
tests:
|
||||||
name: Tests (${{ matrix.name }})
|
name: Tests (${{ matrix.name }})
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
|
timeout-minutes: 45
|
||||||
strategy:
|
strategy:
|
||||||
fail-fast: false
|
fail-fast: false
|
||||||
matrix:
|
matrix:
|
||||||
@@ -40,11 +41,14 @@ jobs:
|
|||||||
- name: Build WASM channels (for integration tests)
|
- name: Build WASM channels (for integration tests)
|
||||||
run: ./scripts/build-wasm-extensions.sh --channels
|
run: ./scripts/build-wasm-extensions.sh --channels
|
||||||
- name: Run Tests
|
- name: Run Tests
|
||||||
run: cargo test ${{ matrix.flags }} -- --nocapture
|
run: |
|
||||||
|
timeout --signal=INT --kill-after=30s 40m \
|
||||||
|
cargo test ${{ matrix.flags }} -- --nocapture
|
||||||
|
|
||||||
heavy-integration-tests:
|
heavy-integration-tests:
|
||||||
name: Heavy Integration Tests
|
name: Heavy Integration Tests
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
|
timeout-minutes: 20
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout repository
|
- name: Checkout repository
|
||||||
uses: actions/checkout@v6
|
uses: actions/checkout@v6
|
||||||
@@ -58,9 +62,13 @@ jobs:
|
|||||||
- name: Build Telegram WASM channel
|
- name: Build Telegram WASM channel
|
||||||
run: cargo build --manifest-path channels-src/telegram/Cargo.toml --target wasm32-wasip2 --release
|
run: cargo build --manifest-path channels-src/telegram/Cargo.toml --target wasm32-wasip2 --release
|
||||||
- name: Run thread scheduling integration tests
|
- name: Run thread scheduling integration tests
|
||||||
run: cargo test --no-default-features --features libsql,integration --test e2e_thread_scheduling -- --nocapture
|
run: |
|
||||||
|
timeout --signal=INT --kill-after=30s 15m \
|
||||||
|
cargo test --no-default-features --features libsql,integration --test e2e_thread_scheduling -- --nocapture
|
||||||
- name: Run Telegram thread-scope regression test
|
- name: Run Telegram thread-scope regression test
|
||||||
run: cargo test --features integration --test telegram_auth_integration test_private_messages_use_chat_id_as_thread_scope -- --exact
|
run: |
|
||||||
|
timeout --signal=INT --kill-after=30s 10m \
|
||||||
|
cargo test --features integration --test telegram_auth_integration test_private_messages_use_chat_id_as_thread_scope -- --exact
|
||||||
|
|
||||||
telegram-tests:
|
telegram-tests:
|
||||||
name: Telegram Channel Tests
|
name: Telegram Channel Tests
|
||||||
@@ -68,6 +76,7 @@ jobs:
|
|||||||
github.event_name != 'pull_request' ||
|
github.event_name != 'pull_request' ||
|
||||||
github.base_ref != 'staging'
|
github.base_ref != 'staging'
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
|
timeout-minutes: 15
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout repository
|
- name: Checkout repository
|
||||||
uses: actions/checkout@v6
|
uses: actions/checkout@v6
|
||||||
@@ -75,7 +84,9 @@ jobs:
|
|||||||
uses: dtolnay/rust-toolchain@stable
|
uses: dtolnay/rust-toolchain@stable
|
||||||
- uses: Swatinem/rust-cache@v2
|
- uses: Swatinem/rust-cache@v2
|
||||||
- name: Run Telegram Channel Tests
|
- name: Run Telegram Channel Tests
|
||||||
run: cargo test --manifest-path channels-src/telegram/Cargo.toml -- --nocapture
|
run: |
|
||||||
|
timeout --signal=INT --kill-after=30s 10m \
|
||||||
|
cargo test --manifest-path channels-src/telegram/Cargo.toml -- --nocapture
|
||||||
|
|
||||||
windows-build:
|
windows-build:
|
||||||
name: Windows Build (${{ matrix.name }})
|
name: Windows Build (${{ matrix.name }})
|
||||||
@@ -110,6 +121,7 @@ jobs:
|
|||||||
github.event_name != 'pull_request' ||
|
github.event_name != 'pull_request' ||
|
||||||
github.base_ref != 'staging'
|
github.base_ref != 'staging'
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
|
timeout-minutes: 30
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout repository
|
- name: Checkout repository
|
||||||
uses: actions/checkout@v6
|
uses: actions/checkout@v6
|
||||||
@@ -125,7 +137,9 @@ jobs:
|
|||||||
- name: Build all WASM extensions against current WIT
|
- name: Build all WASM extensions against current WIT
|
||||||
run: ./scripts/build-wasm-extensions.sh
|
run: ./scripts/build-wasm-extensions.sh
|
||||||
- name: Instantiation test (host linker compatibility)
|
- name: Instantiation test (host linker compatibility)
|
||||||
run: cargo test --all-features wit_compat -- --nocapture
|
run: |
|
||||||
|
timeout --signal=INT --kill-after=30s 20m \
|
||||||
|
cargo test --all-features wit_compat -- --nocapture
|
||||||
|
|
||||||
bench-compile:
|
bench-compile:
|
||||||
name: Benchmark Compilation
|
name: Benchmark Compilation
|
||||||
|
|||||||
+64
-5
@@ -16,6 +16,7 @@ use crate::agent::context_monitor::ContextMonitor;
|
|||||||
use crate::agent::heartbeat::spawn_heartbeat;
|
use crate::agent::heartbeat::spawn_heartbeat;
|
||||||
use crate::agent::routine_engine::{RoutineEngine, spawn_cron_ticker};
|
use crate::agent::routine_engine::{RoutineEngine, spawn_cron_ticker};
|
||||||
use crate::agent::self_repair::{DefaultSelfRepair, RepairResult, SelfRepair};
|
use crate::agent::self_repair::{DefaultSelfRepair, RepairResult, SelfRepair};
|
||||||
|
use crate::agent::session::ThreadState;
|
||||||
use crate::agent::session_manager::SessionManager;
|
use crate::agent::session_manager::SessionManager;
|
||||||
use crate::agent::submission::{Submission, SubmissionParser, SubmissionResult};
|
use crate::agent::submission::{Submission, SubmissionParser, SubmissionResult};
|
||||||
use crate::agent::{HeartbeatConfig as AgentHeartbeatConfig, Router, Scheduler, SchedulerDeps};
|
use crate::agent::{HeartbeatConfig as AgentHeartbeatConfig, Router, Scheduler, SchedulerDeps};
|
||||||
@@ -84,6 +85,15 @@ fn resolve_owner_scope_notification_user(
|
|||||||
trimmed_option(explicit_user).or_else(|| trimmed_option(owner_fallback))
|
trimmed_option(explicit_user).or_else(|| trimmed_option(owner_fallback))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn is_single_message_repl(message: &IncomingMessage) -> bool {
|
||||||
|
message.channel == "repl"
|
||||||
|
&& message
|
||||||
|
.metadata
|
||||||
|
.get("single_message_mode")
|
||||||
|
.and_then(|value| value.as_bool())
|
||||||
|
.unwrap_or(false)
|
||||||
|
}
|
||||||
|
|
||||||
async fn resolve_channel_notification_user(
|
async fn resolve_channel_notification_user(
|
||||||
extension_manager: Option<&Arc<ExtensionManager>>,
|
extension_manager: Option<&Arc<ExtensionManager>>,
|
||||||
channel: Option<&str>,
|
channel: Option<&str>,
|
||||||
@@ -1140,9 +1150,14 @@ impl Agent {
|
|||||||
&& let Submission::UserInput { ref content } = submission
|
&& let Submission::UserInput { ref content } = submission
|
||||||
&& let Some(engine) = self.routine_engine().await
|
&& let Some(engine) = self.routine_engine().await
|
||||||
{
|
{
|
||||||
|
let single_message_repl = is_single_message_repl(message);
|
||||||
// Use post-hook content so that BeforeInbound hooks that rewrite
|
// Use post-hook content so that BeforeInbound hooks that rewrite
|
||||||
// input are respected by event trigger matching.
|
// input are respected by event trigger matching.
|
||||||
let fired = engine.check_event_triggers(message, content).await;
|
let fired = if single_message_repl {
|
||||||
|
engine.check_event_triggers_and_wait(message, content).await
|
||||||
|
} else {
|
||||||
|
engine.check_event_triggers(message, content).await
|
||||||
|
};
|
||||||
if fired > 0 {
|
if fired > 0 {
|
||||||
tracing::debug!(
|
tracing::debug!(
|
||||||
channel = %message.channel,
|
channel = %message.channel,
|
||||||
@@ -1150,10 +1165,16 @@ impl Agent {
|
|||||||
fired,
|
fired,
|
||||||
"Consumed inbound user message with matching event-triggered routine(s)"
|
"Consumed inbound user message with matching event-triggered routine(s)"
|
||||||
);
|
);
|
||||||
return Ok(Some(String::new()));
|
return if single_message_repl {
|
||||||
|
Ok(None)
|
||||||
|
} else {
|
||||||
|
Ok(Some(String::new()))
|
||||||
|
};
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let session_for_empty_exit = Arc::clone(&session);
|
||||||
|
|
||||||
// Process based on submission type
|
// Process based on submission type
|
||||||
let result = match submission {
|
let result = match submission {
|
||||||
Submission::UserInput { content } => {
|
Submission::UserInput { content } => {
|
||||||
@@ -1263,7 +1284,13 @@ impl Agent {
|
|||||||
SubmissionResult::Error { message } => {
|
SubmissionResult::Error { message } => {
|
||||||
Ok(Some(format!("Error: {}", message)))
|
Ok(Some(format!("Error: {}", message)))
|
||||||
}
|
}
|
||||||
_ => Ok(Some(String::new())),
|
_ => {
|
||||||
|
if is_single_message_repl(message) {
|
||||||
|
Ok(None)
|
||||||
|
} else {
|
||||||
|
Ok(Some(String::new()))
|
||||||
|
}
|
||||||
|
}
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
// Authorization checks (including restart channel check) are enforced in handle_system_command
|
// Authorization checks (including restart channel check) are enforced in handle_system_command
|
||||||
@@ -1325,7 +1352,26 @@ impl Agent {
|
|||||||
Ok(Some(content))
|
Ok(Some(content))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
SubmissionResult::Ok { message } => Ok(message),
|
SubmissionResult::Ok {
|
||||||
|
message: output_message,
|
||||||
|
} => {
|
||||||
|
let should_exit =
|
||||||
|
if output_message.as_deref() == Some("") && is_single_message_repl(message) {
|
||||||
|
let sess = session_for_empty_exit.lock().await;
|
||||||
|
sess.threads
|
||||||
|
.get(&thread_id)
|
||||||
|
.map(|thread| thread.state != ThreadState::AwaitingApproval)
|
||||||
|
.unwrap_or(true)
|
||||||
|
} else {
|
||||||
|
false
|
||||||
|
};
|
||||||
|
|
||||||
|
if should_exit {
|
||||||
|
Ok(None)
|
||||||
|
} else {
|
||||||
|
Ok(output_message)
|
||||||
|
}
|
||||||
|
}
|
||||||
SubmissionResult::Error { message } => Ok(Some(format!("Error: {}", message))),
|
SubmissionResult::Error { message } => Ok(Some(format!("Error: {}", message))),
|
||||||
SubmissionResult::Interrupted => Ok(Some("Interrupted.".into())),
|
SubmissionResult::Interrupted => Ok(Some("Interrupted.".into())),
|
||||||
SubmissionResult::NeedApproval { .. } => {
|
SubmissionResult::NeedApproval { .. } => {
|
||||||
@@ -1341,7 +1387,7 @@ impl Agent {
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{
|
use super::{
|
||||||
chat_tool_execution_metadata, resolve_routine_notification_user,
|
chat_tool_execution_metadata, is_single_message_repl, resolve_routine_notification_user,
|
||||||
should_fallback_routine_notification, truncate_for_preview,
|
should_fallback_routine_notification, truncate_for_preview,
|
||||||
};
|
};
|
||||||
use crate::channels::IncomingMessage;
|
use crate::channels::IncomingMessage;
|
||||||
@@ -1503,4 +1549,17 @@ mod tests {
|
|||||||
|
|
||||||
assert!(should_fallback_routine_notification(&error)); // safety: test-only assertion
|
assert!(should_fallback_routine_notification(&error)); // safety: test-only assertion
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn single_message_repl_detection_requires_repl_channel_and_metadata_flag() {
|
||||||
|
let repl = IncomingMessage::new("repl", "owner-scope", "hello")
|
||||||
|
.with_metadata(serde_json::json!({ "single_message_mode": true }));
|
||||||
|
let gateway = IncomingMessage::new("gateway", "owner-scope", "hello")
|
||||||
|
.with_metadata(serde_json::json!({ "single_message_mode": true }));
|
||||||
|
let plain_repl = IncomingMessage::new("repl", "owner-scope", "hello");
|
||||||
|
|
||||||
|
assert!(is_single_message_repl(&repl)); // safety: test-only assertion
|
||||||
|
assert!(!is_single_message_repl(&gateway)); // safety: test-only assertion
|
||||||
|
assert!(!is_single_message_repl(&plain_repl)); // safety: test-only assertion
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+60
-10
@@ -18,6 +18,7 @@ use std::time::Duration;
|
|||||||
use chrono::Utc;
|
use chrono::Utc;
|
||||||
use regex::Regex;
|
use regex::Regex;
|
||||||
use tokio::sync::{RwLock, mpsc};
|
use tokio::sync::{RwLock, mpsc};
|
||||||
|
use tokio::task::JoinHandle;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
use crate::agent::Scheduler;
|
use crate::agent::Scheduler;
|
||||||
@@ -45,6 +46,11 @@ enum EventMatcher {
|
|||||||
System { routine: Routine },
|
System { routine: Routine },
|
||||||
}
|
}
|
||||||
|
|
||||||
|
struct TriggeredRoutine {
|
||||||
|
routine: Routine,
|
||||||
|
detail: String,
|
||||||
|
}
|
||||||
|
|
||||||
/// Distinguishes why sandbox is unavailable so error messages are accurate.
|
/// Distinguishes why sandbox is unavailable so error messages are accurate.
|
||||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
pub enum SandboxReadiness {
|
pub enum SandboxReadiness {
|
||||||
@@ -202,6 +208,44 @@ impl RoutineEngine {
|
|||||||
|
|
||||||
/// Check incoming message against event triggers. Returns number of routines fired.
|
/// Check incoming message against event triggers. Returns number of routines fired.
|
||||||
pub async fn check_event_triggers(&self, message: &IncomingMessage, content: &str) -> usize {
|
pub async fn check_event_triggers(&self, message: &IncomingMessage, content: &str) -> usize {
|
||||||
|
let triggered = self.matching_event_triggers(message, content).await;
|
||||||
|
let fired = triggered.len();
|
||||||
|
for triggered in triggered {
|
||||||
|
std::mem::drop(self.spawn_fire(triggered.routine, "event", Some(triggered.detail)));
|
||||||
|
}
|
||||||
|
fired
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Fire matching event-triggered routines and wait for them to complete.
|
||||||
|
///
|
||||||
|
/// Used by single-message REPL mode so the process does not exit before
|
||||||
|
/// background event-triggered routines finish.
|
||||||
|
pub async fn check_event_triggers_and_wait(
|
||||||
|
&self,
|
||||||
|
message: &IncomingMessage,
|
||||||
|
content: &str,
|
||||||
|
) -> usize {
|
||||||
|
let triggered = self.matching_event_triggers(message, content).await;
|
||||||
|
let fired = triggered.len();
|
||||||
|
let handles: Vec<JoinHandle<()>> = triggered
|
||||||
|
.into_iter()
|
||||||
|
.map(|triggered| self.spawn_fire(triggered.routine, "event", Some(triggered.detail)))
|
||||||
|
.collect();
|
||||||
|
|
||||||
|
for handle in handles {
|
||||||
|
if let Err(e) = handle.await {
|
||||||
|
tracing::warn!(error = %e, "Event-triggered routine task failed");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fired
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn matching_event_triggers(
|
||||||
|
&self,
|
||||||
|
message: &IncomingMessage,
|
||||||
|
content: &str,
|
||||||
|
) -> Vec<TriggeredRoutine> {
|
||||||
let cache = self.event_cache.read().await;
|
let cache = self.event_cache.read().await;
|
||||||
|
|
||||||
// Early return if there are no message matchers at all.
|
// Early return if there are no message matchers at all.
|
||||||
@@ -209,10 +253,9 @@ impl RoutineEngine {
|
|||||||
.iter()
|
.iter()
|
||||||
.any(|m| matches!(m, EventMatcher::Message { .. }))
|
.any(|m| matches!(m, EventMatcher::Message { .. }))
|
||||||
{
|
{
|
||||||
return 0;
|
return Vec::new();
|
||||||
}
|
}
|
||||||
|
let mut triggered = Vec::new();
|
||||||
let mut fired = 0;
|
|
||||||
|
|
||||||
// Collect routine IDs for batch query
|
// Collect routine IDs for batch query
|
||||||
let routine_ids: Vec<Uuid> = cache
|
let routine_ids: Vec<Uuid> = cache
|
||||||
@@ -224,13 +267,13 @@ impl RoutineEngine {
|
|||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
if routine_ids.is_empty() {
|
if routine_ids.is_empty() {
|
||||||
return 0;
|
return Vec::new();
|
||||||
}
|
}
|
||||||
|
|
||||||
// Single batch query instead of N queries
|
// Single batch query instead of N queries
|
||||||
let concurrent_counts = match self.batch_concurrent_counts(&routine_ids).await {
|
let concurrent_counts = match self.batch_concurrent_counts(&routine_ids).await {
|
||||||
Some(counts) => counts,
|
Some(counts) => counts,
|
||||||
None => return 0,
|
None => return Vec::new(),
|
||||||
};
|
};
|
||||||
|
|
||||||
for matcher in cache.iter() {
|
for matcher in cache.iter() {
|
||||||
@@ -285,11 +328,13 @@ impl RoutineEngine {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let detail = truncate(content, 200);
|
let detail = truncate(content, 200);
|
||||||
self.spawn_fire(routine.clone(), "event", Some(detail));
|
triggered.push(TriggeredRoutine {
|
||||||
fired += 1;
|
routine: routine.clone(),
|
||||||
|
detail,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
fired
|
triggered
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Emit a structured event to system-event routines.
|
/// Emit a structured event to system-event routines.
|
||||||
@@ -845,7 +890,12 @@ impl RoutineEngine {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Spawn a fire in a background task.
|
/// Spawn a fire in a background task.
|
||||||
fn spawn_fire(&self, routine: Routine, trigger_type: &str, trigger_detail: Option<String>) {
|
fn spawn_fire(
|
||||||
|
&self,
|
||||||
|
routine: Routine,
|
||||||
|
trigger_type: &str,
|
||||||
|
trigger_detail: Option<String>,
|
||||||
|
) -> JoinHandle<()> {
|
||||||
let run = RoutineRun {
|
let run = RoutineRun {
|
||||||
id: Uuid::new_v4(),
|
id: Uuid::new_v4(),
|
||||||
routine_id: routine.id,
|
routine_id: routine.id,
|
||||||
@@ -882,7 +932,7 @@ impl RoutineEngine {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
execute_routine(engine, routine, run).await;
|
execute_routine(engine, routine, run).await;
|
||||||
});
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
fn check_cooldown(&self, routine: &Routine) -> bool {
|
fn check_cooldown(&self, routine: &Routine) -> bool {
|
||||||
|
|||||||
+51
-9
@@ -431,6 +431,18 @@ impl ReplChannel {
|
|||||||
let _ = execute!(stderr, terminal::Clear(terminal::ClearType::FromCursorDown));
|
let _ = execute!(stderr, terminal::Clear(terminal::ClearType::FromCursorDown));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn finish_single_message_turn(&self) {
|
||||||
|
if self.single_message.is_none() {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
let tx = self.msg_tx.lock().ok().and_then(|mut guard| guard.take());
|
||||||
|
if let Some(tx) = tx {
|
||||||
|
let msg = IncomingMessage::new("repl", &self.user_id, "/quit");
|
||||||
|
let _ = tx.send(msg).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Default for ReplChannel {
|
impl Default for ReplChannel {
|
||||||
@@ -480,7 +492,9 @@ impl Channel for ReplChannel {
|
|||||||
|
|
||||||
async fn start(&self) -> Result<MessageStream, ChannelError> {
|
async fn start(&self) -> Result<MessageStream, ChannelError> {
|
||||||
let (tx, rx) = mpsc::channel(32);
|
let (tx, rx) = mpsc::channel(32);
|
||||||
// Store tx so send_status can inject approval responses directly
|
// Approval prompts inject responses back through this sender.
|
||||||
|
// In single-message mode we keep it until the turn finishes, then
|
||||||
|
// drop it after enqueuing /quit so the receiver stream can close.
|
||||||
if let Ok(mut guard) = self.msg_tx.lock() {
|
if let Ok(mut guard) = self.msg_tx.lock() {
|
||||||
*guard = Some(tx.clone());
|
*guard = Some(tx.clone());
|
||||||
}
|
}
|
||||||
@@ -496,11 +510,10 @@ impl Channel for ReplChannel {
|
|||||||
|
|
||||||
// Single message mode: send it and return
|
// Single message mode: send it and return
|
||||||
if let Some(msg) = single_message {
|
if let Some(msg) = single_message {
|
||||||
let incoming = IncomingMessage::new("repl", &user_id, &msg).with_timezone(&sys_tz);
|
let incoming = IncomingMessage::new("repl", &user_id, &msg)
|
||||||
|
.with_metadata(serde_json::json!({ "single_message_mode": true }))
|
||||||
|
.with_timezone(&sys_tz);
|
||||||
let _ = tx.blocking_send(incoming);
|
let _ = tx.blocking_send(incoming);
|
||||||
// Ensure the agent exits after handling exactly one turn in -m mode,
|
|
||||||
// even when other channels (gateway/http) are enabled.
|
|
||||||
let _ = tx.blocking_send(IncomingMessage::new("repl", &user_id, "/quit"));
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -663,6 +676,7 @@ impl Channel for ReplChannel {
|
|||||||
println!();
|
println!();
|
||||||
println!();
|
println!();
|
||||||
self.stdin_locked.store(false, Ordering::Relaxed);
|
self.stdin_locked.store(false, Ordering::Relaxed);
|
||||||
|
self.finish_single_message_turn().await;
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -681,6 +695,7 @@ impl Channel for ReplChannel {
|
|||||||
println!();
|
println!();
|
||||||
// Unlock stdin so readline can resume
|
// Unlock stdin so readline can resume
|
||||||
self.stdin_locked.store(false, Ordering::Relaxed);
|
self.stdin_locked.store(false, Ordering::Relaxed);
|
||||||
|
self.finish_single_message_turn().await;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -780,6 +795,7 @@ impl Channel for ReplChannel {
|
|||||||
let msg_tx = Arc::clone(&self.msg_tx);
|
let msg_tx = Arc::clone(&self.msg_tx);
|
||||||
let user_id = self.user_id.clone();
|
let user_id = self.user_id.clone();
|
||||||
let lock_flag = Arc::clone(&self.stdin_locked);
|
let lock_flag = Arc::clone(&self.stdin_locked);
|
||||||
|
let single_message_mode = self.single_message.is_some();
|
||||||
tokio::task::spawn_blocking(move || {
|
tokio::task::spawn_blocking(move || {
|
||||||
let action = run_approval_selector(allow_always).unwrap_or("n");
|
let action = run_approval_selector(allow_always).unwrap_or("n");
|
||||||
// Unlock stdin so readline can resume after approval
|
// Unlock stdin so readline can resume after approval
|
||||||
@@ -788,7 +804,12 @@ impl Channel for ReplChannel {
|
|||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
if let Some(tx) = guard.as_ref() {
|
if let Some(tx) = guard.as_ref() {
|
||||||
let msg = IncomingMessage::new("repl", &user_id, action);
|
let msg = if single_message_mode {
|
||||||
|
IncomingMessage::new("repl", &user_id, action)
|
||||||
|
.with_metadata(serde_json::json!({ "single_message_mode": true }))
|
||||||
|
} else {
|
||||||
|
IncomingMessage::new("repl", &user_id, action)
|
||||||
|
};
|
||||||
let _ = tx.blocking_send(msg);
|
let _ = tx.blocking_send(msg);
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
@@ -889,6 +910,7 @@ impl Channel for ReplChannel {
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use futures::StreamExt;
|
use futures::StreamExt;
|
||||||
|
use tokio::time::{Duration, timeout};
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
@@ -897,16 +919,36 @@ mod tests {
|
|||||||
let repl = ReplChannel::with_message("hi".to_string());
|
let repl = ReplChannel::with_message("hi".to_string());
|
||||||
let mut stream = repl.start().await.expect("repl start should succeed");
|
let mut stream = repl.start().await.expect("repl start should succeed");
|
||||||
|
|
||||||
let first = stream.next().await.expect("first message missing");
|
let first = timeout(Duration::from_secs(1), stream.next())
|
||||||
|
.await
|
||||||
|
.expect("timed out waiting for first message")
|
||||||
|
.expect("first message missing");
|
||||||
assert_eq!(first.channel, "repl");
|
assert_eq!(first.channel, "repl");
|
||||||
assert_eq!(first.content, "hi");
|
assert_eq!(first.content, "hi");
|
||||||
|
|
||||||
let second = stream.next().await.expect("quit message missing");
|
assert!(
|
||||||
|
timeout(Duration::from_millis(100), stream.next())
|
||||||
|
.await
|
||||||
|
.is_err(),
|
||||||
|
"single-message mode should wait for the turn to finish before quitting"
|
||||||
|
);
|
||||||
|
|
||||||
|
repl.respond(&first, OutgoingResponse::text("done"))
|
||||||
|
.await
|
||||||
|
.expect("respond should succeed");
|
||||||
|
|
||||||
|
let second = timeout(Duration::from_secs(1), stream.next())
|
||||||
|
.await
|
||||||
|
.expect("timed out waiting for quit message")
|
||||||
|
.expect("quit message missing");
|
||||||
assert_eq!(second.channel, "repl");
|
assert_eq!(second.channel, "repl");
|
||||||
assert_eq!(second.content, "/quit");
|
assert_eq!(second.content, "/quit");
|
||||||
|
|
||||||
assert!(
|
assert!(
|
||||||
stream.next().await.is_none(),
|
timeout(Duration::from_secs(1), stream.next())
|
||||||
|
.await
|
||||||
|
.expect("timed out waiting for stream to close")
|
||||||
|
.is_none(),
|
||||||
"stream should end after /quit"
|
"stream should end after /quit"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -253,6 +253,6 @@ async def test_telegram_hot_activation_transitions_installed_to_active(page):
|
|||||||
assert await card.locator(SEL["ext_pairing_label"]).count() == 0
|
assert await card.locator(SEL["ext_pairing_label"]).count() == 0
|
||||||
|
|
||||||
assert captured_setup_payloads == [
|
assert captured_setup_payloads == [
|
||||||
{"secrets": {"telegram_bot_token": "123456789:ABCdefGhI"}},
|
{"secrets": {"telegram_bot_token": "123456789:ABCdefGhI"}, "fields": {}},
|
||||||
{"secrets": {}},
|
{"secrets": {}, "fields": {}},
|
||||||
]
|
]
|
||||||
|
|||||||
Reference in New Issue
Block a user