From 77e12f7329f3d8e2187b349fb87e45fc01e50309 Mon Sep 17 00:00:00 2001 From: root Date: Fri, 3 Jul 2026 21:31:00 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E7=BD=91=E7=BB=9C=E6=A3=80=E6=B5=8B?= =?UTF-8?q?=E5=A2=9E=E5=BC=BA=20+=20Tracker=20changeset=20=E6=B8=85?= =?UTF-8?q?=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 网络检测 (monitor.go): - TCP 拨测 (8.8.8.8:53 / 1.1.1.1:53 / 208.67.222.222:53) - 延迟阈值检测 (>5s 标记降级) - DNS 多目标检测 (google/baidu/cloudflare) - NetworkCheckResult 新增 TCPReachable + LatencyDegraded Tracker: - 新增 WithKeepChangesets / WithMaxChangesetAge 选项 - Init 时自动清理过期 changeset (默认保留100份/30天) - 先按年龄裁剪,再按数量裁剪 --- cmd/homed/main.go | 5 +- internal/network/monitor.go | 61 +++++++++++++++---- internal/tracker/tracker.go | 117 ++++++++++++++++++++++++++++++------ pkg/types/types.go | 8 ++- 4 files changed, 158 insertions(+), 33 deletions(-) diff --git a/cmd/homed/main.go b/cmd/homed/main.go index 11e9b41..edacefe 100644 --- a/cmd/homed/main.go +++ b/cmd/homed/main.go @@ -150,7 +150,10 @@ func main() { // 变更追踪(overlayfs) // ======================================================================== - trk := tracker.NewTracker(cfg.Daemon.DataDir, agentWorkDir) + trk := tracker.NewTracker(cfg.Daemon.DataDir, agentWorkDir, + tracker.WithKeepChangesets(100), + tracker.WithMaxChangesetAge(30*24*time.Hour), + ) if err := trk.Init(); err != nil { log.Printf("[homed] warning: tracker init: %v", err) } else { diff --git a/internal/network/monitor.go b/internal/network/monitor.go index fa6d8b1..4f7ccd8 100644 --- a/internal/network/monitor.go +++ b/internal/network/monitor.go @@ -12,11 +12,11 @@ import ( ) type Monitor struct { - mu sync.RWMutex - client *http.Client - interval time.Duration + mu sync.RWMutex + client *http.Client + interval time.Duration endpoints []string - status []EndpointStatus + status []EndpointStatus } type EndpointStatus struct { @@ -142,7 +142,11 @@ func (m *Monitor) AggregateResult() types.NetworkCheckResult { m.mu.RLock() defer m.mu.RUnlock() - result := types.NetworkCheckResult{LLMAPIReachable: true, DNSResolving: true} + result := types.NetworkCheckResult{ + LLMAPIReachable: true, + DNSResolving: true, + TCPReachable: true, + } var totalLatency time.Duration checked := 0 @@ -161,14 +165,49 @@ func (m *Monitor) AggregateResult() types.NetworkCheckResult { result.Latency = totalLatency / time.Duration(checked) } - result.DNSResolving = m.checkDNS() + // 延迟阈值检测:平均延迟 > 5s 标记为降级 + if result.Latency > 5*time.Second { + result.LatencyDegraded = true + if result.Error == "" { + result.Error = fmt.Sprintf("high latency: %v", result.Latency) + } + } + + // DNS 多目标检测 + result.DNSResolving = m.checkDNSMulti() + + // TCP 拨测:检测基础网络通畅性 + result.TCPReachable = m.checkTCPReachability() + return result } -func (m *Monitor) checkDNS() bool { - _, err := net.LookupHost("google.com") - if err != nil { - _, err = net.LookupHost("baidu.com") +func (m *Monitor) checkDNSMulti() bool { + targets := []string{"google.com", "baidu.com", "cloudflare.com"} + for _, target := range targets { + _, err := net.LookupHost(target) + if err == nil { + return true + } } - return err == nil + return false +} + +func (m *Monitor) checkTCPReachability() bool { + targets := []struct { + host string + port string + }{ + {"8.8.8.8", "53"}, + {"1.1.1.1", "53"}, + {"208.67.222.222", "53"}, + } + for _, t := range targets { + conn, err := net.DialTimeout("tcp", net.JoinHostPort(t.host, t.port), 3*time.Second) + if err == nil { + conn.Close() + return true + } + } + return false } diff --git a/internal/tracker/tracker.go b/internal/tracker/tracker.go index 5e267e5..e732609 100644 --- a/internal/tracker/tracker.go +++ b/internal/tracker/tracker.go @@ -7,24 +7,39 @@ import ( "os" "os/exec" "path/filepath" + "sort" + "strings" "sync" + "time" ) type Tracker struct { - mu sync.Mutex - dataDir string - workDir string - lowerDir string - upperDir string - mergeDir string - mounted bool - active bool - before *FSState - changeSets []*ChangeSet + mu sync.Mutex + dataDir string + workDir string + lowerDir string + upperDir string + mergeDir string + mounted bool + active bool + before *FSState + changeSets []*ChangeSet + keepChangesets int // 保留最近 N 份 changeset,0 = 不限 + maxChangesetAge time.Duration // changeset 最大保留时长,0 = 不限 } -func NewTracker(dataDir, workDir string) *Tracker { - return &Tracker{ +type TrackerOption func(*Tracker) + +func WithKeepChangesets(n int) TrackerOption { + return func(t *Tracker) { t.keepChangesets = n } +} + +func WithMaxChangesetAge(d time.Duration) TrackerOption { + return func(t *Tracker) { t.maxChangesetAge = d } +} + +func NewTracker(dataDir, workDir string, opts ...TrackerOption) *Tracker { + t := &Tracker{ dataDir: dataDir, workDir: workDir, lowerDir: filepath.Join(workDir, "lower"), @@ -32,6 +47,10 @@ func NewTracker(dataDir, workDir string) *Tracker { mergeDir: filepath.Join(workDir, "merged"), changeSets: make([]*ChangeSet, 0), } + for _, opt := range opts { + opt(t) + } + return t } func (t *Tracker) Init() error { @@ -40,6 +59,7 @@ func (t *Tracker) Init() error { return fmt.Errorf("create overlay dir %s: %w", d, err) } } + t.cleanupChangeSets() log.Printf("[tracker] initialized (work=%s)", t.workDir) return nil } @@ -172,6 +192,65 @@ func (t *Tracker) saveChangeSet(cs *ChangeSet) { } } +func (t *Tracker) cleanupChangeSets() { + dir := filepath.Join(t.dataDir, "changesets") + entries, err := os.ReadDir(dir) + if err != nil { + return + } + + type csFile struct { + name string + info os.FileInfo + } + var files []csFile + for _, e := range entries { + if e.IsDir() || !strings.HasSuffix(e.Name(), ".json") { + continue + } + info, err := e.Info() + if err != nil { + continue + } + files = append(files, csFile{name: e.Name(), info: info}) + } + + // 按修改时间排序 + sort.Slice(files, func(i, j int) bool { + return files[i].info.ModTime().Before(files[j].info.ModTime()) + }) + + now := time.Now() + remaining := make([]csFile, 0, len(files)) + + for _, f := range files { + keep := true + + if t.maxChangesetAge > 0 && now.Sub(f.info.ModTime()) > t.maxChangesetAge { + keep = false + } + + if keep { + remaining = append(remaining, f) + } + } + + // 再按数量裁剪 + if t.keepChangesets > 0 && len(remaining) > t.keepChangesets { + excess := len(remaining) - t.keepChangesets + for i := 0; i < excess; i++ { + path := filepath.Join(dir, remaining[i].name) + os.Remove(path) + } + remaining = remaining[excess:] + } + + if len(files) != len(remaining) { + log.Printf("[tracker] cleanup: removed %d changesets, kept %d", + len(files)-len(remaining), len(remaining)) + } +} + func (t *Tracker) Rollback() error { t.mu.Lock() defer t.mu.Unlock() @@ -213,11 +292,13 @@ func (t *Tracker) Stats() map[string]interface{} { totalChanges += len(cs.Files) } return map[string]interface{}{ - "mounted": t.mounted, - "active": t.active, - "change_sets": len(t.changeSets), - "total_changes": totalChanges, - "merge_dir": t.mergeDir, - "upper_dir": t.upperDir, + "mounted": t.mounted, + "active": t.active, + "change_sets": len(t.changeSets), + "total_changes": totalChanges, + "merge_dir": t.mergeDir, + "upper_dir": t.upperDir, + "keep_changesets": t.keepChangesets, + "max_changeset_age": t.maxChangesetAge.String(), } } diff --git a/pkg/types/types.go b/pkg/types/types.go index 77ce133..e64453f 100644 --- a/pkg/types/types.go +++ b/pkg/types/types.go @@ -46,10 +46,12 @@ type Heartbeat struct { } type NetworkCheckResult struct { - LLMAPIReachable bool `json:"llm_api_reachable"` - DNSResolving bool `json:"dns_resolving"` + LLMAPIReachable bool `json:"llm_api_reachable"` + DNSResolving bool `json:"dns_resolving"` + TCPReachable bool `json:"tcp_reachable"` Latency time.Duration `json:"latency_ms"` - Error string `json:"error,omitempty"` + LatencyDegraded bool `json:"latency_degraded"` + Error string `json:"error,omitempty"` } type SnapshotPolicy struct {