Skip to content

Commit 64e2e7a

Browse files
committed
fix: reconcile persisted state with runtime
1 parent aa2dcf7 commit 64e2e7a

17 files changed

Lines changed: 1435 additions & 222 deletions

‎API.md‎

Lines changed: 36 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -699,6 +699,15 @@ Goja 控制脚本默认只能访问本插件资源。每个插件默认持有一
699699
"page_size": 20,
700700
"total": 5,
701701
"binary_hash": "abc123",
702+
"dataplane": {
703+
"status": "applied",
704+
"desired_generation": 12,
705+
"applied_generation": 12,
706+
"pending_generations": 0
707+
},
708+
"plugins": {
709+
"status": "applied"
710+
},
702711
"workers": [
703712
{
704713
"kind": "kernel",
@@ -735,6 +744,8 @@ Goja 控制脚本默认只能访问本插件资源。每个插件默认持有一
735744
- `draining`
736745
- `error`
737746

747+
`dataplane` 与 `plugins` 是鉴权后的 reconcile 诊断视图。失败时会额外返回 `retry_count`、`last_error`、`last_attempt_at` 和 `last_applied_at`;不要把这些错误详情写入匿名日志或公开状态页。
748+
738749
### KernelRuntimeResponse
739750

740751
`GET /api/kernel/runtime` 是运行时调试视图,常见关键字段:
@@ -819,16 +830,25 @@ Goja 控制脚本默认只能访问本插件资源。每个插件默认持有一
819830
用途:
820831

821832
- 判断本机 Veer 运行时是否 ready,而不只是进程是否完成启动
822-
- 持续检查承载有效规则和范围的 userspace worker、shared proxy,以及所有已启用 Egress NAT 是否已经进入预期的 kernel runtime
833+
- 持续检查控制面 generation 是否已应用、承载有效规则和范围的 userspace worker、shared proxy,以及所有已启用 Egress NAT 是否已经进入预期的 kernel runtime
823834
- 不探测转发后端服务,也不执行端到端转发链路健康检查
824-
- 不需要 Bearer Token
835+
- 不需要 Bearer Token;只返回状态和 generation 计数,不返回内部错误文本。完整错误位于鉴权后的 `GET /api/workers`
825836

826837
启动中或必要运行时组件未就绪时返回 `503`:
827838

828839
```json
829840
{
830841
"status": "starting",
831-
"ready": false
842+
"ready": false,
843+
"dataplane": {
844+
"status": "pending",
845+
"desired_generation": 12,
846+
"applied_generation": 11,
847+
"pending_generations": 1
848+
},
849+
"plugins": {
850+
"status": "applied"
851+
}
832852
}
833853
```
834854

@@ -837,10 +857,21 @@ Goja 控制脚本默认只能访问本插件资源。每个插件默认持有一
837857
```json
838858
{
839859
"status": "ready",
840-
"ready": true
860+
"ready": true,
861+
"dataplane": {
862+
"status": "applied",
863+
"desired_generation": 12,
864+
"applied_generation": 12,
865+
"pending_generations": 0
866+
},
867+
"plugins": {
868+
"status": "applied"
869+
}
841870
}
842871
```
843872

873+
`plugins.status=error` 表示已启用插件的期望状态尚未完整应用,会触发后台退避重试,并使 `/readyz` 保持 `503`,直到运行态收敛。
874+
844875
`GET /metrics`
845876

846877
用途与约束:
@@ -1491,6 +1522,7 @@ Goja 控制脚本默认只能访问本插件资源。每个插件默认持有一
14911522
- `page_size` 最大 `1000`
14921523
- 不传 `page_size` 时返回全部 worker
14931524
- 该接口会把规则 worker、范围 worker、共享站点 worker、kernel worker、egress_nat worker 合并返回
1525+
- 顶层 `dataplane` 与 `plugins` 返回完整 reconcile 状态和错误详情,用于确认数据库期望状态是否已经应用到运行态
14941526

14951527
### 9.2 获取内核运行时
14961528

‎internal/app/api.go‎

Lines changed: 18 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -171,16 +171,27 @@ func buildAPIHandler(cfg *Config, db *sql.DB, pm *ProcessManager) http.Handler {
171171
return
172172
}
173173
w.Header().Set("Cache-Control", "no-store")
174-
ready := pm.isReady()
174+
ready, reconcile, plugins := pm.readinessReconcileSnapshot()
175175
statusCode := http.StatusServiceUnavailable
176176
status := "starting"
177177
if ready {
178178
statusCode = http.StatusOK
179179
status = "ready"
180+
} else if reconcile.Status == "error" || plugins.Status == "error" {
181+
status = "error"
182+
} else if reconcile.Status == "pending" {
183+
status = "reconciling"
180184
}
181185
writeJSON(w, statusCode, map[string]interface{}{
182186
"status": status,
183187
"ready": ready,
188+
"dataplane": map[string]interface{}{
189+
"status": reconcile.Status,
190+
"desired_generation": reconcile.DesiredGeneration,
191+
"applied_generation": reconcile.AppliedGeneration,
192+
"pending_generations": reconcile.PendingGenerations,
193+
},
194+
"plugins": map[string]string{"status": plugins.Status},
184195
})
185196
})
186197
mux.HandleFunc("/metrics", authMiddleware(cfg, func(w http.ResponseWriter, r *http.Request) {
@@ -2548,8 +2559,12 @@ func handleListWorkers(w http.ResponseWriter, r *http.Request, db *sql.DB, pm *P
25482559
kernelRuleIDs := make(map[int64]bool)
25492560
kernelRangeIDs := make(map[int64]bool)
25502561
kernelBinaryHash := ""
2562+
var dataplaneStatus DataplaneReconcileStatus
2563+
var pluginStatus PluginReconcileStatus
25512564
needsSharedSiteCount := false
25522565
pm.mu.Lock()
2566+
dataplaneStatus = pm.dataplaneReconcileStatusLocked()
2567+
pluginStatus = pm.pluginReconcileStatusLocked()
25532568
for id := range pm.kernelRules {
25542569
kernelRuleIDs[id] = true
25552570
}
@@ -2973,6 +2988,8 @@ func handleListWorkers(w http.ResponseWriter, r *http.Request, db *sql.DB, pm *P
29732988
PageSize: pageSize,
29742989
Total: total,
29752990
BinaryHash: kernelBinaryHash,
2991+
Dataplane: dataplaneStatus,
2992+
Plugins: pluginStatus,
29762993
Workers: out,
29772994
}
29782995
writeJSON(w, http.StatusOK, resp)

‎internal/app/api_health_test.go‎

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package app
22

33
import (
4+
"bytes"
45
"encoding/json"
56
"net"
67
"net/http"
@@ -83,6 +84,45 @@ func TestBuildAPIHandlerReadyzReflectsProcessManagerState(t *testing.T) {
8384
}
8485
}
8586

87+
func TestBuildAPIHandlerReadyzRedactsReconcileErrors(t *testing.T) {
88+
db := openTestDB(t)
89+
pm := &ProcessManager{
90+
ready: true,
91+
desiredGeneration: 2,
92+
appliedGeneration: 1,
93+
pluginReconcileLastError: "secret plugin path",
94+
ruleWorkers: map[int]*WorkerInfo{
95+
0: {errored: true, lastError: "secret dataplane path"},
96+
},
97+
}
98+
handler := buildAPIHandler(&Config{WebBind: "127.0.0.1", WebPort: 8080, WebToken: "test-token"}, db, pm)
99+
100+
rec := httptest.NewRecorder()
101+
handler.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, "/readyz", nil))
102+
103+
if rec.Code != http.StatusServiceUnavailable {
104+
t.Fatalf("GET /readyz status = %d, want %d", rec.Code, http.StatusServiceUnavailable)
105+
}
106+
if bytes.Contains(rec.Body.Bytes(), []byte("secret")) || bytes.Contains(rec.Body.Bytes(), []byte("last_error")) {
107+
t.Fatalf("/readyz leaked reconcile detail: %s", rec.Body.String())
108+
}
109+
var resp struct {
110+
Status string `json:"status"`
111+
Dataplane struct {
112+
Status string `json:"status"`
113+
} `json:"dataplane"`
114+
Plugins struct {
115+
Status string `json:"status"`
116+
} `json:"plugins"`
117+
}
118+
if err := json.NewDecoder(rec.Body).Decode(&resp); err != nil {
119+
t.Fatalf("decode /readyz response: %v", err)
120+
}
121+
if resp.Status != "error" || resp.Dataplane.Status != "error" || resp.Plugins.Status != "error" {
122+
t.Fatalf("/readyz response = %+v", resp)
123+
}
124+
}
125+
86126
func TestBuildAPIHandlerReadyzRejectsUnavailableEgressNAT(t *testing.T) {
87127
db := openTestDB(t)
88128
pm := &ProcessManager{

‎internal/app/api_plugin_data.go‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -154,7 +154,15 @@ func handlePluginStateAPI(w http.ResponseWriter, r *http.Request, cfg *Config, d
154154
return
155155
}
156156
if pm != nil {
157-
pm.reconcilePluginsForRuntime()
157+
if _, err := pm.reconcilePluginsForRuntimeWithError(); err != nil {
158+
recordPluginAudit(db, pluginID, "plugin.state", "api", "error", map[string]any{"enabled": *req.Enabled, "error": err.Error()})
159+
writeJSON(w, http.StatusServiceUnavailable, map[string]any{
160+
"error": err.Error(),
161+
"state": pluginStateResponseForID(cfg, db, pm, pluginID),
162+
"stored": true,
163+
})
164+
return
165+
}
158166
pm.redistributeWorkers()
159167
}
160168
recordPluginAudit(db, pluginID, "plugin.state", "api", "success", map[string]any{"enabled": *req.Enabled})

‎internal/app/api_workers_test.go‎

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,9 +4,44 @@ import (
44
"encoding/json"
55
"errors"
66
"net/http/httptest"
7+
"strings"
78
"testing"
89
)
910

11+
func TestHandleListWorkersIncludesDetailedReconcileStatus(t *testing.T) {
12+
db := openTestDB(t)
13+
pm := &ProcessManager{
14+
binaryHash: "deadbeefcafebabe",
15+
ruleWorkers: map[int]*WorkerInfo{},
16+
rangeWorkers: map[int]*WorkerInfo{},
17+
kernelRules: map[int64]bool{},
18+
kernelRanges: map[int64]bool{},
19+
desiredGeneration: 4,
20+
appliedGeneration: 3,
21+
reconcileLastError: "dataplane apply failed",
22+
reconcileRetryCount: 2,
23+
pluginReconcileLastError: "plugin apply failed",
24+
pluginReconcileRetryCount: 3,
25+
}
26+
27+
w := httptest.NewRecorder()
28+
handleListWorkers(w, httptest.NewRequest("GET", "/api/workers", nil), db, pm)
29+
if w.Code != 200 {
30+
t.Fatalf("unexpected status: %d body=%s", w.Code, w.Body.String())
31+
}
32+
var resp WorkerListResponse
33+
if err := json.Unmarshal(w.Body.Bytes(), &resp); err != nil {
34+
t.Fatalf("decode response: %v", err)
35+
}
36+
if resp.Dataplane.Status != "error" || resp.Dataplane.DesiredGeneration != 4 || resp.Dataplane.AppliedGeneration != 3 ||
37+
resp.Dataplane.RetryCount != 2 || !strings.Contains(resp.Dataplane.LastError, "dataplane apply failed") {
38+
t.Fatalf("dataplane reconcile status = %+v", resp.Dataplane)
39+
}
40+
if resp.Plugins.Status != "error" || resp.Plugins.RetryCount != 3 || !strings.Contains(resp.Plugins.LastError, "plugin apply failed") {
41+
t.Fatalf("plugin reconcile status = %+v", resp.Plugins)
42+
}
43+
}
44+
1045
func TestHandleListWorkersIncludesEgressNATWorker(t *testing.T) {
1146
db := openTestDB(t)
1247

0 commit comments

Comments
 (0)