fix: 网络检测增强 + Tracker changeset 清理

网络检测 (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天)
- 先按年龄裁剪,再按数量裁剪
This commit is contained in:
root
2026-07-03 21:31:00 +08:00
parent c125e2afe5
commit 90da05e312
4 changed files with 158 additions and 33 deletions

View File

@ -150,7 +150,10 @@ func main() {
// 变更追踪(overlayfs) // 变更追踪(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 { if err := trk.Init(); err != nil {
log.Printf("[homed] warning: tracker init: %v", err) log.Printf("[homed] warning: tracker init: %v", err)
} else { } else {

View File

@ -12,11 +12,11 @@ import (
) )
type Monitor struct { type Monitor struct {
mu sync.RWMutex mu sync.RWMutex
client *http.Client client *http.Client
interval time.Duration interval time.Duration
endpoints []string endpoints []string
status []EndpointStatus status []EndpointStatus
} }
type EndpointStatus struct { type EndpointStatus struct {
@ -142,7 +142,11 @@ func (m *Monitor) AggregateResult() types.NetworkCheckResult {
m.mu.RLock() m.mu.RLock()
defer m.mu.RUnlock() defer m.mu.RUnlock()
result := types.NetworkCheckResult{LLMAPIReachable: true, DNSResolving: true} result := types.NetworkCheckResult{
LLMAPIReachable: true,
DNSResolving: true,
TCPReachable: true,
}
var totalLatency time.Duration var totalLatency time.Duration
checked := 0 checked := 0
@ -161,14 +165,49 @@ func (m *Monitor) AggregateResult() types.NetworkCheckResult {
result.Latency = totalLatency / time.Duration(checked) 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 return result
} }
func (m *Monitor) checkDNS() bool { func (m *Monitor) checkDNSMulti() bool {
_, err := net.LookupHost("google.com") targets := []string{"google.com", "baidu.com", "cloudflare.com"}
if err != nil { for _, target := range targets {
_, err = net.LookupHost("baidu.com") _, 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
} }

View File

@ -7,24 +7,39 @@ import (
"os" "os"
"os/exec" "os/exec"
"path/filepath" "path/filepath"
"sort"
"strings"
"sync" "sync"
"time"
) )
type Tracker struct { type Tracker struct {
mu sync.Mutex mu sync.Mutex
dataDir string dataDir string
workDir string workDir string
lowerDir string lowerDir string
upperDir string upperDir string
mergeDir string mergeDir string
mounted bool mounted bool
active bool active bool
before *FSState before *FSState
changeSets []*ChangeSet changeSets []*ChangeSet
keepChangesets int // 保留最近 N 份 changeset,0 = 不限
maxChangesetAge time.Duration // changeset 最大保留时长,0 = 不限
} }
func NewTracker(dataDir, workDir string) *Tracker { type TrackerOption func(*Tracker)
return &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, dataDir: dataDir,
workDir: workDir, workDir: workDir,
lowerDir: filepath.Join(workDir, "lower"), lowerDir: filepath.Join(workDir, "lower"),
@ -32,6 +47,10 @@ func NewTracker(dataDir, workDir string) *Tracker {
mergeDir: filepath.Join(workDir, "merged"), mergeDir: filepath.Join(workDir, "merged"),
changeSets: make([]*ChangeSet, 0), changeSets: make([]*ChangeSet, 0),
} }
for _, opt := range opts {
opt(t)
}
return t
} }
func (t *Tracker) Init() error { func (t *Tracker) Init() error {
@ -40,6 +59,7 @@ func (t *Tracker) Init() error {
return fmt.Errorf("create overlay dir %s: %w", d, err) return fmt.Errorf("create overlay dir %s: %w", d, err)
} }
} }
t.cleanupChangeSets()
log.Printf("[tracker] initialized (work=%s)", t.workDir) log.Printf("[tracker] initialized (work=%s)", t.workDir)
return nil 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 { func (t *Tracker) Rollback() error {
t.mu.Lock() t.mu.Lock()
defer t.mu.Unlock() defer t.mu.Unlock()
@ -213,11 +292,13 @@ func (t *Tracker) Stats() map[string]interface{} {
totalChanges += len(cs.Files) totalChanges += len(cs.Files)
} }
return map[string]interface{}{ return map[string]interface{}{
"mounted": t.mounted, "mounted": t.mounted,
"active": t.active, "active": t.active,
"change_sets": len(t.changeSets), "change_sets": len(t.changeSets),
"total_changes": totalChanges, "total_changes": totalChanges,
"merge_dir": t.mergeDir, "merge_dir": t.mergeDir,
"upper_dir": t.upperDir, "upper_dir": t.upperDir,
"keep_changesets": t.keepChangesets,
"max_changeset_age": t.maxChangesetAge.String(),
} }
} }

View File

@ -46,10 +46,12 @@ type Heartbeat struct {
} }
type NetworkCheckResult struct { type NetworkCheckResult struct {
LLMAPIReachable bool `json:"llm_api_reachable"` LLMAPIReachable bool `json:"llm_api_reachable"`
DNSResolving bool `json:"dns_resolving"` DNSResolving bool `json:"dns_resolving"`
TCPReachable bool `json:"tcp_reachable"`
Latency time.Duration `json:"latency_ms"` Latency time.Duration `json:"latency_ms"`
Error string `json:"error,omitempty"` LatencyDegraded bool `json:"latency_degraded"`
Error string `json:"error,omitempty"`
} }
type SnapshotPolicy struct { type SnapshotPolicy struct {