diff --git a/internal/flagutil/metadata.go b/internal/flagutil/metadata.go index 06a08a7b..f15d440b 100644 --- a/internal/flagutil/metadata.go +++ b/internal/flagutil/metadata.go @@ -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) } @@ -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) } diff --git a/internal/flagutil/stdin.go b/internal/flagutil/stdin.go index ed8c89c2..5e526c13 100644 --- a/internal/flagutil/stdin.go +++ b/internal/flagutil/stdin.go @@ -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 @@ -20,6 +20,20 @@ // 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 @@ -27,7 +41,11 @@ import ( "errors" "fmt" "io" + "os" + "sync" "time" + + "github.com/spf13/cobra" ) // StdinReadTimeout bounds how long a command will wait for piped input before @@ -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 @@ -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 +} diff --git a/internal/flagutil/stdin_test.go b/internal/flagutil/stdin_test.go new file mode 100644 index 00000000..33f43486 --- /dev/null +++ b/internal/flagutil/stdin_test.go @@ -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) + } +} diff --git a/internal/output/agentmodecheck.go b/internal/output/agentmodecheck.go new file mode 100644 index 00000000..1b6eea96 --- /dev/null +++ b/internal/output/agentmodecheck.go @@ -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) +} diff --git a/test/stdin-handling.test.ts b/test/stdin-handling.test.ts index a7a25fd3..6bad1498 100644 --- a/test/stdin-handling.test.ts +++ b/test/stdin-handling.test.ts @@ -33,19 +33,24 @@ const ARGS = [ '--token', 'vfp_not_a_real_token', ].join(' '); -/** Runs a real shell pipeline: | 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: | vf agent update --dry-run [flags] */ +function pipeline(producer: string, options: { timeout?: number; flags?: string; env?: Record } = {}) { + 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. @@ -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.