feat: fix stdout and stderr order error
This commit is contained in:
+104
-18
@@ -5,9 +5,12 @@ import (
|
|||||||
"io"
|
"io"
|
||||||
"os"
|
"os"
|
||||||
"os/exec"
|
"os/exec"
|
||||||
|
"runtime"
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/creack/pty"
|
||||||
|
"github.com/engigu/baihu-panel/internal/logger"
|
||||||
"github.com/engigu/baihu-panel/internal/utils"
|
"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 {
|
if len(req.Envs) > 0 {
|
||||||
cmd.Env = append(cmd.Env, req.Envs...)
|
cmd.Env = append(cmd.Env, req.Envs...)
|
||||||
}
|
}
|
||||||
|
// 强制注入终端环境标识及禁用输出缓冲的标志
|
||||||
|
cmd.Env = append(cmd.Env,
|
||||||
|
"TERM=xterm",
|
||||||
|
"PYTHONUNBUFFERED=1",
|
||||||
|
"NODE_NO_WARNINGS=1",
|
||||||
|
)
|
||||||
|
|
||||||
cmd.Stdout = stdout
|
var pipeWriter *os.File
|
||||||
cmd.Stderr = stderr
|
var ptyFile *os.File
|
||||||
|
var copyDone chan struct{}
|
||||||
|
var err error
|
||||||
|
|
||||||
// 使用 cmd.Start() + Wait() 以便在后台处理心跳
|
var started bool
|
||||||
err := cmd.Start()
|
// 尝试开启 PTY 模式(Unix/macOS 且输出合并时)
|
||||||
if err != nil {
|
if runtime.GOOS != "windows" && stdout != nil && (stdout == stderr || stdout == io.Discard) {
|
||||||
// Start 失败的处理
|
// 强制注入终端环境标识及禁用输出缓冲的标志,确保 PTY 模式下最佳实时性能
|
||||||
end := time.Now()
|
cmd.Env = append(cmd.Env,
|
||||||
result := &Result{
|
"TERM=xterm",
|
||||||
Status: "failed",
|
"PYTHONUNBUFFERED=1",
|
||||||
Duration: end.Sub(start).Milliseconds(),
|
"NODE_NO_WARNINGS=1",
|
||||||
ExitCode: 1,
|
)
|
||||||
StartTime: start, // 修正为 start
|
f, ptyErr := pty.Start(cmd)
|
||||||
EndTime: end,
|
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()
|
if !started {
|
||||||
hooks.PostExecute(ctx, logID, result)
|
// 如果 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()
|
err = cmd.Wait()
|
||||||
close(done) // 停止心跳
|
close(done) // 停止心跳
|
||||||
|
|
||||||
|
// PTY 模式下需要显式关闭
|
||||||
|
if ptyFile != nil {
|
||||||
|
ptyFile.Close()
|
||||||
|
}
|
||||||
|
|
||||||
|
// 等待日志复制完成
|
||||||
|
if copyDone != nil {
|
||||||
|
<-copyDone
|
||||||
|
}
|
||||||
|
|
||||||
end := time.Now()
|
end := time.Now()
|
||||||
|
|
||||||
result := &Result{
|
result := &Result{
|
||||||
|
|||||||
@@ -341,16 +341,24 @@ func (s *Scheduler) executeTask(req *ExecutionRequest) (*ExecutionResult, error)
|
|||||||
var combinedBuf safeBuffer
|
var combinedBuf safeBuffer
|
||||||
var stdoutWriter, stderrWriter io.Writer
|
var stdoutWriter, stderrWriter io.Writer
|
||||||
|
|
||||||
if stdout != nil {
|
if stdout != nil && stdout == stderr {
|
||||||
stdoutWriter = io.MultiWriter(&combinedBuf, stdout)
|
// 如果 stdout 和 stderr 是同一个对象,合并成一个 MultiWriter
|
||||||
|
// 这样后面 ExecuteWithHooks 才能识别出它们是同一个,从而开启 PTY 模式
|
||||||
|
mw := io.MultiWriter(&combinedBuf, stdout)
|
||||||
|
stdoutWriter = mw
|
||||||
|
stderrWriter = mw
|
||||||
} else {
|
} else {
|
||||||
stdoutWriter = &combinedBuf
|
if stdout != nil {
|
||||||
}
|
stdoutWriter = io.MultiWriter(&combinedBuf, stdout)
|
||||||
|
} else {
|
||||||
|
stdoutWriter = &combinedBuf
|
||||||
|
}
|
||||||
|
|
||||||
if stderr != nil {
|
if stderr != nil {
|
||||||
stderrWriter = io.MultiWriter(&combinedBuf, stderr)
|
stderrWriter = io.MultiWriter(&combinedBuf, stderr)
|
||||||
} else {
|
} else {
|
||||||
stderrWriter = &combinedBuf
|
stderrWriter = &combinedBuf
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 3. 实际开始执行事件 (经过队列和速率限制之后)
|
// 3. 实际开始执行事件 (经过队列和速率限制之后)
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
"io"
|
"io"
|
||||||
"os"
|
"os"
|
||||||
"sync"
|
"sync"
|
||||||
|
"unicode/utf8"
|
||||||
|
|
||||||
"github.com/engigu/baihu-panel/internal/utils"
|
"github.com/engigu/baihu-panel/internal/utils"
|
||||||
)
|
)
|
||||||
@@ -55,6 +56,7 @@ type TinyLog struct {
|
|||||||
path string
|
path string
|
||||||
writer *bufio.Writer
|
writer *bufio.Writer
|
||||||
subscribers []chan []byte
|
subscribers []chan []byte
|
||||||
|
remainder []byte // Leftover bytes from previous write (partial multi-byte characters)
|
||||||
closed bool
|
closed bool
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -85,17 +87,46 @@ func (l *TinyLog) Write(p []byte) (n int, err error) {
|
|||||||
return 0, os.ErrClosed
|
return 0, os.ErrClosed
|
||||||
}
|
}
|
||||||
|
|
||||||
// 1. Convert to UTF-8 if necessary (common on Windows)
|
// 1. Combine with remainder from previous call
|
||||||
text := utils.ToUTF8(p)
|
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)
|
data := []byte(text)
|
||||||
|
|
||||||
// 2. Write to file buffer
|
// 4. Write to file buffer
|
||||||
_, err = l.writer.Write(data)
|
_, err = l.writer.Write(data)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return 0, err
|
return 0, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// 3. Broadcast to subscribers
|
// 5. Broadcast to subscribers
|
||||||
if len(l.subscribers) > 0 {
|
if len(l.subscribers) > 0 {
|
||||||
for _, ch := range l.subscribers {
|
for _, ch := range l.subscribers {
|
||||||
select {
|
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
|
// Subscribe returns a channel that receives log chunks in real-time
|
||||||
@@ -142,6 +173,22 @@ func (l *TinyLog) Close() error {
|
|||||||
return nil
|
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
|
// Flush buffer to file
|
||||||
if err := l.writer.Flush(); err != nil {
|
if err := l.writer.Flush(); err != nil {
|
||||||
return err
|
return err
|
||||||
|
|||||||
Reference in New Issue
Block a user