Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
9 changes: 0 additions & 9 deletions crates/aster-cli/src/config/env.rs
Original file line number Diff line number Diff line change
Expand Up @@ -262,15 +262,6 @@ fn set(
if let Some(note) = shadow_note(&var, before, local) {
println!("{}", paint(DIM, &note));
}
if !VARS.iter().any(|v| v.var == var) {
println!(
"{}",
paint(
DIM,
&format!("Nothing in Aster reads {var}; `aster env list` shows the ones that are read."),
)
);
}
Ok(())
}

Expand Down
21 changes: 0 additions & 21 deletions crates/aster-cli/src/config/key.rs
Original file line number Diff line number Diff line change
Expand Up @@ -418,17 +418,6 @@ pub(crate) fn set(
if let Some(note) = shadow_note(&var, before, local) {
println!("{}", paint(DIM, &note));
}
if !known(&var) {
println!(
"{}",
paint(
DIM,
&format!(
"Nothing in Aster reads {var}; `aster key list` shows the ones that are read."
),
)
);
}
Ok(())
}

Expand Down Expand Up @@ -538,16 +527,6 @@ fn normalize(var: &str) -> Result<String> {
Ok(var)
}

fn known(var: &str) -> bool {
var == SHARED_KEY_VAR
|| var == "ASTER_JEV_API_KEY"
|| aster_web::KEY_VARS
.iter()
.chain(aster_voice::KEY_VARS)
.any(|(_, known, _)| *known == var)
|| catalog_key_vars().iter().any(|(_, known)| *known == var)
}

fn ask(var: &str) -> Result<String> {
if !crate::picker::is_tty() {
bail!(
Expand Down
12 changes: 0 additions & 12 deletions crates/aster-cli/src/tests/key_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,18 +12,6 @@ fn a_name_is_uppercased_and_checked() {
}
}

#[test]
fn every_web_provider_var_is_known() {
// The catalog `aster key list` reads is the one aster-web resolves from, so
// a provider cannot gain a key without turning up here.
for (_, var, _) in aster_web::KEY_VARS.iter().chain(aster_voice::KEY_VARS) {
assert!(known(var), "{var} should be listed");
}
assert!(known(SHARED_KEY_VAR));
assert!(known("OPENAI_API_KEY"), "catalog vars count as known");
assert!(!known("NOT_A_REAL_KEY"));
}

#[test]
fn an_assignment_is_read_back_with_or_without_quotes() {
assert_eq!(assignment("FOO=bar", "FOO"), Some("bar"));
Expand Down
43 changes: 14 additions & 29 deletions crates/aster-cli/src/tui/chat.rs
Original file line number Diff line number Diff line change
Expand Up @@ -185,8 +185,8 @@ pub async fn run_chat(
app.width = tui.width() as usize;
app.markdown.set_width(app.width);
app.instructions = sync::Arc::new(crate::instructions::discover(&repo_root));
// A cached catalog means the first submit can run now; without one it
// waits for the connect like before.
// A cached catalog gives the first turn its MCP tools. Without one the turn
// starts anyway, and the tools join the turn after the connect lands.
app.mcp = cached;
app.mcp_pending = app.mcp.is_none();
app.limits = limits;
Expand Down Expand Up @@ -270,7 +270,7 @@ pub async fn run_chat(

let mut turn: Option<ChatTurn> = None;
if let Some(seed) = seed.filter(|s| !s.trim().is_empty()) {
turn = app.submit_or_hold(&seed, &[], &mut client, &repo_root);
turn = Some(app.submit(&seed, &[], &mut client, &repo_root));
pane.set_task_running(turn.is_some());
}

Expand Down Expand Up @@ -395,13 +395,11 @@ pub async fn run_chat(
AppEvent::McpReady { runtime, problems } => {
app.mcp = runtime;
app.mcp_pending = false;
// Do not print "MCP connected" anymore.
app.error_box(&problems);
if let Some((text, refs)) = app.held_submit.take() {
if app.flash.as_deref().is_some_and(|f| f.starts_with("MCP servers still connecting")) {
app.flash = None;
turn = Some(app.submit(&text, &refs, &mut client, &repo_root));
pane.set_task_running(true);
}
// Do not print "MCP connected" anymore.
app.error_box(&problems);
}

AppEvent::SetMode(Mode::Yolo) => app.confirm_yolo(&mut pane),
Expand Down Expand Up @@ -432,7 +430,7 @@ pub async fn run_chat(
if turn.is_none() {
let unsent = app.take_unsent();
if !unsent.is_empty() {
turn = app.submit_or_hold(&unsent.join("\n\n"), &[], &mut client, &repo_root);
turn = Some(app.submit(&unsent.join("\n\n"), &[], &mut client, &repo_root));
}
}
pane.set_task_running(turn.is_some());
Expand Down Expand Up @@ -660,7 +658,7 @@ fn on_key(
match pane.handle_key(key, app.width as u16) {
InputResult::Submitted { text, refs } => {
app.flash = None;
*turn = app.submit_or_hold(&text, &refs, client, repo_root);
*turn = Some(app.submit(&text, &refs, client, repo_root));
pane.set_task_running(turn.is_some());
}
InputResult::Command(cmd) => {
Expand All @@ -682,7 +680,7 @@ fn on_key(
// cancelled, so they do not render as part of the new turn.
while events_rx.try_recv().is_ok() {}
app.flash = None;
*turn = app.submit_or_hold(&text, &refs, client, repo_root);
*turn = Some(app.submit(&text, &refs, client, repo_root));
pane.set_task_running(turn.is_some());
}
InputResult::None => {
Expand Down Expand Up @@ -1257,7 +1255,6 @@ struct ChatApp {
instructions: sync::Arc<crate::instructions::Instructions>,
mcp: Option<crate::mcp::McpRuntime>,
mcp_pending: bool,
held_submit: Option<(String, Vec<(String, String)>)>,
turn_injected: Option<sync::Arc<sync::Mutex<Vec<String>>>>,
limits: crate::chat::Limits,
provider_base_url: String,
Expand Down Expand Up @@ -1327,7 +1324,6 @@ impl ChatApp {
instructions: sync::Arc::default(),
mcp: None,
mcp_pending: false,
held_submit: None,
turn_injected: None,
limits: crate::chat::Limits::default(),
provider_base_url: String::new(),
Expand Down Expand Up @@ -1637,28 +1633,17 @@ impl ChatApp {
}
}

fn submit_or_hold(
&mut self,
text: &str,
refs: &[(String, String)],
client: &mut AiClient,
repo_root: &std::path::Path,
) -> Option<ChatTurn> {
if self.mcp_pending {
self.held_submit = Some((text.to_string(), refs.to_vec()));
self.flash = Some("connecting to MCP servers…".into());
return None;
}
Some(self.submit(text, refs, client, repo_root))
}

fn submit(
&mut self,
text: &str,
refs: &[(String, String)],
client: &mut AiClient,
repo_root: &std::path::Path,
) -> ChatTurn {
if self.mcp_pending {
self.flash =
Some("MCP servers still connecting · their tools join your next message".into());
}
// A dismissed session picker leaves no transcript open; start one now
// rather than dropping the conversation on the floor.
if self.recorder.is_none() {
Expand Down Expand Up @@ -1936,7 +1921,7 @@ impl ChatApp {
pane: &mut BottomPane<AppEvent>,
) {
match ev {
// Handled on the run loop, which owns the turn a hold replays into.
// Handled on the run loop, which owns the problems box.
AppEvent::McpReady { .. } => {}
// Entering YOLO goes through `confirm_yolo`, never straight here.
AppEvent::SetMode(Mode::Yolo) => {}
Expand Down
38 changes: 12 additions & 26 deletions crates/aster-cli/src/tui/tests/chat_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -895,45 +895,31 @@ fn record_user_persists_turn() {
assert!(persisted);
}

#[test]
fn a_submit_before_mcp_connects_is_held_not_run() {
let mut client = AiClient::new("http://localhost", "k", "m1");
let mut app = chat_app(client.model.clone());
app.mcp_pending = true;

let turn = app.submit_or_hold("go", &[], &mut client, std::path::Path::new("/tmp"));

assert!(turn.is_none());
assert_eq!(
app.held_submit.as_ref().map(|(t, _)| t.as_str()),
Some("go")
);
assert!(app.history.is_empty());
}

#[tokio::test]
async fn a_submit_after_mcp_connects_runs_straight_away() {
async fn a_submit_before_mcp_connects_runs_straight_away() {
let mut client = AiClient::new("http://localhost", "k", "m1");
let mut app = chat_app(client.model.clone());
app.mcp_pending = true;

let turn = app.submit_or_hold("go", &[], &mut client, std::path::Path::new("/tmp"));
let turn = app.submit("go", &[], &mut client, std::path::Path::new("/tmp"));

assert!(turn.is_some());
assert!(app.held_submit.is_none());
assert!(
app.history
.iter()
.any(|m| m.role == "user" && m.content.text() == "go")
);
turn.unwrap().abort();
assert_eq!(
app.flash.as_deref(),
Some("MCP servers still connecting · their tools join your next message")
);
turn.abort();
}

#[tokio::test]
async fn a_message_typed_mid_turn_queues_into_the_running_turn() {
let mut client = AiClient::new("http://localhost", "k", "m1");
let mut app = chat_app(client.model.clone());
let turn = app.submit_or_hold("go", &[], &mut client, std::path::Path::new("/tmp"));
assert!(turn.is_some());
let turn = app.submit("go", &[], &mut client, std::path::Path::new("/tmp"));

assert!(app.queue_mid_turn("and this too", &[]));

Expand All @@ -945,7 +931,7 @@ async fn a_message_typed_mid_turn_queues_into_the_running_turn() {
assert_eq!(queued, vec!["and this too".to_string()]);
// The scrollback and history wait for the engine's injected event.
assert_eq!(app.history.iter().filter(|m| m.role == "user").count(), 1);
turn.unwrap().abort();
turn.abort();
}

#[test]
Expand All @@ -972,15 +958,15 @@ fn an_injected_event_lands_in_scrollback_and_history() {
async fn unsent_queued_messages_are_reclaimed_when_the_turn_ends() {
let mut client = AiClient::new("http://localhost", "k", "m1");
let mut app = chat_app(client.model.clone());
let turn = app.submit_or_hold("go", &[], &mut client, std::path::Path::new("/tmp"));
let turn = app.submit("go", &[], &mut client, std::path::Path::new("/tmp"));
assert!(app.queue_mid_turn("late arrival", &[]));

let unsent = app.take_unsent();

assert_eq!(unsent, vec!["late arrival".to_string()]);
assert!(app.turn_injected.is_none());
assert!(app.take_unsent().is_empty());
turn.unwrap().abort();
turn.abort();
}

fn app_with_memory(dir: &tempfile::TempDir) -> ChatApp {
Expand Down
109 changes: 109 additions & 0 deletions crates/aster-serve/src/dictation.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
//! The mic button in a browser tab. The server runs `aster dictate` for that
//! tab and relays each NDJSON line as a `dictation` event.

use std::process::Stdio;
use std::sync::Arc;

use serde_json::{Value, json};
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::process::{Child, ChildStdin};

use crate::state::{AppState, Instance};

const STOPPED: &str = "Recording stopped unexpectedly. Try again.";

pub struct Dictation {
id: ulid::Ulid,
child: Child,
stdin: Option<ChildStdin>,
}

pub async fn handle(state: &AppState, instance: &Arc<Instance>, action: &str) {
match action {
"start" => start(state, instance).await,
"stop" => {
let mut slot = instance.dictation.lock().await;
if let Some(mut stdin) = slot.as_mut().and_then(|d| d.stdin.take()) {
let _ = stdin.write_all(b"\n").await;
}
}
"cancel" => {
if let Some(mut dictation) = instance.dictation.lock().await.take() {
let _ = dictation.child.start_kill();
}
}
_ => {}
}
}

async fn start(state: &AppState, instance: &Arc<Instance>) {
let mut cmd = state.cli.command(&["dictate"]);
cmd.stderr(Stdio::null());
let mut child = match cmd.spawn() {
Ok(child) => child,
Err(err) => {
post(
instance,
json!({
"type": "error",
"message": "Couldn't start Aster to listen. Try again.",
"detail": err.to_string(),
}),
);
return;
}
};
let id = ulid::Ulid::new();
let stdin = child.stdin.take();
let stdout = child.stdout.take();
let previous = instance
.dictation
.lock()
.await
.replace(Dictation { id, child, stdin });
if let Some(mut previous) = previous {
let _ = previous.child.start_kill();
}
let Some(stdout) = stdout else {
return;
};

let instance = Arc::clone(instance);
tokio::spawn(async move {
let mut settled = false;
let mut lines = BufReader::new(stdout).lines();
while let Ok(Some(line)) = lines.next_line().await {
let Ok(event) = serde_json::from_str::<Value>(&line) else {
continue;
};
if !is_current(&instance, id).await {
return;
}
settled |= matches!(event["type"].as_str(), Some("transcript" | "error"));
post(&instance, event);
}
let mut slot = instance.dictation.lock().await;
if slot.as_ref().is_some_and(|d| d.id == id) {
*slot = None;
if !settled {
post(
&instance,
json!({ "type": "error", "message": STOPPED, "detail": null }),
);
}
}
});
}

async fn is_current(instance: &Instance, id: ulid::Ulid) -> bool {
instance
.dictation
.lock()
.await
.as_ref()
.is_some_and(|d| d.id == id)
}

fn post(instance: &Instance, event: Value) {
instance.post(json!({ "type": "dictation", "event": event }));
}
Loading
Loading