Skip to content
Open
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
130 changes: 130 additions & 0 deletions cmd/envctl/ask.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
package main

import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"strings"
"time"

"github.com/spf13/cobra"

"github.com/sam-bretz/envctl/internal/daemon"
"github.com/sam-bretz/envctl/internal/workflow"
)

// answerWait covers the coordinator starting the answer and the answer's own
// time limit, with room for a slow start.
const answerWait = 12 * time.Minute

// runGetter is the part of the daemon client awaitAnswer needs, so a test can
// supply the run's progress without a coordinator.
type runGetter interface {
Get(ctx context.Context, id string) (*workflow.Run, error)
}

// awaitAnswer prints the answer to the question just asked. Without --wait it
// reports that the question was recorded and returns.
func awaitAnswer(cmd *cobra.Command, g *globals, client runGetter, asked *workflow.Run, req daemon.ActionRequest, wait bool, timeout time.Duration) error {
id := askedQuestion(asked, req)
if id == "" {
return errors.New("the question was accepted but is not on the run; run envctl run show to find it")
}
if !wait {
return printQuestion(cmd, g, asked.Current().Question(id))
}
ctx, cancel := context.WithTimeout(cmd.Context(), timeout)
defer cancel()
tick := time.NewTicker(questionPoll)
defer tick.Stop()
for {
run, err := client.Get(ctx, asked.ID)
if err != nil {
return err
}
if q := findQuestion(run, id); q != nil && q.Finished() {
return printQuestion(cmd, g, q)
}
select {
case <-ctx.Done():
return fmt.Errorf("no answer within %s; the question stays open, and envctl run show will show the answer when it arrives", timeout)
case <-tick.C:
}
}
}

// questionPoll is how often awaitAnswer checks for an answer.
var questionPoll = 2 * time.Second

// askedQuestion finds the question this request just added: the latest one on
// the stage with that text, which is the one the action appended.
func askedQuestion(run *workflow.Run, req daemon.ActionRequest) string {
rev := run.Revision(req.Revision)
if rev == nil {
rev = run.Current()
}
text := strings.TrimSpace(req.Message)
for i := len(rev.Questions) - 1; i >= 0; i-- {
if q := rev.Questions[i]; q.Node == req.Node && q.Text == text {
return q.ID
}
}
return ""
}

func findQuestion(run *workflow.Run, id string) *workflow.Question {
for i := range run.Revisions {
if q := run.Revisions[i].Question(id); q != nil {
return q
}
}
return nil
}

func printQuestion(cmd *cobra.Command, g *globals, q *workflow.Question) error {
if q == nil {
return errors.New("question not found")
}
if g.jsonOut {
return json.NewEncoder(cmd.OutOrStdout()).Encode(q)
}
out := cmd.OutOrStdout()
switch q.State {
case workflow.QuestionAnswered:
_, err := fmt.Fprintf(out, "%s's supervisor:\n%s\n", q.Node, q.Answer)
return err
case workflow.QuestionFailed:
return fmt.Errorf("the %s supervisor could not answer: %s", q.Node, q.Detail)
}
_, err := fmt.Fprintf(out, "Asked %s's supervisor (%s). envctl run show will show the answer.\n", q.Node, q.ID)
return err
}

// printQuestions lists questions asked on a revision and their answers, which
// is where `envctl run ask --wait=false` says the answer will appear.
func printQuestions(out io.Writer, rev *workflow.Revision) {
if len(rev.Questions) == 0 {
return
}
fmt.Fprintln(out, " questions:")
for _, q := range rev.Questions {
fmt.Fprintf(out, " %s asked %s's supervisor: %s\n", askerName(q.Asker), q.Node, q.Text)
switch q.State {
case workflow.QuestionAnswered:
fmt.Fprintf(out, " answer: %s\n", strings.ReplaceAll(q.Answer, "\n", "\n "))
case workflow.QuestionFailed:
fmt.Fprintf(out, " not answered: %s\n", q.Detail)
default:
fmt.Fprintf(out, " waiting for an answer (%s)\n", q.State)
}
}
}

func askerName(asker string) string {
if asker == "" || asker == "local" {
return "you"
}
return asker
}
114 changes: 114 additions & 0 deletions cmd/envctl/ask_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
package main

import (
"bytes"
"context"
"strings"
"testing"
"time"

"github.com/spf13/cobra"

"github.com/sam-bretz/envctl/internal/daemon"
"github.com/sam-bretz/envctl/internal/workflow"
)

// progress hands back a run whose question moves through the given states,
// one per poll.
type progress struct {
run *workflow.Run
states []workflow.Question
polls int
}

func (p *progress) Get(context.Context, string) (*workflow.Run, error) {
i := min(p.polls, len(p.states)-1)
p.polls++
next := *workflow.Clone(p.run)
next.Current().Questions = []workflow.Question{p.states[i]}
return &next, nil
}

func askedRun(t *testing.T) (*workflow.Run, daemon.ActionRequest) {
t.Helper()
c, err := workflow.Parse([]byte("version: 2\nproject: demo\nrepositories: [{id: app, url: /source}]\nworkflow: {template: feature}\n"))
if err != nil {
t.Fatal(err)
}
r, err := workflow.NewRun("demo", "ship it", "dev", c, time.Now())
if err != nil {
t.Fatal(err)
}
r.Current().Questions = []workflow.Question{{ID: "question_1", Node: "code", Text: "is there a PR?", State: workflow.QuestionPending}}
return r, daemon.ActionRequest{Revision: r.CurrentRevision, Node: "code", Message: " is there a PR? "}
}

func runAsk(t *testing.T, client runGetter, r *workflow.Run, req daemon.ActionRequest, wait bool, timeout time.Duration) (string, error) {
t.Helper()
old := questionPoll
questionPoll = time.Millisecond
t.Cleanup(func() { questionPoll = old })
cmd := &cobra.Command{}
cmd.SetContext(context.Background())
var out bytes.Buffer
cmd.SetOut(&out)
err := awaitAnswer(cmd, &globals{}, client, r, req, wait, timeout)
return out.String(), err
}

func TestRunAskWaitsForTheAnswerAndPrintsIt(t *testing.T) {
r, req := askedRun(t)
p := &progress{run: r, states: []workflow.Question{
{ID: "question_1", Node: "code", State: workflow.QuestionPending},
{ID: "question_1", Node: "code", State: workflow.QuestionAnswering},
{ID: "question_1", Node: "code", State: workflow.QuestionAnswered, Answer: "No PR yet; approved-change opens it."},
}}
out, err := runAsk(t, p, r, req, true, time.Second)
if err != nil {
t.Fatal(err)
}
if !strings.Contains(out, "No PR yet; approved-change opens it.") {
t.Fatalf("answer not printed:\n%s", out)
}
}

func TestRunAskReportsAnUnanswerableQuestionAsAnError(t *testing.T) {
r, req := askedRun(t)
p := &progress{run: r, states: []workflow.Question{{ID: "question_1", Node: "code", State: workflow.QuestionFailed, Detail: "the stage's VM is no longer running"}}}
if _, err := runAsk(t, p, r, req, true, time.Second); err == nil || !strings.Contains(err.Error(), "VM is no longer running") {
t.Fatalf("a failed answer was not reported: %v", err)
}
}

func TestRunAskGivesUpWaitingWithoutLosingTheQuestion(t *testing.T) {
r, req := askedRun(t)
p := &progress{run: r, states: []workflow.Question{{ID: "question_1", Node: "code", State: workflow.QuestionAnswering}}}
_, err := runAsk(t, p, r, req, true, 20*time.Millisecond)
if err == nil || !strings.Contains(err.Error(), "stays open") {
t.Fatalf("a timeout did not say the question is still open: %v", err)
}
}

func TestRunAskWithoutWaitingSaysWhereTheAnswerWillBe(t *testing.T) {
r, req := askedRun(t)
out, err := runAsk(t, &progress{run: r}, r, req, false, time.Second)
if err != nil || !strings.Contains(out, "envctl run show") {
t.Fatalf("no-wait did not point to where the answer appears: %v\n%s", err, out)
}
}

func TestRunShowListsQuestionsAndTheirAnswers(t *testing.T) {
r, _ := askedRun(t)
r.Current().Questions = []workflow.Question{
{Node: "code", Asker: "sam", Text: "is there a PR?", State: workflow.QuestionAnswered, Answer: "not yet"},
{Node: "qa", Text: "why?", State: workflow.QuestionFailed, Detail: "VM released"},
{Node: "plan", Text: "scope?", State: workflow.QuestionAnswering},
}
var out bytes.Buffer
printQuestions(&out, r.Current())
for _, want := range []string{"sam asked code's supervisor: is there a PR?", "answer: not yet", "you asked qa's supervisor: why?", "not answered: VM released", "waiting for an answer (answering)"} {
if !strings.Contains(out.String(), want) {
t.Fatalf("missing %q:\n%s", want, out.String())
}
}
}
16 changes: 15 additions & 1 deletion cmd/envctl/workflow.go
Original file line number Diff line number Diff line change
Expand Up @@ -299,7 +299,7 @@ func runCmd(g *globals) *cobra.Command {
return printRun(cmd, g, run)
}})
}
for _, kind := range []string{"message", "rewind", "cancel", "close", "approve", "priority"} {
for _, kind := range []string{"message", "ask", "rewind", "cancel", "close", "approve", "priority"} {
c.AddCommand(runActionCmd(g, kind))
}
plugins := &cobra.Command{Use: "plugin", Short: "Change invocation plugins and reopen Plan in a new revision"}
Expand Down Expand Up @@ -384,6 +384,7 @@ func runActionCmd(g *globals, action string) *cobra.Command {
req := daemon.ActionRequest{Action: action}
var pluginFile string
var configFile string
wait, waitFor := true, answerWait
c := &cobra.Command{Use: action + " <run>", Short: action + " a workflow", Args: cobra.ExactArgs(1), RunE: func(cmd *cobra.Command, args []string) error {
client, err := connect(cmd.Context(), g)
if err != nil {
Expand Down Expand Up @@ -438,6 +439,9 @@ func runActionCmd(g *globals, action string) *cobra.Command {
if err != nil {
return err
}
if action == "ask" {
return awaitAnswer(cmd, g, client, result, req, wait, waitFor)
}
return printRun(cmd, g, result)
}}
c.Flags().StringVar(&req.OperationID, "operation-id", "", "idempotency key")
Expand All @@ -455,6 +459,15 @@ func runActionCmd(g *globals, action string) *cobra.Command {
_ = c.MarkFlagRequired("text")
c.Flags().StringVar(&req.Recipient, "to", "supervisor", "supervisor or worker")
c.Flags().StringVar(&req.Node, "node", "", "stage to address")
case "ask":
c.Short = "ask a stage's supervisor a question, without changing the work"
c.Flags().StringVar(&req.Node, "node", "", "stage whose supervisor answers")
_ = c.MarkFlagRequired("node")
c.Flags().StringVar(&req.Message, "text", "", "the question")
_ = c.MarkFlagRequired("text")
c.Flags().StringVar(&req.Actor, "actor", os.Getenv("USER"), "who is asking")
c.Flags().BoolVar(&wait, "wait", true, "wait for the answer and print it")
c.Flags().DurationVar(&waitFor, "timeout", answerWait, "how long to wait for an answer")
case "rewind":
c.Flags().StringVar(&req.Node, "to", "", "stage to revisit")
_ = c.MarkFlagRequired("to")
Expand Down Expand Up @@ -488,6 +501,7 @@ func printRun(cmd *cobra.Command, g *globals, r *workflow.Run) error {
}
fmt.Fprintf(out, " VM: %s\n", vm)
printAttention(out, r)
printQuestions(out, rev)
printRuntime(out, "", rev.Runtime)
for node, child := range rev.ChildRuntimes {
if child != nil && (child.Runtime.PreviewURL != "" || len(child.Runtime.Services) > 0) {
Expand Down
2 changes: 1 addition & 1 deletion docs/public/agent-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ Plan must verify downstream requirements before dependent execution. Attach a lo

In Bubble Tea, `[` / `]` browse revisions, `,` / `.` select artifacts, `o` opens them, and `p` attaches a plugin reference file. `T` opens a theme picker; `envctl theme list` and `envctl theme set <name>` do the same from a shell, saving to `~/.config/envctl/config.yaml`. Browsing and detaching do not cancel execution. Historical views are read-only; rewind creates a new revision.

Running attempts expose bounded, redacted `progress` (phase, current check, recent agent activity) in `envctl run show <run-id>` (`--json` for the full object), the Bubble Tea Conversation panel, and MCP run reads. Steer with `envctl run message <run-id> --node <stage> --to worker|supervisor --text ...` or `i` in Bubble Tea. A message reaches a running agent by interrupting its guest job and resuming the same harness session with the message; each message's status (included at start, delivered live, pending, or queued for the next attempt) is shown next to it. Live delivery waits until the harness reports its session, is limited to 8 resumes per role per attempt, and does not apply once an attempt's review has finished. Supervisors see every message for their stage and should reject work that ignores worker steering. Each run shows its token usage in `run show` and the dashboard and stops its agents at `limits.run_tokens` (default 20M counted tokens; cache reads count at `limits.cache_read_weight`, default 10%) or `limits.run_cost_usd`; a run marked needs-attention for its token ceiling needs the user to raise the limit through `run rewind --config` or cancel it. Workers write each output document to a Markdown file in their attempt's `outputs/` directory, and the coordinator reads it from there.
Running attempts expose bounded, redacted `progress` (phase, current check, recent agent activity) in `envctl run show <run-id>` (`--json` for the full object), the Bubble Tea Chat panel, and MCP run reads. Steer with `envctl run message <run-id> --node <stage> --to worker|supervisor --text ...` or `i` in Bubble Tea. To find something out without changing the work, ask instead: `envctl run ask <run-id> --node <stage> --text ...`, `?` in Bubble Tea, or MCP `envctl_action` action `ask`. The stage's supervisor answers through a separate read-only invocation that cannot edit files or run commands, so an attempt's result, review, digest and approval are unchanged and a running worker is not interrupted; the answer appears under the revision's `questions`. Asking needs the stage to have started and the run's VM to still be running, counts toward the token ceiling, and is refused once the run is at that ceiling. A message reaches a running agent by interrupting its guest job and resuming the same harness session with the message; each message's status (included at start, delivered live, pending, or queued for the next attempt) is shown next to it. Live delivery waits until the harness reports its session, is limited to 8 resumes per role per attempt, and does not apply once an attempt's review has finished. Supervisors see every message for their stage and should reject work that ignores worker steering. Each run shows its token usage in `run show` and the dashboard and stops its agents at `limits.run_tokens` (default 20M counted tokens; cache reads count at `limits.cache_read_weight`, default 10%) or `limits.run_cost_usd`; a run marked needs-attention for its token ceiling needs the user to raise the limit through `run rewind --config` or cancel it. Workers write each output document to a Markdown file in their attempt's `outputs/` directory, and the coordinator reads it from there.

An agent silent for its stall window (`limits.stall_seconds`, default 600, overridable per node) is nudged once to report status; a still-silent agent, or a third stall, fails the attempt with evidence and the retry budget applies. A running command extends the window once. If a legitimate step, such as a long test suite, runs quietly for longer than twice the window, raise `stall_seconds` for that node rather than treating the failure as flaky.

Expand Down
11 changes: 10 additions & 1 deletion docs/src/content/docs/cli-reference.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ envctl run message <run-id> --node code --to worker --text "keep the public API

`run show` prints each running attempt's phase (`worker`, `stack`, `checks`, `supervisor`), its current check, and its latest agent activity: messages, tool and command invocations, and check output. `--json` includes the same `progress` object on each attempt, with at most 20 lines of 300 bytes each. Activity comes from the harness's own event stream, redacted by the guest journal and again for configured credentials. It is display state only and never checkpoint evidence. Activity-only updates are written at most every few seconds; phase changes and deliveries are written at once.

`run message` (and `i` in Bubble Tea) records a message against the current revision. `--node` limits it to one stage; without it, the message applies to every stage. `run show` and the Bubble Tea Conversation panel show where each message went:
`run message` (and `i` in Bubble Tea) records a message against the current revision. `--node` limits it to one stage; without it, the message applies to every stage. `run show` and the Bubble Tea Chat panel show where each message went:

- **included when … started**: the message existed when the agent's prompt was frozen.
- **delivered live … (resume N)**: the message reached a running agent.
Expand All @@ -66,6 +66,15 @@ Neither harness accepts input mid-run, so live delivery interrupts and resumes.

Worker-directed messages reach the running worker. The supervisor's prompt includes every message for its stage, with an instruction to reject work that ignores steering addressed to the worker. Messages that arrive while the supervisor is reviewing interrupt and resume the supervisor the same way. Each role allows 8 live resumes per attempt. After that, progress reports the limit, and further messages reach the supervisor review or the next attempt. Messages that arrive after an attempt's review has finished apply to that stage's next attempt, or to a rewind.


### Asking a stage's supervisor

```bash
envctl run ask <run> --node <stage> --text "<question>" # waits and prints the answer
envctl run ask <run> --node <stage> --text "<question>" --wait=false
```

`run ask` puts a question to the stage's supervisor and, by default, waits up to 12 minutes (`--timeout`) for the answer. Unlike `run message`, it never steers: the answer comes from a separate read-only invocation that can read the stage's worktree but cannot edit it or run commands, so the attempt's result, review, digest and approval are unchanged, and a running worker is not interrupted. Claude answers from a branch of the supervisor's session (`--fork-session`), leaving the original transcript as it was; Codex starts a fresh read-only session from the stage's facts. A question needs the stage to have started and the run's VM to still be running, and a run at its token ceiling refuses it; answers count toward that ceiling. `run show` lists questions and answers, and MCP clients ask with `envctl_action` action `ask` and read answers from `envctl_run`.
## Pull requests

```bash
Expand Down
Loading
Loading