AutoK3s v0.9.3 公共基础模块与数据类型超深度分析
目录
- 1. 模块定位与依赖关系
- 2. 模块整体结构
- 3. 核心业务逻辑深度逐行解析
- 4. types/autok3s.go 核心类型详解
- 5. 各 Provider Options 类型定义
- 6. Mermaid 架构图
1. 模块定位与依赖关系
1.1 pkg/common — 公共基础核心包
pkg/common 是 AutoK3s 的 公共基础设施层,承担以下核心职责:
| 全局常量/变量 | common.go | 路径常量、状态常量、全局运行时变量 |
| 数据库存储 | db.go | SQLite + GORM 初始化、AutoMigrate、Schema 创建 |
| 数据模型 | model.go | ClusterState/Template/Credential/Explorer/Setting/SSHKey/Addon 等全部 ORM 模型 |
| Store CRUD | model.go | 所有实体的增删改查方法 |
| 事件广播 | broadcast.go | 发布-订阅模式的事件推送机制 |
| kubeconfig 管理 | file.go | ConfigFileManager 管理 kubeconfig 合并/删除 |
| 日志管理 | log.go | logrus 初始化、集群日志文件管理 |
| Addon 管理 | addon.go | Helm Chart 附加组件的 CRUD |
| Manifest 渲染 | manifest.go | Go template + Sprig 函数的 Helm-style 模板渲染 |
| Helm Dashboard | dashboard.go | helm-dashboard 外部命令管理 |
| Kube-Explorer | explorer.go | kube-explorer 外部命令管理 |
| Prometheus 指标 | metrics.go | 集群/模板计数器、遥测开关 |
| SSH 密钥管理 | sshkey.go | SSH 密钥对的数据库存储与查询 |
| 离线包管理 | package.go | K3s 离线安装包的状态机管理 |
| 设置持久化 | setting_provider.go | DBSettingProvider 实现 settings.Provider 接口 |
| UUID 生成 | uuid.go | 安装实例唯一标识 |
| 默认模板 | default_template.go | 各 Provider 的默认配置模板 |
| Rancher 模板 | rancher_template.go | Rancher Manager 的 HelmChart manifest 模板 |
1.2 被哪些模块依赖
pkg/common 被以下模块依赖:
- pkg/server/ — API 服务器,使用 Store 进行所有 CRUD 操作,使用 Broadcaster 进行 WebSocket 事件推送
- pkg/providers/ — 各云提供商(alibaba/aws/tencent/google/k3d/native),使用 ClusterState、DefaultDB、路径常量
- pkg/cluster/ — 集群生命周期管理,使用 SaveCluster/GetCluster/DeleteCluster
- pkg/settings/ — 设置系统,通过 DBSettingProvider 持久化到数据库
- cmd/ — CLI 命令,使用 IsCLI/Debug/CfgPath 等全局变量
- pkg/metrics/ — 指标系统,通过 common.SetupPrometheusMetrics 初始化
1.3 依赖的外部包
github.com/glebarez/sqlite — 纯 Go SQLite 驱动(无 CGO 依赖)
gorm.io/gorm — ORM 框架
github.com/sirupsen/logrus — 结构化日志
k8s.io/client-go — kubeconfig 管理
github.com/Masterminds/sprig/v3 — Go template 增强函数
helm.sh/helm/v3/pkg/strvals — Helm –set 值解析
github.com/prometheus/client_golang — Prometheus 指标
github.com/pborman/uuid — UUID 生成
github.com/rancher/wrangler/v2 — Rancher API 框架
github.com/rancher/apiserver — Rancher API Server
2. 模块整体结构
2.1 数据库层:SQLite + GORM
数据库配置
文件:db.go
// GetDB 打开并返回数据库连接
func GetDB() (*gorm.DB, error) {
dataSource := GetDataSource() // ~/.autok3s/.db/autok3s.db
config := &gorm.Config{}
if IsCLI && !Debug {
config.Logger = logger.Default.LogMode(logger.Silent) // CLI 模式静默日志
}
return gorm.Open(sqlite.Open(dataSource), config)
}
关键设计:
- 使用 glebarez/sqlite —— 纯 Go 实现,无需 CGO 编译
- 数据库路径:~/.autok3s/.db/autok3s.db
- CLI 模式下静默 GORM 日志,Server 模式或 Debug 模式下开启
- 连接池限制:db.SetMaxOpenConns(1) —— 解决 SQLite “database is locked (SQLITE_BUSY)” 问题(Issue #460)
数据库初始化流程 InitStorage()
func InitStorage(ctx context.Context) error {
// 1. 确保数据库文件存在
dataSource := GetDataSource()
if err := utils.EnsureFileExist(dataSource); return err
// 2. 创建 Store(包含 gorm.DB + Broadcaster)
store, err := NewClusterDB(ctx); return err
// 3. 执行原始 SQL Schema(兼容旧版本)
setup(store.DB)
// 4. GORM AutoMigrate(自动建表/加列)
store.DB.AutoMigrate(
&ClusterState{}, &Template{}, &Package{},
// &Credential{}, ← 注意:注释掉了!兼容 0.5.x 升级问题
&Explorer{}, &Setting{}, &SSHKey{}, &Addon{},
)
// 5. 设置全局 DefaultDB
DefaultDB = store
// 6. 初始化默认 Rancher Addon
rancherAddon := &Addon{
Name: "rancher",
Description: "Default Rancher Manager add-on",
Manifest: []byte(DefaultRancherManifest),
Values: make(types.StringMap),
}
DefaultDB.SaveAddon(rancherAddon)
// 7. 设置 Setting Provider 为 DB 持久化
return settings.SetProvider(&DBSettingProvider{})
}
双重 Schema 策略
AutoK3s 使用双重 Schema 策略确保数据库兼容性:
注意:Credential 表不参与 AutoMigrate(代码注释说明:从 0.5.x 升级时 AutoMigrate 会导致建表错误),仅通过原始 SQL Schema 创建。
原始 SQL Schema 详解
cluster_states 表:
CREATE TABLE IF NOT EXISTS cluster_states (
name TEXT not null, — 集群名称
provider TEXT not null, — 提供商名称
token TEXT, — K3s token
ip TEXT, — 主节点 IP
tls_sans TEXT, — TLS SAN 扩展
cluster_cidr TEXT, — 集群 CIDR
master_extra_args TEXT, — Master 额外参数
worker_extra_args TEXT, — Worker 额外参数
registry TEXT, — 镜像仓库配置
registry_content TEXT, — 镜像仓库内容
data_store TEXT, — 外部数据存储
k3s_version TEXT, — K3s 版本
k3s_channel TEXT, — K3s channel
install_script TEXT, — 安装脚本
mirror TEXT, — 安装镜像
docker_mirror TEXT, — Docker 镜像
docker_script TEXT, — Docker 脚本
network TEXT, — 网络配置
ui bool, — UI 开关(已废弃)
cluster bool, — 集群模式
options BLOB, — Provider 特定选项(JSON)
status TEXT, — 集群状态
master_nodes BLOB, — Master 节点列表(JSON)
worker_nodes BLOB, — Worker 节点列表(JSON)
context_name TEXT, — kubeconfig context 名
master TEXT, — Master 数量
worker TEXT, — Worker 数量
ssh_port TEXT, — SSH 端口
ssh_user TEXT, — SSH 用户
ssh_password TEXT, — SSH 密码
ssh_key_path TEXT, — SSH 密钥路径
ssh_cert TEXT, — SSH 证书
ssh_cert_path TEXT, — SSH 证书路径
ssh_key_passphrase TEXT, — SSH 密钥密码
ssh_agent_auth bool, — SSH Agent 认证
manifests TEXT, — Manifest 文件
enable TEXT, — 启用的组件
standalone bool, — 独立模式
unique (name, provider) — 联合唯一约束
);
templates 表:与 cluster_states 几乎相同,额外增加 is_default bool 字段,无 status/master_nodes/worker_nodes/standalone 字段。
credentials 表:
CREATE TABLE IF NOT EXISTS credentials (
id integer not null primary key autoincrement,
provider TEXT not null,
secrets BLOB
);
2.2 数据模型全景
ClusterState — 集群状态模型
type ClusterState struct {
types.Metadata `json:",inline" mapstructure:",squash" gorm:"embedded"`
Options []byte `json:"options,omitempty" gorm:"type:bytes"`
Status string `json:"status" yaml:"status"`
Standalone bool `json:"standalone" yaml:"standalone" gorm:"type:bool"`
MasterNodes []byte `json:"master-nodes,omitempty" gorm:"type:bytes"`
WorkerNodes []byte `json:"worker-nodes,omitempty" gorm:"type:bytes"`
types.SSH `json:",inline" mapstructure:",squash" gorm:"embedded"`
}
设计要点:
- 嵌入 types.Metadata(gorm:“embedded”)—— 将 Metadata 的所有字段平铺到 cluster_states 表
- 嵌入 types.SSH(gorm:“embedded”)—— SSH 配置同样平铺
- Options 字段:[]byte 存储 Provider 特定选项的 JSON 序列化
- MasterNodes/WorkerNodes:[]byte 存储 Node 列表的 JSON 序列化
- Standalone:标记是否为独立集群(非云端管理)
实现了 ISchemaObject 接口:
- SchemaID() → 返回 apis.Cluster{} 的 schema ID
- ToAPIObject() → 转换为 API 响应对象(调用 ConvertToCluster)
Template — 集群模板模型
type Template struct {
types.Metadata `json:",inline" mapstructure:",squash" gorm:"embedded"`
Options []byte `json:"options,omitempty" gorm:"type:bytes"`
types.SSH `json:",inline" mapstructure:",squash" gorm:"embedded"`
IsDefault bool `json:"is-default" gorm:"type:bool"`
}
与 ClusterState 对比:
- 无 Status/MasterNodes/WorkerNodes/Standalone —— 模板是静态配置
- 有 IsDefault —— 标记默认模板
Credential — 凭证模型
type Credential struct {
ID int `json:"id" gorm:"type:integer;primaryKey;not null;autoIncrement"`
Provider string `json:"provider" gorm:"not null"`
Secrets []byte `json:"secrets,omitempty" gorm:"type:bytes"`
}
- 自增主键 ID(不使用 GORM 默认的 ID 字段名)
- Secrets 存储为 []byte(JSON 序列化的 map[string]string)
- 每个 Provider 仅支持一个凭证(代码中 CreateCredential 会覆盖已有凭证)
Explorer — Kube-Explorer 配置模型
type Explorer struct {
ContextName string `json:"context-name" gorm:"primaryKey;not null"`
Enabled bool `json:"enabled" gorm:"type:bool"`
Port int `json:"port"`
}
- 主键为 ContextName(集群 context 名)
- Port 存储 kube-explorer 监听端口
Setting — 键值设置模型
type Setting struct {
Name string `json:"name" gorm:"primaryKey;not null"`
Value string `json:"value"`
}
- 简单的键值对存储
- 通过 DBSettingProvider 实现 settings.Provider 接口
SSHKey — SSH 密钥模型
type SSHKey struct {
Name string `json:"name" gorm:"primaryKey;not null" wrangler:"required,noupdate"`
GenerateKey bool `json:"generate-key,omitempty" gorm:"-:all" wrangler:"writeOnly,noupdate"`
HasPassword bool `json:"has-password" wrangler:"nocreate,noupdate"`
SSHPassphrase string `json:"ssh-passphrase,omitempty" gorm:"-:all" wrangler:"type=password,nullable"`
Bits int `json:"bits,omitempty" gorm:"-:all" wrangler:"default=2048,nullable"`
SSHCert string `json:"ssh-cert,omitempty" yaml:"ssh-cert,omitempty" wrangler:"nullable"`
SSHKey string `json:"ssh-key,omitempty" yaml:"ssh-key,omitempty" wrangler:"writeOnly,nullable"`
SSHPublicKey string `json:"ssh-key-public,omitempty" yaml:"ssh-key-public,omitempty" wrangler:"nullable"`
}
设计要点:
- gorm:"-:all" 标记的字段不持久化到数据库(GenerateKey/SSHPassphrase/Bits)—— 这些是输入参数,仅在创建时使用
- 持久化字段:Name、HasPassword、SSHCert、SSHKey、SSHPublicKey
- wrangler 标签控制 API 层的读写权限
Addon — 附加组件模型
type Addon struct {
Name string `json:"name" gorm:"primaryKey;not null" wrangler:"required,noupdate"`
Description string `json:"description,omitempty"`
Manifest []byte `json:"manifest" gorm:"type:bytes" wrangler:"required"`
Values types.StringMap `json:"values,omitempty" gorm:"type:stringMap"`
}
- Name 为主键,不可更新(wrangler:"required,noupdate")
- Manifest 存储为 []byte(Helm Chart YAML 模板)
- Values 使用自定义 StringMap 类型(JSON 序列化的 map)
Package — 离线包模型
type Package struct {
Name string `json:"name,omitempty" gorm:"primaryKey;->;<-:create" wrangler:"required,noupdate"`
K3sVersion string `json:"k3sVersion,omitempty" wrangler:"required"`
Archs types.StringArray `json:"archs,omitempty" gorm:"type:text" wrangler:"required"`
FilePath string `json:"filePath,omitempty" wrangler:"nocreate,noupdate"`
State State `json:"state,omitempty" wrangler:"nocreate,noupdate"`
}
- gorm:"primaryKey;->;<-:create" —— 主键,只读(创建后不可更新)
- State 为自定义类型,状态机:Active/OutOfSync/Downloading/Verifying/Validating
event — 内部事件结构(非持久化)
type event struct {
Name string // 事件名:CreateAPIEvent/ChangeAPIEvent/RemoveAPIEvent
Object interface{} // 事件对象
}
type LogEvent struct {
Name string
ContextType string
ContextName string
}
2.3 Store 层:gorm.DB + Broadcaster 的 CRUD 模式
Store 结构
type Store struct {
*gorm.DB // 嵌入 GORM DB
broadcaster *Broadcaster // 事件广播器(私有)
}
Store 是 GORM DB + 事件广播器的组合体,所有 CRUD 操作通过 Store 方法执行,并在操作完成后广播事件。
接口定义
type IIDObject interface {
GetID() string
}
type ISchemaObject interface {
IIDObject
SchemaID() string
ToAPIObject() *apitypes.APIObject
}
- IIDObject:所有持久化模型需实现 GetID() 返回字符串 ID
- ISchemaObject:扩展接口,支持转换为 Rancher API Server 的 APIObject
CRUD 模式统一分析
Cluster CRUD:
- SaveCluster(cluster *types.Cluster) — 查找已存在 → 存在则 Update(Omit name/provider/context_name)→ 不存在则 Create → 更新 Prometheus 指标
- SaveClusterState(state *ClusterState) — 直接 Update(Omit name/provider)
- DeleteCluster(name, provider) — 先 Get(用于广播事件和指标)→ Delete → 广播 RemoveAPIEvent → 递减指标
- ListCluster(provider) — 按 provider 过滤(可选)→ 返回列表
- GetCluster(name, provider) — 按 name+provider 精确查询
- GetClusterByID(contextName) — 按 context_name 查询
- FindCluster(name, provider) — 模糊查询(仅 name 必填)
Template CRUD:
- CreateTemplate(template) — Create → 递增 TemplateCount 指标
- UpdateTemplate(template) — Update(Omit name/provider/context_name)
- DeleteTemplate(name, provider) — 先 Get → Delete → 广播 RemoveAPIEvent → 递减指标
- ListTemplates() — 全量查询
- GetTemplate(name, provider) — 精确查询
Credential CRUD:
- CreateCredential(cred) — 先检查同 provider 是否已有凭证 → 有则 Update(覆盖)→ 无则 Create
- UpdateCredential(cred) — 按 ID Update(Omit id/provider)
- ListCredential() — 全量查询
- GetCredentialByProvider(provider) — 按 provider 查询
- GetCredential(id) — 按 ID 查询
- DeleteCredential(id) — 按 ID 删除
Explorer CRUD:
- SaveExplorer(exp) — 先 Get → 存在则 Update(Omit context_name)→ 不存在则 Create
- GetExplorer(clusterID) — 按 context_name 查询
- DeleteExplorer(clusterID) — 按 context_name 删除
- ListExplorer() — 全量查询
Setting CRUD:
- SaveSetting(s) — 先 Get → 存在则 Update(Omit name)→ 不存在则 Create
- GetSetting(name) — 按 name 查询
- ListSettings() — 全量查询
SSHKey CRUD:
- SaveSSHKey(sshkey) — 先检查存在 → 存在则 Update(Omit name)→ 不存在则 Create
- ListSSHKey(name) — name 为 nil 时全量查询,否则按 name 查询单个
- DeleteSSHKey(name) — 先查询 → 删除
- SSHKeyExists(name) — 返回 bool
Addon CRUD:
- SaveAddon(addon) — 先 Get → 存在则 Update(Omit name)→ 不存在则 Create
- GetAddon(name) — 按 name 查询
- ListAddon() — 全量查询
- DeleteAddon(name) — 先 Get → Delete → 广播 RemoveAPIEvent
Package CRUD:
- ListPackages(name) — name 为 nil 时全量查询,否则按 name 查询单个
- SavePackage(pkg) — 直接 db.Save()(upsert 语义)
- DeletePackage(name) — 先查询 → 删除
- PackageExists(name) — 返回 error
2.4 事件广播:Broadcaster 机制
Broadcaster 架构
type Broadcaster struct {
subs map[Subscriber]subscriberFunc // 订阅者 map:channel → 过滤函数
m sync.RWMutex // 读写锁保护 subs
}
type Subscriber chan interface{} // 订阅者就是一个 channel
type subscriberFunc func(v interface{}) bool // 过滤函数:返回 true 则传递事件
核心方法
Register(sf subscriberFunc) Subscriber — 注册订阅者
- 创建新 channel
- 加锁,将 channel+过滤函数存入 map
- 返回 channel 供消费
Evict(s Subscriber) — 驱逐订阅者
- 加锁,从 map 删除
- close(channel) —— 通知消费者退出
Broadcast(v interface{}) — 广播事件
- 加锁(注意:全局写锁)
- 对每个订阅者启动 goroutine 并行发布
- 使用 WaitGroup 等待所有发布完成
- 发布前调用过滤函数,不匹配则跳过
- 使用 select { case s <- v: } 非阻塞发送(注意:如果 channel 满了会丢弃事件)
Close() — 关闭所有订阅者
- 加锁,遍历所有订阅者
- 逐个 delete + close
Store 中的 GORM Hook 集成
func (d *Store) Register() {
d.DB.Callback().Create().After("gorm:create").Register("gorm:autok3s_create", d.createHandler)
d.DB.Callback().Update().After("gorm:update").Register("gorm:autok3s_update", d.updateHandler)
}
- Create 后 → createHandler → 广播 CreateAPIEvent
- Update 后 → updateHandler → 广播 ChangeAPIEvent
- Delete → 手动广播 RemoveAPIEvent(在 Delete 方法中显式调用)
Watch 机制 — API 事件流
func (d *Store) Watch(apiOp *apitypes.APIRequest, schema *apitypes.APISchema) chan apitypes.APIEvent {
result := make(chan apitypes.APIEvent)
// 注册订阅者,过滤条件为 obj.Type == schema.ID
sub := d.broadcaster.Register(func(v interface{}) bool { … })
go func() {
for {
select {
case v, ok := <-sub:
// 转换并转发事件
result <- getAPIEvent(e, schema)
case <-apiOp.Context().Done():
// 请求上下文取消时,驱逐订阅者并关闭 result channel
d.broadcaster.Evict(sub)
close(result)
return
}
}
}()
return result
}
Log 订阅机制
func (d *Store) Log(apiOp *apitypes.APIRequest, t string, input chan *LogEvent) {
sub := d.broadcaster.Register(func(v interface{}) bool {
event, ok := v.(*LogEvent)
return ok && t == event.ContextType
})
for {
select {
case v, ok := <-sub:
input <- state // 转发 LogEvent
case <-apiOp.Context().Done():
d.broadcaster.Evict(sub)
return
}
}
}
3. 核心业务逻辑深度逐行解析
3.1 InitStorage() — 数据库初始化
文件:db.go 第 49-83 行
func InitStorage(ctx context.Context) error {
// 步骤1: 获取数据库路径 ~/.autok3s/.db/autok3s.db
dataSource := GetDataSource()
// 步骤2: 确保文件存在(touch 操作)
if err := utils.EnsureFileExist(dataSource); err != nil {
return err
}
// 步骤3: 创建 Store 实例
// NewClusterDB 内部:
// – gorm.Open(sqlite.Open(dataSource), config) 打开数据库
// – db.SetMaxOpenConns(1) 解决 SQLite 并发锁问题
// – 创建 Broadcaster 实例
store, err := NewClusterDB(ctx)
if err != nil {
return err
}
// 步骤4: 执行原始 SQL Schema(CREATE TABLE IF NOT EXISTS)
// 对 cluster_states、templates、credentials 三张表
setup(store.DB)
// 步骤5: GORM AutoMigrate
// 自动为 ClusterState/Template/Package/Explorer/Setting/SSHKey/Addon 创建/更新表结构
// Credential 被注释掉 —— 兼容 0.5.x 升级
if err := store.DB.AutoMigrate(
&ClusterState{}, &Template{}, &Package{},
// &Credential{},
&Explorer{}, &Setting{}, &SSHKey{}, &Addon{},
); err != nil {
return err
}
// 步骤6: 设置全局 DefaultDB
DefaultDB = store
// 步骤7: 初始化默认 Rancher Addon
// 检查是否已存在 "rancher" addon
// 不存在则创建默认 addon,包含 DefaultRancherManifest 模板
_, err = DefaultDB.GetAddon("rancher")
if err != nil && err == gorm.ErrRecordNotFound {
rancherAddon := &Addon{
Name: "rancher",
Description: "Default Rancher Manager add-on",
Manifest: []byte(DefaultRancherManifest),
Values: make(types.StringMap),
}
err = DefaultDB.SaveAddon(rancherAddon)
if err != nil {
logrus.Errorf("failed to save default rancher manager add-on template: %v", err)
}
}
// 步骤8: 设置全局 Setting Provider
// DBSettingProvider 实现了 settings.Provider 接口
// 后续所有 settings.Setting 的 Get/Set 操作都会持久化到数据库
return settings.SetProvider(&DBSettingProvider{})
}
3.2 Store CRUD 逐行解析 — 以 SaveCluster 为例
func (d *Store) SaveCluster(cluster *types.Cluster) error {
// 1. 查找已存在的集群状态
state := &ClusterState{}
result := d.DB.Where("name = ? AND provider = ?", cluster.Name, cluster.Provider).Find(state)
// 2. 序列化 Options 为 JSON
opt, err := json.Marshal(cluster.Options)
if err != nil {
return err
}
// 3. 序列化 MasterNodes 和 WorkerNodes 为 JSON
masterNodeBytes, err := json.Marshal(cluster.Status.MasterNodes)
if err != nil {
return err
}
workerNodeBytes, err := json.Marshal(cluster.Status.WorkerNodes)
if err != nil {
return err
}
// 4. 构建新的 ClusterState
state = &ClusterState{
Metadata: cluster.Metadata, // 嵌入 Metadata
Options: opt, // Provider 选项 JSON
Status: cluster.Status.Status,// 状态字符串
MasterNodes: masterNodeBytes, // Master 节点 JSON
WorkerNodes: workerNodeBytes, // Worker 节点 JSON
SSH: cluster.SSH, // 嵌入 SSH 配置
Standalone: cluster.Status.Standalone,
}
// 5. 判断是 Create 还是 Update
if result.RowsAffected == 0 {
// Create:插入新记录
result = d.DB.Create(state)
if result.Error == nil {
// 成功创建后递增 Prometheus 集群计数器
metrics.ClusterCount.With(getLabelsFromMeta(state.Metadata)).Inc()
}
return result.Error
}
// Update:更新已有记录
// Omit("name", "provider", "context_name") —— 不更新这三个字段
// 这确保集群标识不会被意外修改
result = d.DB.Model(state).
Where("name = ? AND provider = ?", cluster.Name, cluster.Provider).
Omit("name", "provider", "context_name").Save(state)
return result.Error
}
3.3 DeleteCluster — 删除与事件广播
func (d *Store) DeleteCluster(name, provider string) error {
// 1. 先获取集群状态(用于广播事件和更新指标)
state, err := d.GetCluster(name, provider)
if err != nil {
return err
}
if state == nil {
return nil // 不存在则直接返回
}
// 2. 执行删除
result := d.DB.Where("name = ? AND provider = ?", name, provider).Delete(&ClusterState{})
// 3. 广播 RemoveAPIEvent 事件
// 即使 DB 删除成功但后续代码也会执行
d.broadcaster.Broadcast(&event{
Name: apitypes.RemoveAPIEvent,
Object: GetAPIObject(state),
})
// 4. 递减 Prometheus 集群计数器
if result.Error == nil {
metrics.ClusterCount.With(getLabelsFromMeta(state.Metadata)).Dec()
}
return result.Error
}
3.4 ConvertToCluster — 状态转换核心函数
func ConvertToCluster(state *ClusterState, nodeInfo bool) types.Cluster {
// 1. 基础转换:从 ClusterState 提取 Metadata、SSH、Status
c := types.Cluster{
Metadata: state.Metadata,
SSH: state.SSH,
Status: types.Status{
Status: state.Status,
},
}
// 2. 获取 Provider 实例,反序列化 Options
p, err := providers.GetProvider(state.Provider)
if err != nil {
logrus.Errorf("failed to get provider by name %s", state.Provider)
return c // 返回基础信息,Options 为空
}
// 3. 调用 Provider 的 GetProviderOptions 反序列化
// 每个 Provider 有自己的 Options 类型(如 alibaba.Options)
opt, err := p.GetProviderOptions(state.Options)
if err != nil {
logrus.Errorf("failed to convert [%s] provider options %s: %v", …)
return c
}
c.Options = opt
// 4. 可选:反序列化节点信息
if nodeInfo {
masterNodes := make([]types.Node, 0)
json.Unmarshal(state.MasterNodes, &masterNodes) // JSON → []Node
workerNodes := make([]types.Node, 0)
json.Unmarshal(state.WorkerNodes, &workerNodes)
c.MasterNodes = masterNodes
c.WorkerNodes = workerNodes
}
return c
}
3.5 ConfigFileManager — kubeconfig 管理
文件:file.go
type ConfigFileManager struct {
mutex sync.RWMutex // 读写锁保护并发 kubeconfig 操作
}
SaveCfg — 保存 kubeconfig
func (c *ConfigFileManager) SaveCfg(context, tempFile string) error {
// 1. 设置 KUBECONFIG 环境变量为主配置文件
kubeConfigPath := filepath.Join(CfgPath, KubeCfgFile) // ~/.autok3s/.kube/config
os.Setenv(clientcmd.RecommendedConfigPathEnvVar, kubeConfigPath)
// 2. 先移除同名的旧 context
err := c.OverwriteCfg(kubeConfigPath, context, c.RemoveCfg)
// 3. 设置 KUBECONFIG 环境变量为合并模式
// KUBECONFIG=~/.autok3s/.kube/config:tempFile
mergeKubeConfigENV := fmt.Sprintf("%s:%s", kubeConfigPath, tempFile)
// Windows 用分号分隔
if runtime.GOOS == "windows" {
mergeKubeConfigENV = fmt.Sprintf("%s;%s", kubeConfigPath, tempFile)
}
os.Setenv(clientcmd.RecommendedConfigPathEnvVar, mergeKubeConfigENV)
// 4. 执行合并并写回主配置文件
return c.OverwriteCfg(filepath.Join(CfgPath, KubeCfgFile), context, c.MergeCfg)
}
RemoveCfg — 移除 context
func (c *ConfigFileManager) RemoveCfg(context string, configAccess clientcmd.ConfigAccess) (*api.Config, error) {
// 1. 获取当前配置
config, err := configAccess.GetStartingConfig()
// 2. 如果当前 context 就是目标 context,清空 CurrentContext
if config.CurrentContext == context {
config.CurrentContext = ""
}
// 3. 删除 Context、Cluster、AuthInfo 条目
delete(config.Contexts, context)
delete(config.Clusters, context)
delete(config.AuthInfos, context)
// 4. 删除关联的 AuthInfo 条目(格式为 user@context)
for key := range config.AuthInfos {
if strings.Contains(key, fmt.Sprintf("@%s", context)) {
delete(config.AuthInfos, key)
}
}
return config, nil
}
MergeCfg — 合并 kubeconfig
func (c *ConfigFileManager) MergeCfg(context string, configAccess clientcmd.ConfigAccess) (*api.Config, error) {
// 1. 加载旧配置
oldConfig, err := clientcmd.LoadFromFile(…)
// 2. 检查 context 是否已存在 —— 存在则报错
if _, ok := oldConfig.Contexts[context]; ok {
return nil, fmt.Errorf("context %s is already exist in kubeconfig file…")
}
// 3. 获取合并后的配置(通过 KUBECONFIG 环境变量)
config, err := configAccess.GetStartingConfig()
return config, err
}
OverwriteCfg — 加锁写入
func (c *ConfigFileManager) OverwriteCfg(path string, context string, cfg func(…) (*api.Config, error)) error {
c.mutex.Lock() // 写锁
defer c.mutex.Unlock()
paOpt := clientcmd.NewDefaultPathOptions()
config, err := cfg(context, paOpt) // 执行 RemoveCfg 或 MergeCfg
if err != nil {
return err
}
return clientcmd.WriteToFile(*config, path) // 写回文件
}
3.6 Manifest 渲染逻辑
文件:manifest.go
GenerateValues — 值合并
func GenerateValues(setValues map[string]string, defaultValues map[string]string) (map[string]interface{}, error) {
values := []string{}
// 1. 先添加默认值(如果 setValues 中没有对应的 key)
for key, value := range defaultValues {
if _, ok := setValues[key]; !ok {
values = append(values, fmt.Sprintf("%s=%s", key, value))
}
}
// 2. 再添加用户设置的值(覆盖默认值)
for key, value := range setValues {
values = append(values, fmt.Sprintf("%s=%s", key, value))
}
// 3. 使用 Helm 的 strvals 解析为嵌套 map
return mergeValues(values)
}
AssembleManifest — 模板渲染
func AssembleManifest(values map[string]interface{}, manifest string, templateFunc template.FuncMap) ([]byte, error) {
// 1. 创建 Go template,注入 Sprig 函数(Helm 同款)
t := template.New("manifest").Funcs(sprig.TxtFuncMap())
// 2. 注入自定义模板函数(如 providerTemplate)
if templateFunc != nil {
t = t.Funcs(templateFunc)
}
// 3. 解析 manifest 模板字符串
t, err := t.Parse(manifest)
// 4. 执行模板渲染,注入 values
var resultContent bytes.Buffer
err = t.Execute(&resultContent, values)
return resultContent.Bytes(), nil
}
mergeValues — Helm 风格值解析
func mergeValues(values []string) (map[string]interface{}, error) {
base := map[string]interface{}{}
for _, value := range values {
// 使用 helm.sh/helm/v3/pkg/strvals.ParseInto
// 支持 key=value 和 key.subkey=value 格式
if err := strvals.ParseInto(value, base); err != nil {
return nil, errors.Wrap(err, "failed parsing –set data")
}
}
return base, nil
}
ValidateName — Addon 名称校验
func ValidateName(name string) error {
// 1. 非空检查
if name == "" {
return errors.New("name is required for addon creation")
}
// 2. 唯一性检查(查数据库)
if _, err := DefaultDB.GetAddon(name); err == nil {
return fmt.Errorf("addon %s is already exist", name)
}
// 3. DNS-1123 子域名校验(小写字母数字连字符)
if errs := validation.IsDNS1123Subdomain(name); len(errs) > 0 {
return fmt.Errorf("name is not validated %s, %v", name, errs)
}
return nil
}
3.7 Dashboard 部署逻辑
文件:dashboard.go
SwitchDashboard — 开关 Helm Dashboard
func SwitchDashboard(ctx context.Context, enabled string) error {
// 1. 启用前检查:集群列表不能为空
if enabled == "true" {
clusters, err := DefaultDB.ListCluster("")
if len(clusters) == 0 {
return errors.New("cannot enable helm-dashboard with empty cluster list")
}
}
// 2. 检查 helm-dashboard 命令是否存在
if err := CheckCommandExist(HelmDashboardCommand); err != nil {
return err
}
// 3. 验证命令可执行
if err := checkDashboardCmd(); err != nil {
return err
}
// 4. 检查当前状态是否需要变更
isEnabled := settings.HelmDashboardEnabled.Get()
if !strings.EqualFold(isEnabled, enabled) {
// 5. 更新设置
settings.HelmDashboardEnabled.Set(enabled)
if enabled == "true" {
enableDashboard(ctx) // 启动
} else {
settings.HelmDashboardEnabled.Set("false")
DashboardCanceled() // 调用 cancel func 停止
}
}
return nil
}
enableDashboard — 启动 Dashboard
func enableDashboard(ctx context.Context) {
// 1. 获取集群列表
clusters, err := DefaultDB.ListCluster("")
// 2. 空集群列表时自动禁用
if len(clusters) == 0 {
settings.HelmDashboardEnabled.Set("false")
return
}
// 3. 获取端口配置,未配置则自动获取空闲端口
dashboardPort := settings.HelmDashboardPort.Get()
if dashboardPort == "" {
freePort, err := k3dutil.GetFreePort() // 使用 k3d 工具获取空闲端口
dashboardPort = strconv.Itoa(freePort)
settings.HelmDashboardPort.Set(dashboardPort) // 持久化端口
}
// 4. 启动 helm-dashboard 进程
dashboardCtx, cancel := context.WithCancel(ctx)
DashboardCanceled = cancel // 保存 cancel 函数供后续停止使用
go func(ctx context.Context, port string) {
StartHelmDashboard(ctx, port)
}(dashboardCtx, dashboardPort)
}
StartHelmDashboard — 执行外部命令
func StartHelmDashboard(ctx context.Context, port string) error {
// 1. 设置 KUBECONFIG 环境变量
os.Setenv(clientcmd.RecommendedConfigPathEnvVar, filepath.Join(CfgPath, KubeCfgFile))
// 2. 支持自定义绑定地址(环境变量 AUTOK3S_HELM_DASHBOARD_ADDRESS)
if os.Getenv("AUTOK3S_HELM_DASHBOARD_ADDRESS") != "" {
dashboardBindAddress = os.Getenv("AUTOK3S_HELM_DASHBOARD_ADDRESS")
}
// 3. 构建命令:helm-dashboard –bind=127.0.0.1 -b –port=XXXX
dashboard := exec.CommandContext(ctx, HelmDashboardCommand,
fmt.Sprintf("–bind=%s", dashboardBindAddress),
"-b",
fmt.Sprintf("–port=%s", port))
// 4. 重定向输出到标准输出/错误
dashboard.Stdout = os.Stdout
dashboard.Stderr = os.Stderr
// 5. 启动并等待
dashboard.Start()
return dashboard.Wait()
}
3.8 Kube-Explorer 部署逻辑
文件:explorer.go
EnableExplorer — 启动指定集群的 Explorer
func EnableExplorer(ctx context.Context, config string) (int, error) {
// 1. 检查是否已启动
if _, ok := ExplorerWatchers[config]; ok {
return 0, fmt.Errorf("kube-explorer for cluster %s has already started", config)
}
// 2. 检查命令存在性和可执行性
CheckCommandExist(KubeExplorerCommand)
checkExplorerCmd()
// 3. 查询数据库中的 Explorer 配置
exp, err := DefaultDB.GetExplorer(config)
// 4. 如果未配置或未启用,分配端口并保存
if exp == nil || !exp.Enabled {
var port int
if exp == nil {
port, err = k3dutil.GetFreePort() // 新建:自动获取空闲端口
} else {
port = exp.Port // 已存在:复用之前的端口
}
exp = &Explorer{
ContextName: config,
Port: port,
Enabled: true,
}
DefaultDB.SaveExplorer(exp)
}
// 5. 启动 kube-explorer 进程
explorerCtx, cancel := context.WithCancel(ctx)
ExplorerWatchers[config] = cancel // 保存 cancel 函数
go func(ctx context.Context, config string, port int) {
StartKubeExplorer(ctx, config, port)
}(explorerCtx, config, exp.Port)
return exp.Port, nil
}
DisableExplorer — 停止指定集群的 Explorer
func DisableExplorer(config string) error {
// 1. 检查是否在运行
if _, ok := ExplorerWatchers[config]; !ok {
return fmt.Errorf("cann't disable unactive kube-explorer for cluster %s", config)
}
// 2. 更新数据库状态为禁用
exp, err := DefaultDB.GetExplorer(config)
if exp == nil || exp.Enabled {
DefaultDB.SaveExplorer(&Explorer{
ContextName: config,
Port: port,
Enabled: false,
})
}
// 3. 调用 cancel 停止进程
ExplorerWatchers[config]()
delete(ExplorerWatchers, config)
return nil
}
InitExplorer — 启动时恢复所有 Explorer
func InitExplorer(ctx context.Context) {
// 1. 查询所有 Explorer 配置
expList, err := DefaultDB.ListExplorer()
// 2. 遍历已启用的 Explorer,逐个启动
for _, exp := range expList {
if exp.Enabled {
go func(ctx context.Context, name string) {
EnableExplorer(ctx, name)
}(ctx, exp.ContextName)
}
}
}
StartKubeExplorer — 执行外部命令
func StartKubeExplorer(ctx context.Context, config string, port int) error {
// 构建命令:kube-explorer
// –kubeconfig=~/.autok3s/.kube/config
// –context=<config>
// –http-listen-port=<port>
// –https-listen-port=0
explorer := exec.CommandContext(ctx, KubeExplorerCommand,
fmt.Sprintf("–kubeconfig=%s", filepath.Join(CfgPath, KubeCfgFile)),
fmt.Sprintf("–context=%s", config),
fmt.Sprintf("–http-listen-port=%d", port),
"–https-listen-port=0")
explorer.Stdout = os.Stdout
explorer.Stderr = os.Stderr
explorer.Start()
return explorer.Wait()
}
3.9 Metrics 采集
文件:metrics.go
SetupPrometheusMetrics — 初始化指标
func SetupPrometheusMetrics(version string) {
// 1. 设置活跃指标(带 UUID + version 标签)
labels := uuidLabels()
labels["version"] = version
metrics.Active.With(labels).Set(1)
// 2. 从数据库恢复集群计数
clusters, err := DefaultDB.ListCluster("")
for _, cluster := range clusters {
metrics.ClusterCount.With(getLabelsFromMeta(cluster.Metadata)).Add(1)
}
// 3. 从数据库恢复模板计数
templates, err := DefaultDB.ListTemplates()
for _, template := range templates {
metrics.TemplateCount.With(getLabelsFromMeta(template.Metadata)).Add(1)
}
// 4. 设置指标启用函数(动态判断遥测是否开启)
metrics.SetupEnableFunc(func() bool {
enable := GetTelemetryEnable()
return enable != nil && *enable
})
}
getLabelsFromMeta — 标签生成
func getLabelsFromMeta(meta types.Metadata) prometheus.Labels {
version := meta.K3sVersion
if version == "" {
version = "unknown"
}
uuid := GetUUID()
if uuid == "" {
uuid = "unknown"
}
return prometheus.Labels{
"provider": meta.Provider,
"k3sversion": version,
"install_uuid": uuid,
}
}
三个标签维度:provider(云提供商)、k3sversion(K3s 版本)、install_uuid(安装实例 ID)。
MetricsPrompt — 遥测提示
func MetricsPrompt(cmd *cobra.Command) {
// 1. 跳过特定命令
if cmd.Use == "version" || cmd.Use == "serve" || … {
return
}
// 2. 非终端环境跳过
if !utils.IsTerm() {
return
}
// 3. 已设置过则跳过
if should := GetTelemetryEnable(); should != nil {
return
}
// 4. 首次使用时询问用户
rtn := uti


