Files
ops-mcp/internal/tool/docker.go
T
yangzhaohan 795ad7660c feat(tool): 添加 Docker 日志实时跟踪功能
- 扩展 DockerTool 支持实时日志跟踪模式
- 新增 follow 参数用于启用持续日志输出
- 实现 SSE 流式推送日志到客户端
- 添加后台协程处理日志流和错误输出
- 集成 server 通知机制推送实时日志行
- 优化命令执行和上下文取消逻辑
2026-07-06 11:43:19 +08:00

184 lines
4.8 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package tool
import (
"bufio"
"bytes"
"context"
"fmt"
"log/slog"
"os/exec"
"strconv"
"github.com/mark3labs/mcp-go/mcp"
"github.com/mark3labs/mcp-go/server"
)
// DockerTool 提供 docker logs 查询与实时跟踪。
type DockerTool struct{}
func NewDockerTool() *DockerTool {
return &DockerTool{}
}
func (d *DockerTool) Name() string { return "docker" }
func (d *DockerTool) Description() string { return "Docker 容器日志查询与实时跟踪" }
func (d *DockerTool) Initialize(_ context.Context) error {
if _, err := exec.LookPath("docker"); err != nil {
return fmt.Errorf("docker CLI not found: %w", err)
}
slog.Info("docker CLI ready")
return nil
}
func (d *DockerTool) Shutdown(_ context.Context) error { return nil }
func (d *DockerTool) HealthCheck(ctx context.Context) error {
return exec.CommandContext(ctx, "docker", "version").Run()
}
func (d *DockerTool) Register(mcpServer *server.MCPServer) error {
mcpServer.AddTool(mcp.NewTool("docker_logs",
mcp.WithDescription("获取 Docker 容器日志,等价于 docker logs --tail N <container>follow 模式通过 SSE 实时推送日志流"),
mcp.WithString("container",
mcp.Required(),
mcp.Description("容器名称或 ID"),
),
mcp.WithNumber("tail",
mcp.Description("返回最后 N 行日志,默认 100"),
),
mcp.WithBoolean("follow",
mcp.Description("是否持续跟踪日志输出(类似 docker logs -f),默认 false。SSE 模式下日志行通过 notifications/docker_logs/stream 推送"),
),
), d.handleLogs)
return nil
}
func (d *DockerTool) handleLogs(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) {
args := getArgs(req)
container, _ := args["container"].(string)
if container == "" {
return mcp.NewToolResultError("container 参数必填"), nil
}
tail := 100
if t, ok := args["tail"].(float64); ok && t > 0 {
tail = int(t)
}
follow := false
if f, ok := args["follow"].(bool); ok {
follow = f
}
if !follow {
return d.oneShotLogs(ctx, container, tail)
}
return d.followLogs(ctx, container, tail)
}
func (d *DockerTool) oneShotLogs(ctx context.Context, container string, tail int) (*mcp.CallToolResult, error) {
cmd := exec.CommandContext(ctx, "docker", "logs",
"--tail", strconv.Itoa(tail),
"--timestamps",
container,
)
var stdout, stderr bytes.Buffer
cmd.Stdout = &stdout
cmd.Stderr = &stderr
if err := cmd.Run(); err != nil {
return mcp.NewToolResultError(
fmt.Sprintf("docker logs 失败: %v\n%s", err, stderr.String()),
), nil
}
output := stdout.String()
if output == "" {
output = "(容器没有日志输出)"
}
return mcp.NewToolResultText(output), nil
}
func (d *DockerTool) followLogs(ctx context.Context, container string, tail int) (*mcp.CallToolResult, error) {
srv := server.ServerFromContext(ctx)
followCtx, cancelFollow := context.WithCancel(context.Background())
cmd := exec.CommandContext(followCtx, "docker", "logs",
"--tail", strconv.Itoa(tail),
"--timestamps",
"-f",
container,
)
stdout, err := cmd.StdoutPipe()
if err != nil {
cancelFollow()
return mcp.NewToolResultError(fmt.Sprintf("创建管道失败: %v", err)), nil
}
stderr, err := cmd.StderrPipe()
if err != nil {
cancelFollow()
return mcp.NewToolResultError(fmt.Sprintf("创建管道失败: %v", err)), nil
}
if err := cmd.Start(); err != nil {
cancelFollow()
return mcp.NewToolResultError(fmt.Sprintf("docker logs 启动失败: %v", err)), nil
}
// 后台读取 stderr
go func() {
scanner := bufio.NewScanner(stderr)
for scanner.Scan() {
slog.Warn("docker logs stderr", "container", container, "line", scanner.Text())
}
if err := scanner.Err(); err != nil {
slog.Warn("docker logs stderr scanner error", "container", container, "err", err)
}
}()
// 后台流式推送 stdout
go func() {
defer cancelFollow()
defer func() {
if err := cmd.Wait(); err != nil {
slog.Warn("docker logs -f exited with error", "container", container, "err", err)
}
}()
scanner := bufio.NewScanner(stdout)
for scanner.Scan() {
line := scanner.Text()
if srv == nil {
continue
}
if err := srv.SendNotificationToClient(ctx, "notifications/docker_logs/stream", map[string]any{
"container": container,
"line": line,
}); err != nil {
slog.Debug("docker logs notification failed, stopping stream", "err", err)
return
}
}
if err := scanner.Err(); err != nil {
slog.Warn("docker logs stdout scanner error", "container", container, "err", err)
}
if srv != nil {
_ = srv.SendNotificationToClient(ctx, "notifications/docker_logs/stream", map[string]any{
"container": container,
"line": "--- 日志流结束 ---",
})
}
}()
return mcp.NewToolResultText(
fmt.Sprintf("开始跟踪容器 %s 日志(tail=%d, follow=true)…", container, tail),
), nil
}