package actuator import ( "bufio" "context" "fmt" "io" "time" "golang.org/x/crypto/ssh" ) // RunStreaming runs a command on an established SSH client and forwards output // chunks to sink as they arrive. A nil sink collects output silently. Returns // the full combined output and any command error. func RunStreaming(ctx context.Context, client *ssh.Client, command string, sink func(stream string, chunk []byte), timeout time.Duration) (string, error) { session, err := client.NewSession() if err != nil { return "", fmt.Errorf("create session: %w", err) } defer session.Close() outPipe, err := session.StdoutPipe() if err != nil { return "", fmt.Errorf("stdout pipe: %w", err) } errPipe, err := session.StderrPipe() if err != nil { return "", fmt.Errorf("stderr pipe: %w", err) } type streamResult struct { out string err error } resultCh := make(chan streamResult, 1) go func() { var combined []byte done := make(chan struct{}, 2) readStream := func(stream string, r io.Reader) { sc := bufio.NewScanner(r) for sc.Scan() { line := sc.Bytes() chunk := make([]byte, len(line)) copy(chunk, line) if sink != nil { sink(stream, chunk) } if stream == "stdout" || stream == "" { if len(combined) > 0 { combined = append(combined, '\n') } combined = append(combined, chunk...) } } done <- struct{}{} } go readStream("stdout", outPipe) go readStream("stderr", errPipe) runErr := session.Run(command) <-done <-done resultCh <- streamResult{out: string(combined), err: runErr} }() if timeout > 0 { var cancel context.CancelFunc ctx, cancel = context.WithTimeout(ctx, timeout) defer cancel() } select { case <-ctx.Done(): session.Close() return "", ctx.Err() case res := <-resultCh: if res.err != nil { return res.out, fmt.Errorf("command: %w", res.err) } return res.out, nil } }