Skip to content
Closed
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
13 changes: 7 additions & 6 deletions internal/flagutil/metadata.go
Original file line number Diff line number Diff line change
Expand Up @@ -251,10 +251,11 @@ func BuildRequest[T any](cmd *cobra.Command, meta []FlagMeta, bodyFieldPath stri
// flag) have no body for stdin to fill — consuming piped JSON there would
// both surprise pipelines and re-relax required path/query params via the
// bodyPrePopulated relaxation below.
if !bodyPrePopulated && (bodyFieldPath != "" || bodyFlagName != "") && HasStdinInput(cmd) {
if !bodyPrePopulated && (bodyFieldPath != "" || bodyFlagName != "") {
// Bounded: an open pipe that never sends data and never closes would
// otherwise block here forever. See internal/flagutil/stdin.go.
stdinData, err := readStdinBounded(cmd.InOrStdin())
// otherwise block here forever, and in an agent's shell a silent socket
// is no body at all. See internal/flagutil/stdin.go.
stdinData, err := readStdinBody(cmd)
if err != nil {
return nil, fmt.Errorf("failed to read stdin: %w", err)
}
Expand Down Expand Up @@ -331,9 +332,9 @@ func BuildRequestBody[T any](cmd *cobra.Command, flagName string, annotations st

if FlagChanged(cmd, flagName) {
requestData, _ = GetStringFlag(cmd, flagName)
} else if HasStdinInput(cmd) {
// Bounded, same reasoning as BuildRequest above.
stdin, err := readStdinBounded(cmd.InOrStdin())
} else {
// Same rules as BuildRequest above.
stdin, err := readStdinBody(cmd)
if err != nil {
return nil, fmt.Errorf("failed to read stdin: %w", err)
}
Expand Down
94 changes: 88 additions & 6 deletions internal/flagutil/stdin.go
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
// This file is not generated by Speakeasy — it bounds the stdin read that
// BuildRequest performs when a command accepts a request body.
// This file is not generated by Speakeasy — it decides whether stdin holds a
// request body for BuildRequest and BuildRequestBody, and bounds the read.
//
// The problem it solves: `HasStdinInput` reports true for any non-TTY stdin,
// including a pipe that is open but will never carry data or close. io.ReadAll
Expand All @@ -20,14 +20,32 @@
// curl, node and python3 pipelines every single time, which is silent data loss
// dressed up as a fix. Reading with a deadline is slower in the pathological
// case and correct in every other one.
//
// One shape is an exception: a socket, when an AI coding agent runs vf. Claude
// Code runs any Bash command that contains a heredoc with stdin connected to a
// Unix socket that it never writes to and never closes. Every body command in
// such a call waited the full deadline and then failed, even when its flags
// held the whole request. Waiting there protects nothing. Shells hand vf a pipe
// or a file, for `producer | vf`, `vf < file` and heredocs alike, never a
// socket. A socket comes from a program that spawned vf, and a program that
// means to send a body writes it when it starts vf: Node's first byte lands
// within a millisecond, even with a hundred spawns at once. So in agent mode a
// socket that delivers nothing within agentSocketWait is taken to carry no
// body. Pipes and files keep the full deadline in every mode, so a slow
// `curl ... | vf` is still read. A program that must write late can pass the
// body with --body instead.

package flagutil

import (
"errors"
"fmt"
"io"
"os"
"sync"
"time"

"github.com/spf13/cobra"
)

// StdinReadTimeout bounds how long a command will wait for piped input before
Expand All @@ -36,34 +54,84 @@ import (
// dropping that input silently would be worse than the hang this replaces.
const StdinReadTimeout = 10 * time.Second

// agentSocketWait is how long, in agent mode, a socket on stdin has to deliver
// anything before vf treats it as carrying no body.
const agentSocketWait = 250 * time.Millisecond

// ErrStdinTimeout is returned when stdin stayed open past StdinReadTimeout
// without delivering data or closing.
var ErrStdinTimeout = errors.New("timed out reading stdin")

// readStdinBounded reads r to EOF, giving up after StdinReadTimeout.
// isAgentMode reports whether an AI coding agent is running vf. The output
// package owns agent detection and installs its IsAgentMode here when the
// program starts, because output imports flagutil and the reverse import would
// be a cycle. Asking each time, rather than keeping a copy, means flagutil
// always sees the mode as finally resolved.
var isAgentMode = func() bool { return false }

// SetAgentModeCheck installs the function that reports agent mode.
func SetAgentModeCheck(check func() bool) {
isAgentMode = check
}

// silenceMeansNoBody reports whether stdin of this type, when it stays silent
// for agentSocketWait, carries no body rather than a body still on its way.
func silenceMeansNoBody(isAgentMode bool, stdin os.FileMode) bool {
return isAgentMode && stdin&os.ModeSocket != 0
}

// readStdinBody returns the request body piped to cmd, or nil when there is
// none.
func readStdinBody(cmd *cobra.Command) ([]byte, error) {
if !HasStdinInput(cmd) {
return nil, nil
}
in := cmd.InOrStdin()
if in == os.Stdin {
if stat, err := os.Stdin.Stat(); err == nil && silenceMeansNoBody(isAgentMode(), stat.Mode()) {
return readStdinBounded(in, agentSocketWait)
}
}
return readStdinBounded(in, 0)
}

// readStdinBounded reads r to EOF, giving up after StdinReadTimeout. When
// silenceLimit is positive and r delivers nothing at all within it, not even
// EOF, it returns no data and no error: there is no body.
//
// On timeout the read goroutine is abandoned while it is still parked in
// On either timeout the read goroutine is abandoned while it is still parked in
// read(2). That leaks one goroutine and its OS thread for the remaining life of
// the process, which is acceptable here and nowhere else: this runs once per
// invocation of a short-lived CLI that is about to exit. The alternative —
// putting the file descriptor into non-blocking mode — mutates the open file
// description that the parent shell shares, so it would corrupt the caller's
// terminal on any exit path that failed to restore it.
func readStdinBounded(r io.Reader) ([]byte, error) {
func readStdinBounded(r io.Reader, silenceLimit time.Duration) ([]byte, error) {
type result struct {
data []byte
err error
}
done := make(chan result, 1) // buffered so the abandoned goroutine can finish and exit
first := &firstReadSignal{r: r, arrived: make(chan struct{})}

go func() {
data, err := io.ReadAll(r)
data, err := io.ReadAll(first)
done <- result{data: data, err: err}
}()

timer := time.NewTimer(StdinReadTimeout)
defer timer.Stop()

if silenceLimit > 0 {
silence := time.NewTimer(silenceLimit)
defer silence.Stop()
select {
case <-first.arrived:
case <-silence.C:
return nil, nil
}
}

select {
case res := <-done:
return res.data, res.err
Expand All @@ -76,3 +144,17 @@ func readStdinBounded(r io.Reader) ([]byte, error) {
)
}
}

// firstReadSignal closes arrived once the first read from r returns, whether
// with data, EOF or an error.
type firstReadSignal struct {
r io.Reader
once sync.Once
arrived chan struct{}
}

func (f *firstReadSignal) Read(p []byte) (int, error) {
n, err := f.r.Read(p)
f.once.Do(func() { close(f.arrived) })
return n, err
}
97 changes: 97 additions & 0 deletions internal/flagutil/stdin_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
package flagutil

import (
"io"
"os"
"strings"
"testing"
"time"
)

func TestSilenceMeansNoBodyOnlyForASocketInAgentMode(t *testing.T) {
cases := []struct {
name string
isAgentMode bool
stdin os.FileMode
want bool
}{
{"socket in agent mode", true, os.ModeSocket, true},
{"socket outside agent mode", false, os.ModeSocket, false},
// A pipe is what a shell pipeline gives vf; a slow producer must be read.
{"pipe in agent mode", true, os.ModeNamedPipe, false},
{"file in agent mode", true, 0, false},
}
for _, tc := range cases {
if got := silenceMeansNoBody(tc.isAgentMode, tc.stdin); got != tc.want {
t.Errorf("%s: silenceMeansNoBody = %v, want %v", tc.name, got, tc.want)
}
}
}

// silentReader returns a reader that delivers nothing and never closes until
// the test ends.
func silentReader(t *testing.T) io.Reader {
t.Helper()
r, w := io.Pipe()
t.Cleanup(func() { w.Close() })
return r
}

// writeAfter returns a reader that delivers body after delay, then closes.
func writeAfter(delay time.Duration, body string) io.Reader {
r, w := io.Pipe()
go func() {
time.Sleep(delay)
io.WriteString(w, body)
w.Close()
}()
return r
}

func TestReadStdinBoundedTreatsSilenceAsNoBody(t *testing.T) {
started := time.Now()
data, err := readStdinBounded(silentReader(t), 20*time.Millisecond)
if err != nil || data != nil {
t.Fatalf("got (%q, %v), want no body and no error", data, err)
}
if waited := time.Since(started); waited > time.Second {
t.Errorf("waited %s on a silent reader, want about the silence limit", waited)
}
}

func TestReadStdinBoundedReadsABodyThatStartsInTime(t *testing.T) {
data, err := readStdinBounded(writeAfter(10*time.Millisecond, `{"prompt":"x"}`), time.Second)
if err != nil || string(data) != `{"prompt":"x"}` {
t.Fatalf("got (%q, %v), want the body", data, err)
}
}

// The limit is on silence, not on the whole read: once anything arrives, the
// rest may take as long as StdinReadTimeout allows.
func TestReadStdinBoundedKeepsReadingAfterTheFirstBytes(t *testing.T) {
r, w := io.Pipe()
go func() {
io.WriteString(w, `{"prompt":`)
time.Sleep(100 * time.Millisecond)
io.WriteString(w, `"x"}`)
w.Close()
}()
data, err := readStdinBounded(r, 20*time.Millisecond)
if err != nil || string(data) != `{"prompt":"x"}` {
t.Fatalf("got (%q, %v), want the whole body", data, err)
}
}

func TestReadStdinBoundedReturnsAtOnceOnEOF(t *testing.T) {
data, err := readStdinBounded(strings.NewReader(""), time.Second)
if err != nil || len(data) != 0 {
t.Fatalf("got (%q, %v), want an empty body", data, err)
}
}

func TestReadStdinBoundedWithoutALimitWaitsForASlowBody(t *testing.T) {
data, err := readStdinBounded(writeAfter(300*time.Millisecond, `{"prompt":"slow"}`), 0)
if err != nil || string(data) != `{"prompt":"slow"}` {
t.Fatalf("got (%q, %v), want the slow body", data, err)
}
}
15 changes: 15 additions & 0 deletions internal/output/agentmodecheck.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
// This file is not generated by Speakeasy. It lets flagutil ask whether agent
// mode is on: in agent mode a silent socket on stdin is no request body (see
// internal/flagutil/stdin.go), and flagutil cannot import this package to ask,
// since this package imports flagutil.
//
// flagutil calls IsAgentMode itself instead of keeping a copy of the answer, so
// it follows however agent mode is finally resolved, flags included.

package output

import "github.com/voiceflow/cli/internal/flagutil"

func init() {
flagutil.SetAgentModeCheck(IsAgentMode)
}
70 changes: 63 additions & 7 deletions test/stdin-handling.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,19 +33,24 @@ const ARGS = [
'--token', 'vfp_not_a_real_token',
].join(' ');

/** Runs a real shell pipeline: <producer> | vf agent update --dry-run */
function pipeline(producer: string, options: { timeout?: number } = {}) {
return execa({ reject: false, timeout: options.timeout ?? 20_000, stdin: 'ignore' })(
'sh', ['-c', `${producer} | ${VF} ${ARGS}`],
/** Runs a real shell pipeline: <producer> | vf agent update --dry-run [flags] */
function pipeline(producer: string, options: { timeout?: number; flags?: string; env?: Record<string, string> } = {}) {
return execa({ reject: false, timeout: options.timeout ?? 20_000, stdin: 'ignore', env: options.env })(
'sh', ['-c', `${producer} | ${VF} ${ARGS} ${options.flags ?? ''}`],
);
}

/** The dry-run output contains the request body; pull the prompt back out. */
function sentPrompt(output: string): string | null {
const m = output.match(/"prompt":\s*"([^"]*)"/);
/** The dry-run output contains the request body; pull a string field back out. */
function sentField(output: string, field: string): string | null {
const m = output.match(new RegExp(`"${field}":\\s*"([^"]*)"`));
return m ? m[1] : null;
}

const sentPrompt = (output: string) => sentField(output, 'prompt');

/** Puts vf in agent mode the way Claude Code's shell does. */
const AGENT = { CLAUDECODE: '1' };

describe('stdin from cold producers', () => {
// Each of these is a producer that does real work before writing. They are
// the shapes that a readiness check drops, so they are the regression guard.
Expand Down Expand Up @@ -97,6 +102,57 @@ describe('stdin that never delivers', () => {
});
});

describe("stdin in an agent's shell", () => {
// Claude Code runs any Bash command that contains a heredoc with stdin
// connected to a Unix socket that it never writes to and never closes. Every
// body command in such a call used to wait the full 10s and then fail, flags
// or not. execa's stdin: 'pipe' hands vf the same thing: Node spawns children
// over a socketpair, and nothing is written unless the test writes it.
const agentShell = (args: string[]) =>
execa({ reject: false, timeout: 20_000, stdin: 'pipe', env: AGENT })(VF, args);

// Well under the 10s the silent socket used to cost, with room for a slow CI
// machine to start the binary.
const NO_WAIT_MS = 5_000;

it('takes the body from the flags without waiting on a silent socket', async () => {
const started = Date.now();
const result = await agentShell([...ARGS.split(' '), '--prompt', 'from-flag']);

expect(Date.now() - started, 'vf waited on a socket that never writes').toBeLessThan(NO_WAIT_MS);
expect(result.exitCode, result.stderr).toBe(0);
expect(sentPrompt(result.stderr + result.stdout)).toBe('from-flag');
});

it('does not wait either when no body flag is given', async () => {
const started = Date.now();
const result = await agentShell(ARGS.split(' '));

expect(Date.now() - started, 'vf waited on a socket that never writes').toBeLessThan(NO_WAIT_MS);
expect(result.exitCode, result.stderr).toBe(0);
});

it('still reads a body the socket delivers as vf starts', async () => {
const result = await execa({ reject: false, timeout: 20_000, input: '{"prompt":"from-socket"}', env: AGENT })(
VF, ARGS.split(' '),
);
expect(sentPrompt(result.stderr + result.stdout)).toBe('from-socket');
});

// Only a socket is cut short. A shell pipe is what `producer | vf` gives vf,
// in agent mode too, so a producer slower than the socket's wait still has
// its body read, and merged with the flags as the README documents.
it('still reads a slow shell producer, merged with the flags', async () => {
const result = await pipeline(`sh -c 'sleep 1; echo "{\\"prompt\\":\\"from-slow-pipe\\"}"'`, {
flags: '--instructions from-flag',
env: AGENT,
});
expect(result.timedOut).toBe(false);
expect(sentPrompt(result.stderr + result.stdout), 'slow body dropped in agent mode').toBe('from-slow-pipe');
expect(sentField(result.stderr + result.stdout, 'instructions')).toBe('from-flag');
});
});

describe('backward compatibility', () => {
// execa's `input:` is how the rest of the suite pipes stdin. It must keep
// working, but it is deliberately NOT counted as coverage for the cases above.
Expand Down
Loading