From 27769c27026fb8a7793be581f32f9d2e56bd15ab Mon Sep 17 00:00:00 2001 From: engigu Date: Tue, 10 Feb 2026 11:57:43 +0800 Subject: [PATCH] feat: fix stdout and stderr order error --- internal/executor/executor.go | 122 ++++++++++++++++++++++++---- internal/executor/scheduler.go | 24 ++++-- internal/services/tasks/tiny_log.go | 57 +++++++++++-- 3 files changed, 172 insertions(+), 31 deletions(-) diff --git a/internal/executor/executor.go b/internal/executor/executor.go index 686f75e..a3378c5 100644 --- a/internal/executor/executor.go +++ b/internal/executor/executor.go @@ -5,9 +5,12 @@ import ( "io" "os" "os/exec" + "runtime" "strings" "time" + "github.com/creack/pty" + "github.com/engigu/baihu-panel/internal/logger" "github.com/engigu/baihu-panel/internal/utils" ) @@ -107,28 +110,101 @@ func ExecuteWithHooks(ctx context.Context, req Request, stdout, stderr io.Writer if len(req.Envs) > 0 { cmd.Env = append(cmd.Env, req.Envs...) } + // 强制注入终端环境标识及禁用输出缓冲的标志 + cmd.Env = append(cmd.Env, + "TERM=xterm", + "PYTHONUNBUFFERED=1", + "NODE_NO_WARNINGS=1", + ) - cmd.Stdout = stdout - cmd.Stderr = stderr + var pipeWriter *os.File + var ptyFile *os.File + var copyDone chan struct{} + var err error - // 使用 cmd.Start() + Wait() 以便在后台处理心跳 - err := cmd.Start() - if err != nil { - // Start 失败的处理 - end := time.Now() - result := &Result{ - Status: "failed", - Duration: end.Sub(start).Milliseconds(), - ExitCode: 1, - StartTime: start, // 修正为 start - EndTime: end, + var started bool + // 尝试开启 PTY 模式(Unix/macOS 且输出合并时) + if runtime.GOOS != "windows" && stdout != nil && (stdout == stderr || stdout == io.Discard) { + // 强制注入终端环境标识及禁用输出缓冲的标志,确保 PTY 模式下最佳实时性能 + cmd.Env = append(cmd.Env, + "TERM=xterm", + "PYTHONUNBUFFERED=1", + "NODE_NO_WARNINGS=1", + ) + f, ptyErr := pty.Start(cmd) + if ptyErr == nil { + logger.Infof("[Executor] 任务 #%d 启动于 PTY 模式", logID) + ptyFile = f + started = true + copyDone = make(chan struct{}) + go func() { + defer close(copyDone) + // io.Copy 对于 PTY 来说是最稳健且即时的流式拷贝 + io.Copy(stdout, f) + f.Close() + }() + } else { + logger.Errorf("[Executor] 任务 #%d PTY 启动失败: %v", logID, ptyErr) } - // 执行后钩子 - if hooks != nil { - result.Output += "\n[System Error] " + err.Error() - hooks.PostExecute(ctx, logID, result) + } + + if !started { + // 如果 stdout 和 stderr 指针不一致,但在逻辑上我们知道它们是同一个 MultiWriter, + // 这里会显示为 Pipe 模式。 + if stdout != stderr && stdout != io.Discard { + logger.Debugf("[Executor] 任务 #%d stdout (%p) and stderr (%p) are different, falling back to Pipe mode.", logID, stdout, stderr) } - return result, err + logger.Infof("[Executor] 任务 #%d 启动于 Pipe 模式", logID) + if stdout != nil && stdout == stderr { + pr, pw, err := os.Pipe() + if err == nil { + cmd.Stdout = pw + cmd.Stderr = pw + pipeWriter = pw + copyDone = make(chan struct{}) + go func() { + io.Copy(stdout, pr) + pr.Close() + close(copyDone) + }() + } else { + cmd.Stdout = stdout + cmd.Stderr = stderr + } + } else { + cmd.Stdout = stdout + cmd.Stderr = stderr + } + + // 使用 cmd.Start() + Wait() 以便在后台处理心跳 + err = cmd.Start() + if err != nil { + if pipeWriter != nil { + pipeWriter.Close() + } + // Start 失败的处理 + end := time.Now() + result := &Result{ + Status: "failed", + Duration: end.Sub(start).Milliseconds(), + ExitCode: 1, + StartTime: start, // 修正为 start + EndTime: end, + } + // 执行后钩子 + if hooks != nil { + result.Output += "\n[System Error] " + err.Error() + hooks.PostExecute(ctx, logID, result) + } + return result, err + } + + // 在父进程中关闭写端,这样子进程退出后 pr 才会收到 EOF + if pipeWriter != nil { + pipeWriter.Close() + } + } else { + // PTY 模式下 cmd.Start() 已经在 pty.Start(cmd) 中调用过了 } // 启动心跳协程 @@ -153,6 +229,16 @@ func ExecuteWithHooks(ctx context.Context, req Request, stdout, stderr io.Writer err = cmd.Wait() close(done) // 停止心跳 + // PTY 模式下需要显式关闭 + if ptyFile != nil { + ptyFile.Close() + } + + // 等待日志复制完成 + if copyDone != nil { + <-copyDone + } + end := time.Now() result := &Result{ diff --git a/internal/executor/scheduler.go b/internal/executor/scheduler.go index 77c7b59..de56ca6 100644 --- a/internal/executor/scheduler.go +++ b/internal/executor/scheduler.go @@ -341,16 +341,24 @@ func (s *Scheduler) executeTask(req *ExecutionRequest) (*ExecutionResult, error) var combinedBuf safeBuffer var stdoutWriter, stderrWriter io.Writer - if stdout != nil { - stdoutWriter = io.MultiWriter(&combinedBuf, stdout) + if stdout != nil && stdout == stderr { + // 如果 stdout 和 stderr 是同一个对象,合并成一个 MultiWriter + // 这样后面 ExecuteWithHooks 才能识别出它们是同一个,从而开启 PTY 模式 + mw := io.MultiWriter(&combinedBuf, stdout) + stdoutWriter = mw + stderrWriter = mw } else { - stdoutWriter = &combinedBuf - } + if stdout != nil { + stdoutWriter = io.MultiWriter(&combinedBuf, stdout) + } else { + stdoutWriter = &combinedBuf + } - if stderr != nil { - stderrWriter = io.MultiWriter(&combinedBuf, stderr) - } else { - stderrWriter = &combinedBuf + if stderr != nil { + stderrWriter = io.MultiWriter(&combinedBuf, stderr) + } else { + stderrWriter = &combinedBuf + } } // 3. 实际开始执行事件 (经过队列和速率限制之后) diff --git a/internal/services/tasks/tiny_log.go b/internal/services/tasks/tiny_log.go index 0ff5971..7c0dbb5 100644 --- a/internal/services/tasks/tiny_log.go +++ b/internal/services/tasks/tiny_log.go @@ -8,6 +8,7 @@ import ( "io" "os" "sync" + "unicode/utf8" "github.com/engigu/baihu-panel/internal/utils" ) @@ -55,6 +56,7 @@ type TinyLog struct { path string writer *bufio.Writer subscribers []chan []byte + remainder []byte // Leftover bytes from previous write (partial multi-byte characters) closed bool } @@ -85,17 +87,46 @@ func (l *TinyLog) Write(p []byte) (n int, err error) { return 0, os.ErrClosed } - // 1. Convert to UTF-8 if necessary (common on Windows) - text := utils.ToUTF8(p) + // 1. Combine with remainder from previous call + originalInputLen := len(p) + payload := p + if len(l.remainder) > 0 { + payload = append(l.remainder, p...) + l.remainder = nil + } + + // 2. Identify trailing partial UTF-8 sequence + lastSafe := len(payload) + // UTF-8 characters are max 4 bytes. Check the last few bytes. + for i := len(payload) - 1; i >= 0 && i >= len(payload)-4; i-- { + if utf8.RuneStart(payload[i]) { + if !utf8.FullRune(payload[i:]) { + // Indeed a partial rune at the end + lastSafe = i + l.remainder = make([]byte, len(payload)-i) + copy(l.remainder, payload[i:]) + } + break + } + } + + // If the entire payload is partial and not longer than a max UTF-8 char, + // keep it all for the next call. + if lastSafe == 0 && len(l.remainder) > 0 { + return originalInputLen, nil + } + + // 3. Convert only the complete part to UTF-8 + text := utils.ToUTF8(payload[:lastSafe]) data := []byte(text) - // 2. Write to file buffer + // 4. Write to file buffer _, err = l.writer.Write(data) if err != nil { return 0, err } - // 3. Broadcast to subscribers + // 5. Broadcast to subscribers if len(l.subscribers) > 0 { for _, ch := range l.subscribers { select { @@ -106,7 +137,7 @@ func (l *TinyLog) Write(p []byte) (n int, err error) { } } - return len(p), nil + return originalInputLen, nil } // Subscribe returns a channel that receives log chunks in real-time @@ -142,6 +173,22 @@ func (l *TinyLog) Close() error { return nil } + // Process any remaining bytes + if len(l.remainder) > 0 { + text := utils.ToUTF8(l.remainder) + data := []byte(text) + _, _ = l.writer.Write(data) + + // Also notify subscribers of the last bit + for _, ch := range l.subscribers { + select { + case ch <- data: + default: + } + } + l.remainder = nil + } + // Flush buffer to file if err := l.writer.Flush(); err != nil { return err