containerd
# 1.简介
# 1.1.运行时
- 容器运行时能够管理容器运行的整个生命周期具体一点就是镜像制作、镜像格式定义、镜像管理、镜像分发、容器运行及容器实例创建。容器运行时的行业标准叫做
OCI规范,分为容器运行时规范和容器镜像规范。容器运行时规范描述基于bundle(目录)运行容器,目录保存容器的规范文件config.json和rootfs,rootfs包含容器运行时所需的操作系统文件;容器镜像规范定义镜像打包及展开为bundle。 
kubelet调用容器运行时创建容器,容器运行时实现CRI接口就能接入使用。早期的kubelet默认支持docker作为容器运行时,中间层docker shim由kubelet维护。随着越来越多运行时实现支持CRI,docker shim已从kubelet移除。
# 1.2.containerd
containerd源于docker,最初作为docker的内部组件负责容器的生命周期管理。随着容器技术的普及,促使coontainerd向一个更加模块化、轻量级且安全的容器运行时环境演变。2017年,containerd被捐赠给``Cloud Native Computing Foundation(CNCF),成为独立的开源项目继续发展,标志着containerd作为docker`的内部组件成为容器技术生态中的重要角色。--- 客户端与服务端通信 containerd通过gRPC协议实现客户端与服务端间通信,提供灵活高效的远程过程调用方案 --- 容器运行时 containerd内置runc,runc是一个符合OCI标准的轻量级容器运行时,用于直接在操作系统上运行容器 --- 镜像管理 containerd提供镜像拉取、存储和分发能力,支持OCI镜像规范,可以和docker镜像仓库及其它兼容OCI标准的镜像仓库交互 --- 存储和快照管理 containerd利用快照机制管理容器的文件系统状态,支持层次化存储和增量更新,提供容器镜像和文件系统管理效率和性能 --- 插件系统 containerd设计插件系统,支持功能的扩展和定制,可以集成额外的存储、网络或其他资源管理功能1
2
3
4
5
6
7
8
9
10
11
12
13
14
# 1.3.交互模式
containerd有多种客户端,利用namespace隔离不同客户端的容器和镜像,未指定时docker的命名空间是moby,kubernetes的命名空间是k8s.io。其中,container代表一个容器的元数据,containerd中的Task用于获取容器对象及转换为操作系统可运行的进程。
# 1.4.核心模块
containerd/ ├── api/ │ ├── services/ # gRPC服务接口定义(contentService、tasksService) │ └── types/ # 通用数据类型(Protobuf生成的数据结构) ├── cmd/ │ ├── containerd/ # 主服务入口(启动daemon和插件) │ ├── ctr/ # 客户端工具(直接操作containerd) │ ├── containerd-shim*/ # shim实现(v1/v2版本,管理容器生命周期) ├── content/ # 镜像内容存储核心逻辑 │ ├── local/ # 本地存储实现(boltdb元数据+blob存储) ├── contrib/ # 第三方扩展 │ ├── seccomp/ # 默认seccomp规则 │ ├── nvidia/ # NVIDIA GPU支持 │ ├── apparmor/ # AppArmor配置模板 ├── images/ # 镜像元数据管理 │ ├── image.go # 镜像对象定义 │ ├── importexport.go # 镜像导入/导出逻辑 ├── leases/ # 分布式锁管理 │ ├── lease.go # 租约机制实现 ├── metadata/ # 元数据存储(基于BoltDB) │ ├── bolt.go # 数据库操作核心逻辑 ├── mount/ # 跨平台挂载操作 │ ├── mount_linux.go # Linux挂载实现 │ ├── mount_windows.go # Windows挂载实现 ├── oci/ # OCI规范实现 │ ├── spec_opts_linux.go # 生成Linux容器的config.json ├── plugin/ # 插件框架 │ ├── plugin.go # 插件注册接口定义 ├── remotes/ │ └── docker/ # Docker镜像仓库交互逻辑 ├── runtime/ │ ├── v2/ # shim v2实现(默认OCI运行时接口) │ ├── linux/ # Linux平台特有运行时逻辑 ├── services/ │ ├── tasks/ # 容器任务管理服务(调用runtime/执行操作) │ ├── content/ # 镜像内容存储服务(管理layer blob) │ ├── snapshots/ # 存储快照服务(驱动接口定义) ├── snapshots/ # 存储驱动实现 │ ├── overlay/ # OverlayFS驱动 │ ├── btrfs/ # Btrfs驱动 │ ├── devmapper/ # DeviceMapper驱动 ├── sys/ # 系统调用封装 │ ├── oom_linux.go # OOM事件监控(Linux)1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
# 2.容器运行时
# 2.1.启动
containerd启动时会初始化工具App工具,这是containerd的命令行入口,初始化后会启动App服务。func main() { // 命令行工具 app := command.App() // 启动containerd if err := app.Run(os.Args); err != nil { ... } }1
2
3
4
5
6
7
8初始化
App,会通过server.New()启动containerd及命令行工具,加载所有处理器插件注册到gRPC、tcp和ttrpc服务,实现对外提供基于containerd操作容器能力。// App returns a *cli.App instance. func App() *cli.App { ... app.Action = func(context *cli.Context) error { ... go func() { defer close(chsrv) // 启动gRPC、tcp、ttrpc服务及注册插件 server, err := server.New(ctx, config) ... }() ... } return app }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15server.New()是containerd运行的主逻辑,通过加载所有init注册的Registration,调用p.Init()实例化所有处理器,将处理器以插件的形式注册到对应的gRPC、ttrpc和tcp服务。// 初始化cri处理器 func init() { ... // 注册到plugin.register.r plugin.Register(&plugin.Registration{ Type: plugin.GRPCPlugin, ID: "cri", Config: &config, Requires: []plugin.Type{ plugin.EventPlugin, plugin.ServicePlugin, plugin.WarningPlugin, plugin.SnapshotPlugin, }, // 实例化函数 InitFn: initCRIService, }) } // New creates and initializes a new containerd server func New(ctx context.Context, config *srvconfig.Config) (*Server, error) { ... // 加载init初始化的Registration(根据requires定义的插件依赖顺序排序) plugins, err := LoadPlugins(ctx, config) ... // 遍历Registration for _, p := range plugins { ... // 调用Init()实例化处理器 result := p.Init(initContext) ... instance, err := result.Instance() // gRPC类型插件 if src, ok := instance.(grpcService); ok { grpcServices = append(grpcServices, src) } // ttrpc类型插件 if src, ok := instance.(ttrpcService); ok { ttrpcServices = append(ttrpcServices, src) } // tcp类型插件 if service, ok := instance.(tcpService); ok { tcpServices = append(tcpServices, service) } s.plugins = append(s.plugins, result) } ... // 注册插件至gRPC服务 for _, service := range grpcServices { if err := service.Register(grpcServer); err != nil { return nil, err } } // 注册插件至ttrpc服务 for _, service := range ttrpcServices { if err := service.RegisterTTRPC(ttrpcServer); err != nil { return nil, err } } // 注册插件至tcp服务 for _, service := range tcpServices { if err := service.RegisterTCP(tcpServer); err != nil { return nil, err } } ... return s, nil } // gRPCServer的实现 func (c *criService) register(s *grpc.Server) error { // 容器运行时/镜像运行时的统一实现 instrumented := newInstrumentedService(c) // 注册容器运行时 runtime.RegisterRuntimeServiceServer(s, instrumented) // 注册镜像运行时 runtime.RegisterImageServiceServer(s, instrumented) ... return nil } // 注册服务 func (s *Server) RegisterService(sd *ServiceDesc, ss any) { if ss != nil { // 接口签名校验,ss必须实现serviceDesc规定的接口 ht := reflect.TypeOf(sd.HandlerType).Elem() st := reflect.TypeOf(ss) ... } // 注册 s.register(sd, ss) } func (s *Server) register(sd *ServiceDesc, ss any) { s.mu.Lock() defer s.mu.Unlock() ... // 已经注册过 if _, ok := s.services[sd.ServiceName]; ok { ... } // 存储接口签名及实现 info := &serviceInfo{ serviceImpl: ss, methods: make(map[string]*MethodDesc), streams: make(map[string]*StreamDesc), mdata: sd.Metadata, } // 记录方法名<-->方法映射 for i := range sd.Methods { d := &sd.Methods[i] info.methods[d.MethodName] = d } ... // 记录接口名对应的handler s.services[sd.ServiceName] = info }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
# 2.2.sandbox创建
criService注册到gRPC服务后,kubelet创建sandbox请求会转发至接口描述的处理函数_RuntimeService_RunPodSandbox_Handler,处理函数会基于方法实现调用RunSandbox完成sandbox创建。func _RuntimeService_RunPodSandbox_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(RunPodSandboxRequest) // 解码请求入参 if err := dec(in); err != nil { return nil, err } ... // 基于接口实现的处理handle handler := func(ctx context.Context, req interface{}) (interface{}, error) { return srv.(RuntimeServiceServer).RunPodSandbox(ctx, req.(*RunPodSandboxRequest)) } // 执行处理handle return interceptor(ctx, in, info, handler) }1
2
3
4
5
6
7
8
9
10
11
12
13
14RunSandbox远程调用最终由criService.RunPodSandbox()处理,sandbox的创建结果会返回给kubelet,sandbox创建涉及pause容器及网络。// RunPodSandbox creates and starts a pod-level sandbox. Runtimes should ensure // the sandbox is in ready state. func (c *criService) RunPodSandbox(ctx context.Context, r *runtime.RunPodSandboxRequest) (_ *runtime.RunPodSandboxResponse, retErr error) { ... // 生成唯一的sandbox ID id := util.GenerateID() ... // 生成sandbox名称 name := makeSandboxName(metadata) ... // 唯一性校验,name<-->sandbox ID唯一 if err := c.sandboxNameIndex.Reserve(name, id); err != nil { return nil, fmt.Errorf("failed to reserve sandbox name %q: %w", name, err) } ... // 初始化sandbox对象,状态为unknown sandbox := sandboxstore.NewSandbox( ... sandboxstore.Status{ State: sandboxstore.StateUnknown, CreatedAt: time.Now().UTC(), }, ) // 拉取镜像,sandbox中运行的是infra pause容器 image, err := c.ensureImageExists(ctx, c.config.SandboxImage, config) ... // 镜像转为containerd理解的对象 containerdImage, err := c.toContainerdImage(ctx, *image) ... // 获取容器运行时支持,默认为runc ociRuntime, err := c.getSandboxRuntime(config, r.GetRuntimeHandler()) ... // 生成满足OCI规范的sandbox定义 spec, err := c.sandboxContainerSpec(id, config, &image.ImageSpec.Config, "", ociRuntime.PodAnnotations) ... // 创建containerd层面的容器对象,保存在boltdb(gRPC调用containerd的容器服务) container, err := c.client.NewContainer(ctx, id, opts...) ... // 创建sandbox文件系统目录 sandboxRootDir := c.getSandboxRootDir(id) c.os.MkdirAll(sandboxRootDir, 0755) ... volatileSandboxRootDir := c.getVolatileSandboxRootDir(id) c.os.MkdirAll(volatileSandboxRootDir, 0755) ... // 安装hostname、resolv.conf文件 if err = c.setupSandboxFiles(id, config); err != nil { return nil, fmt.Errorf("failed to setup sandbox files: %w", err) } ... // 非hostnet网络模式,调用CNI安装网络 if podNetwork { ... // 创建容器网络命名空间(挂载点关联网络命名空间,避免无引用时被Linux回收) sandbox.NetNS, err = netns.NewNetNS(netnsMountDir) ... // 调用CNI安装网络,顺序执行CNI插件 if err := c.setupPodNetwork(ctx, &sandbox); err != nil { return nil, fmt.Errorf("failed to setup network for sandbox %q: %w", id, err) } ... } ... // 创建Task(gRPC调用containerd services中的task服务,最终根据构造的容器运行时信息、配置、根文件系统、挂载等调用runc创建) task, err := container.NewTask(ctx, containerdio.NullIO, taskOpts...) ... // 异步等待容器进程结束 exitCh, err := task.Wait(ctrdutil.NamespacedContext()) ... // 节点资源插件 nric, err := nri.New() ... if nric != nil { ... // 运行相关插件,可能是设备插件 if _, err := nric.InvokeWithSandbox(ctx, task, v1.Create, nriSB); err != nil { return nil, fmt.Errorf("nri invoke: %w", err) } } // 启动Task,runc拉起容器进程,此时容器真正启动 if err := task.Start(ctx); err != nil { return nil, fmt.Errorf("failed to start sandbox container task %q: %w", id, err) } // 更新sandbox对象,状态为ready,记录进程pid if err := sandbox.Status.Update(func(status sandboxstore.Status) (sandboxstore.Status, error) { status.Pid = task.Pid() status.State = sandboxstore.StateReady status.CreatedAt = info.CreatedAt return status, nil }) ... // 存储sandbox到缓存,后续加入业务容器、清理、状态判断使用 if err := c.sandboxStore.Add(sandbox); err != nil { return nil, fmt.Errorf("failed to add sandbox %+v into store: %w", sandbox, err) } ... return &runtime.RunPodSandboxResponse{PodSandboxId: id}, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98c.setupPodNetwork用于生成网络配置相关参数,利用CNI的可执行二进制安装sandbox网络。// setupPodNetwork setups up the network for a pod func (c *criService) setupPodNetwork(ctx context.Context, sandbox *sandboxstore.Sandbox) error { ... // 网络安装 result, err := netPlugin.Setup(ctx, id, path, opts...) ... if err != nil { networkPluginOperationsErrors.WithValues(networkSetUpOp).Inc() return err } // 回填sandboxIP和结果 if configs, ok := result.Interfaces[defaultIfName]; ok && len(configs.IPConfigs) > 0 { sandbox.IP, sandbox.AdditionalIPs = selectPodIPs(ctx, configs.IPConfigs, c.config.IPPreference) sandbox.CNIResult = result return nil } return fmt.Errorf("failed to find network info for sandbox %q", id) } // Setup setups the network in the namespace and returns a Result func (c *libcni) Setup(ctx context.Context, id string, path string, opts ...NamespaceOpts) (*Result, error) { // CNI初始化状态检查 if err := c.Status(); err != nil { return nil, err } // 生成网络命名空间对象 ns, err := newNamespace(id, path, opts...) ... // sandbox加入网络 result, err := c.attachNetworks(ctx, ns) ... return c.createResult(result) } func (c *libcni) attachNetworks(ctx context.Context, ns *Namespace) ([]*types100.Result, error) { ... // 顺序遍历Plugin执行CNI,执行网络安装 for i, network := range c.Networks() { wg.Add(1) go asynchAttach(ctx, i, network, ns, &wg, rc) } ... return results, firstError } func asynchAttach(ctx context.Context, index int, n *Network, ns *Namespace, wg *sync.WaitGroup, rc chan asynchAttachResult) { ... r, err := n.Attach(ctx, ns) rc <- asynchAttachResult{index: index, res: r, err: err} } func (n *Network) Attach(ctx context.Context, ns *Namespace) (*types100.Result, error) { // 执行CNI ADD操作 r, err := n.cni.AddNetworkList(ctx, n.config, ns.config(n.ifName)) ... return types100.NewResultFromResult(r) } // AddNetworkList executes a sequence of plugins with the ADD command func (c *CNIConfig) AddNetworkList(ctx context.Context, list *NetworkConfigList, rt *RuntimeConf) (types.Result, error) { ... // 遍历每个CNI Plugin for _, net := range list.Plugins { // ADD网络操作 result, err = c.addNetwork(ctx, list.Name, list.CNIVersion, net, result, rt) ... } // 容器网络安装结果写入文件 if err = c.cacheAdd(result, list.Bytes, list.Name, rt); err != nil { return nil, fmt.Errorf("failed to set network %q cached result: %w", list.Name, err) } return result, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
# 2.3.sandbox停止
kubelet通过gRPC调用请求停止sandbox,根据请求方法名步入_RuntimeService_StopPodSandbox_Handler处理器。func _RuntimeService_RemovePodSandbox_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(RemovePodSandboxRequest) // 解码请求入参 if err := dec(in); err != nil { return nil, err } ... // 基于容器运行时实现的处理函数 handler := func(ctx context.Context, req interface{}) (interface{}, error) { return srv.(RuntimeServiceServer).StopPodSandbox(ctx, req.(*StopPodSandboxRequest)) } // 执行处理函数 return interceptor(ctx, in, info, handler) }1
2
3
4
5
6
7
8
9
10
11
12
13
14srv是容器运行时接口实现,这里是criService,最终请求会步入criService.StopPodSandbox处理。// StopPodSandbox stops the sandbox. If there are any running containers in the // sandbox, they should be forcibly terminated. func (c *criService) StopPodSandbox(ctx context.Context, r *runtime.StopPodSandboxRequest) (*runtime.StopPodSandboxResponse, error) { // 获取缓存中的sandbox对象 sandbox, err := c.sandboxStore.Get(r.GetPodSandboxId()) ... // 停止sandbox if err := c.stopPodSandbox(ctx, sandbox); err != nil { return nil, err } return &runtime.StopPodSandboxResponse{}, nil } func (c *criService) stopPodSandbox(ctx context.Context, sandbox sandboxstore.Sandbox) error { ... // 停止sandbox关联所有容器 containers := c.containerStore.List() for _, container := range containers { if container.SandboxID != id { continue } // 停止每个容器(业务容器,调用runc停止容器进程,unknown状态容器主动删除注册的Task对象,正常状态由defer的事件监听器删除) if err := c.stopContainer(ctx, container, 0); err != nil { return fmt.Errorf("failed to stop container %q: %w", container.ID, err) } } // 清理相关文件(hostname/resolv.conf) if err := c.cleanupSandboxFiles(id, sandbox.Config); err != nil { return fmt.Errorf("failed to cleanup sandbox files: %w", err) } ... // sandbox处于ready或unknown状态 if state == sandboxstore.StateReady || state == sandboxstore.StateUnknown { // 停止sandbox(pause容器) if err := c.stopSandboxContainer(ctx, sandbox); err != nil { return fmt.Errorf("failed to stop sandbox container %q in %q state: %w", id, state, err) } } // 清理网络相关资源 if sandbox.NetNS != nil { ... // 基于CNI插件清理网络 if err := c.teardownPodNetwork(ctx, sandbox); err != nil { return fmt.Errorf("failed to destroy network for sandbox %q: %w", id, err) } // 删除网络命名空间对象 if err := sandbox.NetNS.Remove(); err != nil { return fmt.Errorf("failed to remove network namespace for sandbox %q: %w", id, err) } } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52stopSandboxContainer用于停止sandbox容器,其实就是获取sandbox pause容器,调用runc停止pause容器进程。// stopSandboxContainer kills the sandbox container. // `task.Delete` is not called here because it will be called when // the event monitor handles the `TaskExit` event. func (c *criService) stopSandboxContainer(ctx context.Context, sandbox sandboxstore.Sandbox) error { ... // 获取pause容器 container := sandbox.Container ... // 获取pause容器关联Task task, err := container.Task(ctx, nil) if err != nil { ... // task获取失败的unknown pause容器进一步尝试清理 if state == sandboxstore.StateUnknown { return cleanupUnknownSandbox(ctx, id, sandbox, c) } return nil } // unknown状态的pause容器 if state == sandboxstore.StateUnknown { ... // 异步监听Task停止 exitCh, err := task.Wait(waitCtx) if err != nil { ... // 监听异常,尝试清理unknown pause容器 return cleanupUnknownSandbox(ctx, id, sandbox, c) } ... } // 调用runc停止Task(Task对象删除由defer的事件监听器执行) if err = task.Kill(ctx, syscall.SIGKILL); err != nil && !errdefs.IsNotFound(err) { return fmt.Errorf("failed to kill sandbox container: %w", err) } // 基于Task.wait()的异步管道监听结束事件 return c.waitSandboxStop(ctx, sandbox) } // cleanupUnknownSandbox会清理runc的容器进程及containerd注册的Task对象,内部调用的是handleSandboxExit func handleSandboxExit(ctx context.Context, e *eventtypes.TaskExit, sb sandboxstore.Sandbox, c *criService) error { // 获取pause容器关联的Task进程对象 task, err := sb.Container.Task(ctx, nil) if err != nil { ... } else { // 根据获取的Task调用runc停止容器进程(task.client.TaskService().Delete) if _, err = task.Delete(ctx, WithNRISandboxDelete(sb.ID), containerd.WithProcessKill); err != nil { ... } } // runc的Task已经退出 if errdefs.IsNotFound(err) { // 尝试清理注册到containerd的Task对象,避免ID冲突,容器进程停止但Task注册还在的内存泄露 _, err = c.client.TaskService().Delete(ctx, &apitasks.DeleteTaskRequest{ContainerID: sb.Container.ID()}) ... } // 更新sandbox状态为notReady err = sb.Status.Update(func(status sandboxstore.Status) (sandboxstore.Status, error) { status.State = sandboxstore.StateNotReady status.Pid = 0 return status, nil }) ... // 通知sandbox channel已停止 sb.Stop() return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68teardownPodNetwork会调用CNI Plugin清理安装的网络,覆盖路由、iptables规则清理甚至网卡对的清理。// teardownPodNetwork removes the network from the pod func (c *criService) teardownPodNetwork(ctx context.Context, sandbox sandboxstore.Sandbox) error { // 获取网络插件(/etc/cni/net.d下的定义文件) netPlugin := c.getNetworkPlugin(sandbox.RuntimeHandler) ... // 根据配置的网络插件移除网络资源 err = netPlugin.Remove(ctx, id, path, opts...) ... return nil } // Remove removes the network config from the namespace func (c *libcni) Remove(ctx context.Context, id string, path string, opts ...NamespaceOpts) error { // 检查插件初始化状态 if err := c.Status(); err != nil { return err } // 获取网络命名空间对象 ns, err := newNamespace(id, path, opts...) ... for _, network := range c.Networks() { if err := network.Remove(ctx, ns); err != nil { // 网络资源不存在说明删除成功 if (path == "" && strings.Contains(err.Error(), "no such file or directory")) || strings.Contains(err.Error(), "not found") { continue } return err } } return nil } func (n *Network) Remove(ctx context.Context, ns *Namespace) error { // CNI接口调用 return n.cni.DelNetworkList(ctx, n.config, ns.config(n.ifName)) } // DelNetworkList executes a sequence of plugins with the DEL command func (c *CNIConfig) DelNetworkList(ctx context.Context, list *NetworkConfigList, rt *RuntimeConf) error { ... // CNI支持的版本检查 if gtet, err := version.GreaterThanOrEqualTo(list.CNIVersion, "0.4.0"); err != nil { return err }else if gtet { // 读取容器网络ADD时的结果 cachedResult, err = c.getCachedResult(list.Name, list.CNIVersion, rt) ... } ... // 倒序执行CNI插件的DEL(例如flannel指定的bridge和ipam) for i := len(list.Plugins) - 1; i >= 0; i-- { net := list.Plugins[i] // 执行插件的DEL,由CNI完成资源回收 if err := c.delNetwork(ctx, list.Name, list.CNIVersion, net, cachedResult, rt); err != nil { return fmt.Errorf("plugin %s failed (delete): %w", pluginDescription(net.Network), err) } } // 删除ADD时保存的容器网络文件 _ = c.cacheDel(list.Name, rt) return nil } // func (c *CNIConfig) delNetwork(ctx context.Context, name, cniVersion string, net *NetworkConfig, prevResult types.Result, rt *RuntimeConf) error { // 检查执行器(负责运行插件二进制)已初始化 c.ensureExec() // 查找CNI插件对应二进制文件 pluginPath, err := c.exec.FindInPath(net.Network.Type, c.Path) ... // 构造插件执行配置 newConf, err := buildOneConfig(name, cniVersion, net, prevResult, rt) ... // 执行二进制插件的DEL命令 return invoke.ExecPluginWithoutResult(ctx, pluginPath, newConf.Bytes, c.args("DEL", rt), c.exec) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
# 2.4.sandbox移除
kubelet GC负责调用containerd删除Pod,根据请求方法名步入_RuntimeService_RemovePodSandbox_Handler处理器。func _RuntimeService_RemovePodSandbox_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(RemovePodSandboxRequest) // 解码请求入参 if err := dec(in); err != nil { return nil, err } ... // 基于容器运行时实现的处理函数 handler := func(ctx context.Context, req interface{}) (interface{}, error) { return srv.(RuntimeServiceServer).RemovePodSandbox(ctx, req.(*RemovePodSandboxRequest)) } // 执行处理函数 return interceptor(ctx, in, info, handler) }1
2
3
4
5
6
7
8
9
10
11
12
13
14srv是容器运行时接口实现,这里是criService,最终请求会步入criService.RemovePodSandbox处理。// RemovePodSandbox removes the sandbox. If there are running containers in the // sandbox, they should be forcibly removed. func (c *criService) RemovePodSandbox(ctx context.Context, r *runtime.RemovePodSandboxRequest) (*runtime.RemovePodSandboxResponse, error) { ... // 根据sandboxID获取缓存中的sandbox sandbox, err := c.sandboxStore.Get(r.GetPodSandboxId()) ... // 停止sandbox if err := c.stopPodSandbox(ctx, sandbox); err != nil { return nil, fmt.Errorf("failed to forcibly stop sandbox %q: %w", id, err) } ... // 网络删除状态判断 ... // 移除sandbox里的所有容器 cntrs := c.containerStore.List() for _, cntr := range cntrs { if cntr.SandboxID != id { continue } // 移除boltdb中containerd容器,移除容器status文件,移除内存cri容器对象 _, err = c.RemoveContainer(ctx, &runtime.RemoveContainerRequest{ContainerId: cntr.ID}) ... } // 清理sandbox目录 sandboxRootDir := c.getSandboxRootDir(id) ensureRemoveAll(ctx, sandboxRootDir) ... volatileSandboxRootDir := c.getVolatileSandboxRootDir(id) ensureRemoveAll(ctx, volatileSandboxRootDir) ... // 删除sandbox pause容器(boltdb) sandbox.Container.Delete(ctx, containerd.WithSnapshotCleanup) ... // 删除containerd缓存的sandbox对象 c.sandboxStore.Delete(id) // 释放sandboxName<-->sandboxID映射 c.sandboxNameIndex.ReleaseByKey(id) return &runtime.RemovePodSandboxResponse{}, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
# 2.5.container创建
kubelet通过gRPC调用请求创建init容器或业务容器,根据请求方法名步入_RuntimeService_CreateContainer_Handler处理器。func _RuntimeService_CreateContainer_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(CreateContainerRequest) // 解码请求入参 if err := dec(in); err != nil { return nil, err } ... // 基于接口实现的处理函数 handler := func(ctx context.Context, req interface{}) (interface{}, error) { return srv.(RuntimeServiceServer).CreateContainer(ctx, req.(*CreateContainerRequest)) } // 执行处理函数 return interceptor(ctx, in, info, handler) }1
2
3
4
5
6
7
8
9
10
11
12
13
14srv是容器运行时接口实现,这里是criService,最终请求会步入criService.CreateContainer处理。// CreateContainer creates a new container in the given PodSandbox. func (c *criService) CreateContainer(ctx context.Context, r *runtime.CreateContainerRequest) (_ *runtime.CreateContainerResponse, retErr error) { ... // 获取缓存的sandbox(请求会传入sandbox ID) sandbox, err := c.sandboxStore.Get(r.GetPodSandboxId()) ... // 获取sandbox中pause容器关联的Task s, err := sandbox.Container.Task(ctx, nil) ... // Task代表的容器进程pid sandboxPid := s.Pid() // 生成唯一的容器ID id := util.GenerateID() ... // 传入的容器名及生成的containerd层面容器名 containerName := metadata.Name name := makeContainerName(metadata, sandboxConfig.GetMetadata()) // 检查容器名称及ID是否唯一,映射name<-->containerID if err = c.containerNameIndex.Reserve(name, id); err != nil { return nil, fmt.Errorf("failed to reserve container name %q: %w", name, err) } ... // 初始化container存储对象 meta := containerstore.Metadata{ ID: id, Name: name, SandboxID: sandboxID, Config: config, } // 从镜像存储解析镜像内容,kubelet创建容器前已经请求过拉取镜像 image, err := c.localResolve(config.GetImage().GetImage()) ... // 镜像转为containerd理解的对象 containerdImage, err := c.toContainerdImage(ctx, image) ... // 生成容器root目录和临时目录 containerRootDir := c.getContainerRootDir(id) c.os.MkdirAll(containerRootDir, 0755) ... volatileContainerRootDir := c.getVolatileContainerRootDir(id) c.os.MkdirAll(volatileContainerRootDir, 0755) ... // 处理挂载点及卷 volumeMounts = c.volumeMounts(containerRootDir, config.GetMounts(), &image.ImageSpec.Config) ... // 容器卷挂载 mounts := c.containerMounts(sandboxID, config) // 获取容器运行时支持,默认runc ociRuntime, err := c.getSandboxRuntime(sandboxConfig, sandbox.Metadata.RuntimeHandler) ... // 生成满足OCI的容器定义 spec, err := c.containerSpec(id, sandboxID, sandboxPid, sandbox.NetNSPath, containerName, containerdImage.Name(), config, sandboxConfig, &image.ImageSpec.Config, append(mounts, volumeMounts...), ociRuntime) ... // 初始化容器IO,用于日志输出到宿主机指定位置 containerIO, err := cio.NewContainerIO(id...) ... // 创建containerd层面容器对象,存储到boltdb if cntr, err = c.client.NewContainer(ctx, id, opts...); err != nil { return nil, fmt.Errorf("failed to create containerd container: %w", err) } ... // 创建cri层面容器对象(containerd容器的包装),补充k8s需要的管理信息 container, err := containerstore.NewContainer(...) if err != nil { return nil, fmt.Errorf("failed to create internal container object for %q: %w", id, err) } ... // 存储cri层面容器对象到缓存 if err := c.containerStore.Add(container); err != nil { return nil, fmt.Errorf("failed to add container %q into store: %w", id, err) } return &runtime.CreateContainerResponse{ContainerId: id}, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
# 2.6.container启动
kubelet通过gRPC调用启动init容器或业务容器,根据请求方法名步入_RuntimeService_StartContainer_Handler处理器。func _RuntimeService_StartContainer_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(StartContainerRequest) // 解码请求入参 if err := dec(in); err != nil { return nil, err } ... // 基于容器运行时实现调用的处理函数 handler := func(ctx context.Context, req interface{}) (interface{}, error) { return srv.(RuntimeServiceServer).StartContainer(ctx, req.(*StartContainerRequest)) } // 执行处理函数 return interceptor(ctx, in, info, handler) }1
2
3
4
5
6
7
8
9
10
11
12
13
14srv是容器运行时接口实现,这里是criService,最终请求会步入criService.StartContainer处理。// StartContainer starts the container. func (c *criService) StartContainer(ctx context.Context, r *runtime.StartContainerRequest) (retRes *runtime.StartContainerResponse, retErr error) { ... // 获取缓存的容器对象 cntr, err := c.containerStore.Get(r.GetContainerId()) ... // 容器对象状态设置为starting if err := setContainerStarting(cntr); err != nil { return nil, fmt.Errorf("failed to set starting state for container %q: %w", id, err) } ... // 获取boltdb存储的sandbox对象 sandbox, err := c.sandboxStore.Get(meta.SandboxID) ... // 基于容器IO(stdout/stderr)设置输出管道,日志输出到指定位置 ioCreation := func(id string) (_ containerdio.IO, err error) { stdoutWC, stderrWC, err := c.createContainerLoggers(meta.LogPath, config.GetTty()) ... cntr.IO.AddOutput("log", stdoutWC, stderrWC) cntr.IO.Pipe() return cntr.IO, nil } ... // 获取容器运行时支持 ociRuntime, err := c.getSandboxRuntime(sandbox.Config, sandbox.Metadata.RuntimeHandler) ... // 请求runc创建Task,其实创建的是实际的容器进程对象,此时进程未启动 task, err := container.NewTask(ctx, ioCreation, taskOpts...) ... // 异步监听容器退出 exitCh, err := task.Wait(ctrdutil.NamespacedContext()) ... // 初始化节点资源插件 nric, err := nri.New() ... if nric != nil { ... // 执行设备插件 if _, err := nric.InvokeWithSandbox(ctx, task, v1.Create, nriSB); err != nil { return nil, fmt.Errorf("nri invoke: %w", err) } } ... // 调用runc启动Task,此时容器正式启动 if err := task.Start(ctx); err != nil { return nil, fmt.Errorf("failed to start containerd task %q: %w", id, err) } // 同步更新容器状态,记录进程pid、启动时间 if err := cntr.Status.UpdateSync(func(status containerstore.Status) (containerstore.Status, error) { status.Pid = task.Pid() status.StartedAt = time.Now().UnixNano() return status, nil }) ... return &runtime.StartContainerResponse{}, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
# 2.7.container停止
kubelet通过gRPC调用停止init容器或业务容器,根据请求方法名步入_RuntimeService_StopContainer_Handler处理器。func _RuntimeService_StopContainer_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(StopContainerRequest) // 解码请求入参 if err := dec(in); err != nil { return nil, err } ... // 基于容器运行时调用的处理函数 handler := func(ctx context.Context, req interface{}) (interface{}, error) { return srv.(RuntimeServiceServer).StopContainer(ctx, req.(*StopContainerRequest)) } // 执行处理函数 return interceptor(ctx, in, info, handler) }1
2
3
4
5
6
7
8
9
10
11
12
13
14srv是容器运行时接口实现,这里是criService,最终请求会步入criService.StopContainer处理。// StopContainer stops a running container with a grace period (i.e., timeout). func (c *criService) StopContainer(ctx context.Context, r *runtime.StopContainerRequest) (*runtime.StopContainerResponse, error) { start := time.Now() // 获取缓存的容器 container, err := c.containerStore.Get(r.GetContainerId()) ... // 停止容器 if err := c.stopContainer(ctx, container, time.Duration(r.GetTimeout())*time.Second); err != nil { return nil, err } ... return &runtime.StopContainerResponse{}, nil } // stopContainer stops a container based on the container metadata. func (c *criService) stopContainer(ctx context.Context, container containerstore.Container, timeout time.Duration) error { ... // 获取容器关联Task,其实就是容器进程对象 task, err := container.Container.Task(ctx, nil) ... // unknown状态 if state == runtime.ContainerState_CONTAINER_UNKNOWN { ... // 异步监听容器退出 exitCh, err := task.Wait(waitCtx) ... stopCh := c.eventMonitor.startContainerExitMonitor(exitCtx, id, task.Pid(), exitCh) defer func() { ... <-stopCh }() } // 终止Task,最终的Task删除会在event处理器监听容器退出完成 if timeout > 0 { stopSignal := "SIGTERM" // 容器配了停止信号 if container.StopSignal != "" { stopSignal = container.StopSignal } else { // 镜像存储读取镜像信息 image, err := c.imageStore.Get(container.ImageRef) ... // 容器没有配停止信号,设置为镜像定义的停止信号 if image.ImageSpec.Config.StopSignal != "" { stopSignal = image.ImageSpec.Config.StopSignal } } ... if container.IsStopSignaledWithTimeout == nil { sswt = true } else { // 原子CAS操作,保证stop信号只发送一次 sswt = atomic.CompareAndSwapUint32(container.IsStopSignaledWithTimeout, 0, 1) } if sswt { // 向runc发送停止信号 if err = task.Kill(ctx, sig); err != nil && !errdefs.IsNotFound(err) { return fmt.Errorf("failed to stop container %q: %w", id, err) } } ... // 等待容器进程退出 err = c.waitContainerStop(sigTermCtx, container) ... } // 超时未退出,强制停止容器 if err = task.Kill(ctx, syscall.SIGKILL); err != nil && !errdefs.IsNotFound(err) { return fmt.Errorf("failed to kill container %q: %w", id, err) } // 等待容器退出 return c.waitContainerStop(ctx, container) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72注意
stopContainer只负责停止容器进程,containerd的事件监听器event monitor捕获容器退出事件会移除containerd的内部task实体,kubelet重建Pod会gRPC调用容器删除旧的init容器。其余情况下,pod删除后sandbox和container清理由kubelet GC调用containerd完成。
# 2.8.container移除
kubelet重启Pod或新建Pod时,会尝试通过gRPC移除旧的容器,根据请求方法名步入_RuntimeService_RemoveContainer_Handler处理器。func _RuntimeService_RemoveContainer_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(RemoveContainerRequest) // 解码请求入参 if err := dec(in); err != nil { return nil, err } ... // 基于容器运行时实现的处理函数 handler := func(ctx context.Context, req interface{}) (interface{}, error) { return srv.(RuntimeServiceServer).RemoveContainer(ctx, req.(*RemoveContainerRequest)) } // 执行处理函数 return interceptor(ctx, in, info, handler) }1
2
3
4
5
6
7
8
9
10
11
12
13
14srv是容器运行时接口实现,这里是criService,最终请求会步入criService.RemoveContainer处理。// RemoveContainer removes the container. func (c *criService) RemoveContainer(ctx context.Context, r *runtime.RemoveContainerRequest) (_ *runtime.RemoveContainerResponse, retErr error) { ... // 利用containerID获取缓存的container container, err := c.containerStore.Get(ctrID) ... // 获取容器信息 i, err := container.Container.Info(ctx) if err != nil { ... // 获取不到容器信息删除存储中的容器对象,释放name<-->containerID映射 c.containerStore.Delete(ctrID) c.containerNameIndex.ReleaseByKey(ctrID) return &runtime.RemoveContainerResponse{}, nil } ... // running或unknown状态停止container if state == runtime.ContainerState_CONTAINER_RUNNING || state == runtime.ContainerState_CONTAINER_UNKNOWN { if err := c.stopContainer(ctx, container, 0); err != nil { return nil, fmt.Errorf("failed to forcibly stop container %q: %w", id, err) } } // 设置容器状态为正在移除 if err := setContainerRemoving(container); err != nil { return nil, fmt.Errorf("failed to set removing state for container %q: %w", id, err) } ... // 删除containerd层面的容器对象及快照 if err := container.Container.Delete(ctx, containerd.WithSnapshotCleanup); err != nil { ... } // 删除cri层面的容器元数据(metadata) if err := container.Delete(); err != nil { return nil, fmt.Errorf("failed to delete container checkpoint for %q: %w", id, err) } // 删除容器相关目录 containerRootDir := c.getContainerRootDir(id) ensureRemoveAll(ctx, containerRootDir) ... volatileContainerRootDir := c.getVolatileContainerRootDir(id) ensureRemoveAll(ctx, volatileContainerRootDir) ... // 清理cri层面容器对象 c.containerStore.Delete(id) // 释放name<-->containerID映射 c.containerNameIndex.ReleaseByKey(id) return &runtime.RemoveContainerResponse{}, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
# 3.Task
# 3.1.简介
containerd调用runc创建的容器是静态的,需要通过Task启动容器。Task对象是容器运行时的核心概念之一,代表容器进程的实际执行实例,负责管理容器的生命周期和运行时状态,本质上关联一个管理容器生命周期的shim进程。// 管理容器中单个进程的实现 type process struct { id string // 进程唯一标识 task *task // 所属父task pid uint32 // 进程的实际PID io cio.IO // 进程的IO } // 管理容器主进程及附属进程实现 type task struct { client *Client // 创建runc gRPC连接的客户端 c Container // 关联的容器对象(静态配置) io cio.IO // 容器IO id string // Task的唯一标识 pid uint32 // 容器进程PID } // 容器运行时接口 type Task interface { Process // 继承Process接口(容器进程基础操作方法start/exec/wait) Pause(context.Context) error // 暂停任务 Resume(context.Context) error // 恢复任务 Exec(context.Context, string, *specs.Process, cio.Creator) (Process, error) // Task中创建新进程 Pids(context.Context) ([]ProcessInfo, error) // Task内所有进程列表 Checkpoint(context.Context, ...CheckpointTaskOpts) (Image, error) // 保存Task状态到OCI镜像 Update(context.Context, ...UpdateTaskOpts) error // 更新Task配置(资源限制、OCI镜像定义) LoadProcess(context.Context, string, cio.Attach) (Process, error) // 加载已存在的exec进程 Metrics(context.Context) (*types.Metric, error) // 运行时指标(CPU/Mem) Spec(context.Context) (*oci.Spec, error) // 获取OCI运行时配置 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30补充
容器是一个静态的配置对象,包含镜像文件系统、环境变量、挂载点等信息,不涉及运行时操作。
Task解耦容器定义和运行时,作为容器的运行时实例负责执行容器内的应用程序。
# 3.2.task创建
container是一个静态的配置对象,容器内应用程序的执行由Task负责,Task在容器创建期间创建,维护容器的静态配置及动态资源,管理容器运行的进程及子进程信息。func (c *container) NewTask(ctx context.Context, ioCreate cio.Creator...) (_ Task, err error) { // 创建容器的IO管道 i, err := ioCreate(c.id) ... // 构建CreateTaskRequest request := &tasks.CreateTaskRequest{ ContainerID: c.id, Terminal: cfg.Terminal, Stdin: cfg.Stdin, Stdout: cfg.Stdout, Stderr: cfg.Stderr, } // 获取容器运行时信息 r, err := c.get(ctx) ... // 容器有快照(挂载持久化存储),获取挂载点并添加到请求 if r.SnapshotKey != "" { ... // 获取快照挂载点 mounts, err := s.Mounts(ctx, r.SnapshotKey) ... for _, m := range mounts { ... // 根文件系统挂载加入请求 request.Rootfs = append(...) } } info := TaskInfo{ runtime: r.Runtime.Name, } // 根据传入的可选项修改Task信息(运行时路径、镜像路径) for _, o := range opts { if err := o(ctx, c.client, &info); err != nil { return nil, err } } if info.RootFS != nil { // 传入的根文件系统加入请求 for _, m := range info.RootFS { request.Rootfs = append(...) } } // 运行时路径加入请求 request.RuntimePath = info.RuntimePath ... // 初始化Task t := &task{ client: c.client, io: i, id: c.id, c: c, } ... // gRPC调用创建Task response, err := c.client.TaskService().Create(ctx, request) ... // 记录主进程PID t.pid = response.Pid return t, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60c.client.TaskService().Create会基于gRPC调用containerd/services的task服务创建任务,task创建其实就是启动shim进程,利用shim进程调用runc创建挂起的容器进程init.pid。func (c *tasksClient) Create(ctx context.Context, in *CreateTaskRequest, opts ...grpc.CallOption) (*CreateTaskResponse, error) { out := new(CreateTaskResponse) // 基于gRPC调用task服务创建任务 err := c.cc.Invoke(ctx, "/containerd.services.tasks.v1.Tasks/Create", in, out, opts...) ... return out, nil } // 最终调用到taskServer注册的gRPC服务 func (s *service) Create(ctx context.Context, r *api.CreateTaskRequest) (*api.CreateTaskResponse, error) { return s.local.Create(ctx, r) } // 创建shim进程,封装shim client,监听shim退出推送event func (l *local) Create(ctx context.Context, r *api.CreateTaskRequest, _ ...grpc.CallOption) (*api.CreateTaskResponse, error) { // 从boltdb获取容器信息 container, err := l.getContainer(ctx, r.ContainerID) ... // 返回容器恢复路径(运行中的容器压缩Tar后直接启动,用于docker、podman等客户端容器迁移) checkpointPath, err := getRestorePath(container.Runtime.Name, r.Options) ... // 通过checkpoint镜像启动的任务 if checkpointPath == "" && r.Checkpoint != nil { // 创建临时文件 checkpointPath, err = os.MkdirTemp(os.Getenv("XDG_RUNTIME_DIR"), "ctrd-checkpoint") ... // content IO reader, err := l.store.ReaderAt(...) ... // 解压OCI Tar到临时文件 _, err = archive.Apply(ctx, checkpointPath, content.NewReader(reader)) reader.Close() ... } // 构造创建请求 opts := runtime.CreateOpts{ Spec: container.Spec, IO: runtime.IO{...}, Checkpoint: checkpointPath, Runtime: container.Runtime.Name, ... } ... // 获取容器运行时的shim接口(containerd-shim-runc-v2) rtime, err := l.getRuntime(container.Runtime.Name) ... // 检测Task任务是否存在(map缓存中) _, err = rtime.Get(ctx, r.ContainerID) ... // 创建任务(启动shim进程) c, err := rtime.Create(ctx, r.ContainerID, opts) ... labels := map[string]string{"runtime": container.Runtime.Name} // 注册任务到metric collector及OOM监听器,监听任务退出 // 任务退出时,会触发m.trigger回调,通知向event服务发送事件 if err := l.monitor.Monitor(c, labels); err != nil { return nil, fmt.Errorf("monitor task: %w", err) } pid, err := c.PID(ctx) ... // 返回容器对应的PID,用于后续跟踪 return &api.CreateTaskResponse{ ContainerID: r.ContainerID, Pid: pid, }, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66rtime.Create用于启动shim进程,基于shim进程连接runc创建任务,shim进程常驻直到容器删除,相当于容器的守护进程,针对容器的操作偶会基于shim进程调用runc进行。// Create launches new shim instance and creates new task func (m *TaskManager) Create(ctx context.Context, taskID string, opts runtime.CreateOpts) (runtime.Task, error) { // 启动shim进程 process, err := m.manager.Start(ctx, taskID, opts) ... shim := process.(*shimTask) // 调用shim创建任务(shim调用runc) t, err := shim.Create(ctx, opts) ... return t, nil }1
2
3
4
5
6
7
8
9
10
11m.manager.Start用于启动containerd-shim-runc-v2进程,创建连接shim进程的ttrpc客户端,注册shimTask任务到shimManager。// 启动shimTask及注册shimTask func (m *ShimManager) Start(ctx context.Context, id string, opts runtime.CreateOpts) (_ ShimProcess, retErr error) { // 初始化bundle,记录容器信息(OCI规范的容器数据,包括容器应用的启动命令) bundle, err := NewBundle(ctx, m.root, m.state, id, opts.Spec.Value) ... // 启动shim进程 shim, err := m.startShim(ctx, bundle, id, opts) ... // 初始化shimTask对象 shimTask := &shimTask{ shim: shim, task: task.NewTaskClient(shim.client), } // 注册shimTask任务 if err := m.shims.Add(ctx, shimTask); err != nil { return nil, fmt.Errorf("failed to add task: %w", err) } return shimTask, nil } // 启动shim进程 func (m *ShimManager) startShim(ctx context.Context, bundle *Bundle, id string, opts runtime.CreateOpts) (*shim, error) { // 初始化ns ns, err := namespaces.NamespaceRequired(ctx) ... // 将容器运行时runc路径解析为shim可执行的文件路径(io.containerd.runc.v2--> /usr/local/bin/containerd-shim-runc-v2) runtimePath, err := m.resolveRuntimePath(opts.Runtime) ... // 构建shim运行对象(bundle记录管理的容器数据) b := shimBinary(bundle, shimBinaryConfig{ runtime: runtimePath, // shim向containerd上报事件/注册服务地址 address: m.containerdAddress, // shim响应containerd指令调用地址,用于快速事件通知/状态同步 ttrpcAddress: m.containerdTTRPCAddress, // 是否启用核心调度 schedCore: m.schedCore, }) // 启动shim进程注册关闭回调 shim, err := b.Start(ctx, topts, func() { ... m.shims.Delete(ctx, id) }) ... return shim, nil } // 启动shim进程注册关闭回调 func (b *binary) Start(ctx context.Context, opts *types.Any, onClose func()) (_ *shim, err error) { // 准备启动命令参数 args := []string{"-id", b.bundle.ID} ... args = append(args, "start") // 启动shim进程 cmd, err := client.Command(...) if err != nil { return nil, err } ... // 打开shim日志管道 f, err := openShimLog(shimCtx, b.bundle, client.AnonDialer) ... go func() { defer f.Close() // shim的stderr拷贝到containerd的stderr _, err := io.Copy(os.Stderr, f) ... }() // 阻塞同步shim执行结果 out, err := cmd.CombinedOutput() ... // 解析shim暴露的ttrpc socket地址 address := strings.TrimSpace(string(out)) // 连接shim服务 conn, err := client.Connect(address, client.AnonDialer) ... // 保存shim runtime路径--->/usr/bin/containerd-shim-runc-v2(containerd重启恢复runtime) = os.WriteFile(filepath.Join(b.bundle.Path, "shim-binary-path"), []byte(b.runtime), 0600) ... // 创建ttrpc客户端 client := ttrpc.NewClient(conn, ttrpc.WithOnClose(onCloseWithShimLog)) return &shim{ bundle: b.bundle, client: client, }, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86shim.Create调用shim进程,间接调用runc运行时创建新的容器任务,本质上是创建init.pid进程文件,但不会真正启动容器应用进程。func (s *shimTask) Create(ctx context.Context, opts runtime.CreateOpts) (runtime.Task, error) { ... // 创建Task的gRPC请求 request := &task.CreateTaskRequest{ ID: s.ID(), Bundle: s.bundle.Path, Stdin: opts.IO.Stdin, Stdout: opts.IO.Stdout, Stderr: opts.IO.Stderr, Terminal: opts.IO.Terminal, Checkpoint: opts.Checkpoint, Options: topts, } ... // gRPC调用shim服务 _, err := s.task.Create(ctx, request) ... return s, nil } // s.task.Create的调用处理 func (s *service) Create(ctx context.Context, r *taskAPI.CreateTaskRequest) (_ *taskAPI.CreateTaskResponse, err error) { ... // 创建由shim管理的,基于runc的container实例 container, err := runc.NewContainer(ctx, s.platform, r) if err != nil { return nil, err } // 注册容器到shim缓存(后续根据ID操作) s.containers[r.ID] = container // 通知containerd容器的TaskCreate事件 s.send(&eventstypes.TaskCreate{ ContainerID: r.ID, Bundle: r.Bundle, Rootfs: r.Rootfs, IO: &eventstypes.TaskIO{ Stdin: r.Stdin, Stdout: r.Stdout, Stderr: r.Stderr, Terminal: r.Terminal, }, Checkpoint: r.Checkpoint, Pid: uint32(container.Pid()), }) // 获取容器进程对象(init进程) proc, _ := container.Process("") // 更新任务状态 handleStarted(container, proc) return &taskAPI.CreateTaskResponse{ Pid: uint32(container.Pid()), }, nil } // NewContainer returns a new runc container func NewContainer(ctx context.Context, platform stdio.Platform, r *task.CreateTaskRequest) (_ *Container, retErr error) { ... // 构造容器创建配置 config := &process.CreateConfig{ ID: r.ID, Bundle: r.Bundle, Runtime: opts.BinaryName, Rootfs: mounts, Terminal: r.Terminal, Stdin: r.Stdin, Stdout: r.Stdout, Stderr: r.Stderr, Checkpoint: r.Checkpoint, ParentCheckpoint: r.ParentCheckpoint, Options: r.Options, } // 向bundle写入runtime名称,options配置(记录容器的runtime配置,用于恢复) WriteOptions(r.Bundle, opts) WriteRuntime(r.Bundle, opts.BinaryName) ... // 挂载点挂到rootfs for _, rm := range mounts { ... m.Mount(rootfs) } // 创建init进程对象,封装runc相关信息(启动命令) p, err := newInit( ctx, r.Bundle, filepath.Join(r.Bundle, "work"), ns, platform, config, &opts, rootfs, ) ... // 执行runc create在容器命名空间初始化init进程,但不执行CMD(init.pid文件,此时init进程挂起) if err := p.Create(ctx, config); err != nil { return nil, errdefs.ToGRPC(err) } // 组装返回container对象 container := &Container{ ID: r.ID, Bundle: r.Bundle, // init进程 process: p, // 进程映射表 processes: make(map[string]process.Process), // 预留子进程表 reservedProcess: make(map[string]struct{}), } pid := p.Pid() // init进程有效进一步加载cgroup控制器,用于后续资源限制 if pid > 0 { ... container.cgroup = cg } return container, nil } // init进程创建 func (p *Init) Create(ctx context.Context, r *CreateConfig) error { var ( err error socket *runc.Socket pio *processIO // init进程文件(<bundle>/init.pid) pidFile = newPidFile(p.Bundle) ) if r.Terminal { // 终端容器,创建console socket socket, err = runc.NewTempConsoleSocket(); err != nil { ... } else { // 创建pipe IO pio, err = createIO(ctx, p.id, p.IoUID, p.IoGID, p.stdio) ... p.io = pio } // 容器恢复流程 if r.Checkpoint != "" { return p.createCheckpointedState(r, pidFile) } ... // runc create --bundle <r.Bundle> --pid-file <pidFile> <r.ID> // 创建并挂起init进程 if err := p.runtime.Create(ctx, r.ID, r.Bundle, opts); err != nil { return p.runtimeError(err, "OCI runtime create failed") } ... // 读取容器init pid保存 pid, err := pidFile.Read() ... p.pid = pid return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153NewTask流程分析
1.创建
containerd层面的容器I/O,用于后续shim/容器进程IO转发2.生成
Task配置,gRPC调用taskManager创建容器任务3.
taskManager调用shimManager拉起shim进程,创建shimTask任务,shimTask与shim进程绑定,负责容器交互与状态通知4.
shimTask基于shim client调用shim创建容器进程,底层由shim调用runc创建挂起的容器进程5.获取
shim返回的容器进程,监听shim是否挂掉、容器进程是否退出6.返回
Task对象,代表容器进程
# 3.3.task退出
Task.wait()用于创建带缓冲管道,启动协程异步等待容器任务退出,通过管道返回任务退出状态。func (t *task) Wait(ctx context.Context) (<-chan ExitStatus, error) { c := make(chan ExitStatus, 1) go func() { ... // gRPC调用runc请求监听该Task退出,阻塞调用 r, err := t.client.TaskService().Wait(ctx, &tasks.WaitRequest{ ContainerID: t.id, }) ... // 返回Task退出状态 c <- ExitStatus{ code: r.ExitStatus, exitedAt: r.ExitedAt, } }() return c, nil } func (c *tasksClient) Wait(ctx context.Context, in *WaitRequest, opts ...grpc.CallOption) (*WaitResponse, error) { out := new(WaitResponse) // gRPC调用containerd的task服务 err := c.cc.Invoke(ctx, "/containerd.services.tasks.v1.Tasks/Wait", in, out, opts...) if err != nil { return nil, err } return out, nil } func (l *local) Wait(ctx context.Context, r *api.WaitRequest, _ ...grpc.CallOption) (*api.WaitResponse, error) { // 获取注册的shimTask t, err := l.getTask(ctx, r.ContainerID) ... // 转为Process接口 p := runtime.Process(t) if r.ExecID != "" { ... // 不是init进程,获取exec进程 p = t.Process(ctx, r.ExecID) } // 阻塞等待 exit, err := p.Wait(ctx) ... return &api.WaitResponse{ ExitStatus: exit.Status, ExitedAt: exit.Timestamp, }, nil } func (s *shimTask) Wait(ctx context.Context) (*runtime.Exit, error) { // 获取容器进程 taskPid, err := s.PID(ctx) ... // ttrpc调用wait response, err := s.task.Wait(ctx, &task.WaitRequest{ ID: s.ID(), }) ... return &runtime.Exit{ Pid: taskPid, Timestamp: response.ExitedAt, Status: response.ExitStatus, }, nil } // Wait for a process to exit func (s *service) Wait(ctx context.Context, r *taskAPI.WaitRequest) (*taskAPI.WaitResponse, error) { // 获取runc容器 container, err := s.getContainer(r.ID) ... // 获取进程对象 p, err := container.Process(r.ExecID) ... // 阻塞调用wait p.Wait() return &taskAPI.WaitResponse{ ExitStatus: uint32(p.ExitStatus()), ExitedAt: p.ExitedAt(), }, nil } // Wait for the process to exit func (p *Init) Wait() { // 退出管道有数据才结束(checkProcess会将进程退出状态更新到Init,同时关闭管道返回) <-p.waitBlock }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85补充
1.
task.wait底层会从shimTask获取runc container,进一步拿到容器进程,基于容器进程的p.wait方法内的管道实现退出监听2.监听管道会共享给事件监听器
eventMonitor,用于容器退出后更新维护的容器状态
# 3.4.task启动
c.client.TaskService().Create调用TaskManager创建shimTask,启动shim进程,创建容器init进程,维护Task与容器信息的关系,但此时容器应用并未真正启动,task.start会真正调用shim-->runc启动容器应用。func (t *task) Start(ctx context.Context) error { // gRPC调用TaskManager r, err := t.client.TaskService().Start(ctx, &tasks.StartRequest{ ContainerID: t.id, }) ... // 记录容器进程PID,这是init进程的子进程 t.pid = r.Pid return nil }1
2
3
4
5
6
7
8
9
10t.client.TaskService().Start基于gRPC调用到conatinerd TaskManager启动容器进程。func (c *tasksClient) Start(ctx context.Context, in *StartRequest, opts ...grpc.CallOption) (*StartResponse, error) { out := new(StartResponse) err := c.cc.Invoke(ctx, "/containerd.services.tasks.v1.Tasks/Start", in, out, opts...) if err != nil { return nil, err } return out, nil } // 真正处理函数 func (l *local) Start(ctx context.Context, r *api.StartRequest, _ ...grpc.CallOption) (*api.StartResponse, error) { // 获取container关联的任务句柄(shim侧的Task代理) t, err := l.getTask(ctx, r.ContainerID) ... // 强转为Process接口(代表init进程) p := runtime.Process(t) // 定位具体进程 if r.ExecID != "" { // 进一步获取exec进程 p, err = t.Process(ctx, r.ExecID) ... } // 启动进程 p.Start(ctx) ... // 获取进程状态(容器恢复、containerd重启重连shim,之前的pid可能就是旧的,这里会再更新一次) state, err := p.State(ctx) ... return &api.StartResponse{ Pid: state.Pid, }, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32l.getTask获取的任务是shimTask,p.start本质上是调用shimTask.Start启动容器进程,也就是真正将容器应用启动。// shimTask注册时,task给的是TaskClient func (s *shimTask) Start(ctx context.Context) error { // 调用启动容器任务 _, err := s.task.Start(ctx, &task.StartRequest{ ID: s.ID(), }) ... return nil } // 进程启动(ttrpc调用) func (c *taskClient) Start(ctx context.Context, req *StartRequest) (*StartResponse, error) { var resp StartResponse // ttrpc调用containerd的shim服务 if err := c.client.Call(ctx, "containerd.task.v2.Task", "Start", req, &resp); err != nil { return nil, err } return &resp, nil } // task实现 func (s *service) Start(ctx context.Context, r *taskAPI.StartRequest) (*taskAPI.StartResponse, error) { // 获取runc container container, err := s.getContainer(r.ID) ... // 启动后的回调函数 handleStarted, cleanup := s.preStart(cinit) ... // 启动容器进程 p, err := container.Start(ctx, r) ... handleStarted(container, p) return &taskAPI.StartResponse{ Pid: uint32(p.Pid()), }, nil } // Start a container process func (c *Container) Start(ctx context.Context, r *task.StartRequest) (process.Process, error) { // 根据execID获取启动的进程对象 p, err := c.Process(r.ExecID) ... // 启动进程 if err := p.Start(ctx); err != nil { return p, err } // 缓存容器的cgroups句柄 if c.Cgroup() == nil && p.Pid() > 0 { ... c.cgroup = cg } return p, nil } // Start the init process func (p *Init) Start(ctx context.Context) error { ... // 启动进程 return p.initState.Start(ctx) } // 启动进程 func (s *createdState) Start(ctx context.Context) error { s.p.start(ctx) ... return s.transition("running") } // 调用runc启动 func (p *Init) start(ctx context.Context) error { err := p.runtime.Start(ctx, p.id) return p.runtimeError(err, "OCI runtime start failed") } // Start will start an already created container func (r *Runc) Start(context context.Context, id string) error { return r.runOrError(r.command(context, "start", id)) } // 最终还是二进制调用runc,运行容器应用的启动命令(runc create,runc会启动init.pid记录的进程,该进程对应容器应用启动命令) func (r *Runc) command(context context.Context, args ...string) *exec.Cmd { command := r.Command ... cmd := exec.CommandContext(context, command, append(r.args(), args...)...) cmd.SysProcAttr = &syscall.SysProcAttr{ Setpgid: r.Setpgid, } cmd.Env = filterEnv(os.Environ(), "NOTIFY_SOCKET") if r.PdeathSignal != 0 { cmd.SysProcAttr.Pdeathsig = r.PdeathSignal } return cmd }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93补充
1.
runc create启动容器时,会从容器的bundle目录找config.json,runc会按照这份配置将容器的第一个进程拉起2.第一个进程的启动命令记录在
config.json内的spec.Process.Args列表3.容器的内存限制体现在
cgroups组,容器在对应ns下启动后cgroups会加入容器所在命令空间,以限制资源占用{ "Process": { "Args": ["/manager", "--port=8080"], "Env" : ["PATH=/usr/local/sbin:..."], ... } }1
2
3
4
5
6
7
# 3.5.task停止
容器停止时会调用
container.Container.Task()获取容器关联的task实例,通过task代理停止容器进程。// 获取containerd container关联的task实例 func (c *container) Task(ctx context.Context, attach cio.Attach) (Task, error) { return c.loadTask(ctx, attach) } func (c *container) loadTask(ctx context.Context, ioAttach cio.Attach) (Task, error) { // RPC调用TaskManager获取task实例需要容器进程数据 response, err := c.client.TaskService().Get(ctx, &tasks.GetRequest{ ContainerID: c.id, }) ... // 拷贝容器IO var i cio.IO if ioAttach != nil && response.Process.Status != tasktypes.StatusUnknown { // Do not attach IO for task in unknown state, because there // are no fifo paths anyway. attachExistingIO(response, ioAttach) ... } // 初始化task实例 t := &task{ client: c.client, io: i, id: response.Process.ID, pid: response.Process.Pid, c: c, } return t, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29c.client.TaskService().Get()是任务获取的核心函数,本质上会调用/containerd.services.tasks.v1.Tasks/Get的gRPC接口获取容器进程信息。func (c *tasksClient) Get(ctx context.Context, in *GetRequest, opts ...grpc.CallOption) (*GetResponse, error) { out := new(GetResponse) err := c.cc.Invoke(ctx, "/containerd.services.tasks.v1.Tasks/Get", in, out, opts...) if err != nil { return nil, err } return out, nil }1
2
3
4
5
6
7
8/containerd.services.tasks.v1.Tasks/Get是由containerd的io.containerd.grpc.v1.tasks插件提供的服务。// containerd\services\tasks\service.go func (s *service) Get(ctx context.Context, r *api.GetRequest) (*api.GetResponse, error) { return s.local.Get(ctx, r) }1
2
3
4插件实例最终会调用
local.Get()方法获取任务,插件注册时local实例指向io.containerd.service.v1.task-service实现。func (l *local) Get(ctx context.Context, r *api.GetRequest, _ ...grpc.CallOption) (*api.GetResponse, error) { // 获取容器关联shimTask实例 task, err := l.getTask(ctx, r.ContainerID) ... // 转为Process接口 p := runtime.Process(task) if r.ExecID != "" { // exec进程,重新获取exec进程对象 p, err = task.Process(ctx, r.ExecID) ... } // 关联进程状态 t, err := getProcessState(ctx, p) ... return &api.GetResponse{ Process: t, }, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18l.getTask()会获取boltdb存储的containerd容器实例,进一步根据containerID获取shimManager注册的shimTask。// 获取task实例 func (l *local) getTask(ctx context.Context, id string) (runtime.Task, error) { // 从boltdb加载容器数据 container, err := l.getContainer(ctx, id) ... // 基于containerID获取shimTask注册的task实例 return l.getTaskFromContainer(ctx, container) } func (l *local) getTaskFromContainer(ctx context.Context, container *containers.Container) (runtime.Task, error) { // 获取容器运行时(不同运行时存储各自的shimTask) runtime, err := l.getRuntime(container.Runtime.Name) ... // 获取创建时注册的task实例 t, err := runtime.Get(ctx, container.ID) ... return t, nil } // Get a specific task(v2) func (m *TaskManager) Get(ctx context.Context, id string) (runtime.Task, error) { // 获取shimManager.TaskList列表注册的shimTask实例 return m.manager.shims.Get(ctx, id) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24获取的容器进程对应
shimTask实例会转为更上层的进程接口,调用getProcessState()进一步包装为containerd层面的task代理返回。func getProcessState(ctx context.Context, p runtime.Process) (*task.Process, error) { ... // 调用shim获取容器状态 state, err := p.State(ctx) ... // 返回带状态的进程对象 return &task.Process{ ID: p.ID(), Pid: state.Pid, Status: status, Stdin: state.Stdin, Stdout: state.Stdout, Stderr: state.Stderr, Terminal: state.Terminal, ExitStatus: state.ExitStatus, ExitedAt: state.ExitedAt, }, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18p.State()会调用/containerd.task.v2.Task/State的ttrpc接口查询容器进程状态信息,该接口是由containerd-shim-runc-v2提供的ttrpc服务。func (s *shimTask) State(ctx context.Context) (runtime.State, error) { // taskClient调用 response, err := s.task.State(ctx, &task.StateRequest{ ID: s.ID(), }) ... return runtime.State{ Pid: response.Pid, Status: status, Stdin: response.Stdin, Stdout: response.Stdout, Stderr: response.Stderr, Terminal: response.Terminal, ExitStatus: response.ExitStatus, ExitedAt: response.ExitedAt, }, nil } func (c *taskClient) State(ctx context.Context, req *StateRequest) (*StateResponse, error) { var resp StateResponse // ttrpc服务调用 if err := c.client.Call(ctx, "containerd.task.v2.Task", "State", req, &resp); err != nil { return nil, err } return &resp, nil } // State returns runtime state information for a process func (s *service) State(ctx context.Context, r *taskAPI.StateRequest) (*taskAPI.StateResponse, error) { // 获取runc container container, err := s.getContainer(r.ID) ... // 获取进程对象 p, err := container.Process(r.ExecID) ... // 获取进程状态 st, err := p.Status(ctx) ... sio := p.Stdio() return &taskAPI.StateResponse{ ID: p.ID(), Bundle: container.Bundle, Pid: uint32(p.Pid()), Status: status, Stdin: sio.Stdin, Stdout: sio.Stdout, Stderr: sio.Stderr, Terminal: sio.Terminal, ExitStatus: uint32(p.ExitStatus()), ExitedAt: p.ExitedAt(), }, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52容器进程代理
task实例获取后,stopContainer()会调用task.kill()停止容器进程。func (t *task) Kill(ctx context.Context, s syscall.Signal, opts ...KillOpts) error { ... // gRPC调用停止容器进程 _, err := t.client.TaskService().Kill(ctx, &tasks.KillRequest{ Signal: uint32(s), ContainerID: t.id, ExecID: i.ExecID, All: i.All, }) ... return nil }1
2
3
4
5
6
7
8
9
10
11
12t.client.TaskService().Kill会调用/containerd.services.tasks.v1.Tasks/Kill的gRPC接口kill容器进程,containerd提供该服务的是io.containerd.grpc.v1.tasks插件。func (c *tasksClient) Kill(ctx context.Context, in *KillRequest, opts ...grpc.CallOption) (*types1.Empty, error) { out := new(types1.Empty) err := c.cc.Invoke(ctx, "/containerd.services.tasks.v1.Tasks/Kill", in, out, opts...) if err != nil { return nil, err } return out, nil } // io.containerd.grpc.v1.tasks插件实现 func (s *service) Kill(ctx context.Context, r *api.KillRequest) (*ptypes.Empty, error) { return s.local.Kill(ctx, r) } // service插件底层实现 func (l *local) Kill(ctx context.Context, r *api.KillRequest, _ ...grpc.CallOption) (*ptypes.Empty, error) { // 调用注册的shimTask实例 t, err := l.getTask(ctx, r.ContainerID) ... // 转为上层进程接口 p := runtime.Process(t) if r.ExecID != "" { t.Process(ctx, r.ExecID) ... } // 调用shim停止容器进程 p.Kill(ctx, r.Signal, r.All) ... return empty, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30p.Kill()会调用/containerd.task.v2.Task/Kill的ttrpc接口停止容器进程,该接口是由containerd-shim-runc-v2提供的ttrpc服务。// p.Kill()--->p.task.Kill()--->taskClient.Kill() func (c *taskClient) Kill(ctx context.Context, req *KillRequest) (*types1.Empty, error) { var resp types1.Empty if err := c.client.Call(ctx, "containerd.task.v2.Task", "Kill", req, &resp); err != nil { return nil, err } return &resp, nil } // Kill a process with the provided signal func (s *service) Kill(ctx context.Context, r *taskAPI.KillRequest) (*ptypes.Empty, error) { // 获取runc container container, err := s.getContainer(r.ID) ... // 调用runc kill if err := container.Kill(ctx, r); err != nil { return nil, errdefs.ToGRPC(err) } return empty, nil } // Kill a process func (c *Container) Kill(ctx context.Context, r *task.KillRequest) error { // 获取进程实例(init/exec) p, err := c.Process(r.ExecID) ... // 停止进程 return p.Kill(ctx, r.Signal, r.All) } // Kill the init process func (p *Init) Kill(ctx context.Context, signal uint32, all bool) error { ... return p.initState.Kill(ctx, signal, all) } // 不同进程状态都是走到init.kill func (p *Init) kill(ctx context.Context, signal uint32, all bool) error { // 调用runc kill err := p.runtime.Kill(ctx, p.id, int(signal), &runc.KillOpts{ All: all, }) return checkKillError(err) } // Kill sends the specified signal to the container func (r *Runc) Kill(context context.Context, id string, sig int, opts *KillOpts) error { args := []string{ "kill", } ... // 调用runc可执行文件执行runc kill return r.runOrError(r.command(context, append(args, id, strconv.Itoa(sig))...)) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
# 3.6.task删除
containerd启动容器进程前,调用task.wait获取exitCh启动containerExitMonitor,用于容器进程退出时回收相关资源,触发taskManager注册的shim资源回调。// startContainerExitMonitor starts an exit monitor for a given container. func (em *eventMonitor) startContainerExitMonitor(ctx context.Context, id string, pid uint32, exitCh <-chan containerd.ExitStatus) <-chan struct{} { stopCh := make(chan struct{}) go func() { defer close(stopCh) select { // task.wait()获取的exitCh case exitRes := <-exitCh: ... e := &eventtypes.TaskExit{ ContainerID: id, ... } ... err = func() error { ... // 获取缓存的cri container实例 cntr, err := em.c.containerStore.Get(e.ID) if err == nil { // 进程退出的核心清理函数 handleContainerExit(dctx, e, cntr, em.c) ... return nil } ... return nil }() ... return case <-ctx.Done(): } }() return stopCh } // handleContainerExit handles TaskExit event for container. func handleContainerExit(ctx context.Context, e *eventtypes.TaskExit, cntr containerstore.Container, c *criService) error { // gRPC调用taskManager获取task实例(shimTask代理) task, err := cntr.Container.Task(ctx,...) ... // 清理容器进程 task.Delete(ctx, WithNRISandboxDelete(cntr.SandboxID), containerd.WithProcessKill) ... // shimTask找不到,根据containerID远调清理 if errdefs.IsNotFound(err) { _, err = c.client.TaskService().Delete(ctx, &apitasks.DeleteTaskRequest{ContainerID: cntr.Container.ID()}) ... } // 更新缓存的容器状态 err = cntr.Status.UpdateSync(func(status containerstore.Status) (containerstore.Status, error) { ... return status, nil }) ... // 标识容器停止 cntr.Stop() return nil } // c.client.TaskService().Delete()--->taskManager func (l *local) Delete(ctx context.Context, r *api.DeleteTaskRequest, _ ...grpc.CallOption) (*api.DeleteResponse, error) { // 获取boltdb中container container实例 container, err := l.getContainer(ctx, r.ContainerID) ... // 回收容器进程及shim侧资源 exit, err := rtime.Delete(ctx, r.ContainerID) ... return &api.DeleteResponse{ ExitStatus: exit.Status, ExitedAt: exit.Timestamp, Pid: exit.Pid, }, nil } // Delete deletes the task and shim instance func (m *TaskManager) Delete(ctx context.Context, taskID string) (*runtime.Exit, error) { // 获取注册的shimTask item, err := m.manager.shims.Get(ctx, taskID) ... // 清理shim侧资源 exit, err := shimTask.delete(ctx, func(ctx context.Context, id string) { m.manager.shims.Delete(ctx, id) }) ... return exit, nil } func (s *shimTask) delete(ctx context.Context, removeTask func(ctx context.Context, id string)) (*runtime.Exit, error) { // 删除容器进程及附属资源 response, shimErr := s.task.Delete(ctx, &task.DeleteRequest{ ID: s.ID(), }) ... // 清理shimManager.TaskList注册的shimTask if shimErr == nil { removeTask(ctx, s.ID()) } // 同步socket、eventCh、epoll fd资源(启动时注册的回调) s.waitShutdown(ctx) ... // 回收shim进程及附属资源 s.shim.delete(ctx) ... // 清理shimManager.TaskList注册的shimTask removeTask(ctx, s.ID()) ... return &runtime.Exit{ ... Pid: response.Pid, }, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111s.task.Delete()调用shim删除容器进程、umount根文件系统,实现容器资源回收目的。// Delete the initial process and container func (s *service) Delete(ctx context.Context, r *taskAPI.DeleteRequest) (*taskAPI.DeleteResponse, error) { // 获取runc container container, err := s.getContainer(r.ID) ... // 进程删除/umount根文件系统 p, err := container.Delete(ctx, r) ... // init进程删除,清除runc conainer实例 if r.ExecID == "" { ... delete(s.containers, r.ID) ... } return &taskAPI.DeleteResponse{ ... Pid: uint32(p.Pid()), }, nil } // Delete the container or a process by id func (c *Container) Delete(ctx context.Context, r *task.DeleteRequest) (process.Process, error) { // 获取进程实例 p, err := c.Process(r.ExecID) ... // 执行回收 if err := p.Delete(ctx); err != nil { return nil, err } // exec进程,从缓存移除 if r.ExecID != "" { c.ProcessRemove(r.ExecID) } return p, nil } // init进程删除(exec删除实现仅移除pidFile) func (p *Init) delete(ctx context.Context) error { ... // 调用runc delete err := p.runtime.Delete(ctx, p.id, nil) ... // umount根文件系统 mount.UnmountAll(p.Rootfs, 0) ... return err }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47s.shim.delete用于回收shim侧资源,关闭context上下文,同步等待shim进程退出,解绑根文件系统及清除container bundle持久化数据。func (s *shim) delete(ctx context.Context) error { ... // 关闭shim client context s.Close() ... // 等待shim进程退出(shim进程回收后会关闭userCloseWaitCh,触发该方法返回) s.client.UserOnCloseWait(ctx) ... // 清理container bundle信息 s.bundle.Delete() ... return result.ErrorOrNil() } // Delete a bundle atomically func (b *Bundle) Delete() error { // 获取b.Path/work的符号链接指向的实际路径 work, werr := os.Readlink(filepath.Join(b.Path, "work")) // 容器根文件系统路径 rootfs := filepath.Join(b.Path, "rootfs") ... // 卸载容器文件系统挂载点 mount.UnmountAll(rootfs, 0) ... // 移除根文件系统 os.Remove(rootfs) ... return fmt.Errorf("failed to remove both bundle and workdir locations: %v: %w", err2, err) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29回到
container.NewTask(),task实例创建时拉起containerd-shim-runc-v2进程,注册shim进程的资源清理回调,用于容器进程关闭上游触发shim侧资源清理。func (m *ShimManager) startShim(ctx context.Context, bundle *Bundle, id string, opts runtime.CreateOpts) (*shim, error) { ... // 启动shim进程 shim, err := b.Start(ctx, topts, func() { // 退出时的资源清理 cleanupAfterDeadShim(context.Background(), id, ns, m.shims, m.events, b) // 清理shimManager.TaskList注册的shimTask实例 m.shims.Delete(ctx, id) }) ... return shim, nil } // shim进程启动 func (b *binary) Start(ctx context.Context, opts *types.Any, onClose func()) (_ *shim, err error) { ... // 退出回调 onCloseWithShimLog := func() { onClose() ... } ... // 创建shim client(ttrpc) client := ttrpc.NewClient(conn, ttrpc.WithOnClose(onCloseWithShimLog)) return &shim{ bundle: b.bundle, client: client, }, nil } func NewClient(conn net.Conn, opts ...ClientOpts) *Client { ctx, cancel := context.WithCancel(context.Background()) c := &Client{ ... // context关闭的监听管道 closed: cancel, ... // 退出清理回调 userCloseFunc: func() {}, ... } ... for _, o := range opts { o(c) } // 开启监听 go c.run() return c } func (c *Client) run() { ... // Sender goroutine // Receives calls from dispatch, adds them to the set of active calls, and sends them // to the server. go func() { var streamID uint32 = 1 for { select { case <-c.ctx.Done(): return case call := <-c.calls: id := streamID // 客户端发送的streamID总是奇数 streamID += 2 // 注册调用者到waiters if err := waiters.set(id, call); err != nil { call.errs <- err continue } // 发送请求到shim服务 if err := c.send(id, messageTypeRequest, call.req); err != nil { call.errs <- err // 发送失败移除请求 waiters.get(id) } } } }() go func() { defer close(receiverDone) for { select { // context退出 case <-c.ctx.Done(): c.setError(c.ctx.Err()) return default: // 接收shim进程上报的数据(socket) mh, p, err := c.channel.recv() ... // msg消息 msg := &message{ messageHeader: mh, p: p[:mh.Length], err: err, } // 找到调用方 call, ok, err := waiters.get(mh.StreamID) ... // 反序列化响应 call.errs <- c.recv(call.resp, msg) } } }() defer func() { // 关闭连接 c.conn.Close() // 退出清理回调,停止shim进程,回收资源,停止日志拷贝 c.userCloseFunc() // 关闭等待管道 close(c.userCloseWaitCh) }() for { select { // 接收者退出 case <-receiverDone: c.Close() // context关闭 case <-c.ctx.Done(): // 终止所有挂起的调用 waiters.abort(c.error()) return } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127eventMonitor执行回收时,会关闭shim.client.context,触发select...case...退出执行资源清理,依次为shim ttrpc conn关闭、执行回调、关闭userCloseWaitCh管道。func cleanupAfterDeadShim(ctx context.Context, id, ns string, rt *runtime.TaskList, events *exchange.Exchange, binaryCall *binary) { ... // 删除shim进程 response, err := binaryCall.Delete(ctx) ... // 获取shimManager.TaskList注册的shimTask,获取不到说明删除成功 if _, err := rt.Get(ctx, id); err != nil { return } ... } func (b *binary) Delete(ctx context.Context) (*runtime.Exit, error) { ... // 构建containerd-shim-runc-v2 delete命令 cmd, err := client.Command(ctx, &client.CommandConfig{ Runtime: b.runtime, ... Args: []string{ "-id", b.bundle.ID, "-bundle", b.bundle.Path, "delete", }, }) // 执行delete命令 cmd.Run() ... // 清理container bundle信息 b.bundle.Delete() return &runtime.Exit{ ... Pid: response.Pid, }, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35containerd-shim-runc-v2调用runc创建容器进程对象后,会调用handleStarted更新shim的内部状态,处理收到的进程退出事件,登记正在运行的进程到s.running,以补偿处理未及时登记的进程退出事件。// 处理进程状态 handleStarted = func(c *runc.Container, p process.Process) { ... // 获取刚创建的容器进程 pid = p.Pid() ... // 检查是否属于init进程 _, init := p.(*process.Init) ... if !init { // 更新待处理的exec数量 s.pendingExecs[c]-- // 获取退出的init进程 iExits, initExited := exits[c.Pid()] // init进程在exec进程前退出,handleProcessExit()暂未处理init进程退出事件 if initExited && s.pendingExecs[c] == 0 { // 清理记录的runc container退出进程 delete(s.pendingExecs, c) // 记录退出的init进程 initExits = iExits ... // 更新runc container正在运行的init进程 if len(skipped) == 0 { delete(s.running, initPid) } else { s.running[initPid] = skipped } } } // 获取当前进程的退出事件 ees, exited := exits[pid] // 清理当前容器关联的退出事件订阅 delete(s.exitSubscribers, &exits) // GC回收 exits = nil // 当前进程退出或已经没有进程 if pid == 0 || exited { ... // 处理当前进程退出事件 for _, ee := range ees { s.handleProcessExit(ee, c, p) } // 处理init进程退出事件 for _, eee := range initExits { for _, cp := range initCps { s.handleProcessExit(eee, cp.Container, cp.Process) } } } else { // 当前进程正常运行,记录 s.running[pid] = append(s.running[pid], containerProcess{ Container: c, Process: p, }) ... } } // s.mu must be locked when calling handleProcessExit func (s *service) handleProcessExit(e runcC.Exit, c *runc.Container, p process.Process) { // 当前进程是init主进程 if ip, ok := p.(*process.Init); ok { // 配置了KillAllOnExit,需要kill所有子进程 if runc.ShouldKillAllOnExit(s.context, c.Bundle) { // runc kill所有子进程 ip.KillAll(s.context) ... } } // 标记进程退出(exec进程主动退出) p.SetExited(e.Status) ... }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
注意
1.容器进程刚创建后执行
handleStarted,用于更新容器进程状态,清理提前退出进程2.
task.wait()返回exitCh,内部协程阻塞调用shim.wait()--->process--->runc.wait(),进程退出后exit状态推入exitCh通知eventMonitor执行回收3.
task.wait()阻塞调用异常时,containerExitMonitor()会获取异常事件,主动查询task及相关进程状态,处理容器退出
# 3.7.事件推送
eventMonitor是containerd用于事件监听和处理的核心组件,负责containerd的事件系统中订阅事件——exit/oom,将事件分发给订阅者。// Register CRI service plugin func init() { ... plugin.Register(&plugin.Registration{ Type: plugin.GRPCPlugin, ... // 初始化及启动cri插件服务 InitFn: initCRIService, }) } func initCRIService(ic *plugin.InitContext) (interface{}, error) { ... // 初始化cri s, err := server.NewCRIService(c, client, warn) ... go func() { // 启动cri服务 s.Run(ready) ... }() return s, nil } // NewCRIService returns a new instance of CRIService func NewCRIService(config criconfig.Config,client *containerd.Client, warn warning.Service) (CRIService,error) { ... // 初始化事件监听器 c.eventMonitor = newEventMonitor(c) ... return c, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32criService.Run()启动运行时服务,内部会基于eventMonitor订阅默认的oom/image事件,同时启动eventMonitor。// Run starts the CRI service. func (c *criService) Run(ready func()) error { // 事件订阅 c.eventMonitor.subscribe(c.client) // 启动eventMonitor eventMonitorErrCh := c.eventMonitor.start() ... return nil }1
2
3
4
5
6
7
8
9c.eventMonitor.subscribe(client)基于gRPC订阅oom/image事件,shim进程发布的事件过滤后会从server端发送回来。// subscribe starts to subscribe containerd events. func (em *eventMonitor) subscribe(subscriber events.Subscriber) { filters := []string{ `topic=="/tasks/oom"`, `topic~="/images/"`, } em.ch, em.errCh = subscriber.Subscribe(em.ctx, filters...) } // The subscriber can stop receiving events by canceling the provided context. // The errs channel will be closed and return a nil error. func (c *Client) Subscribe(ctx context.Context, filters ...string) (ch <-chan *events.Envelope, errs <-chan error) { return c.EventService().Subscribe(ctx, filters...) } func (e *eventRemote) Subscribe(ctx context.Context, filters ...string) (ch <-chan *events.Envelope, errs <-chan error) { ... // 创建eventsSubscribeClient session, err := e.client.Subscribe(ctx, &eventsapi.SubscribeRequest{ Filters: filters, }) ... go func() { ... for { // 收到事件 ev, err := session.Recv() ... select { // 推至ch管道 case evq <- &events.Envelope{ Timestamp: ev.Timestamp, Namespace: ev.Namespace, Topic: ev.Topic, Event: ev.Event, }: // context cancel case <-ctx.Done(): ... return } } }() return ch, errs } // e.client.Subscribe func (c *eventsClient) Subscribe(ctx context.Context, in *SubscribeRequest, opts ...grpc.CallOption) (Events_SubscribeClient, error) { // gRPC调用客户端 stream, err := c.cc.NewStream(ctx, &_Events_serviceDesc.Streams[0], "/containerd.services.events.v1.Events/Subscribe", opts...) ... x := &eventsSubscribeClient{stream} // 发送事件订阅请求 x.ClientStream.SendMsg(in) ... // 订阅请求后禁用send能力 x.ClientStream.CloseSend() return x, nil } // c.cc.NewStream()最终会调用containerd events服务的Subscribe()订阅 // /containerd.services.events.v1.Events/Subscribe func _Events_Subscribe_Handler(srv interface{}, stream grpc.ServerStream) error { m := new(SubscribeRequest) // 解析请求 stream.RecvMsg(m) ... return srv.(EventsServer).Subscribe(m, &eventsSubscribeServer{stream}) } // 订阅事件 func (s *service) Subscribe(req *api.SubscribeRequest, srv api.Events_SubscribeServer) error { ctx, cancel := context.WithCancel(srv.Context()) defer cancel() // 订阅事件 eventq, errq := s.events.Subscribe(ctx, req.Filters...) for { select { // 有事件产生 case ev := <-eventq: // 通过streamClient发送回消费端 srv.Send(toProto(ev)) ... } } } // NewQueue returns a queue to the provided Sink dst. func NewQueue(dst Sink) *Queue { eq := Queue{ // channel dst: dst, events: list.New(), } eq.cond = sync.NewCond(&eq.mu) // 启动queue go eq.run() return &eq } // run is the main goroutine to flush events to the target sink. func (eq *Queue) run() { for { // 获取queue中事件 event := eq.next() ... // 调用channel.Write(event)--->channel.C eq.dst.Write(event) ... } } // s.events.Subscribe(ctx, req.Filters...)注册订阅 func (e *Exchange) Subscribe(ctx context.Context, fs ...string) (ch <-chan *events.Envelope, errs <-chan error) { var ( ... channel = goevents.NewChannel(0) queue = goevents.NewQueue(channel) dst goevents.Sink = queue ) ... // 带fileter条件订阅 if len(fs) > 0 { ... dst = goevents.NewFilter(queue, goevents.MatcherFunc(func(gev goevents.Event) bool { return filter.Match(adapt(gev)) })) } // 消费端信息保存到广播表(b.adds-->b.sinks,用于发送) e.broadcaster.Add(dst) go func() { ... loop: for { select { // shim推送事件 case ev := <-channel.C: env, ok := ev.(*events.Envelope) ... select { // event推至ch case evch <- env: case <-ctx.Done(): break loop } case <-ctx.Done(): break loop } } ... }() return } // 增加订阅者 func (b *Broadcaster) Add(sink Sink) error { return b.configure(b.adds, sink) } func (b *Broadcaster) configure(ch chan configureRequest, sink Sink) error { response := make(chan error, 1) for { select { case ch <- configureRequest{ sink: sink, response: response}: // 置空,确保请求只发送一次 ch = nil ... } } } // containerd调用exchange.NewExchange()注册exchange插件 func init() { plugin.Register(&plugin.Registration{ Type: plugin.EventPlugin, ID: "exchange", InitFn: func(ic *plugin.InitContext) (interface{}, error) { // TODO: In 2.0, create exchange since ic.Events will be removed return exchange.NewExchange(), nil }, }) } // NewExchange returns a new event Exchange func NewExchange() *Exchange { return &Exchange{ // 初始化broadcaster broadcaster: goevents.NewBroadcaster(), } } // NewBroadcaster appends one or more sinks to the list of sinks. func NewBroadcaster(sinks ...Sink) *Broadcaster { b := Broadcaster{ // 订阅链接 sinks: sinks, // 事件 events: make(chan Event), adds: make(chan configureRequest), removes: make(chan configureRequest), shutdown: make(chan struct{}), closed: make(chan struct{}), } // Start the broadcaster go b.run() return &b } // run is the main broadcast loop, started when the broadcaster is created. // Under normal conditions, it waits for events on the event channel. After // Close is called, this goroutine will exit. func (b *Broadcaster) run() { ... for { select { // 收到事件 case event := <-b.events: // 按照订阅列表转发 for _, sink := range b.sinks { // queue.Write() sink.Write(event) ... } // e.broadcaster.Add(dst)的请求 case request := <-b.adds: ... // 追加到订阅列表 if !found { b.sinks = append(b.sinks, request.sink) } ... // e.broadcaster.Remove(dst) case request := <-b.removes: remove(request.sink) ... case <-b.shutdown: // close all the underlying sinks for _, sink := range b.sinks { sink.Close() ... } return } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252c.eventMonitor.start()启动eventMonitor,将streamClient接收到的订阅事件分发处理。func (em *eventMonitor) start() <-chan error { ... go func() { ... for { select { // 收到订阅事件 case e := <-em.ch: ... // 事件转换 id, evt, err := convertEvent(e.Event) ... // 事件分发 if err := em.handleEvent(evt); err != nil { // 处理失败的事件加入退避队列 em.backOff.enBackOff(id, evt) } ... case <-backOffCheckCh: // 获取所有过期的订阅queue ids := em.backOff.getExpiredIDs() for _, id := range ids { queue := em.backOff.deBackOff(id) // 处理queue中的已有事件 for i, any := range queue.events { if err := em.handleEvent(any); err != nil { // 处理失败的事件丢弃 em.backOff.reBackOff(id, queue.events[i:], queue.duration) break } } } } } }() return errCh } // handleEvent handles a containerd event. func (em *eventMonitor) handleEvent(any interface{}) error { ... switch e := any.(type) { // 容器异常退出 case *eventtypes.TaskExit: ... // 获取缓存容器实例 cntr, err := em.c.containerStore.Get(e.ID) if err == nil { ... // 容器退出事件处理 handleContainerExit(ctx, e, cntr, em.c) return nil } // sandbox pause容器退出事件 sb, err := em.c.sandboxStore.Get(e.ID) if err == nil { ... handleSandboxExit(ctx, e, sb, em.c) return nil } return nil case *eventtypes.TaskOOM: // 获取缓存容器 cntr, err := em.c.containerStore.Get(e.ContainerID) ... // 更新持久化的容器状态 err = cntr.Status.UpdateSync(func(status containerstore.Status) (containerstore.Status, error) { status.Reason = oomExitReason return status, nil }) ... // 更新镜像缓存、ID、digest case *eventtypes.ImageCreate: return em.c.updateImage(ctx, e.Name) case *eventtypes.ImageUpdate: return em.c.updateImage(ctx, e.Name) case *eventtypes.ImageDelete: return em.c.updateImage(ctx, e.Name) } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81containerd-shim-runc-v2是containerd管理容器生命周期的核心组件,shim进程启动会注册publisher插件、ttrpc服务插件及子进程退出监听。子进程退出监听依赖Linux的PR_SET_CHILD_SUBREAPER机制,shim通过SIGCHLD感知到子进程退出。func run(ctx context.Context, manager Manager, initFunc Init, name string, config Config) error { ... // 信号监听 signals, err := setupSignals(config) ... // 未禁用时,启动subreaper确保孤儿进程回收 if !config.NoSubreaper { subreaper() ... } ... // 初始化事件发布器,用于向containerd回发事件(内部开启失败事件重试) publisher, err := NewPublisher(ttrpcAddress) ... // Register event plugin plugin.Register(&plugin.Registration{ Type: plugin.EventPlugin, ID: "publisher", InitFn: func(ic *plugin.InitContext) (interface{}, error) { return publisher, nil }, }) ... // 加载插件 plugins := plugin.Graph(func(*plugin.Registration) bool { return false }) for _, p := range plugins { ... // 初始化插件 result := p.Init(initContext) ... // 获取插件实例 instance, err := result.Instance() ... // ttrpc插件 if src, ok := instance.(ttrpcService); ok { ttrpcServices = append(ttrpcServices, src) } } ... // 初始化ttrpc服务 server, err := newServer(ttrpc.WithUnaryServerInterceptor(unaryInterceptor)) ... // 注册ttrpc服务 for _, srv := range ttrpcServices { srv.RegisterTTRPC(server) ... } // 处理信号监听的退出事件(阻塞) serve(ctx, server, signals, sd.Shutdown) ... // shim服务停止,移除socket if address, err := ReadAddress("address"); err == nil { _ = RemoveSocket(address) } ... }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56setupSignals()用于设置进程的信号处理机制,shim进程订阅特定的Linux信号,保证容器生命周期的正确管理和资源清理。// setupSignals creates a new signal handler for all signals and sets the shim as a // sub-reaper so that the container processes are reparented func setupSignals(config Config) (chan os.Signal, error) { signals := make(chan os.Signal, 32) // 订阅自身程序生命周期的相关信号 smp := []os.Signal{unix.SIGTERM, unix.SIGINT, unix.SIGPIPE} // 订阅子进程状态变化信号 if !config.NoReaper { smp = append(smp, unix.SIGCHLD) } // 订阅信号绑定signals管道 signal.Notify(signals, smp...) return signals, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14serve()会启动ttrpc服务、获取signals订阅的系统信号,处理优雅退出及子进程回收。// serve serves the ttrpc API over a unix socket in the current working directory // and blocks until the context is canceled func serve(ctx context.Context, server *ttrpc.Server, signals chan os.Signal, shutdown func()) error { ... // 创建socket,用于shim与containerd通信 l, err := serveListener(socketFlag) ... go func() { defer l.Close() // 启动ttrpc监听服务 server.Serve(ctx, l) ... }() ... // 自身进程退出时的资源清理及释放 go handleExitSignals(ctx, logger, shutdown) // 监听SIGCHLD信号回收子进程 return reap(ctx, logger, signals) } func reap(ctx context.Context, logger *logrus.Entry, signals chan os.Signal) error { for { select { case <-ctx.Done(): return ctx.Err() case s := <-signals: switch s { // 子进程退出 case unix.SIGCHLD: // 子进程回收 reaper.Reap() ... case unix.SIGPIPE: } } } } // Reap should be called when the process receives an SIGCHLD. Reap will reap // all exited processes and close their wait channels func Reap() error { ... // 基于unix.Wait4()系统调用收集所有退出子进程直至没有更多退出子进程 exits, err := reap(false) for _, e := range exits { // 阻塞通知订阅者子进程退出 done := Default.notify(runc.Exit{ Timestamp: now, Pid: e.Pid, Status: e.Status, }) select { case <-done: case <-time.After(1 * time.Second): } } return err } // Default.notify func (m *Monitor) notify(e runc.Exit) chan struct{} { ... go func() { defer close(done) for { ... subscribers = m.getSubscribers() for _, s := range subscribers { s.do(func() { ... select { // 事件推入subscriber.exitCh case s.c <- e: success[s.c] = struct{}{} case <-timer.C: recv = false failed++ } }) } ... } }() return done }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85shim进程启动时会注册io.containerd.ttrpc.v1/task插件对外提供ttrpc服务,插件注册时NewTaskService作为初始化方法,内部会注册subscriber,替换runc.Monitor为全局默认Monitor,同时启动forward()转发及processExits()退出进程回收。func init() { plugin.Register(&plugin.Registration{ Type: plugin.TTRPCPlugin, ID: "task", InitFn: func(ic *plugin.InitContext) (interface{}, error) { pp, err := ic.GetByID(plugin.EventPlugin, "publisher") ... ss, err := ic.GetByID(plugin.InternalPlugin, "shutdown") ... return task.NewTaskService(ic.Context, pp.(shim.Publisher), ss.(shutdown.Service)) }, }) } // NewTaskService creates a new instance of a task service func NewTaskService(ctx context.Context, publisher shim.Publisher, sd shutdown.Service) (taskAPI.TaskService, error) { ... // OOM事件监听(oom是由内核直接kill进程,系统调用无法感知,因此单独处理) ep, err = oomv2.New(publisher) ... // 异步监听OOM事件推送到monitor go ep.Run(ctx) s := &service{ ... events: make(chan interface{}, 128), // 注册默认的subscriber,绑定插件的ec管道 ec: reaper.Default.Subscribe(), ep: ep, ... } // 子进程回收(killAll)及事件转发到events go s.processExits() // 替换runc模块的全局Monitor runcC.Monitor = reaper.Default ... // events的事件转发给订阅者 go s.forward(ctx, publisher) // 关闭事件转发回调 sd.RegisterCallback(func(context.Context) error { close(s.events) return nil }) // 移除socket回调 if address, err := shim.ReadAddress("address"); err == nil { sd.RegisterCallback(func(context.Context) error { return shim.RemoveSocket(address) }) } return s, nil } // Subscribe to process exit changes // 注册及返回事件管道 func (m *Monitor) Subscribe() chan runc.Exit { c := make(chan runc.Exit, bufferSize) m.Lock() m.subscribers[c] = &subscriber{ c: c, } m.Unlock() return c }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62shim调用go-runc运行runc命令时,会执行Monitor.start()获取事件管道,Monitor.wait()消费事件管道,用以监听runc进程的退出事件。Monitor.wait()正常触发后,会取消subscriber订阅及返回命令执行状态。// Start starts the command a registers the process with the reaper func (m *Monitor) Start(c *exec.Cmd) (chan runc.Exit, error) { // 订阅退出事件 ec := m.Subscribe() // 执行runc命令 if err := c.Start(); err != nil { m.Unsubscribe(ec) return nil, err } return ec, nil } // Wait blocks until a process is signal as dead. // User should rely on the value of the exit status to determine if the // command was successful or not. func (m *Monitor) Wait(c *exec.Cmd, ec chan runc.Exit) (int, error) { // 消费到退出事件 for e := range ec { // 退出进程是runc操作进程 if e.Pid == c.Process.Pid { // 再次获取退出进程 c.Wait() m.Unsubscribe(ec) return e.Status, nil } } // return no such process if the ec channel is closed and no more exit // events will be sent return -1, ErrNoSuchProcess }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30s.processExits()执行子进程回收,配合handleStarted处理退出事件。handleStarted()先执行会更新容器进程状态,回收退出进程资源,解注册exitSubscribers订阅,后续s.processExits()不会向handleStarted()发布事件。s.processExits()先执行,会转发退出事件,此时未执行handleStarted()的进程还未加入running列表,s.processExits()不会处理这些事件。func (s *service) processExits() { for e := range s.ec { // 退出事件转发到exitSubscriber(handleStarted)处理 for subscriber := range s.exitSubscribers { (*subscriber)[e.Pid] = append((*subscriber)[e.Pid], e) } ... // 当前进程在running列表 for _, cp := range s.running[e.Pid] { _, init := cp.Process.(*process.Init) // init进程+存在未结束子进程就跳过,交给handleStarted处理 if init && s.pendingExecs[cp.Container] != 0 { skipped = append(skipped, cp) // 加入待回收列表 } else { cps = append(cps, cp) } } // 更新running列表 if len(skipped) > 0 { s.running[e.Pid] = skipped } else { delete(s.running, e.Pid) } ... // 退出进程的回收及事件转存(s.events) for _, cp := range cps { s.handleProcessExit(e, cp.Container, cp.Process) } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31containerd调用shim.NewContainer()方法时会初始化container实例关联的cgroupManager,该cgroupManager作为后续oom事件监听工具。// NewContainer returns a new runc container func NewContainer(ctx context.Context, platform stdio.Platform, r *task.CreateTaskRequest) (_ *Container, retErr error) { ... pid := p.Pid() if pid > 0 { var cg interface{} // cgroup v2模式 if cgroups.Mode() == cgroups.Unified { // 容器进程的cgroup相对路径 g, err := cgroupsv2.PidGroupPath(pid) ... // 构造cgroupManager cg, err = cgroupsv2.LoadManager("/sys/fs/cgroup", g) ... // cgroup v1模式 } else { // 初始化cgroupManager cg, err = cgroups.Load(cgroups.V1, cgroups.PidPath(pid)) ... } // cg关联容器 container.cgroup = cg } return container, nil } // 启动容器进程时,会将container.cgroup加入oomWatcher func (s *service) Start(ctx context.Context, r *taskAPI.StartRequest) (*taskAPI.StartResponse, error) { ... switch r.ExecID { // init进程 case "": // 获取cgroup s.ep.Add(container.ID, cg) ... // 事件发送 s.send(&eventstypes.TaskStart{ ContainerID: container.ID, Pid: uint32(p.Pid()), }) default: // 事件发送(s.events--->publisher.Publish) s.send(&eventstypes.TaskExecStarted{ ContainerID: container.ID, ExecID: r.ExecID, Pid: uint32(p.Pid()), }) } // 更新容器进程状态 handleStarted(container, p) return &taskAPI.StartResponse{ Pid: uint32(p.Pid()), }, nil } // Add cgroups.Cgroup to the epoll monitor func (w *watcher) Add(id string, cgx interface{}) error { // 获取cg cg, ok := cgx.(*cgroupsv2.Manager) ... eventCh, errCh := cg.EventChan() go func() { for { i := item{id: id} select { // 事件推送 case ev := <-eventCh: i.ev = ev w.itemCh <- i case err := <-errCh: // channel is closed when cgroup gets deleted if err != nil { i.err = err w.itemCh <- i } return } } }() return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81cg.EventChan()异步监听进程相关的oomEvent,相关事件推入eventCh管道,watcher.Add()开启协程处理及并转发到watcher.itemCh。func (c *Manager) EventChan() (<-chan Event, <-chan error) { ec := make(chan Event) errCh := make(chan error, 1) // 开启监听协程 go c.waitForEvents(ec, errCh) return ec, errCh } func (c *Manager) waitForEvents(ec chan<- Event, errCh chan<- error) { defer close(errCh) // 获取cgroup监听文件描述符 fd, _, err := c.MemoryEventFD() ... for { buffer := make([]byte, syscall.SizeofInotifyEvent*10) // 阻塞读取事件(memory.enents) bytesRead, err := syscall.Read(fd, buffer) ... // 读到完整事件 if bytesRead >= syscall.SizeofInotifyEvent { out := make(map[string]interface{}) // 读取memry.events中的事件 readKVStatsFile(c.path, "memory.events", out) ... // 解析事件 e, err := parseMemoryEvents(out) ... // 将事件发出去 ec <- e // cgroup即将销毁或进程全部退出 if c.isCgroupEmpty() { return } } } } // MemoryEventFD returns inotify file descriptor and 'memory.events' inotify watch descriptor // cgroup v2的内存相关事件会记录在memory.events和cgroup.events两个文件 func (c *Manager) MemoryEventFD() (int, uint32, error) { // 创建事件监听的主文件描述符 fd, err := syscall.InotifyInit() ... // 监听memory.events的文件修改事件(memory.max/memory.high) wd, err := syscall.InotifyAddWatch(fd, fpath, unix.IN_MODIFY) ... // 监听cgroup.events文件修改事件(容器退出) syscall.InotifyAddWatch(fd, evpath, unix.IN_MODIFY) ... return fd, uint32(wd), nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52shim注册io.containerd.ttrpc.v1/task插件会启动oomWatcher,watcher.itemCh的事件通过publisher.Publish()发送到containerd监听的ttrpc socket。// Run the loop func (w *watcher) Run(ctx context.Context) { lastOOMMap := make(map[string]uint64) // key: id, value: ev.OOM for { select { case <-ctx.Done(): w.Close() return case i := <-w.itemCh: if i.err != nil { delete(lastOOMMap, i.id) continue } // 获取上一次的oom次数 lastOOM := lastOOMMap[i.id] // 当前是更新的oomKiil事件 if i.ev.OOMKill > lastOOM { // 发布事件 w.publisher.Publish(ctx, runtime.TaskOOMEventTopic, &eventstypes.TaskOOM{ ContainerID: i.id, }) ... } // 记录更新的oom事件 if i.ev.OOMKill > 0 { lastOOMMap[i.id] = i.ev.OOMKill } } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30容器退出的相关事件最终转发到
s.events,s.events中保存的事件会在s.forward()调用publisher.Publish()上报给containerd。func (s *service) forward(ctx context.Context, publisher shim.Publisher) { ... for e := range s.events { publisher.Publish(ctx, runc.GetTopic(e), e) ... } // close(s.events)后关闭publisher连接及资源 publisher.Close() }1
2
3
4
5
6
7
8
9containerd支持命令行传入ttrpc服务地址,未指定则会给定默认值。基于命令行启动shim进程时,ttrpc地址会以env形式传送到shim,供shim上报事件时初始化publisher。// NewPublisher creates a new remote events publisher func NewPublisher(address string) (*RemoteEventsPublisher, error) { // 建立ttrpc连接 client, err := ttrpcutil.NewClient(address) ... // 初始化publisher实例 l := &RemoteEventsPublisher{ client: client, closed: make(chan struct{}), requeue: make(chan *item, queueSize), } // 重试失败管道中的事件,超出最大重试删除 go l.processQueue() return l, nil } // Publish publishes the event by forwarding it to the configured ttrpc server func (l *RemoteEventsPublisher) Publish(ctx context.Context, topic string, event events.Event) error { // 获取ns ns, err := namespaces.NamespaceRequired(ctx) ... // 编码事件 any, err := typeurl.MarshalAny(event) ... // 构造投递的事件 i := &item{...} // 转发事件 if err := l.forwardRequest(i.ctx, &v1.ForwardRequest{Envelope: i.ev}); err != nil { l.queue(i) return err } return nil } // 事件上报 func (l *RemoteEventsPublisher) forwardRequest(ctx context.Context, req *v1.ForwardRequest) error { // 获取连接containerd的ttrpc客户端 service, err := l.client.EventsService() if err == nil { ... // 请求转发 _, err = service.Forward(fCtx, req) ... } ... // 重连 l.client.Reconnect() ... // 再次获取连接客户端 ... // 重新发送 _, err = service.Forward(fCtx, req) ... return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56publisher上报的事件由containerd.ttrpcService转发到已注册的exchange插件,由exchange模块调用handleEvent分发处理。func (s *ttrpcService) Forward(ctx context.Context, r *api.ForwardRequest) (*ptypes.Empty, error) { if err := s.events.Forward(ctx, fromTProto(r.Envelope)); err != nil { return nil, errdefs.ToGRPC(err) } return &ptypes.Empty{}, nil } // Forward accepts an envelope to be directly distributed on the exchange. func (e *Exchange) Forward(ctx context.Context, envelope *events.Envelope) (err error) { ... return e.broadcaster.Write(envelope) } // 事件转入broadcaster模块的events,由broadcaster的异步协程调用sink.Write()--->queue.Write()--->channel.Write()处理 func (b *Broadcaster) Write(event Event) error { select { case b.events <- event: case <-b.closed: return ErrSinkClosed } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
# 3.8.shim进程
客户端通过调用
containerd启动容器时,containerd不会直接交互底层容器运行时,而是调用io.containerd.grpc.v1.tasks插件创建containerd-shim-runc-v2进程。containerd-shim-runc-v2进程的父进程是systemd(1)。这样shim进程就脱离containerd管理,避免containerd退出影响到容器运行。后续的容器操作,都会基于shim运行runc二进制操作容器,runc运行后直接退出,shim进程就会成为容器进程的父进程,负责收集容器进程状态上报给containerd,以及pid=1的进程退出后接管容器的进程进行清理,避免出现僵尸进程。
注意
1.
containerd-shim-runc-v2进程启动后会挂在systemd(1)进程,脱离containerd进程管理2.
runc init创建init进程后,会根据OCI Container规范,将init进程替换为container启动命令进程3.
containerd-shim-runc-v2进程会执行PR_SET_CHILD_SUBREAPER系统调用,将runc start启动的init进程变为孤儿进程时接管